fix(timing): harden live active projections

This commit is contained in:
Bryan Helmkamp 2026-07-25 23:43:43 -04:00
parent c4971b93d3
commit dd9f75fb05
No known key found for this signature in database
17 changed files with 523 additions and 86 deletions

View file

@ -73,6 +73,7 @@ import {
type RunFileSelection,
type RunGraphDirection,
} from "./query-keys";
import { isTerminalRunStatus } from "./run-actions";
const immutableOptions: SWRConfiguration = {
revalidateIfStale: false,
@ -179,10 +180,19 @@ export function useRunsPage(opts: RunsPageOptions = {}, enabled = true) {
);
}
export function useRun(id: string | undefined) {
export function useRun(id: string | undefined, refreshInterval?: number) {
return useSWR<Run | null>(
id ? queryKeys.runs.detail(id) : null,
() => apiNullableData(() => runsApi.retrieveRun(id!)),
refreshInterval
? {
refreshInterval: (run) =>
run?.timestamps.started_at &&
!isTerminalRunStatus(run.lifecycle.status.kind)
? refreshInterval
: 0,
}
: undefined,
);
}

View file

@ -87,8 +87,6 @@ describe("queryKeys", () => {
test("agent activity events invalidate per-stage resources", () => {
for (const event of [
"stage.prompt",
"agent.tool.started",
"agent.tool.completed",
"command.started",
"command.completed",
]) {
@ -97,8 +95,19 @@ describe("queryKeys", () => {
queryKeys.runs.stageContextWindow("run-1", "stage-1"),
]);
}
for (const event of ["agent.tool.started", "agent.tool.completed"]) {
expect(queryKeysForRunEvent("run-1", event, "stage-1")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
queryKeys.runs.stageEvents("run-1", "stage-1"),
queryKeys.runs.stageContextWindow("run-1", "stage-1"),
]);
}
expect(queryKeysForRunEvent("run-1", "agent.message", "stage-1")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
queryKeys.runs.stageEvents("run-1", "stage-1"),
queryKeys.runs.stageContextWindow("run-1", "stage-1"),
]);
@ -106,7 +115,9 @@ describe("queryKeys", () => {
test("agent message without a node_id still invalidates projected state", () => {
expect(queryKeysForRunEvent("run-1", "agent.message")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
]);
});
});

View file

@ -85,6 +85,8 @@ describe("queryKeysForRunEvent", () => {
test("interrupt settlement invalidates projected control state and stage activity", () => {
expect(queryKeysForRunEvent("run-1", "agent.round.interrupted", "nap@1")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.billing("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.events("run-1", 1000),
queryKeys.runs.stageEvents("run-1", "nap@1"),
@ -134,22 +136,40 @@ describe("queryKeysForRunEvent", () => {
"agent.error",
]) {
expect(queryKeysForRunEvent("run-1", event, "code@1")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
queryKeys.runs.stageEvents("run-1", "code@1"),
]);
}
expect(
queryKeysForRunEvent("run-1", "agent.message", "code@1"),
).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
queryKeys.runs.stageEvents("run-1", "code@1"),
queryKeys.runs.stageContextWindow("run-1", "code@1"),
]);
expect(queryKeysForRunEvent("run-1", "agent.session.ended")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
]);
});
test("tool timing events invalidate live summaries and stage resources", () => {
for (const event of ["agent.tool.started", "agent.tool.completed"]) {
expect(queryKeysForRunEvent("run-1", event, "code@1")).toEqual([
queryKeys.runs.detail("run-1"),
queryKeys.runs.state("run-1"),
queryKeys.runs.billing("run-1"),
queryKeys.runs.stageEvents("run-1", "code@1"),
queryKeys.runs.stageContextWindow("run-1", "code@1"),
]);
}
});
test("watchdog timeout refreshes the stage events for that stage", () => {
expect(
queryKeysForRunEvent("run-1", "watchdog.timeout", "code@1"),

View file

@ -118,6 +118,10 @@ const INFERENCE_EVENTS = new Set([
"agent.error",
"agent.session.ended",
]);
const TOOL_TIMING_EVENTS = new Set([
"agent.tool.started",
"agent.tool.completed",
]);
// Todo / task mutation events refresh `getRunState` consumers (so per-stage
// todo projections update live) and the run events list.
const TODO_EVENTS = new Set([
@ -184,6 +188,12 @@ export function queryKeysForRunEvent(
if (AGENT_CONTROL_STATE_EVENTS.has(event)) {
keys.unshift(queryKeys.runs.state(runId));
}
if (event === "agent.round.interrupted") {
keys.unshift(
queryKeys.runs.detail(runId),
queryKeys.runs.billing(runId),
);
}
if (stageId) {
keys.push(queryKeys.runs.stageEvents(runId, stageId));
keys.push(queryKeys.runs.stageContextWindow(runId, stageId));
@ -192,7 +202,11 @@ export function queryKeysForRunEvent(
}
if (INFERENCE_EVENTS.has(event)) {
const keys: Key[] = [queryKeys.runs.state(runId)];
const keys: Key[] = [
queryKeys.runs.detail(runId),
queryKeys.runs.state(runId),
queryKeys.runs.billing(runId),
];
if (stageId) {
keys.push(queryKeys.runs.stageEvents(runId, stageId));
if (event === "agent.message") {
@ -202,6 +216,19 @@ export function queryKeysForRunEvent(
return keys;
}
if (TOOL_TIMING_EVENTS.has(event)) {
const keys: Key[] = [
queryKeys.runs.detail(runId),
queryKeys.runs.state(runId),
queryKeys.runs.billing(runId),
];
if (stageId) {
keys.push(queryKeys.runs.stageEvents(runId, stageId));
keys.push(queryKeys.runs.stageContextWindow(runId, stageId));
}
return keys;
}
if (event === "watchdog.timeout") {
return stageId ? [queryKeys.runs.stageEvents(runId, stageId)] : [];
}

View file

@ -69,6 +69,8 @@ import {
export const handle = { hideHeader: true };
const RUN_TIMING_REFRESH_INTERVAL_MS = 30_000;
type LifecycleTrigger = () => Promise<LifecycleMutationResult | undefined>;
export function meta({ data }: any) {
@ -78,7 +80,7 @@ export function meta({ data }: any) {
export default function RunDetail({ params }: { params: { id: string } }) {
const demoMode = useDemoMode();
const runQuery = useRun(params.id);
const runQuery = useRun(params.id, RUN_TIMING_REFRESH_INTERVAL_MS);
const runStateQuery = useRunState(params.id);
const summary = runQuery.data;
const run = summary ? buildRunDetailRun(summary) : null;
@ -116,7 +118,7 @@ export default function RunDetail({ params }: { params: { id: string } }) {
childrenCount,
});
const steerBarRef = useRef<SteerBarHandle | null>(null);
const now = useTickingNow(30_000);
const now = useTickingNow(RUN_TIMING_REFRESH_INTERVAL_MS);
const { fullHeight, hideSteerBar } = childRouteLayoutFlags(matches);
useRunEvents(params.id);

View file

@ -155,8 +155,9 @@ git diff --check
- Inference time is Fabro-observed LLM request/stream elapsed time, not
provider-reported model-only compute time.
- LLM retry backoff, queueing outside a request/stream, human waits, steering
waits, and scheduler gaps are wall time but not active time.
- Queueing outside a request/stream, human waits, steering waits, and scheduler
gaps are wall time but not active time. Retry delay inside an open LLM request
bracket follows the executor stopwatch and counts as inference time.
- ~~Active timing is finalized-event based in v1; live active-time ticking can
be added later if it becomes necessary.~~ **Superseded 2026-07-25.** It became
necessary: a run parked in one long agent stage reported ~12% of its wall time
@ -164,7 +165,6 @@ git diff --check
accumulate inference and tool brackets from the event log and expose
`StageProjection::live_timing(now)`, the active-time twin of
`live_wall_time_ms`. Finalized values remain authoritative and still replace
the live estimate at terminal events. See
`.ai/plans/live-active-time-accumulation.md`.
the live estimate at terminal events. Implemented in PR #647.
- No compatibility layer is required for existing API clients or stored run
event data.

View file

@ -10772,9 +10772,16 @@ components:
duplicated completion in a replayed log cannot drain the batch early.
type: object
required:
- session_id
- started_at
- open_call_ids
properties:
session_id:
type: string
description: >
Root agent session that dispatched the batch. Transitions are
gated on it so delayed events from a replaced session cannot
mutate the current batch.
started_at:
type: string
format: date-time
@ -10783,6 +10790,8 @@ components:
other calls were outstanding.
open_call_ids:
type: array
minItems: 1
uniqueItems: true
items:
type: string
description: Calls dispatched but not yet completed, by tool call id.

View file

@ -11,7 +11,7 @@ use axum_extra::extract::Query as ExtraQuery;
use base64::Engine as _;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use bytes::Bytes;
use chrono::Utc;
use chrono::{DateTime, Utc};
use fabro_api::types::{
BoardColumn, RunManifest, SubmitAnswerRequest, UpdateRunParentRequest, UpdateRunRequest,
};
@ -23,7 +23,7 @@ use fabro_store::{
};
use fabro_types::settings::ResolveError;
use fabro_types::{
AutomationRef, Principal, RunClientProvenance, RunId, RunProvenance, RunServerProvenance,
AutomationRef, Principal, Run, RunClientProvenance, RunId, RunProvenance, RunServerProvenance,
RunStatusKind, StageContextWindow, StageContextWindowStaleness,
StageContextWindowUnavailableReason, StageHandler, StageModelUsage, StageProjection,
SystemActorKind, WorkflowSettings, parse_blob_ref,
@ -298,7 +298,7 @@ async fn validate_parent_link(
}
async fn updated_run_response(state: &AppState, run_id: &RunId) -> Response {
match state.stores.run_summaries.get(run_id, Utc::now()).await {
match run_summary_at(state, run_id, Utc::now()).await {
Ok(Some(summary)) => (
StatusCode::OK,
Json(state.decorate_run_summary(summary).await),
@ -311,6 +311,28 @@ async fn updated_run_response(state: &AppState, run_id: &RunId) -> Response {
}
}
/// Read the durable summary and overlay its timing from the live projection.
///
/// The SQLite read model stores active timing as of the most recent event.
/// An open inference or tool bracket keeps accruing between events, so detail
/// reads need the projection's current estimate while the run is non-terminal.
async fn run_summary_at(
state: &AppState,
run_id: &RunId,
now: DateTime<Utc>,
) -> fabro_store::Result<Option<Run>> {
let Some(mut summary) = state.stores.run_summaries.get(run_id, now).await? else {
return Ok(None);
};
if summary.timestamps.completed_at.is_none() {
let cached = state.stores.runs.get_cached_run(run_id).await?;
if let Some(timing) = cached.and_then(|cached| cached.projection.live_run_timing(now)) {
summary.timing = Some(timing);
}
}
Ok(Some(summary))
}
async fn list_runs(
_auth: RequiredRunManagementActor,
State(state): State<Arc<AppState>>,
@ -934,7 +956,7 @@ async fn get_run_status(
RequireRunManagementTarget(id, _actor): RequireRunManagementTarget,
State(state): State<Arc<AppState>>,
) -> Response {
match state.stores.run_summaries.get(&id, Utc::now()).await {
match run_summary_at(&state, &id, Utc::now()).await {
Ok(Some(run)) => {
(StatusCode::OK, Json(state.decorate_run_summary(run).await)).into_response()
}

View file

@ -5503,6 +5503,50 @@ async fn list_run_stages_exposes_execution_identity_for_resumed_stage() {
assert_eq!(second["resumed_from_stage_id"], "work@1");
}
#[tokio::test]
async fn run_billing_includes_live_stage_timing_in_rows_and_totals() {
let state = test_app_state_with_isolated_storage();
let app = crate::test_support::build_test_router(Arc::clone(&state));
let run_id = RunId::new();
create_durable_run_with_events(&state, run_id, &[
workflow_event::Event::RunSubmitted {
definition_blob: None,
},
workflow_event::Event::RunStarting,
workflow_event::Event::RunRunning,
workflow_run_started_event(run_id),
])
.await;
append_scoped_stage_event(
&state,
run_id,
"work",
1,
&stage_started_event("work", "command"),
)
.await;
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
let response = app
.oneshot(
Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/billing")))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
let body = response_json!(response, StatusCode::OK).await;
let stages = body["stages"].as_array().unwrap();
assert_eq!(stages.len(), 1);
let row_timing = &stages[0]["timing"];
assert!(row_timing["active_time_ms"].as_u64().unwrap() > 0);
assert_eq!(row_timing["tool_time_ms"], row_timing["active_time_ms"]);
assert_eq!(&body["totals"]["timing"], row_timing);
}
/// `checkpoint.completed_nodes` records every visit, so a looped node appears
/// once per re-entry. Billing must dedup so a retried node renders as one row
/// and `runtime_secs` is summed across all visits exactly once.
@ -7800,6 +7844,51 @@ async fn get_run_status_returns_status() {
assert!(body["labels"].is_object());
}
#[tokio::test]
async fn get_run_status_advances_live_active_timing_between_events() {
let state = test_app_state_with_isolated_storage();
let app = crate::test_support::build_test_router(Arc::clone(&state));
let run_id = RunId::new();
create_durable_run_with_events(&state, run_id, &[
workflow_event::Event::RunSubmitted {
definition_blob: None,
},
workflow_event::Event::RunStarting,
workflow_event::Event::RunRunning,
workflow_run_started_event(run_id),
])
.await;
append_scoped_stage_event(
&state,
run_id,
"work",
1,
&stage_started_event("work", "command"),
)
.await;
// The SQLite summary stores timing at the StageStarted event. A later
// detail read must overlay the in-flight command's active time from the
// projection even though no newer event has arrived.
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
let response = app
.oneshot(
Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}")))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
let body = response_json!(response, StatusCode::OK).await;
let timing = &body["timing"];
assert!(timing["active_time_ms"].as_u64().unwrap() > 0);
assert_eq!(timing["tool_time_ms"], timing["active_time_ms"]);
assert!(timing["wall_time_ms"].as_u64().unwrap() >= timing["active_time_ms"].as_u64().unwrap());
}
#[tokio::test]
async fn get_run_status_not_found() {
let app = test_app_with();

View file

@ -19,6 +19,7 @@ use fabro_types::{
SandboxProviderKind, StageCompletion, StageHandler, StageId, StageInferenceProjection,
StageModelUsage, StageOutcome, StageProjection, StageState, StartRecord, SubAgentProjection,
SubAgentStatus, TodoListKind, TodoListProjection, TodoProjection, WorkflowRef, first_event_seq,
timing,
};
use fabro_util::error::render_compact_with_causes;
@ -405,7 +406,7 @@ impl RunProjectionReducer for RunProjection {
};
stage.response = response;
stage.completion = Some(completion);
stage.timing = Some(props.timing);
stage.set_authoritative_timing(props.timing);
if let Some(billing) = &props.billing {
stage.usage.replace_with_billed_usage(billing);
stage.model = Some(billing.model().clone());
@ -428,7 +429,7 @@ impl RunProjectionReducer for RunProjection {
failure_reason,
timestamp: ts,
});
stage.timing = Some(props.timing);
stage.set_authoritative_timing(props.timing);
if let Some(billing) = &props.billing {
stage.usage.replace_with_billed_usage(billing);
stage.model = Some(billing.model().clone());
@ -481,7 +482,7 @@ impl RunProjectionReducer for RunProjection {
close_inference_bracket(self, stored, props.visit, event.seq, ts);
}
EventBody::AgentSessionEnded(_) => {
close_inference_brackets_for_session(self, stored, ts);
close_active_brackets_for_session(self, stored, ts);
}
EventBody::AgentSessionActivated(props) => {
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
@ -746,7 +747,11 @@ impl RunProjectionReducer for RunProjection {
});
}
EventBody::AgentToolStarted(props) => {
let is_root_session = stored.parent_session_id.is_none();
let root_session_id = if stored.parent_session_id.is_none() {
stored.session_id.clone()
} else {
None
};
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
else {
return Ok(());
@ -770,19 +775,22 @@ impl RunProjectionReducer for RunProjection {
// A subagent's tools run inside the root session's tool call,
// so the root batch already covers them. Timing them again
// would double-count that span.
if is_root_session {
stage.open_tool_call(props.tool_call_id.clone(), ts);
if let Some(session_id) = root_session_id {
stage.open_tool_call(session_id, props.tool_call_id.clone(), ts);
}
}
EventBody::AgentToolCompleted(props) => {
if stored.parent_session_id.is_some() {
return Ok(());
}
let Some(session_id) = stored.session_id.as_deref() else {
return Ok(());
};
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
else {
return Ok(());
};
stage.close_tool_call(&props.tool_call_id, ts);
stage.close_tool_call(session_id, &props.tool_call_id, ts);
}
_ => {}
}
@ -1123,16 +1131,10 @@ fn close_bracket_on_stage(stage: &mut StageProjection, ts: DateTime<Utc>) {
let Some(inference) = stage.inference.take() else {
return;
};
stage.accumulate_inference_ms(elapsed_ms(inference.started_at, ts));
stage.accumulate_inference_ms(timing::elapsed_ms(inference.started_at, ts));
}
/// Non-negative milliseconds between two instants, saturating at zero so a
/// clock skew or an out-of-order replay cannot produce a negative span.
fn elapsed_ms(from: DateTime<Utc>, to: DateTime<Utc>) -> u64 {
u64::try_from(to.signed_duration_since(from).num_milliseconds().max(0)).unwrap_or(0)
}
/// Close every bracket opened by the session that just ended.
/// Close every active bracket opened by the session that just ended.
///
/// `agent.session.ended` is the only ordering-safe backstop for terminal
/// cancel and wall-clock timeout, which tear the session down through
@ -1147,7 +1149,7 @@ fn elapsed_ms(from: DateTime<Utc>, to: DateTime<Utc>) -> u64 {
/// opened. Implemented as a normal stage lookup it would find no target and
/// silently no-op, leaving the bracket open forever on exactly the path it
/// exists to cover.
fn close_inference_brackets_for_session(
fn close_active_brackets_for_session(
state: &mut RunProjection,
stored: &RunEvent,
ts: DateTime<Utc>,
@ -1166,6 +1168,7 @@ fn close_inference_brackets_for_session(
if opened_here {
close_bracket_on_stage(stage, ts);
}
stage.close_tool_batch_for_session(session_id, ts);
}
}
@ -1475,14 +1478,15 @@ fn finalize_unfinished_stages_after_run_failed(
StageState::Failed
};
for (_, stage) in state.iter_stages_mut() {
for (_, stage) in state.iter_stages_unordered_mut() {
if stage.state.is_terminal() {
continue;
}
// Close any bracket still open so its span is not dropped on the
// Close any brackets still open so their spans are not dropped on the
// floor when the live estimate is frozen into `timing` below.
close_bracket_on_stage(stage, timestamp);
stage.close_open_tool_batch(timestamp);
// Freeze the live estimate before flipping to a terminal state:
// `live_timing` reads `effective_state` and would return wall-only
@ -1490,7 +1494,9 @@ fn finalize_unfinished_stages_after_run_failed(
let frozen = stage.live_timing(timestamp);
stage.state = terminal_state;
if stage.timing.is_none() && stage.started_at.is_some() {
stage.timing = Some(frozen);
stage.set_authoritative_timing(frozen);
} else {
stage.clear_live_timing();
}
}
}
@ -1641,12 +1647,16 @@ mod tests {
StageId::new("plan", 1)
}
fn agent_event(seq: u32, ts: &str, body: EventBody) -> EventEnvelope {
fn session_event(seq: u32, ts: &str, session_id: &str, body: EventBody) -> EventEnvelope {
let mut event = test_stage_event_at(seq, ts, body, stage_id());
event.event.session_id = Some("session-1".to_string());
event.event.session_id = Some(session_id.to_string());
event
}
fn agent_event(seq: u32, ts: &str, body: EventBody) -> EventEnvelope {
session_event(seq, ts, "session-1", body)
}
/// An event from a sub-agent session nested under the root session.
fn child_event(seq: u32, ts: &str, body: EventBody) -> EventEnvelope {
let mut event = agent_event(seq, ts, body);
@ -1855,6 +1865,81 @@ mod tests {
assert!(stage(&state).tool_batch.is_some());
}
#[test]
fn a_foreign_session_completion_does_not_mutate_the_open_batch() {
let mut state = started_state();
state
.apply_event(&agent_event(
2,
"2026-04-07T12:00:00Z",
tool_started("call-a"),
))
.unwrap();
state
.apply_event(&session_event(
3,
"2026-04-07T12:00:05Z",
"session-2",
tool_completed("call-a"),
))
.unwrap();
let batch = stage(&state).tool_batch.as_ref().unwrap();
assert_eq!(batch.session_id, "session-1");
assert!(batch.open_call_ids.contains("call-a"));
assert_eq!(stage(&state).live_tool_ms, 0);
}
#[test]
fn a_replacement_session_starts_a_separate_tool_batch() {
let mut state = started_state();
state
.apply_event(&agent_event(
2,
"2026-04-07T12:00:00Z",
tool_started("old-call"),
))
.unwrap();
state
.apply_event(&session_event(
3,
"2026-04-07T12:00:05Z",
"session-2",
tool_started("new-call"),
))
.unwrap();
assert_eq!(stage(&state).live_tool_ms, 5_000);
let batch = stage(&state).tool_batch.as_ref().unwrap();
assert_eq!(batch.session_id, "session-2");
assert_eq!(
batch.open_call_ids,
["new-call".to_string()].into_iter().collect()
);
// A delayed completion from the old session cannot close the new
// session's batch even when call ids happen to collide.
state
.apply_event(&agent_event(
4,
"2026-04-07T12:00:07Z",
tool_completed("new-call"),
))
.unwrap();
assert!(stage(&state).tool_batch.is_some());
state
.apply_event(&session_event(
5,
"2026-04-07T12:00:09Z",
"session-2",
tool_completed("new-call"),
))
.unwrap();
assert_eq!(stage(&state).live_tool_ms, 9_000);
assert!(stage(&state).tool_batch.is_none());
}
#[test]
fn subagent_tool_calls_do_not_double_count_against_the_root_batch() {
let mut state = started_state();
@ -1892,14 +1977,21 @@ mod tests {
}
#[test]
fn session_end_accumulates_rather_than_discarding_the_bracket() {
fn session_end_accumulates_every_open_active_bracket() {
let mut state = started_state();
state
.apply_event(&agent_event(2, "2026-04-07T12:00:05Z", llm_started()))
.unwrap();
state
.apply_event(&agent_event(
3,
"2026-04-07T12:00:07Z",
tool_started("call-a"),
))
.unwrap();
let mut ended = test_stage_event_at(
3,
4,
"2026-04-07T12:00:20Z",
EventBody::AgentSessionEnded(AgentSessionEndedProps {}),
stage_id(),
@ -1908,7 +2000,9 @@ mod tests {
state.apply_event(&ended).unwrap();
assert_eq!(stage(&state).live_inference_ms, 15_000);
assert_eq!(stage(&state).live_tool_ms, 13_000);
assert!(stage(&state).inference.is_none());
assert!(stage(&state).tool_batch.is_none());
}
#[test]
@ -1986,6 +2080,48 @@ mod tests {
stage.timing.unwrap(),
"a terminal stage reports its finalized breakdown, not a live estimate"
);
assert_eq!(stage.live_inference_ms, 0);
assert_eq!(stage.live_tool_ms, 0);
assert!(stage.inference.is_none());
assert!(stage.tool_batch.is_none());
}
#[test]
fn run_failure_freezes_open_work_and_clears_live_bookkeeping() {
let mut state = started_state();
state.status = RunStatus::Running;
state
.apply_event(&agent_event(2, "2026-04-07T12:00:01Z", llm_started()))
.unwrap();
state
.apply_event(&agent_event(3, "2026-04-07T12:00:04Z", agent_message()))
.unwrap();
state
.apply_event(&agent_event(
4,
"2026-04-07T12:00:05Z",
tool_started("call-a"),
))
.unwrap();
let mut failed = test_event(
5,
EventBody::RunFailed(run_failed_props(FailureReason::WorkflowError)),
None,
);
failed.event.ts = test_dt("2026-04-07T12:00:10Z");
state.apply_event(&failed).unwrap();
let stage = stage(&state);
assert_eq!(
stage.timing,
Some(fabro_types::StageTiming::new(10_000, 3_000, 5_000))
);
assert_eq!(stage.state, StageState::Failed);
assert_eq!(stage.live_inference_ms, 0);
assert_eq!(stage.live_tool_ms, 0);
assert!(stage.inference.is_none());
assert!(stage.tool_batch.is_none());
}
}

View file

@ -361,6 +361,11 @@ fn main() {
"fabro_types::StageInferenceProjection",
&[],
),
(
"StageToolBatchProjection",
"fabro_types::StageToolBatchProjection",
&[],
),
("LlmOutputKind", "fabro_types::LlmOutputKind", &[]),
("PermissionLevel", "fabro_types::PermissionLevel", &[]),
(

View file

@ -69,10 +69,10 @@ pub mod types {
StageContextWindowCategory, StageContextWindowCountMethod, StageContextWindowProjection,
StageContextWindowStaleness, StageContextWindowUnavailableReason,
StageContextWindowWarning, StageHandler, StageId, StageInferenceProjection,
StageModelUsage, StageOutcome, StageProjection, StageState, SubAgentProjection,
SubAgentStatus, SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse,
TodoListProjection, TurnId, UpdateVariableRequest, UserPrincipal, Variable,
VariableListResponse, WorkflowSettings,
StageModelUsage, StageOutcome, StageProjection, StageState, StageToolBatchProjection,
SubAgentProjection, SubAgentStatus, SystemActorKind, SystemIntegrationStatus,
SystemIntegrationsResponse, TodoListProjection, TurnId, UpdateVariableRequest,
UserPrincipal, Variable, VariableListResponse, WorkflowSettings,
};
pub use crate::generated::types::*;

View file

@ -18,6 +18,7 @@ use fabro_api::types::{
StageContextWindowUnavailableReason as ApiStageContextWindowUnavailableReason,
StageContextWindowWarning as ApiStageContextWindowWarning,
StageInferenceProjection as ApiStageInferenceProjection, StageProjection as ApiStageProjection,
StageToolBatchProjection as ApiStageToolBatchProjection,
SubAgentProjection as ApiSubAgentProjection, SubAgentStatus as ApiSubAgentStatus,
TodoListProjection as ApiTodoListProjection,
};
@ -29,8 +30,8 @@ use fabro_types::{
ParallelBranchResult, PermissionLevel, SkillsProjection, StageContextWindow,
StageContextWindowBreakdownItem, StageContextWindowCategory, StageContextWindowCountMethod,
StageContextWindowProjection, StageContextWindowStaleness, StageContextWindowUnavailableReason,
StageContextWindowWarning, StageInferenceProjection, StageProjection, SubAgentProjection,
SubAgentStatus, TodoListKind, TodoListProjection,
StageContextWindowWarning, StageInferenceProjection, StageProjection, StageToolBatchProjection,
SubAgentProjection, SubAgentStatus, TodoListKind, TodoListProjection,
};
use serde_json::json;
@ -68,9 +69,32 @@ fn stage_projection_reuses_nested_agent_state_types() {
);
assert_same_type::<ApiStageContextWindowWarning, StageContextWindowWarning>();
assert_same_type::<ApiStageInferenceProjection, StageInferenceProjection>();
assert_same_type::<ApiStageToolBatchProjection, StageToolBatchProjection>();
assert_same_type::<ApiLlmOutputKind, LlmOutputKind>();
}
#[test]
fn stage_tool_batch_projection_matches_openapi_json_shape() {
let batch = StageToolBatchProjection {
session_id: "ses_root".to_string(),
started_at: "2026-04-29T12:34:00Z".parse().unwrap(),
open_call_ids: ["call_1".to_string(), "call_2".to_string()]
.into_iter()
.collect(),
};
let value = serde_json::to_value(&batch).unwrap();
assert_eq!(
value,
json!({
"session_id": "ses_root",
"started_at": "2026-04-29T12:34:00Z",
"open_call_ids": ["call_1", "call_2"]
})
);
let api_batch: ApiStageToolBatchProjection = serde_json::from_value(value).unwrap();
assert_eq!(api_batch, batch);
}
#[test]
fn stage_inference_projection_matches_openapi_json_shape() {
let inference = StageInferenceProjection {
@ -148,6 +172,9 @@ fn stage_projection_without_inference_round_trips() {
let stage: StageProjection = serde_json::from_value(value.clone()).unwrap();
assert!(stage.inference.is_none());
assert!(stage.tool_batch.is_none());
assert_eq!(stage.live_inference_ms, 0);
assert_eq!(stage.live_tool_ms, 0);
assert_eq!(serde_json::to_value(stage).unwrap(), value);
}

View file

@ -125,7 +125,7 @@ pub use run_projection::{
StageContextWindowBreakdownItem, StageContextWindowCategory, StageContextWindowCountMethod,
StageContextWindowProjection, StageContextWindowStaleness, StageContextWindowUnavailableReason,
StageContextWindowWarning, StageInferenceProjection, StageModelUsage, StageProjection,
SubAgentProjection, SubAgentStatus, first_event_seq,
StageToolBatchProjection, SubAgentProjection, SubAgentStatus, first_event_seq,
};
pub use run_sandbox::{
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,

View file

@ -12,7 +12,7 @@ use crate::{
AgentToolSummary, BilledTokenCounts, Checkpoint, Conclusion, InterviewQuestionRecord,
InvalidTransition, LlmOutputKind, ModelRef, PermissionLevel, PullRequestLink, RunApproval,
RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus, RunTiming, StageCompletion,
StageHandler, StageId, StageState, StageTiming, StartRecord, TodoListProjection,
StageHandler, StageId, StageState, StageTiming, StartRecord, TodoListProjection, timing,
};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
@ -433,6 +433,10 @@ fn is_zero_ms(value: &u64) -> bool {
/// duplicated log must not let a repeated completion drain the batch early.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct StageToolBatchProjection {
/// Root agent session that dispatched the batch. Later transitions are
/// gated on it so delayed events from a replaced session cannot mutate
/// the current session's batch.
pub session_id: String,
/// When the batch opened — the first `agent.tool.started` observed while
/// no other calls were outstanding.
pub started_at: DateTime<Utc>,
@ -603,10 +607,9 @@ impl StageProjection {
state,
StageState::Running | StageState::Retrying | StageState::Pending
) {
return self.started_at.map(|started| {
u64::try_from(now.signed_duration_since(started).num_milliseconds().max(0))
.unwrap_or(0)
});
return self
.started_at
.map(|started| timing::elapsed_ms(started, now));
}
self.timing.map(|timing| timing.wall_time_ms)
}
@ -645,9 +648,6 @@ impl StageProjection {
}
let wall_time_ms = self.live_wall_time_ms(now).unwrap_or(0);
let elapsed = |since: DateTime<Utc>| {
u64::try_from(now.signed_duration_since(since).num_milliseconds().max(0)).unwrap_or(0)
};
// `handler` is absent on projections built from events written before
// stage execution identity. Treat those as agent stages, matching
@ -660,11 +660,11 @@ impl StageProjection {
let open_inference = self
.inference
.as_ref()
.map_or(0, |inference| elapsed(inference.started_at));
.map_or(0, |inference| timing::elapsed_ms(inference.started_at, now));
let open_tool = self
.tool_batch
.as_ref()
.map_or(0, |batch| elapsed(batch.started_at));
.map_or(0, |batch| timing::elapsed_ms(batch.started_at, now));
(
self.live_inference_ms.saturating_add(open_inference),
self.live_tool_ms.saturating_add(open_tool),
@ -693,9 +693,26 @@ impl StageProjection {
}
/// Record a dispatched tool call, opening a batch if none is outstanding.
pub fn open_tool_call(&mut self, tool_call_id: String, started_at: DateTime<Utc>) {
///
/// If a replacement root session starts work before the old session's end
/// event arrives, freeze the old batch at this boundary before opening
/// the new one. This keeps the sessions separate without dropping time.
pub fn open_tool_call(
&mut self,
session_id: String,
tool_call_id: String,
started_at: DateTime<Utc>,
) {
let replaces_open_batch = self
.tool_batch
.as_ref()
.is_some_and(|batch| batch.session_id != session_id);
if replaces_open_batch {
self.close_open_tool_batch(started_at);
}
self.tool_batch
.get_or_insert_with(|| StageToolBatchProjection {
session_id,
started_at,
open_call_ids: BTreeSet::new(),
})
@ -705,21 +722,52 @@ impl StageProjection {
/// Retire a tool call. Folds the batch into the live accumulator once the
/// last outstanding call reports, so concurrent calls count once.
pub fn close_tool_call(&mut self, tool_call_id: &str, now: DateTime<Utc>) {
pub fn close_tool_call(&mut self, session_id: &str, tool_call_id: &str, now: DateTime<Utc>) {
let Some(batch) = self.tool_batch.as_mut() else {
return;
};
batch.open_call_ids.remove(tool_call_id);
if !batch.open_call_ids.is_empty() {
if batch.session_id != session_id
|| !batch.open_call_ids.remove(tool_call_id)
|| !batch.open_call_ids.is_empty()
{
return;
}
let elapsed = u64::try_from(
now.signed_duration_since(batch.started_at)
.num_milliseconds()
.max(0),
)
.unwrap_or(0);
self.live_tool_ms = self.live_tool_ms.saturating_add(elapsed);
self.close_open_tool_batch(now);
}
/// Close a tool batch only when it belongs to `session_id`.
pub fn close_tool_batch_for_session(&mut self, session_id: &str, now: DateTime<Utc>) {
let opened_here = self
.tool_batch
.as_ref()
.is_some_and(|batch| batch.session_id == session_id);
if opened_here {
self.close_open_tool_batch(now);
}
}
/// Fold any open tool batch into the live accumulator.
pub fn close_open_tool_batch(&mut self, now: DateTime<Utc>) {
let Some(batch) = self.tool_batch.take() else {
return;
};
self.live_tool_ms = self
.live_tool_ms
.saturating_add(timing::elapsed_ms(batch.started_at, now));
}
/// Install a worker-provided terminal timing and discard transient live
/// bookkeeping that is no longer authoritative.
pub fn set_authoritative_timing(&mut self, timing: StageTiming) {
self.timing = Some(timing);
self.clear_live_timing();
}
/// Discard transient timing accumulators and open brackets.
pub fn clear_live_timing(&mut self) {
self.live_inference_ms = 0;
self.live_tool_ms = 0;
self.inference = None;
self.tool_batch = None;
}
@ -1138,20 +1186,15 @@ mod iter_stages_tests {
#[cfg(test)]
mod live_timing_tests {
use std::collections::HashMap;
use std::num::NonZeroU32;
use chrono::{DateTime, TimeZone, Utc};
use super::{RunProjection, StageToolBatchProjection};
use crate::{
Graph, ModelRef, RunId, RunSpec, StageHandler, StageInferenceProjection, StageProjection,
StageState, StageTiming, StartRecord, WorkflowSettings, test_support,
StageState, StageTiming, StartRecord, WorkflowSettings, first_event_seq, test_support,
};
fn seq(n: u32) -> NonZeroU32 {
NonZeroU32::new(n).unwrap()
}
fn at(seconds: i64) -> DateTime<Utc> {
Utc.timestamp_opt(1_700_000_000 + seconds, 0).unwrap()
}
@ -1180,7 +1223,7 @@ mod live_timing_tests {
/// In-flight stage that started at `at(0)`.
fn running(handler: StageHandler) -> StageProjection {
let mut stage = StageProjection::new(seq(1));
let mut stage = StageProjection::new(first_event_seq(1));
stage.handler = Some(handler);
stage.started_at = Some(at(0));
stage.state = StageState::Running;
@ -1225,6 +1268,7 @@ mod live_timing_tests {
// Inference open for 20s, tools open for 10s, at t=120s.
stage.inference = Some(open_bracket(at(100)));
stage.tool_batch = Some(StageToolBatchProjection {
session_id: "session-1".to_string(),
started_at: at(110),
open_call_ids: ["call-1".to_string()].into_iter().collect(),
});
@ -1330,17 +1374,17 @@ mod live_timing_tests {
base_sha: None,
});
let baseline = projection.stage_entry("baseline", 1, seq(1));
let baseline = projection.stage_entry("baseline", 1, first_event_seq(1));
baseline.handler = Some(StageHandler::Command);
baseline.state = StageState::Succeeded;
baseline.timing = Some(StageTiming::new(42_666, 0, 42_663));
let assess = projection.stage_entry("assess", 1, seq(2));
let assess = projection.stage_entry("assess", 1, first_event_seq(2));
assess.handler = Some(StageHandler::Agent);
assess.state = StageState::Succeeded;
assess.timing = Some(StageTiming::new(86_025, 78_230, 7_588));
let plan = projection.stage_entry("plan", 1, seq(3));
let plan = projection.stage_entry("plan", 1, first_event_seq(3));
plan.handler = Some(StageHandler::Agent);
plan.started_at = Some(at(146));
plan.state = StageState::Running;
@ -1367,7 +1411,8 @@ mod live_timing_tests {
});
for (index, node) in ["branch-a", "branch-b", "branch-c"].iter().enumerate() {
let stage = projection.stage_entry(node, 1, seq(u32::try_from(index).unwrap() + 1));
let stage =
projection.stage_entry(node, 1, first_event_seq(u32::try_from(index).unwrap() + 1));
stage.handler = Some(StageHandler::Agent);
stage.state = StageState::Succeeded;
stage.timing = Some(StageTiming::new(60_000, 60_000, 0));

View file

@ -18,6 +18,16 @@
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
/// Non-negative milliseconds between two instants.
///
/// Clock skew or an out-of-order replay can put `end` before `start`; those
/// spans contribute zero rather than wrapping into a large unsigned value.
#[must_use]
pub fn elapsed_ms(start: DateTime<Utc>, end: DateTime<Utc>) -> u64 {
u64::try_from(end.signed_duration_since(start).num_milliseconds().max(0))
.expect("non-negative chrono millisecond durations fit in u64")
}
/// Timing breakdown for one stage visit.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct StageTiming {
@ -74,16 +84,18 @@ impl StageTiming {
/// legitimately sum past run wall time.
#[must_use]
pub fn clamped_to_wall(&self) -> Self {
if self.active_time_ms <= self.wall_time_ms {
let active_time_ms = u128::from(self.inference_time_ms) + u128::from(self.tool_time_ms);
if active_time_ms <= u128::from(self.wall_time_ms) {
return *self;
}
// Preserve the split rather than truncating one side, so a clamped
// stage still shows where its time went. Widen for the multiply: the
// quotient is bounded by `wall_time_ms` because `active_time_ms`
// exceeds it here, so it always fits back into u64.
let scaled = u128::from(self.inference_time_ms) * u128::from(self.wall_time_ms)
/ u128::from(self.active_time_ms);
let inference_time_ms = u64::try_from(scaled).unwrap_or(self.wall_time_ms);
// quotient is bounded by `wall_time_ms` because the exact, widened
// active total exceeds it here, so it always fits back into u64.
let scaled =
u128::from(self.inference_time_ms) * u128::from(self.wall_time_ms) / active_time_ms;
let inference_time_ms =
u64::try_from(scaled).expect("scaled inference time is bounded by wall time");
let tool_time_ms = self.wall_time_ms.saturating_sub(inference_time_ms);
Self::new(self.wall_time_ms, inference_time_ms, tool_time_ms)
}
@ -166,8 +178,7 @@ impl RunTiming {
/// Milliseconds elapsed from `start` to `now`, clamped at zero.
#[must_use]
pub fn wall_time_ms_since(start: DateTime<Utc>, now: DateTime<Utc>) -> u64 {
u64::try_from(now.signed_duration_since(start).num_milliseconds().max(0))
.expect("non-negative milliseconds fit in u64")
elapsed_ms(start, now)
}
}
@ -184,7 +195,9 @@ impl From<StageTiming> for RunTiming {
#[cfg(test)]
mod tests {
use super::{RunTiming, StageTiming};
use chrono::{TimeZone, Utc};
use super::{RunTiming, StageTiming, elapsed_ms};
#[test]
fn stage_timing_new_derives_active_as_sum_of_inference_and_tool() {
@ -215,6 +228,23 @@ mod tests {
assert_eq!(sum.active_time_ms, 175);
}
#[test]
fn stage_timing_clamp_uses_the_unsaturated_active_total() {
let timing = StageTiming::new(u64::MAX, u64::MAX, u64::MAX).clamped_to_wall();
assert_eq!(timing.inference_time_ms, u64::MAX / 2);
assert_eq!(timing.tool_time_ms, u64::MAX.saturating_sub(u64::MAX / 2));
assert_eq!(timing.active_time_ms, u64::MAX);
}
#[test]
fn elapsed_ms_clamps_out_of_order_instants_to_zero() {
let later = Utc.timestamp_opt(100, 0).unwrap();
let earlier = Utc.timestamp_opt(99, 0).unwrap();
assert_eq!(elapsed_ms(later, earlier), 0);
}
#[test]
fn run_timing_wall_only_zeroes_breakdown_and_active() {
let timing = RunTiming::wall_only(1500);

View file

@ -18,6 +18,10 @@
* One open tool batch: tool calls dispatched together that have not all reported completion. `open_call_ids` is a set rather than a count so a duplicated completion in a replayed log cannot drain the batch early.
*/
export interface StageToolBatchProjection {
/**
* Root agent session that dispatched the batch. Transitions are gated on it so delayed events from a replaced session cannot mutate the current batch.
*/
'session_id': string;
/**
* When the batch opened — the first dispatched call observed while no other calls were outstanding.
*/
@ -25,5 +29,5 @@ export interface StageToolBatchProjection {
/**
* Calls dispatched but not yet completed, by tool call id.
*/
'open_call_ids': Array<string>;
'open_call_ids': Set<string>;
}