diff --git a/apps/fabro-web/app/lib/queries.ts b/apps/fabro-web/app/lib/queries.ts index 3c8eb8a1d..aa7328b5e 100644 --- a/apps/fabro-web/app/lib/queries.ts +++ b/apps/fabro-web/app/lib/queries.ts @@ -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( 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, ); } diff --git a/apps/fabro-web/app/lib/query-keys.test.ts b/apps/fabro-web/app/lib/query-keys.test.ts index ac441b998..52a5ccef3 100644 --- a/apps/fabro-web/app/lib/query-keys.test.ts +++ b/apps/fabro-web/app/lib/query-keys.test.ts @@ -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"), ]); }); }); diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index e984be74d..4fce18b4d 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -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"), diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index d0d4f4b4d..523774e18 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -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)] : []; } diff --git a/apps/fabro-web/app/routes/run-detail.tsx b/apps/fabro-web/app/routes/run-detail.tsx index 97ce00702..ad2d78986 100644 --- a/apps/fabro-web/app/routes/run-detail.tsx +++ b/apps/fabro-web/app/routes/run-detail.tsx @@ -69,6 +69,8 @@ import { export const handle = { hideHeader: true }; +const RUN_TIMING_REFRESH_INTERVAL_MS = 30_000; + type LifecycleTrigger = () => Promise; 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(null); - const now = useTickingNow(30_000); + const now = useTickingNow(RUN_TIMING_REFRESH_INTERVAL_MS); const { fullHeight, hideSteerBar } = childRouteLayoutFlags(matches); useRunEvents(params.id); diff --git a/docs/plans/2026-05-21-wall-and-active-time-metrics-plan.md b/docs/plans/2026-05-21-wall-and-active-time-metrics-plan.md index 3683e5879..c0980ed01 100644 --- a/docs/plans/2026-05-21-wall-and-active-time-metrics-plan.md +++ b/docs/plans/2026-05-21-wall-and-active-time-metrics-plan.md @@ -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. diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 568aa30e5..f2460a69b 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -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. diff --git a/lib/apps/fabro-server/src/server/handler/runs.rs b/lib/apps/fabro-server/src/server/handler/runs.rs index 0e2c8cbcb..4b9acdf1c 100644 --- a/lib/apps/fabro-server/src/server/handler/runs.rs +++ b/lib/apps/fabro-server/src/server/handler/runs.rs @@ -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, +) -> fabro_store::Result> { + 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>, @@ -934,7 +956,7 @@ async fn get_run_status( RequireRunManagementTarget(id, _actor): RequireRunManagementTarget, State(state): State>, ) -> 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() } diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 71b4dfebf..790089ab1 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -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(); diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index df05f84b3..14e86124a 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -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) { 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, to: DateTime) -> 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, to: DateTime) -> 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, @@ -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()); } } diff --git a/lib/foundation/fabro-api/build.rs b/lib/foundation/fabro-api/build.rs index fc943e99d..740585545 100644 --- a/lib/foundation/fabro-api/build.rs +++ b/lib/foundation/fabro-api/build.rs @@ -361,6 +361,11 @@ fn main() { "fabro_types::StageInferenceProjection", &[], ), + ( + "StageToolBatchProjection", + "fabro_types::StageToolBatchProjection", + &[], + ), ("LlmOutputKind", "fabro_types::LlmOutputKind", &[]), ("PermissionLevel", "fabro_types::PermissionLevel", &[]), ( diff --git a/lib/foundation/fabro-api/src/lib.rs b/lib/foundation/fabro-api/src/lib.rs index 4312acb91..b9a6a7185 100644 --- a/lib/foundation/fabro-api/src/lib.rs +++ b/lib/foundation/fabro-api/src/lib.rs @@ -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::*; diff --git a/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs b/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs index 5faed1bd3..22a544e69 100644 --- a/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs +++ b/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs @@ -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::(); assert_same_type::(); + assert_same_type::(); assert_same_type::(); } +#[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); } diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index 68f5644a2..b7f3b0474 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -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, diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 0576d32f6..e2b9ad3f6 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -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, @@ -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| { - 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) { + /// + /// 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, + ) { + 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) { + pub fn close_tool_call(&mut self, session_id: &str, tool_call_id: &str, now: DateTime) { 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) { + 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) { + 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.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)); diff --git a/lib/foundation/fabro-types/src/timing.rs b/lib/foundation/fabro-types/src/timing.rs index 4fe45a36f..812a72ef2 100644 --- a/lib/foundation/fabro-types/src/timing.rs +++ b/lib/foundation/fabro-types/src/timing.rs @@ -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, end: DateTime) -> 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, now: DateTime) -> 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 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); diff --git a/lib/packages/fabro-api-client/src/models/stage-tool-batch-projection.ts b/lib/packages/fabro-api-client/src/models/stage-tool-batch-projection.ts index a6bd7718c..c382b5dd1 100644 --- a/lib/packages/fabro-api-client/src/models/stage-tool-batch-projection.ts +++ b/lib/packages/fabro-api-client/src/models/stage-tool-batch-projection.ts @@ -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; + 'open_call_ids': Set; }