From c4971b93d32a950bb971fa3c5902e083c8905b2e Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 25 Jul 2026 14:54:38 -0400 Subject: [PATCH] fix(timing): accumulate active time for in-flight stages MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `active_time_ms` was only ever computed from terminal stage events, so a stage still running contributed zero to the run rollup. A run parked in one long agent stage reported 2m 8s of active time against 16m 53s of wall clock — the two finished stages — while the running stage had been doing continuous inference and tool work for over 14 minutes. `live_run_timing` summed `filter_map(|stage| stage.timing)`, and `stage.timing` is only written at finalization. Wall time ticked live off `start_time`; active time did not tick at all. Stage projections now accumulate brackets from the event log: - Closing an inference bracket folds its span into `live_inference_ms` instead of discarding it, including across retries, matching the in-process stopwatch. - Tool calls open a batch on the first outstanding call and close it when the last one drains, so tools running concurrently within a turn count once — the same span `execute_tool_calls` is bracketed by. Summing per-call durations would over-count parallel tool use. Subagent tool events are excluded; they run inside the root call's span already. - `StageProjection::live_timing(now)` composes accumulators with any open bracket, per handler: agent stages use the brackets, prompt and command stages count elapsed time as inference and tool respectively, and handlers that wait on a human, timer, condition, or child branches report zero. Active is clamped to wall per stage. A worker killed mid-turn leaves its bracket open forever, and without the clamp it would tick up unbounded. The clamp does not need to detect the dead worker: a stage cannot have been active longer than it has existed. `watchdog.timeout` remains the authority on whether a run is stuck. The clamp is deliberately not applied at run level, where concurrent branches can legitimately sum past run wall time. Timing is derived from events rather than emitted by the worker, so this needs no event-schema change and applies to runs already stored. `StageProjection.timing` keeps its terminal-only meaning, and the authoritative breakdown still replaces the live estimate at terminal events. The billing endpoint had the same hole behind its `wall_only` fallback: running stages reported zero inference/tool/active. Not visible in the product, which renders only `wall_time_ms`, but wrong for any other consumer of `GET /runs/{id}/billing`. Parallel branch stages lose their breakdown permanently, even after completion, because `parallel.branch.completed` carries only `duration_ms`. That is a separate data-loss bug, tracked in #644. Co-Authored-By: Claude Opus 5 (1M context) --- .../app/routes/run-detail/header.tsx | 9 +- ...05-21-wall-and-active-time-metrics-plan.md | 11 +- docs/public/api-reference/fabro-api.yaml | 66 ++- .../src/server/handler/billing.rs | 16 +- lib/components/fabro-store/src/run_state.rs | 475 +++++++++++++++++- .../fabro-types/src/run_projection.rs | 462 ++++++++++++++++- lib/foundation/fabro-types/src/timing.rs | 26 + .../src/.openapi-generator/FILES | 1 + .../fabro-api-client/src/models/index.ts | 1 + .../fabro-api-client/src/models/run-timing.ts | 2 +- .../src/models/stage-projection.ts | 12 + .../src/models/stage-timing.ts | 2 +- .../src/models/stage-tool-batch-projection.ts | 29 ++ 13 files changed, 1060 insertions(+), 52 deletions(-) create mode 100644 lib/packages/fabro-api-client/src/models/stage-tool-batch-projection.ts diff --git a/apps/fabro-web/app/routes/run-detail/header.tsx b/apps/fabro-web/app/routes/run-detail/header.tsx index efc68f09b..3b7081596 100644 --- a/apps/fabro-web/app/routes/run-detail/header.tsx +++ b/apps/fabro-web/app/routes/run-detail/header.tsx @@ -326,6 +326,7 @@ function DurationPopover({ }) { const endMs = completedAt != null ? Date.parse(completedAt) : now; const sinceCreatedMs = Math.max(0, endMs - Date.parse(createdAt)); + const isRunning = completedAt == null; return ( <> Duration @@ -335,8 +336,14 @@ function DurationPopover({
{formatDurationMs(sinceCreatedMs)}
-
Active (inference + tools)
+
+ Active (inference + tools){isRunning ? " — estimated" : ""} +
{formatDurationMs(timing.active_time_ms)}
+
+ {formatDurationMs(timing.inference_time_ms)} inference ·{" "} + {formatDurationMs(timing.tool_time_ms)} tools +
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 90438d3b5..3683e5879 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 @@ -157,7 +157,14 @@ git diff --check 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. -- Active timing is finalized-event based in v1; live active-time ticking can be - added later if it becomes necessary. +- ~~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 + as active, because in-flight stages contributed nothing. Stage projections now + 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`. - 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 b028d89c0..568aa30e5 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -10673,7 +10673,38 @@ components: - type: "null" description: | Per-attempt timing breakdown for the latest terminal attempt: - wall time plus the active inference/tool breakdown. + wall time plus the active inference/tool breakdown. Null while the + stage is still in flight; the live estimate is derived from + `live_inference_ms`, `live_tool_ms`, and any open bracket. + live_inference_ms: + type: integer + format: uint64 + minimum: 0 + default: 0 + description: | + Inference time accumulated from closed brackets during the current + attempt. Live estimate only — the authoritative value arrives with + the terminal event and lands in `timing`. Excludes the currently + open bracket, whose span is measured from `inference.started_at`. + example: 78230 + live_tool_ms: + type: integer + format: uint64 + minimum: 0 + default: 0 + description: | + Tool time accumulated from closed tool batches during the current + attempt. A batch spans the first dispatched call through the + completion that drains the last outstanding one, so tools running + concurrently within a turn are counted once. + example: 7588 + tool_batch: + oneOf: + - $ref: "#/components/schemas/StageToolBatchProjection" + - type: "null" + description: | + Open tool batch: when the batch started and which calls have not + yet reported completion. usage: $ref: "#/components/schemas/BilledTokenCounts" model: @@ -10734,6 +10765,28 @@ components: $ref: "#/components/schemas/StageState" description: Lifecycle state of the stage projection. + StageToolBatchProjection: + description: > + 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. + type: object + required: + - started_at + - open_call_ids + properties: + started_at: + type: string + format: date-time + description: > + When the batch opened — the first dispatched call observed while no + other calls were outstanding. + open_call_ids: + type: array + items: + type: string + description: Calls dispatched but not yet completed, by tool call id. + StageInferenceProjection: description: > One open inference bracket: a dispatched LLM request that has not yet @@ -12244,6 +12297,12 @@ components: observed LLM request/stream elapsed time; `tool_time_ms` is tool or command execution elapsed time; `active_time_ms` equals `inference_time_ms + tool_time_ms`. + + For a terminal stage these come from the worker's own stopwatch and are + authoritative. For a stage still in flight they are a live estimate + reconstructed from the event log, and `active_time_ms` is clamped to + `wall_time_ms`. The estimate is replaced by the authoritative + breakdown when the stage reaches a terminal event. type: object required: - wall_time_ms @@ -12278,6 +12337,11 @@ components: Timing rollup for an entire run. Active fields sum work across stage visits, so `active_time_ms` can exceed `wall_time_ms` when parallel branches run concurrently. + + For a running run, stages still in flight contribute a live estimate + rather than nothing, so wall and active both advance continuously. + Unlike `StageTiming`, active is not clamped to wall here — concurrent + branches can legitimately sum past run wall time. type: object required: - wall_time_ms diff --git a/lib/apps/fabro-server/src/server/handler/billing.rs b/lib/apps/fabro-server/src/server/handler/billing.rs index a1888713c..6cc4cbdde 100644 --- a/lib/apps/fabro-server/src/server/handler/billing.rs +++ b/lib/apps/fabro-server/src/server/handler/billing.rs @@ -175,7 +175,7 @@ fn live_billing_rows(projection: &RunProjection, now: DateTime) -> Vec= row.latest_visit { @@ -188,20 +188,6 @@ fn live_billing_rows(projection: &RunProjection, now: DateTime) -> Vec) -> StageTiming { - if let Some(timing) = stage.timing { - return timing; - } - if let Some(live_wall) = stage.live_wall_time_ms(now) { - return StageTiming::wall_only(live_wall); - } - StageTiming::default() -} - fn stage_has_billing_row(stage: &StageProjection) -> bool { stage.completion.is_some() || stage.timing.is_some() diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index 9ab1e2b2d..df05f84b3 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -449,7 +449,7 @@ impl RunProjectionReducer for RunProjection { context_window.event_seq = Some(event.seq); stage.context_window = Some(context_window); } - close_inference_bracket(self, stored, props.visit, event.seq); + close_inference_bracket(self, stored, props.visit, event.seq, ts); } EventBody::AgentLlmStarted(props) => { open_inference_bracket(self, stored, props, event.seq, ts); @@ -478,10 +478,10 @@ impl RunProjectionReducer for RunProjection { inference.first_output_kind = None; } EventBody::AgentError(props) => { - close_inference_bracket(self, stored, props.visit, event.seq); + close_inference_bracket(self, stored, props.visit, event.seq, ts); } EventBody::AgentSessionEnded(_) => { - close_inference_brackets_for_session(self, stored); + close_inference_brackets_for_session(self, stored, ts); } EventBody::AgentSessionActivated(props) => { let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) @@ -497,7 +497,7 @@ impl RunProjectionReducer for RunProjection { return Ok(()); }; stage.agent_control = AgentControlState::WaitingForSteer; - close_inference_bracket(self, stored, props.visit, event.seq); + close_inference_bracket(self, stored, props.visit, event.seq, ts); } EventBody::AgentSteeringInjected(props) => { let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) @@ -746,6 +746,7 @@ impl RunProjectionReducer for RunProjection { }); } EventBody::AgentToolStarted(props) => { + let is_root_session = stored.parent_session_id.is_none(); let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) else { return Ok(()); @@ -766,6 +767,22 @@ impl RunProjectionReducer for RunProjection { projection.invoked = true; } } + // 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); + } + } + EventBody::AgentToolCompleted(props) => { + if stored.parent_session_id.is_some() { + 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); } _ => {} } @@ -1044,6 +1061,22 @@ fn matching_inference_slot<'a>( visit: u32, seq: u32, ) -> Option<&'a mut Option> { + Some(&mut matching_inference_stage(state, stored, visit, seq)?.inference) +} + +/// Resolve the stage owning a bracket this event is allowed to mutate, +/// borrowing the whole projection so the caller can also fold elapsed time +/// into the stage's live accumulators. +/// +/// Same gating as [`matching_inference_slot`]: `None` for child-session +/// events, for a stage with no open bracket, and for a bracket belonging to a +/// different session (which is what post-failover events look like). +fn matching_inference_stage<'a>( + state: &'a mut RunProjection, + stored: &RunEvent, + visit: u32, + seq: u32, +) -> Option<&'a mut StageProjection> { if stored.parent_session_id.is_some() { return None; } @@ -1053,7 +1086,7 @@ fn matching_inference_slot<'a>( .inference .as_ref() .is_some_and(|inference| inference.session_id == session_id); - opened_here.then_some(&mut stage.inference) + opened_here.then_some(stage) } /// Resolve the open inference bracket this event is allowed to mutate. @@ -1066,11 +1099,37 @@ fn matching_inference_bracket<'a>( matching_inference_slot(state, stored, visit, seq)?.as_mut() } -/// Close the bracket on a stage-addressed terminal event. -fn close_inference_bracket(state: &mut RunProjection, stored: &RunEvent, visit: u32, seq: u32) { - if let Some(slot) = matching_inference_slot(state, stored, visit, seq) { - *slot = None; - } +/// Close the bracket on a stage-addressed terminal event, folding its elapsed +/// time into the stage's live inference accumulator. +fn close_inference_bracket( + state: &mut RunProjection, + stored: &RunEvent, + visit: u32, + seq: u32, + ts: DateTime, +) { + let Some(stage) = matching_inference_stage(state, stored, visit, seq) else { + return; + }; + close_bracket_on_stage(stage, ts); +} + +/// Take the open bracket and add its span to `live_inference_ms`. +/// +/// Retries inside the bracket are deliberately included: the in-process +/// stopwatch counts a retried attempt's elapsed time as inference, and +/// `agent.llm.retry` keeps the bracket open rather than reopening it. +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)); +} + +/// 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. @@ -1088,7 +1147,11 @@ fn close_inference_bracket(state: &mut RunProjection, stored: &RunEvent, visit: /// 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(state: &mut RunProjection, stored: &RunEvent) { +fn close_inference_brackets_for_session( + state: &mut RunProjection, + stored: &RunEvent, + ts: DateTime, +) { if stored.parent_session_id.is_some() { return; } @@ -1101,7 +1164,7 @@ fn close_inference_brackets_for_session(state: &mut RunProjection, stored: &RunE .as_ref() .is_some_and(|inference| inference.session_id == session_id); if opened_here { - stage.inference = None; + close_bracket_on_stage(stage, ts); } } } @@ -1417,18 +1480,17 @@ fn finalize_unfinished_stages_after_run_failed( continue; } + // Close any bracket still open so its span is not dropped on the + // floor when the live estimate is frozen into `timing` below. + close_bracket_on_stage(stage, timestamp); + + // Freeze the live estimate before flipping to a terminal state: + // `live_timing` reads `effective_state` and would return wall-only + // once the stage no longer looks in-flight. + let frozen = stage.live_timing(timestamp); stage.state = terminal_state; - if stage.timing.is_none() { - if let Some(started_at) = stage.started_at { - let wall_time_ms = u64::try_from( - timestamp - .signed_duration_since(started_at) - .num_milliseconds() - .max(0), - ) - .expect("non-negative milliseconds fit in u64"); - stage.timing = Some(fabro_types::StageTiming::wall_only(wall_time_ms)); - } + if stage.timing.is_none() && stage.started_at.is_some() { + stage.timing = Some(frozen); } } } @@ -1560,6 +1622,373 @@ mod tests { use super::{RunProjection, RunProjectionReducer, build_summary}; use crate::{Error, EventEnvelope, StageId}; + /// Live accumulation of inference and tool time while a stage is in + /// flight. The finalized breakdown still arrives with the terminal event + /// and replaces these; these exist so a long-running stage is not reported + /// as doing no work. + mod live_active_accumulation { + use fabro_types::run_event::{ + AgentLlmFirstOutputProps, AgentLlmRetryProps, AgentLlmStartedProps, + AgentToolCompletedProps, AgentToolStartedProps, + }; + use fabro_types::{ + LlmOutputKind, LlmRetryPhase, ModelRef, Speed, StageOutcome, StageProjection, + }; + + use super::*; + + fn stage_id() -> StageId { + StageId::new("plan", 1) + } + + fn agent_event(seq: u32, ts: &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 + } + + /// 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); + event.event.session_id = Some("session-child".to_string()); + event.event.parent_session_id = Some("session-1".to_string()); + event + } + + fn llm_started() -> EventBody { + EventBody::AgentLlmStarted(AgentLlmStartedProps { + requested_model: ModelRef { + provider: "anthropic".parse().unwrap(), + model_id: "claude-fable-5".into(), + speed: Some(Speed::Fast), + }, + visit: 1, + }) + } + + fn tool_started(tool_call_id: &str) -> EventBody { + EventBody::AgentToolStarted(AgentToolStartedProps { + tool_name: "Bash".to_string(), + tool_call_id: tool_call_id.to_string(), + arguments: json!({}), + visit: 1, + tool_call: None, + turn_id: None, + parent_message_id: None, + }) + } + + fn tool_completed(tool_call_id: &str) -> EventBody { + EventBody::AgentToolCompleted(AgentToolCompletedProps { + tool_name: "Bash".to_string(), + tool_call_id: tool_call_id.to_string(), + output: json!("ok"), + is_error: false, + visit: 1, + tool_result: None, + turn_id: None, + }) + } + + fn agent_message() -> EventBody { + EventBody::AgentMessage(live_agent_message_props(live_counts(10, 5))) + } + + fn started_state() -> RunProjection { + let mut state = initialized_projection(); + state + .apply_event(&test_stage_event_at( + 1, + "2026-04-07T12:00:00Z", + EventBody::StageStarted(started_props()), + stage_id(), + )) + .unwrap(); + state + } + + fn stage(state: &RunProjection) -> &StageProjection { + state.stage(&stage_id()).unwrap() + } + + #[test] + fn closing_an_inference_bracket_accumulates_its_span() { + 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:06Z", + EventBody::AgentLlmFirstOutput(AgentLlmFirstOutputProps { + kind: LlmOutputKind::Text, + visit: 1, + }), + )) + .unwrap(); + state + .apply_event(&agent_event(4, "2026-04-07T12:00:12Z", agent_message())) + .unwrap(); + + // 12:00:05 -> 12:00:12; first_output is a marker, not the close. + assert_eq!(stage(&state).live_inference_ms, 7_000); + assert!(stage(&state).inference.is_none()); + } + + #[test] + fn concurrent_tool_calls_count_once_not_per_call() { + let mut state = started_state(); + for (seq, id) in [(2, "call-a"), (3, "call-b"), (4, "call-c")] { + state + .apply_event(&agent_event(seq, "2026-04-07T12:00:00Z", tool_started(id))) + .unwrap(); + } + // All three finish 10s later. Summing per-call spans would report + // 30s; the batch actually occupied 10s of wall time. + for (seq, id) in [(5, "call-a"), (6, "call-b"), (7, "call-c")] { + state + .apply_event(&agent_event( + seq, + "2026-04-07T12:00:10Z", + tool_completed(id), + )) + .unwrap(); + } + + assert_eq!(stage(&state).live_tool_ms, 10_000); + assert!(stage(&state).tool_batch.is_none()); + } + + #[test] + fn a_batch_stays_open_until_its_last_call_reports() { + let mut state = started_state(); + state + .apply_event(&agent_event( + 2, + "2026-04-07T12:00:00Z", + tool_started("call-a"), + )) + .unwrap(); + state + .apply_event(&agent_event( + 3, + "2026-04-07T12:00:02Z", + tool_started("call-b"), + )) + .unwrap(); + state + .apply_event(&agent_event( + 4, + "2026-04-07T12:00:05Z", + tool_completed("call-a"), + )) + .unwrap(); + + assert_eq!( + stage(&state).live_tool_ms, + 0, + "batch must not close while call-b is outstanding" + ); + + state + .apply_event(&agent_event( + 5, + "2026-04-07T12:00:09Z", + tool_completed("call-b"), + )) + .unwrap(); + + // Measured from the batch open, not from the last call's start. + assert_eq!(stage(&state).live_tool_ms, 9_000); + } + + #[test] + fn successive_batches_accumulate() { + let mut state = started_state(); + for (seq, ts, body) in [ + (2, "2026-04-07T12:00:00Z", tool_started("call-a")), + (3, "2026-04-07T12:00:04Z", tool_completed("call-a")), + (4, "2026-04-07T12:00:10Z", tool_started("call-b")), + (5, "2026-04-07T12:00:16Z", tool_completed("call-b")), + ] { + state.apply_event(&agent_event(seq, ts, body)).unwrap(); + } + + assert_eq!(stage(&state).live_tool_ms, 10_000); + } + + #[test] + fn a_duplicate_completion_does_not_drain_the_batch_early() { + let mut state = started_state(); + state + .apply_event(&agent_event( + 2, + "2026-04-07T12:00:00Z", + tool_started("call-a"), + )) + .unwrap(); + state + .apply_event(&agent_event( + 3, + "2026-04-07T12:00:00Z", + tool_started("call-b"), + )) + .unwrap(); + // call-a reports twice, as a replayed or duplicated log can. + state + .apply_event(&agent_event( + 4, + "2026-04-07T12:00:03Z", + tool_completed("call-a"), + )) + .unwrap(); + state + .apply_event(&agent_event( + 5, + "2026-04-07T12:00:04Z", + tool_completed("call-a"), + )) + .unwrap(); + + assert_eq!(stage(&state).live_tool_ms, 0); + assert!(stage(&state).tool_batch.is_some()); + } + + #[test] + fn subagent_tool_calls_do_not_double_count_against_the_root_batch() { + let mut state = started_state(); + state + .apply_event(&agent_event( + 2, + "2026-04-07T12:00:00Z", + tool_started("root-call"), + )) + .unwrap(); + // The sub-agent's own tools run inside the root call's span. + state + .apply_event(&child_event( + 3, + "2026-04-07T12:00:01Z", + tool_started("child-call"), + )) + .unwrap(); + state + .apply_event(&child_event( + 4, + "2026-04-07T12:00:02Z", + tool_completed("child-call"), + )) + .unwrap(); + state + .apply_event(&agent_event( + 5, + "2026-04-07T12:00:08Z", + tool_completed("root-call"), + )) + .unwrap(); + + assert_eq!(stage(&state).live_tool_ms, 8_000); + } + + #[test] + fn session_end_accumulates_rather_than_discarding_the_bracket() { + let mut state = started_state(); + state + .apply_event(&agent_event(2, "2026-04-07T12:00:05Z", llm_started())) + .unwrap(); + + let mut ended = test_stage_event_at( + 3, + "2026-04-07T12:00:20Z", + EventBody::AgentSessionEnded(AgentSessionEndedProps {}), + stage_id(), + ); + ended.event.session_id = Some("session-1".to_string()); + state.apply_event(&ended).unwrap(); + + assert_eq!(stage(&state).live_inference_ms, 15_000); + assert!(stage(&state).inference.is_none()); + } + + #[test] + fn a_foreign_session_close_leaves_the_bracket_and_accumulator_alone() { + let mut state = started_state(); + state + .apply_event(&agent_event(2, "2026-04-07T12:00:05Z", llm_started())) + .unwrap(); + + // Post-failover: a new session emits the message, so the old + // bracket is not this event's to close or bill. + let mut foreign = + test_stage_event_at(3, "2026-04-07T12:00:20Z", agent_message(), stage_id()); + foreign.event.session_id = Some("session-2".to_string()); + state.apply_event(&foreign).unwrap(); + + assert_eq!(stage(&state).live_inference_ms, 0); + assert!(stage(&state).inference.is_some()); + } + + #[test] + fn a_retry_keeps_accumulating_within_one_bracket() { + let mut state = started_state(); + state + .apply_event(&agent_event(2, "2026-04-07T12:00:00Z", llm_started())) + .unwrap(); + state + .apply_event(&agent_event( + 3, + "2026-04-07T12:00:04Z", + EventBody::AgentLlmRetry(AgentLlmRetryProps { + provider: "anthropic".to_string(), + model: "claude-fable-5".to_string(), + attempt: 0, + delay_secs: 0.0, + error: json!({ "kind": "stream" }), + phase: Some(LlmRetryPhase::Consume), + visit: 1, + }), + )) + .unwrap(); + state + .apply_event(&agent_event(4, "2026-04-07T12:00:11Z", agent_message())) + .unwrap(); + + // The whole bracket counts, retry included, matching the + // in-process stopwatch. + assert_eq!(stage(&state).live_inference_ms, 11_000); + assert_eq!(stage(&state).inference, None); + } + + #[test] + fn stage_completion_replaces_the_live_estimate_with_finalized_timing() { + let mut state = started_state(); + state + .apply_event(&agent_event(2, "2026-04-07T12:00:00Z", llm_started())) + .unwrap(); + state + .apply_event(&agent_event(3, "2026-04-07T12:00:09Z", agent_message())) + .unwrap(); + assert_eq!(stage(&state).live_inference_ms, 9_000); + + state + .apply_event(&test_stage_event_at( + 4, + "2026-04-07T12:00:10Z", + EventBody::StageCompleted(completed_props(10_000, StageOutcome::Succeeded)), + stage_id(), + )) + .unwrap(); + + let stage = stage(&state); + assert_eq!( + stage.live_timing(test_dt("2026-04-07T12:30:00Z")), + stage.timing.unwrap(), + "a terminal stage reports its finalized breakdown, not a live estimate" + ); + } + } + fn test_event(seq: u32, body: EventBody, node_id: Option<&str>) -> EventEnvelope { let event = RunEvent { id: format!("evt-{seq}"), diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 027cce2fc..0576d32f6 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -1,5 +1,5 @@ use std::borrow::Cow; -use std::collections::{BTreeMap, HashMap}; +use std::collections::{BTreeMap, BTreeSet, HashMap}; use std::num::NonZeroU32; use chrono::{DateTime, Utc}; @@ -354,10 +354,31 @@ pub struct StageProjection { /// immutable projections under their own `StageId`s. /// /// `None` for stages still in flight (`started_at` is set but no terminal - /// event has been observed yet). For live wall-time ticking, the UI uses - /// `started_at`; once terminal this carries the finalized breakdown. + /// event has been observed yet). For a live breakdown while in flight, use + /// [`StageProjection::live_timing`]; once terminal this carries the + /// finalized, authoritative breakdown. #[serde(default, skip_serializing_if = "Option::is_none")] pub timing: Option, + /// Inference time accumulated from closed brackets during this attempt. + /// + /// Live estimate only: the authoritative value arrives with the terminal + /// event and lands in `timing`. Excludes the currently-open bracket, which + /// [`StageProjection::live_timing`] adds from `inference.started_at`. + #[serde(default, skip_serializing_if = "is_zero_ms")] + pub live_inference_ms: u64, + /// Tool time accumulated from closed tool batches during this attempt. + /// + /// A batch spans the first `agent.tool.started` with no outstanding calls + /// through the `agent.tool.completed` that drains the last one, so tools + /// running concurrently within a turn are counted once. This matches how + /// the in-process stopwatch brackets `execute_tool_calls`; summing + /// per-call durations would over-count parallel tool use. + #[serde(default, skip_serializing_if = "is_zero_ms")] + pub live_tool_ms: u64, + /// Open tool batch for this stage: when the current batch started, and the + /// `tool_call_id`s that have not yet reported completion. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub tool_batch: Option, #[serde(default)] pub usage: BilledTokenCounts, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -395,6 +416,30 @@ pub struct StageProjection { pub state: StageState, } +/// Serde guard so zero-valued live accumulators stay off the wire. +#[allow( + clippy::trivially_copy_pass_by_ref, + reason = "serde skip_serializing_if predicates receive fields by reference" +)] +fn is_zero_ms(value: &u64) -> bool { + *value == 0 +} + +/// One open tool batch: tool calls dispatched together that have not all +/// reported completion. +/// +/// `open_call_ids` is a set rather than a count because `agent.tool.completed` +/// identifies its call by id, and a projection replaying a truncated or +/// duplicated log must not let a repeated completion drain the batch early. +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct StageToolBatchProjection { + /// When the batch opened — the first `agent.tool.started` observed while + /// no other calls were outstanding. + pub started_at: DateTime, + /// Calls dispatched but not yet completed, by `tool_call_id`. + pub open_call_ids: BTreeSet, +} + /// One open inference bracket: a dispatched LLM request that has not yet /// produced a message, error, or interrupt. #[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] @@ -506,6 +551,9 @@ impl StageProjection { response: None, completion: None, timing: None, + live_inference_ms: 0, + live_tool_ms: 0, + tool_batch: None, usage: BilledTokenCounts::default(), model: None, root_agent_todos: None, @@ -563,6 +611,118 @@ impl StageProjection { self.timing.map(|timing| timing.wall_time_ms) } + /// Live timing breakdown in milliseconds — the active-time twin of + /// [`Self::live_wall_time_ms`]. + /// + /// Once terminal, returns the stored `timing` unchanged: the finalized + /// breakdown comes from the worker's own stopwatch and is authoritative. + /// + /// While in flight, returns an estimate reconstructed from the event log: + /// accumulated closed brackets plus whatever bracket is open right now. + /// The estimate is per-handler, because only agent stages emit brackets at + /// all: + /// + /// - `Agent` — accumulated inference and tool brackets, plus the open + /// inference bracket and open tool batch. + /// - `Prompt` — one inference call spanning the stage, so elapsed time + /// since `started_at` counts as inference. Matches the finalized + /// `active_only(inference, 0)`. + /// - `Command` — the command *is* the work, so elapsed time counts as tool. + /// Matches the finalized `active_only(0, duration_ms)`. + /// - Everything else — zero. Waiting on a human, a timer, a condition, or + /// child branches is wall time, not active time. + /// + /// Active is clamped to wall. A worker killed mid-turn leaves its bracket + /// open forever (see [`StageInferenceProjection`]), and without the clamp + /// that bracket would tick up without bound. The clamp does not need to + /// know the worker died: a stage cannot have been active longer than it + /// has existed. `watchdog.timeout` remains the authority on whether a run + /// is actually stuck. + #[must_use] + pub fn live_timing(&self, now: DateTime) -> StageTiming { + if let Some(timing) = self.timing { + return timing; + } + + 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 + // `StageHandler::from_handler_type`: the accumulators below are only + // ever populated by agent events, so a legacy non-agent stage still + // reads zero rather than being credited work it never did. + let handler = self.handler.unwrap_or(StageHandler::Agent); + let (inference_time_ms, tool_time_ms) = match handler { + StageHandler::Agent => { + let open_inference = self + .inference + .as_ref() + .map_or(0, |inference| elapsed(inference.started_at)); + let open_tool = self + .tool_batch + .as_ref() + .map_or(0, |batch| elapsed(batch.started_at)); + ( + self.live_inference_ms.saturating_add(open_inference), + self.live_tool_ms.saturating_add(open_tool), + ) + } + StageHandler::Prompt => (wall_time_ms, 0), + StageHandler::Command => (0, wall_time_ms), + // Waiting on a human, a timer, a condition, or child branches is + // wall time, not active time. + StageHandler::Human + | StageHandler::Wait + | StageHandler::Conditional + | StageHandler::Parallel + | StageHandler::ParallelFanIn + | StageHandler::StackManagerLoop + | StageHandler::Start + | StageHandler::Exit => (0, 0), + }; + + StageTiming::new(wall_time_ms, inference_time_ms, tool_time_ms).clamped_to_wall() + } + + /// Fold a closed inference bracket into the live accumulator. + pub fn accumulate_inference_ms(&mut self, elapsed_ms: u64) { + self.live_inference_ms = self.live_inference_ms.saturating_add(elapsed_ms); + } + + /// 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) { + self.tool_batch + .get_or_insert_with(|| StageToolBatchProjection { + started_at, + open_call_ids: BTreeSet::new(), + }) + .open_call_ids + .insert(tool_call_id); + } + + /// 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) { + 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() { + 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.tool_batch = None; + } + /// Begin a new automatic attempt within this stage execution: clear every /// per-attempt field so prior-attempt data does not leak, then record /// `started_at` and `state = Running`. Preserves `first_event_seq` @@ -709,10 +869,13 @@ impl RunProjection { /// terminal conclusion yet. /// /// Run-level wall time ticks from `run.started` to `now`. Active time sums - /// inference and tool timing from stages that have already emitted a - /// terminal stage event. Stage projections do not currently track live - /// inference/tool time while a stage is still running, so active time steps - /// forward when each stage completes while wall time advances continuously. + /// [`StageProjection::live_timing`] across every stage, so an in-flight + /// stage contributes its live estimate rather than nothing — both halves + /// advance continuously. Terminal stages contribute their finalized, + /// authoritative breakdown. + /// + /// Active is not clamped to run wall time here: concurrent branches can + /// legitimately sum past it. The clamp applies per stage. #[must_use] pub fn live_run_timing(&self, now: DateTime) -> Option { let start = self.start.as_ref()?; @@ -720,7 +883,7 @@ impl RunProjection { let active = self .stages .values() - .filter_map(|stage| stage.timing) + .map(|stage| stage.live_timing(now)) .fold(RunTiming::default(), |acc, timing| { acc.saturating_add(&RunTiming::from(timing)) }); @@ -971,3 +1134,286 @@ 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, + }; + + fn seq(n: u32) -> NonZeroU32 { + NonZeroU32::new(n).unwrap() + } + + fn at(seconds: i64) -> DateTime { + Utc.timestamp_opt(1_700_000_000 + seconds, 0).unwrap() + } + + fn projection() -> RunProjection { + RunProjection::new( + "Test run".to_string(), + RunSpec { + run_id: RunId::new(), + settings: WorkflowSettings::default(), + graph: Graph::new("test"), + graph_source: None, + workflow_slug: None, + automation: None, + source_directory: None, + labels: HashMap::default(), + provenance: test_support::test_run_provenance(), + manifest_blob: None, + definition_blob: None, + git: None, + fork_source_ref: None, + }, + at(0), + ) + } + + /// In-flight stage that started at `at(0)`. + fn running(handler: StageHandler) -> StageProjection { + let mut stage = StageProjection::new(seq(1)); + stage.handler = Some(handler); + stage.started_at = Some(at(0)); + stage.state = StageState::Running; + stage + } + + fn open_bracket(started_at: DateTime) -> StageInferenceProjection { + StageInferenceProjection { + session_id: "session-1".to_string(), + started_at, + requested_model: ModelRef { + provider: "anthropic".parse().unwrap(), + model_id: "claude-sonnet-5".into(), + speed: None, + }, + first_output_at: None, + first_output_kind: None, + retries: 0, + } + } + + #[test] + fn terminal_stage_returns_stored_timing_unchanged() { + let mut stage = running(StageHandler::Agent); + stage.state = StageState::Succeeded; + stage.timing = Some(StageTiming::new(90_000, 78_000, 7_000)); + // Live accumulators are stale leftovers; the finalized value wins. + stage.live_inference_ms = 5; + stage.live_tool_ms = 5; + + assert_eq!( + stage.live_timing(at(600)), + StageTiming::new(90_000, 78_000, 7_000) + ); + } + + #[test] + fn agent_stage_sums_accumulators_and_open_brackets() { + let mut stage = running(StageHandler::Agent); + stage.live_inference_ms = 30_000; + stage.live_tool_ms = 5_000; + // Inference open for 20s, tools open for 10s, at t=120s. + stage.inference = Some(open_bracket(at(100))); + stage.tool_batch = Some(StageToolBatchProjection { + started_at: at(110), + open_call_ids: ["call-1".to_string()].into_iter().collect(), + }); + + assert_eq!( + stage.live_timing(at(120)), + StageTiming::new(120_000, 50_000, 15_000) + ); + } + + #[test] + fn agent_stage_without_brackets_reports_only_accumulators() { + let mut stage = running(StageHandler::Agent); + stage.live_inference_ms = 30_000; + stage.live_tool_ms = 5_000; + + assert_eq!( + stage.live_timing(at(120)), + StageTiming::new(120_000, 30_000, 5_000) + ); + } + + #[test] + fn prompt_stage_counts_elapsed_as_inference() { + let stage = running(StageHandler::Prompt); + + assert_eq!( + stage.live_timing(at(45)), + StageTiming::new(45_000, 45_000, 0) + ); + } + + #[test] + fn command_stage_counts_elapsed_as_tool() { + let stage = running(StageHandler::Command); + + assert_eq!( + stage.live_timing(at(45)), + StageTiming::new(45_000, 0, 45_000) + ); + } + + #[test] + fn waiting_handlers_report_wall_time_with_zero_active() { + for handler in [ + StageHandler::Human, + StageHandler::Wait, + StageHandler::Conditional, + StageHandler::Parallel, + StageHandler::ParallelFanIn, + StageHandler::StackManagerLoop, + StageHandler::Start, + StageHandler::Exit, + ] { + let stage = running(handler); + let timing = stage.live_timing(at(600)); + + assert_eq!( + timing, + StageTiming::new(600_000, 0, 0), + "{handler} should report wall time only" + ); + } + } + + #[test] + fn open_bracket_from_a_killed_worker_is_clamped_to_wall() { + let mut stage = running(StageHandler::Agent); + // Bracket opened before the stage even started — the pathological + // shape a killed worker leaves behind. Without the clamp this would + // report 700s of inference against 600s of wall. + stage.inference = Some(open_bracket(at(-100))); + + let timing = stage.live_timing(at(600)); + + assert_eq!(timing.wall_time_ms, 600_000); + assert_eq!(timing.active_time_ms, 600_000); + } + + #[test] + fn clamping_preserves_the_inference_tool_split() { + let mut stage = running(StageHandler::Agent); + // 3:1 inference:tool, totalling 200s of active against 100s of wall. + stage.live_inference_ms = 150_000; + stage.live_tool_ms = 50_000; + + let timing = stage.live_timing(at(100)); + + assert_eq!(timing.wall_time_ms, 100_000); + assert_eq!(timing.active_time_ms, 100_000); + assert_eq!(timing.inference_time_ms, 75_000); + assert_eq!(timing.tool_time_ms, 25_000); + } + + #[test] + fn live_run_timing_counts_in_flight_stages_not_just_terminal_ones() { + // The shape that motivated this change: two finished stages and one + // long-running agent stage that had been active nearly the whole run. + let mut projection = projection(); + projection.start = Some(StartRecord { + start_time: at(0), + run_branch: None, + base_sha: None, + }); + + let baseline = projection.stage_entry("baseline", 1, 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)); + 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)); + plan.handler = Some(StageHandler::Agent); + plan.started_at = Some(at(146)); + plan.state = StageState::Running; + plan.live_inference_ms = 700_000; + plan.live_tool_ms = 150_000; + + let timing = projection.live_run_timing(at(1_013)).unwrap(); + + assert_eq!(timing.wall_time_ms, 1_013_000); + assert_eq!(timing.inference_time_ms, 778_230); + assert_eq!(timing.tool_time_ms, 200_251); + // Before this change the in-flight stage contributed nothing and the + // run reported 128,481 ms of active time against 1,013,000 ms of wall. + assert_eq!(timing.active_time_ms, 978_481); + } + + #[test] + fn live_run_timing_may_exceed_run_wall_when_branches_overlap() { + let mut projection = projection(); + projection.start = Some(StartRecord { + start_time: at(0), + run_branch: None, + base_sha: None, + }); + + 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)); + stage.handler = Some(StageHandler::Agent); + stage.state = StageState::Succeeded; + stage.timing = Some(StageTiming::new(60_000, 60_000, 0)); + } + + let timing = projection.live_run_timing(at(60)).unwrap(); + + assert_eq!(timing.wall_time_ms, 60_000); + assert_eq!( + timing.active_time_ms, 180_000, + "concurrent branches legitimately sum past run wall time" + ); + } +} + +#[cfg(test)] +mod live_timing_legacy_tests { + use chrono::{DateTime, TimeZone, Utc}; + + use crate::{StageProjection, StageState, StageTiming, first_event_seq}; + + fn at(seconds: i64) -> DateTime { + Utc.timestamp_opt(1_700_000_000 + seconds, 0).unwrap() + } + + #[test] + fn a_legacy_stage_without_a_recorded_handler_uses_its_accumulators() { + let mut stage = StageProjection::new(first_event_seq(1)); + stage.handler = None; + stage.started_at = Some(at(0)); + stage.state = StageState::Running; + stage.live_inference_ms = 30_000; + + assert_eq!( + stage.live_timing(at(120)), + StageTiming::new(120_000, 30_000, 0) + ); + } + + #[test] + fn a_legacy_stage_with_no_accumulators_reports_no_active_time() { + let mut stage = StageProjection::new(first_event_seq(1)); + stage.handler = None; + stage.started_at = Some(at(0)); + stage.state = StageState::Running; + + assert_eq!(stage.live_timing(at(120)), StageTiming::new(120_000, 0, 0)); + } +} diff --git a/lib/foundation/fabro-types/src/timing.rs b/lib/foundation/fabro-types/src/timing.rs index 0d1e555e8..4fe45a36f 100644 --- a/lib/foundation/fabro-types/src/timing.rs +++ b/lib/foundation/fabro-types/src/timing.rs @@ -62,6 +62,32 @@ impl StageTiming { Self::new(0, inference_time_ms, tool_time_ms) } + /// Scale the breakdown down so `active_time_ms` does not exceed + /// `wall_time_ms`, preserving the inference/tool ratio. + /// + /// Only meaningful for live estimates of a single in-flight stage, where + /// an open bracket left behind by a killed worker would otherwise tick up + /// without bound. Finalized timings come from the worker's stopwatch and + /// already satisfy the invariant. + /// + /// Deliberately *not* applied at run level: concurrent branches can + /// legitimately sum past run wall time. + #[must_use] + pub fn clamped_to_wall(&self) -> Self { + if self.active_time_ms <= 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); + 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) + } + /// Sum two timings field-by-field. Used to aggregate visits of one node /// and to accumulate run-level rollups. #[must_use] diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index e3a1a7c0e..4ae265b76 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -487,6 +487,7 @@ models/stage-projection.ts models/stage-state.ts models/stage-summary.ts models/stage-timing.ts +models/stage-tool-batch-projection.ts models/start-record.ts models/start-run-request.ts models/steer-run-request.ts diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 5cd29601d..029a2d3ab 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -457,6 +457,7 @@ export * from './stage-projection'; export * from './stage-state'; export * from './stage-summary'; export * from './stage-timing'; +export * from './stage-tool-batch-projection'; export * from './start-record'; export * from './start-run-request'; export * from './steer-run-request'; diff --git a/lib/packages/fabro-api-client/src/models/run-timing.ts b/lib/packages/fabro-api-client/src/models/run-timing.ts index fb585b0e9..d4fa62262 100644 --- a/lib/packages/fabro-api-client/src/models/run-timing.ts +++ b/lib/packages/fabro-api-client/src/models/run-timing.ts @@ -15,7 +15,7 @@ /** - * Timing rollup for an entire run. Active fields sum work across stage visits, so `active_time_ms` can exceed `wall_time_ms` when parallel branches run concurrently. + * Timing rollup for an entire run. Active fields sum work across stage visits, so `active_time_ms` can exceed `wall_time_ms` when parallel branches run concurrently. For a running run, stages still in flight contribute a live estimate rather than nothing, so wall and active both advance continuously. Unlike `StageTiming`, active is not clamped to wall here — concurrent branches can legitimately sum past run wall time. */ export interface RunTiming { 'wall_time_ms': number; diff --git a/lib/packages/fabro-api-client/src/models/stage-projection.ts b/lib/packages/fabro-api-client/src/models/stage-projection.ts index 57fd342ca..8cb88c0ac 100644 --- a/lib/packages/fabro-api-client/src/models/stage-projection.ts +++ b/lib/packages/fabro-api-client/src/models/stage-projection.ts @@ -60,6 +60,9 @@ import type { StageState } from './stage-state'; import type { StageTiming } from './stage-timing'; // May contain unused imports in some cases // @ts-ignore +import type { StageToolBatchProjection } from './stage-tool-batch-projection'; +// May contain unused imports in some cases +// @ts-ignore import type { SubAgentProjection } from './sub-agent-projection'; // May contain unused imports in some cases // @ts-ignore @@ -96,6 +99,15 @@ export interface StageProjection { */ 'started_at'?: string | null; 'timing'?: StageTiming | null; + /** + * Inference time accumulated from closed brackets during the current attempt. Live estimate only — the authoritative value arrives with the terminal event and lands in `timing`. Excludes the currently open bracket, whose span is measured from `inference.started_at`. + */ + 'live_inference_ms'?: number; + /** + * Tool time accumulated from closed tool batches during the current attempt. A batch spans the first dispatched call through the completion that drains the last outstanding one, so tools running concurrently within a turn are counted once. + */ + 'live_tool_ms'?: number; + 'tool_batch'?: StageToolBatchProjection | null; 'usage': BilledTokenCounts; 'model'?: BillingModelRef | null; 'todos'?: TodoListProjection | null; diff --git a/lib/packages/fabro-api-client/src/models/stage-timing.ts b/lib/packages/fabro-api-client/src/models/stage-timing.ts index d9915c409..5bdd92bb4 100644 --- a/lib/packages/fabro-api-client/src/models/stage-timing.ts +++ b/lib/packages/fabro-api-client/src/models/stage-timing.ts @@ -15,7 +15,7 @@ /** - * Timing breakdown for one stage visit. Fields are all milliseconds. `wall_time_ms` is elapsed clock time; `inference_time_ms` is Fabro- observed LLM request/stream elapsed time; `tool_time_ms` is tool or command execution elapsed time; `active_time_ms` equals `inference_time_ms + tool_time_ms`. + * Timing breakdown for one stage visit. Fields are all milliseconds. `wall_time_ms` is elapsed clock time; `inference_time_ms` is Fabro- observed LLM request/stream elapsed time; `tool_time_ms` is tool or command execution elapsed time; `active_time_ms` equals `inference_time_ms + tool_time_ms`. For a terminal stage these come from the worker\'s own stopwatch and are authoritative. For a stage still in flight they are a live estimate reconstructed from the event log, and `active_time_ms` is clamped to `wall_time_ms`. The estimate is replaced by the authoritative breakdown when the stage reaches a terminal event. */ export interface StageTiming { 'wall_time_ms': number; 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 new file mode 100644 index 000000000..a6bd7718c --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/stage-tool-batch-projection.ts @@ -0,0 +1,29 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * 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 { + /** + * When the batch opened — the first dispatched call observed while no other calls were outstanding. + */ + 'started_at': string; + /** + * Calls dispatched but not yet completed, by tool call id. + */ + 'open_call_ids': Array; +}