diff --git a/apps/fabro-web/app/components/stage-sidebar.tsx b/apps/fabro-web/app/components/stage-sidebar.tsx index 2d3122602..e2cfe88f8 100644 --- a/apps/fabro-web/app/components/stage-sidebar.tsx +++ b/apps/fabro-web/app/components/stage-sidebar.tsx @@ -1,4 +1,4 @@ -import { useState, useEffect, useRef, type ComponentType } from "react"; +import { useEffect, useRef, type ComponentType } from "react"; import { Link } from "react-router"; import type { StageState } from "@qltysh/fabro-api-client"; import { @@ -12,6 +12,7 @@ import { import { Bars3BottomLeftIcon, DocumentTextIcon, MapIcon } from "@heroicons/react/24/outline"; import { formatDurationSecs } from "../lib/format"; import { ACTIVE_STAGE_STATES, formatStageLabel } from "../lib/stage-sidebar"; +import { useTickingNow } from "../lib/time"; export interface Stage { id: string; @@ -43,7 +44,6 @@ interface StageSidebarProps { export function StageSidebar({ stages, runId, selectedStageId, activeLink }: StageSidebarProps) { // Track when we first observed each running stage (for ticking timer) const runningStartRef = useRef>(new Map()); - const [, setTick] = useState(0); // Track start times for running stages useEffect(() => { @@ -63,16 +63,13 @@ export function StageSidebar({ stages, runId, selectedStageId, activeLink }: Sta }, [stages]); // Tick every second while any stage is running - useEffect(() => { - if (!stages.some((s) => ACTIVE_STAGE_STATES.has(s.status))) return; - const interval = setInterval(() => setTick((t) => t + 1), 1000); - return () => clearInterval(interval); - }, [stages]); + const hasActive = stages.some((s) => ACTIVE_STAGE_STATES.has(s.status)); + const now = useTickingNow(hasActive); function stageDuration(stage: Stage): string { if (ACTIVE_STAGE_STATES.has(stage.status)) { const start = runningStartRef.current.get(stage.id); - if (start) return formatDurationSecs(Math.floor((Date.now() - start) / 1000)); + if (start) return formatDurationSecs(Math.floor((now - start) / 1000)); return "0s"; } return stage.duration; diff --git a/apps/fabro-web/app/lib/query-keys.test.ts b/apps/fabro-web/app/lib/query-keys.test.ts index 5d7b3a53a..a06acaedb 100644 --- a/apps/fabro-web/app/lib/query-keys.test.ts +++ b/apps/fabro-web/app/lib/query-keys.test.ts @@ -22,6 +22,7 @@ describe("queryKeys", () => { ]); expect(queryKeysForRunEvent("run-1", "stage.completed", "stage-1")).toEqual([ queryKeys.runs.stages("run-1"), + queryKeys.runs.billing("run-1"), queryKeys.runs.events("run-1", 1000), queryKeys.runs.graph("run-1", "LR"), queryKeys.runs.graph("run-1", "TB"), @@ -48,4 +49,4 @@ describe("queryKeys", () => { test("agent activity events without a node_id invalidate nothing", () => { expect(queryKeysForRunEvent("run-1", "agent.message")).toEqual([]); }); -}); \ No newline at end of file +}); diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 14db0e3f2..432791352 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -50,12 +50,16 @@ describe("queryKeysForRunEvent", () => { ]); }); - test("stage.retrying invalidates the same keys as other stage events", () => { - const keys = queryKeysForRunEvent("run-1", "stage.retrying", "verify@2"); - expect(keys).toContain(queryKeys.runs.stages("run-1")); - expect(keys).toContain(queryKeys.runs.events("run-1", 1000)); - expect(keys).toContain(queryKeys.runs.detail("run-1")); - expect(keys).toContain(queryKeys.runs.stageEvents("run-1", "verify@2")); + test("stage.retrying invalidates stages, billing, events, graph, detail, and stage events", () => { + expect(queryKeysForRunEvent("run-1", "stage.retrying", "verify@2")).toEqual([ + queryKeys.runs.stages("run-1"), + queryKeys.runs.billing("run-1"), + queryKeys.runs.events("run-1", 1000), + queryKeys.runs.graph("run-1", "LR"), + queryKeys.runs.graph("run-1", "TB"), + queryKeys.runs.detail("run-1"), + queryKeys.runs.stageEvents("run-1", "verify@2"), + ]); }); }); diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index 1546d1578..63855296e 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -109,6 +109,7 @@ export function queryKeysForRunEvent( if (STAGE_EVENTS.has(event)) { const keys = [ queryKeys.runs.stages(runId), + queryKeys.runs.billing(runId), queryKeys.runs.events(runId, 1000), queryKeys.runs.graph(runId, "LR"), queryKeys.runs.graph(runId, "TB"), diff --git a/apps/fabro-web/app/lib/stage-sidebar.ts b/apps/fabro-web/app/lib/stage-sidebar.ts index e44e5f4a5..84a5de2de 100644 --- a/apps/fabro-web/app/lib/stage-sidebar.ts +++ b/apps/fabro-web/app/lib/stage-sidebar.ts @@ -1,13 +1,22 @@ -import type { PaginatedRunStageList, StageState } from "@qltysh/fabro-api-client"; +import { StageState } from "@qltysh/fabro-api-client"; +import type { PaginatedRunStageList } from "@qltysh/fabro-api-client"; import type { Stage } from "../components/stage-sidebar"; import { isVisibleStage } from "../data/runs"; import { formatDurationSecs } from "./format"; -export const ACTIVE_STAGE_STATES: ReadonlySet = new Set(["running", "retrying"]); +export const ACTIVE_STAGE_STATES: ReadonlySet = new Set([ + StageState.RUNNING, + StageState.RETRYING, +]); +export const IN_FLIGHT_STAGE_STATES: ReadonlySet = new Set([ + StageState.PENDING, + StageState.RUNNING, + StageState.RETRYING, +]); export const SUCCEEDED_STAGE_STATES: ReadonlySet = new Set([ - "succeeded", - "partially_succeeded", + StageState.SUCCEEDED, + StageState.PARTIALLY_SUCCEEDED, ]); /** diff --git a/apps/fabro-web/app/lib/time.ts b/apps/fabro-web/app/lib/time.ts index 1d8a818d4..fb6f9b87f 100644 --- a/apps/fabro-web/app/lib/time.ts +++ b/apps/fabro-web/app/lib/time.ts @@ -1,3 +1,21 @@ +import { useEffect, useState } from "react"; + +/** + * Re-renders the calling component every `intervalMs` milliseconds while + * `active` is true, returning the current `Date.now()` value at each tick. + * Returns the captured value when paused, so renders are stable. + */ +export function useTickingNow(active: boolean, intervalMs = 1000): number { + const [now, setNow] = useState(() => Date.now()); + useEffect(() => { + if (!active) return; + setNow(Date.now()); + const interval = setInterval(() => setNow(Date.now()), intervalMs); + return () => clearInterval(interval); + }, [active, intervalMs]); + return now; +} + function relativeTime(seconds: number, past: boolean): string { if (seconds < 60) return past ? "just now" : "in <1m"; const minutes = Math.floor(seconds / 60); @@ -21,4 +39,3 @@ export function timeAgo(iso: string): string { export function timeUntil(iso: string): string { return relativeTime(Math.floor((new Date(iso).getTime() - Date.now()) / 1000), false); } - diff --git a/apps/fabro-web/app/routes/run-billing.test.tsx b/apps/fabro-web/app/routes/run-billing.test.tsx index ed9e3040c..c35ab52cb 100644 --- a/apps/fabro-web/app/routes/run-billing.test.tsx +++ b/apps/fabro-web/app/routes/run-billing.test.tsx @@ -75,12 +75,14 @@ describe("RunBilling", () => { model: null, billing: zeroBilling(), runtime_secs: 0, + state: "succeeded", }, { stage: { id: "command", name: "command" }, model: null, billing: zeroBilling(), runtime_secs: 61, + state: "succeeded", }, ], totals: { @@ -96,7 +98,7 @@ describe("RunBilling", () => { expect(text).toMatch(/—\s*\/\s*—/); expect(text).toContain("1m 1s"); expect(text).not.toContain("By model"); - expect(text).not.toContain("No completed stages yet"); + expect(text).not.toContain("No stages yet"); }); test("renders mixed LLM and non-LLM rows while counting only LLM rows by model", () => { @@ -108,6 +110,7 @@ describe("RunBilling", () => { model: null, billing: zeroBilling(), runtime_secs: 0, + state: "succeeded", }, { stage: { id: "agent", name: "agent" }, @@ -119,6 +122,7 @@ describe("RunBilling", () => { total_usd_micros: 240000, }), runtime_secs: 42, + state: "succeeded", }, ], totals: { @@ -155,11 +159,61 @@ describe("RunBilling", () => { expect(textFromInstance(byModelFooterCells[1])).toBe("1"); }); - test("keeps the empty state for runs with no completed stages", () => { + test("keeps the empty state for runs with no stages", () => { const renderer = renderBilling(billing()); const text = textFromNode(renderer.toJSON()); - expect(text).toContain("No completed stages yet"); - expect(text).toContain("Stages will appear once the run produces completed nodes."); + expect(text).toContain("No stages yet"); + expect(text).toContain("Stages will appear as soon as the run starts executing."); + }); + + test("renders an in-flight row with live runtime and includes its elapsed time in the footer", () => { + const originalNow = Date.now; + // Pin "now" to 30s after the in-flight row started. + const startedAt = "2026-04-29T12:00:00.000Z"; + const fakeNow = new Date("2026-04-29T12:00:30.000Z").getTime(); + Date.now = () => fakeNow; + + try { + const renderer = renderBilling( + billing({ + stages: [ + { + stage: { id: "in-flight", name: "in-flight" }, + model: null, + // Server reports 0 runtime / no billing; the row is still being executed. + billing: zeroBilling(), + runtime_secs: 0, + started_at: startedAt, + state: "running", + }, + ], + // Server total is 0 because the in-flight row hasn't been finalized. + totals: { + runtime_secs: 0, + ...zeroBilling(), + }, + }), + ); + + const text = textFromNode(renderer.toJSON()); + // Empty-state must NOT show — the table should appear as soon as the + // first stage starts. + expect(text).not.toContain("No stages yet"); + expect(text).toContain("in-flight"); + + // Both the row's runtime cell and the footer total should reflect + // ~30s elapsed since started_at. + expect(text).toContain("30s"); + + const footers = renderer.root.findAll((node) => node.type === "tfoot"); + const footerCells = footers[0].findAll((node) => node.type === "td"); + // The Run time column in the footer is index 3 (Total / [empty Model] / + // Tokens / Run time / Billing). + const footerRuntime = textFromInstance(footerCells[3]); + expect(footerRuntime).toContain("30s"); + } finally { + Date.now = originalNow; + } }); }); diff --git a/apps/fabro-web/app/routes/run-billing.tsx b/apps/fabro-web/app/routes/run-billing.tsx index 6bb6ab763..be0d52386 100644 --- a/apps/fabro-web/app/routes/run-billing.tsx +++ b/apps/fabro-web/app/routes/run-billing.tsx @@ -1,7 +1,11 @@ +import { useMemo } from "react"; + import { EmptyState } from "../components/state"; import { formatDurationSecs } from "../lib/format"; import { useRunBilling } from "../lib/queries"; -import type { RunBilling } from "@qltysh/fabro-api-client"; +import { IN_FLIGHT_STAGE_STATES } from "../lib/stage-sidebar"; +import { useTickingNow } from "../lib/time"; +import type { RunBilling, RunBillingStage } from "@qltysh/fabro-api-client"; const EMPTY_VALUE = "—"; @@ -14,78 +18,103 @@ function formatUsdMicros(usdMicros?: number | null) { return usdMicros == null ? EMPTY_VALUE : `$${(usdMicros / 1_000_000).toFixed(2)}`; } -function mapBilling(billing: RunBilling | undefined) { - if (!billing) { - return { - stages: [], - totalRuntime: formatDurationSecs(0), - totalUsdMicros: undefined, - totalInput: null, - totalOutput: null, - modelBreakdown: [], - modelStageCount: 0, - }; - } +function isInFlight(stage: RunBillingStage): boolean { + return stage.state != null && IN_FLIGHT_STAGE_STATES.has(stage.state); +} - const stages = billing.stages.map((stage) => { - const hasModel = stage.model != null; - return { - stage: stage.stage.name, - model: stage.model?.id ?? null, - inputTokens: hasModel ? stage.billing.input_tokens : null, - outputTokens: hasModel - ? stage.billing.output_tokens + stage.billing.reasoning_tokens - : null, - runtime: formatDurationSecs(stage.runtime_secs), - totalUsdMicros: stage.billing.total_usd_micros, - }; - }); - const totalRuntime = formatDurationSecs(billing.totals.runtime_secs); - const hasLlmStages = billing.by_model.length > 0; - const totalInput = hasLlmStages ? billing.totals.input_tokens : null; - const totalOutput = hasLlmStages - ? billing.totals.output_tokens + billing.totals.reasoning_tokens - : null; - const totalUsdMicros = billing.totals.total_usd_micros; - const modelBreakdown = billing.by_model - .map((entry) => ({ - model: entry.model.id, - stages: entry.stages, - inputTokens: entry.billing.input_tokens, - outputTokens: entry.billing.output_tokens + entry.billing.reasoning_tokens, - totalUsdMicros: entry.billing.total_usd_micros, - })) - .sort((a, b) => (b.totalUsdMicros ?? -1) - (a.totalUsdMicros ?? -1)); - const modelStageCount = modelBreakdown.reduce((sum, row) => sum + row.stages, 0); +interface MappedStageRow { + stage: string; + model: string | null; + inputTokens: number | null; + outputTokens: number | null; + runtimeSecs: number; + totalUsdMicros: number | null | undefined; +} + +function liveRuntimeSecs(stage: RunBillingStage, now: number): number { + if (stage.started_at) { + const startedMs = new Date(stage.started_at).getTime(); + if (Number.isFinite(startedMs)) { + return Math.max(0, (now - startedMs) / 1000); + } + } + return stage.runtime_secs; +} + +function mapStageRow(stage: RunBillingStage, runtimeSecs: number): MappedStageRow { + const hasModel = stage.model != null; return { - stages, - totalRuntime, - totalUsdMicros, - totalInput, - totalOutput, - modelBreakdown, - modelStageCount, + stage: stage.stage.name, + model: stage.model?.id ?? null, + inputTokens: hasModel ? stage.billing.input_tokens : null, + outputTokens: hasModel + ? stage.billing.output_tokens + stage.billing.reasoning_tokens + : null, + runtimeSecs, + totalUsdMicros: stage.billing.total_usd_micros, }; } export default function RunBilling({ params }: { params: { id: string } }) { const billingQuery = useRunBilling(params.id); - const { - stages, - totalRuntime, - totalUsdMicros, - totalInput, - totalOutput, - modelBreakdown, - modelStageCount, - } = mapBilling(billingQuery.data); + const billing = billingQuery.data; + const hasInFlight = billing?.stages.some(isInFlight) ?? false; - if (!stages.length) { + // Tick once per second only while a stage is in-flight. + const now = useTickingNow(hasInFlight); + + // Completed rows don't depend on `now`; memoize them by `billing` so we + // don't reallocate them every tick. + const completedRows = useMemo(() => { + if (!billing) return []; + return billing.stages.map((stage) => mapStageRow(stage, stage.runtime_secs)); + }, [billing]); + + // The model breakdown is server-derived and stable across ticks too. + const modelBreakdown = useMemo(() => { + if (!billing) return []; + return billing.by_model + .map((entry) => ({ + model: entry.model.id, + stages: entry.stages, + inputTokens: entry.billing.input_tokens, + outputTokens: entry.billing.output_tokens + entry.billing.reasoning_tokens, + totalUsdMicros: entry.billing.total_usd_micros, + })) + .sort((a, b) => (b.totalUsdMicros ?? -1) - (a.totalUsdMicros ?? -1)); + }, [billing]); + + // Re-derive only the in-flight rows on each tick; everything else stays put. + const rows = useMemo(() => { + if (!billing) return []; + if (!hasInFlight) return completedRows; + return billing.stages.map((stage, idx) => + isInFlight(stage) + ? mapStageRow(stage, liveRuntimeSecs(stage, now)) + : completedRows[idx], + ); + }, [billing, completedRows, hasInFlight, now]); + + // While ticking, sum the displayed row runtimes so the footer updates in + // lock-step. Otherwise trust the server's authoritative total. + const totalRuntimeSecs = hasInFlight + ? rows.reduce((sum, row) => sum + row.runtimeSecs, 0) + : (billing?.totals.runtime_secs ?? 0); + + const hasLlmStages = (billing?.by_model.length ?? 0) > 0; + const totalInput = hasLlmStages ? (billing?.totals.input_tokens ?? null) : null; + const totalOutput = hasLlmStages && billing + ? billing.totals.output_tokens + billing.totals.reasoning_tokens + : null; + const totalUsdMicros = billing?.totals.total_usd_micros; + const modelStageCount = modelBreakdown.reduce((sum, row) => sum + row.stages, 0); + + if (!rows.length) { return (
); @@ -105,7 +134,7 @@ export default function RunBilling({ params }: { params: { id: string } }) { - {stages.map((row) => ( + {rows.map((row) => ( {row.stage} @@ -115,7 +144,9 @@ export default function RunBilling({ params }: { params: { id: string } }) { {formatTokens(row.inputTokens)} /{" "} {formatTokens(row.outputTokens)} - {row.runtime} + + {formatDurationSecs(row.runtimeSecs)} + {formatUsdMicros(row.totalUsdMicros)} @@ -131,7 +162,7 @@ export default function RunBilling({ params }: { params: { id: string } }) { {formatTokens(totalOutput)} - {totalRuntime} + {formatDurationSecs(totalRuntimeSecs)} {formatUsdMicros(totalUsdMicros)} diff --git a/apps/fabro-web/app/routes/run-stages.tsx b/apps/fabro-web/app/routes/run-stages.tsx index a8fc69168..2b01fc5f3 100644 --- a/apps/fabro-web/app/routes/run-stages.tsx +++ b/apps/fabro-web/app/routes/run-stages.tsx @@ -40,6 +40,7 @@ import type { Stage } from "../components/stage-sidebar"; import { EmptyState } from "../components/state"; import { CopyButton } from "../components/ui"; import { formatDurationSecs } from "../lib/format"; +import { useTickingNow } from "../lib/time"; import { fetchRunCommandLog, useRunStageEvents, useRunStages } from "../lib/queries"; import { STAGE_ACTIVITY_EVENT_TYPES, type StageActivityEventType } from "../lib/run-events"; import { ACTIVE_STAGE_STATES, formatStageLabel, mapRunStagesToSidebarStages } from "../lib/stage-sidebar"; @@ -563,7 +564,6 @@ function RunningStageDuration({ const [startedAt, setStartedAt] = useState(() => isRunning ? Date.now() : null, ); - const [, setTick] = useState(0); useEffect(() => { setStartedAt((current) => { @@ -572,14 +572,10 @@ function RunningStageDuration({ }); }, [isRunning]); - useEffect(() => { - if (!isRunning) return; - const interval = setInterval(() => setTick((tick) => tick + 1), 1000); - return () => clearInterval(interval); - }, [isRunning]); + const now = useTickingNow(isRunning); if (isRunning && startedAt) { - return formatDurationSecs(Math.floor((Date.now() - startedAt) / 1000)); + return formatDurationSecs(Math.floor((now - startedAt) / 1000)); } return duration; } diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 4319f6774..0b499e7fa 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -5314,6 +5314,20 @@ components: oneOf: - $ref: "#/components/schemas/CommandTermination" - type: "null" + started_at: + type: ["string", "null"] + format: date-time + description: Wall-clock time the latest attempt of this stage started, if known. + duration_ms: + type: ["integer", "null"] + format: uint64 + minimum: 0 + description: Wall-clock duration of the stage's latest terminal attempt, if known. + state: + oneOf: + - $ref: "#/components/schemas/StageState" + - type: "null" + description: Lifecycle state of the stage projection. InterviewOption: description: Option stored with an interview question in the event log. @@ -6333,6 +6347,11 @@ components: minimum: 1 description: 1-based visit count; bumped each time the workflow re-enters this node. example: 2 + started_at: + type: ["string", "null"] + format: date-time + description: Wall-clock time the latest attempt of this stage started, if known. + example: "2026-04-29T12:34:56Z" # ── File Diff Schemas ────────────────────────────────────────────── @@ -6526,6 +6545,16 @@ components: type: number description: Wall-clock runtime in seconds, summed across every visit of this node. example: 154.0 + started_at: + type: ["string", "null"] + format: date-time + description: Wall-clock time the latest attempt of this stage started, if known. + example: "2026-04-29T12:34:56Z" + state: + oneOf: + - $ref: "#/components/schemas/StageState" + - type: "null" + description: Lifecycle state of the stage. Use to detect in-flight rows for client-side runtime ticking. RunBillingTotals: description: Aggregate billing totals across all stages of a run. diff --git a/lib/crates/fabro-api/tests/run_billing_stage_round_trip.rs b/lib/crates/fabro-api/tests/run_billing_stage_round_trip.rs index c160a840a..52fb4ece1 100644 --- a/lib/crates/fabro-api/tests/run_billing_stage_round_trip.rs +++ b/lib/crates/fabro-api/tests/run_billing_stage_round_trip.rs @@ -1,4 +1,5 @@ use fabro_api::types::RunBillingStage; +use fabro_types::StageState; use serde_json::json; #[test] @@ -28,3 +29,59 @@ fn run_billing_stage_model_accepts_required_null() { assert!(encoded.get("model").is_some()); assert!(encoded["model"].is_null()); } + +#[test] +fn run_billing_stage_round_trips_terminal_row_with_started_at_and_state() { + let value = json!({ + "stage": { + "id": "build", + "name": "build" + }, + "model": { "id": "claude-sonnet-4-5" }, + "billing": { + "input_tokens": 12, + "output_tokens": 34, + "total_tokens": 46, + "reasoning_tokens": 0, + "cache_read_tokens": 0, + "cache_write_tokens": 0 + }, + "runtime_secs": 5.5, + "started_at": "2026-04-29T12:34:56Z", + "state": "succeeded" + }); + + let stage: RunBillingStage = + serde_json::from_value(value.clone()).expect("terminal stage row should deserialize"); + assert!(stage.started_at.is_some()); + assert_eq!(stage.state, Some(StageState::Succeeded)); + assert_eq!(serde_json::to_value(stage).unwrap(), value); +} + +#[test] +fn run_billing_stage_round_trips_in_flight_row() { + let value = json!({ + "stage": { + "id": "build", + "name": "build" + }, + "model": null, + "billing": { + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + "reasoning_tokens": 0, + "cache_read_tokens": 0, + "cache_write_tokens": 0 + }, + "runtime_secs": 1.25, + "started_at": "2026-04-29T12:34:56Z", + "state": "running" + }); + + let stage: RunBillingStage = + serde_json::from_value(value.clone()).expect("in-flight stage row should deserialize"); + assert!(stage.model.is_none()); + assert_eq!(stage.state, Some(StageState::Running)); + assert_eq!(serde_json::to_value(stage).unwrap(), value); +} diff --git a/lib/crates/fabro-api/tests/stage_projection_round_trip.rs b/lib/crates/fabro-api/tests/stage_projection_round_trip.rs index 0523197cf..ad4dcb579 100644 --- a/lib/crates/fabro-api/tests/stage_projection_round_trip.rs +++ b/lib/crates/fabro-api/tests/stage_projection_round_trip.rs @@ -28,7 +28,10 @@ fn stage_projection_round_trips_representative_json() { "parallel_results": [{ "branch": 0, "status": "succeeded" }], "stdout": "ok", "stderr": "", - "termination": "exited" + "termination": "exited", + "started_at": "2026-04-29T12:34:00Z", + "duration_ms": 56000, + "state": "succeeded" }); let state: StageProjection = serde_json::from_value(value.clone()).unwrap(); diff --git a/lib/crates/fabro-server/src/demo/mod.rs b/lib/crates/fabro-server/src/demo/mod.rs index c900c8b91..539731c82 100644 --- a/lib/crates/fabro-server/src/demo/mod.rs +++ b/lib/crates/fabro-server/src/demo/mod.rs @@ -1214,30 +1214,35 @@ mod runs { "Detect Drift", StageState::Succeeded, Some(72.0), + None, ), run_stage_from_stage_id( &StageId::new("propose-changes", 1), "Propose Changes", StageState::Succeeded, Some(154.0), + None, ), run_stage_from_stage_id( &StageId::new("review-changes", 1), "Review Changes", StageState::Succeeded, Some(45.0), + None, ), run_stage_from_stage_id( &StageId::new("apply-changes", 1), "Apply Changes", StageState::Succeeded, Some(118.0), + None, ), run_stage_from_stage_id( &StageId::new("apply-changes", 2), "Apply Changes", StageState::Running, None, + None, ), ] } @@ -1374,6 +1379,8 @@ mod runs { total_usd_micros: Some(480_000), }, runtime_secs: 72.0, + started_at: None, + state: Some(StageState::Succeeded), }, RunBillingStage { stage: BillingStageRef { @@ -1393,6 +1400,8 @@ mod runs { total_usd_micros: Some(720_000), }, runtime_secs: 154.0, + started_at: None, + state: Some(StageState::Succeeded), }, RunBillingStage { stage: BillingStageRef { @@ -1412,6 +1421,8 @@ mod runs { total_usd_micros: Some(190_000), }, runtime_secs: 45.0, + started_at: None, + state: Some(StageState::Succeeded), }, RunBillingStage { stage: BillingStageRef { @@ -1431,6 +1442,8 @@ mod runs { total_usd_micros: Some(870_000), }, runtime_secs: 118.0, + started_at: None, + state: Some(StageState::Running), }, ], totals: RunBillingTotals { diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 965989c8c..276374d2f 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -570,6 +570,7 @@ pub(crate) fn run_stage_from_stage_id( name: impl Into, status: StageState, duration_secs: Option, + started_at: Option>, ) -> RunStage { RunStage { id: stage_id.to_string(), @@ -579,6 +580,7 @@ pub(crate) fn run_stage_from_stage_id( node_id: stage_id.node_id().to_string(), visit: std::num::NonZeroU32::new(stage_id.visit()) .expect("StageId stores a non-zero visit"), + started_at, } } diff --git a/lib/crates/fabro-server/src/server/handler/billing.rs b/lib/crates/fabro-server/src/server/handler/billing.rs index 193ebd78d..5c000869c 100644 --- a/lib/crates/fabro-server/src/server/handler/billing.rs +++ b/lib/crates/fabro-server/src/server/handler/billing.rs @@ -1,13 +1,14 @@ +use std::collections::HashMap; use std::sync::Arc; -use fabro_store::RunProjectionReducer; -use fabro_types::{EventBody, RunProjection, StageId}; +use chrono::{DateTime, Utc}; +use fabro_types::{RunProjection, StageProjection, StageState}; use super::super::{ - ApiError, AppState, BillingByModel, BillingStageRef, EventEnvelope, HashMap, IntoResponse, - Json, ListResponse, ModelReference, PaginationParams, Path, Query, RequiredUser, Response, - Router, RunBilling, RunBillingStage, RunBillingTotals, RunId, StageState, State, StatusCode, - get, parse_run_id_path, run_stage_from_stage_id, + ApiError, AppState, BillingByModel, BillingStageRef, IntoResponse, Json, ListResponse, + ModelReference, PaginationParams, Path, Query, RequiredUser, Response, Router, RunBilling, + RunBillingStage, RunBillingTotals, RunId, State, StatusCode, get, parse_run_id_path, + run_stage_from_stage_id, }; pub(super) fn routes() -> Router> { @@ -16,41 +17,6 @@ pub(super) fn routes() -> Router> { .route("/runs/{id}/billing", get(get_run_billing)) } -/// Map a `stage.*` lifecycle event body to the [`StageState`] it implies. -/// Returns `None` for any other variant. -fn stage_state_from_lifecycle(body: &EventBody) -> Option { - match body { - EventBody::StageStarted(_) => Some(StageState::Running), - EventBody::StageRetrying(_) => Some(StageState::Retrying), - EventBody::StageFailed(props) => Some(if props.will_retry { - StageState::Retrying - } else { - StageState::Failed - }), - EventBody::StageCompleted(props) => Some(StageState::from(props.status)), - _ => None, - } -} - -/// Single-pass scan over `events` building the latest [`StageState`] for each -/// [`StageId`] from lifecycle events (started/retrying/completed/failed). Each -/// later lifecycle event overwrites earlier ones, leaving the latest as the -/// stored value — equivalent to "scan in reverse, take first match" but in O(E) -/// for the whole list rather than O(stages × events). -fn latest_stage_states(events: &[EventEnvelope]) -> HashMap { - let mut states = HashMap::new(); - for envelope in events { - let Some(stage_id) = envelope.event.stage_id.as_ref() else { - continue; - }; - let Some(state) = stage_state_from_lifecycle(&envelope.event.body) else { - continue; - }; - states.insert(stage_id.clone(), state); - } - states -} - async fn list_run_stages( _auth: RequiredUser, State(state): State>, @@ -62,42 +28,30 @@ async fn list_run_stages( Err(response) => return response, }; - let events = match state.store.open_run_reader(&id).await { - Ok(run_store) => run_store.list_events().await.unwrap_or_default(), - Err(_) => return ApiError::not_found("Run not found.").into_response(), + let Ok(run_store) = state.store.open_run_reader(&id).await else { + return ApiError::not_found("Run not found.").into_response(); }; - - let projection = match RunProjection::apply_events(&events) { - Ok(projection) => projection, + let projection = match run_store.state().await { + Ok(state) => state, Err(err) => { - tracing::warn!( - run_id = %id, - error = %err, - "Failed to build run projection; returning empty stages list", - ); - RunProjection::default() + return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) + .into_response(); } }; - let stage_durations = fabro_workflow::extract_stage_durations_by_stage_id(&events); - let lifecycle_states = latest_stage_states(&events); - let mut stages = Vec::new(); - for (stage_id, stage_projection) in projection.iter_stages() { - // Prefer the latest lifecycle event; fall back to the projection's - // stored completion (e.g. for runs recovered from snapshot only). - let status = lifecycle_states.get(stage_id).copied().unwrap_or_else(|| { - stage_projection - .completion - .as_ref() - .map_or(StageState::Pending, |c| StageState::from(c.outcome)) - }); - stages.push(run_stage_from_stage_id( - stage_id, - stage_id.node_id().to_string(), - status, - stage_durations.get(stage_id).map(|ms| *ms as f64 / 1000.0), - )); - } + let now = Utc::now(); + let stages = projection + .iter_stages() + .map(|(stage_id, stage)| { + run_stage_from_stage_id( + stage_id, + stage_id.node_id().to_string(), + stage.effective_state(), + stage.runtime_secs(now), + stage.started_at, + ) + }) + .collect::>(); (StatusCode::OK, Json(ListResponse::new(stages))).into_response() } @@ -121,6 +75,7 @@ async fn get_run_billing( .into_response(); } }; + let rollup = fabro_workflow::billing_rollup_from_projection(&projection); let by_model = rollup .by_model @@ -133,20 +88,33 @@ async fn get_run_billing( stages: model.stages, }) .collect::>(); - let stages = rollup + + let rollup_by_node = rollup .stages .iter() - .map(|stage| RunBillingStage { - billing: stage.billing.clone(), - model: stage - .model_id - .as_ref() - .map(|id| ModelReference { id: id.clone() }), - runtime_secs: stage.duration_ms as f64 / 1000.0, - stage: BillingStageRef { - id: stage.node_id.clone(), - name: stage.node_id.clone(), - }, + .map(|stage| (stage.node_id.as_str(), stage)) + .collect::>(); + let live_rows = live_billing_rows(&projection, Utc::now()); + let runtime_secs = live_rows.iter().map(|row| row.runtime_secs).sum::(); + let stages = live_rows + .into_iter() + .map(|row| { + let rollup_stage = rollup_by_node.get(row.node_id.as_str()); + RunBillingStage { + billing: rollup_stage + .map(|stage| stage.billing.clone()) + .unwrap_or_default(), + model: rollup_stage + .and_then(|stage| stage.model_id.as_ref()) + .map(|id| ModelReference { id: id.clone() }), + runtime_secs: row.runtime_secs, + stage: BillingStageRef { + id: row.node_id.clone(), + name: row.node_id, + }, + started_at: row.started_at, + state: row.state, + } }) .collect::>(); @@ -154,16 +122,80 @@ async fn get_run_billing( by_model, stages, totals: RunBillingTotals { - cache_read_tokens: rollup.totals.cache_read_tokens, + cache_read_tokens: rollup.totals.cache_read_tokens, cache_write_tokens: rollup.totals.cache_write_tokens, - input_tokens: rollup.totals.input_tokens, - output_tokens: rollup.totals.output_tokens, - reasoning_tokens: rollup.totals.reasoning_tokens, - runtime_secs: rollup.runtime_ms as f64 / 1000.0, - total_tokens: rollup.totals.total_tokens, - total_usd_micros: rollup.totals.total_usd_micros, + input_tokens: rollup.totals.input_tokens, + output_tokens: rollup.totals.output_tokens, + reasoning_tokens: rollup.totals.reasoning_tokens, + runtime_secs, + total_tokens: rollup.totals.total_tokens, + total_usd_micros: rollup.totals.total_usd_micros, }, }; (StatusCode::OK, Json(response)).into_response() } + +struct LiveBillingRow { + node_id: String, + runtime_secs: f64, + started_at: Option>, + state: Option, + latest_visit: u32, +} + +fn live_billing_rows(projection: &RunProjection, now: DateTime) -> Vec { + let mut row_indices = HashMap::::new(); + let mut rows = Vec::::new(); + + for (stage_id, stage) in projection.iter_stages() { + let node_id = stage_id.node_id(); + if is_exit_stage(projection, node_id) || !stage_has_billing_row(stage) { + continue; + } + + let index = *row_indices.entry(node_id.to_string()).or_insert_with(|| { + let index = rows.len(); + rows.push(LiveBillingRow { + node_id: node_id.to_string(), + runtime_secs: 0.0, + started_at: None, + state: None, + latest_visit: 0, + }); + index + }); + let row = &mut rows[index]; + row.runtime_secs += billing_runtime_secs(stage, now).unwrap_or(0.0); + + if stage_id.visit() >= row.latest_visit { + row.latest_visit = stage_id.visit(); + row.started_at = stage.started_at; + row.state = Some(stage.effective_state()); + } + } + + rows +} + +fn billing_runtime_secs(stage: &StageProjection, now: DateTime) -> Option { + stage + .duration_ms + .map(|ms| ms as f64 / 1000.0) + .or_else(|| stage.runtime_secs(now)) +} + +fn stage_has_billing_row(stage: &StageProjection) -> bool { + stage.completion.is_some() + || stage.duration_ms.is_some() + || stage.usage.is_some() + || stage.started_at.is_some() + || stage.state.is_some() +} + +fn is_exit_stage(projection: &RunProjection, node_id: &str) -> bool { + projection + .spec() + .and_then(|spec| spec.graph().nodes.get(node_id)) + .is_some_and(|node| node.handler_type() == Some("exit")) +} diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs index 08196dbdd..aadc30c8e 100644 --- a/lib/crates/fabro-server/src/server/tests.rs +++ b/lib/crates/fabro-server/src/server/tests.rs @@ -2878,6 +2878,198 @@ async fn list_run_stages_shows_retrying_when_failed_will_retry() { assert_eq!(stage_status(&body, "work@1"), "retrying"); } +#[tokio::test] +async fn run_billing_retried_node_then_succeeded_emits_one_row_with_final_attempt_duration() { + 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_event::Event::StageStarted { + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + handler_type: "command".to_string(), + attempt: 1, + max_attempts: 3, + }, + workflow_event::Event::StageFailed { + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + failure: FailureDetail::new("transient", FailureCategory::TransientInfra), + will_retry: true, + duration_ms: 10, + billing: None, + actor: None, + }, + workflow_event::Event::StageRetrying { + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + attempt: 2, + max_attempts: 3, + delay_ms: 0, + }, + workflow_event::Event::StageStarted { + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + handler_type: "command".to_string(), + attempt: 2, + max_attempts: 3, + }, + workflow_event::Event::StageCompleted { + node_id: "work".to_string(), + name: "Work".to_string(), + index: 0, + duration_ms: 25, + status: "succeeded".to_string(), + preferred_label: None, + suggested_next_ids: Vec::new(), + billing: None, + failure: None, + notes: None, + files_touched: Vec::new(), + context_updates: None, + jump_to_node: None, + context_values: None, + node_visits: None, + loop_failure_signatures: None, + restart_failure_signatures: None, + response: None, + attempt: 2, + max_attempts: 3, + }, + ]) + .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, "retry collapses to one row per node_id"); + let row = &stages[0]; + assert_eq!(row["stage"]["id"], "work"); + assert_eq!( + row["state"], "succeeded", + "final state mirrors the latest StageCompleted" + ); + let runtime = row["runtime_secs"].as_f64().unwrap(); + assert!( + (runtime - 0.025).abs() < f64::EPSILON, + "runtime should equal final attempt's 25ms, got {runtime}" + ); +} + +fn revisit_test_started(node_id: &str) -> workflow_event::Event { + workflow_event::Event::StageStarted { + node_id: node_id.to_string(), + name: node_id.to_string(), + index: 0, + handler_type: "command".to_string(), + attempt: 1, + max_attempts: 1, + } +} + +fn revisit_test_completed_with_visit( + node_id: &str, + duration_ms: u64, + visit: usize, +) -> workflow_event::Event { + let mut node_visits = std::collections::BTreeMap::new(); + node_visits.insert(node_id.to_string(), visit); + workflow_event::Event::StageCompleted { + node_id: node_id.to_string(), + name: node_id.to_string(), + index: 0, + duration_ms, + status: "succeeded".to_string(), + preferred_label: None, + suggested_next_ids: Vec::new(), + billing: None, + failure: None, + notes: None, + files_touched: Vec::new(), + context_updates: None, + jump_to_node: None, + context_values: None, + node_visits: Some(node_visits), + loop_failure_signatures: None, + restart_failure_signatures: None, + response: None, + attempt: 1, + max_attempts: 1, + } +} + +#[tokio::test] +async fn run_billing_revisited_node_collapses_to_two_rows_with_summed_visit_duration() { + 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, + // A → B → A loop. Per-visit `node_visits` payload steers the reducer + // to attribute each StageCompleted to the right visit. + revisit_test_started("a"), + revisit_test_completed_with_visit("a", 1, 1), + revisit_test_started("b"), + revisit_test_completed_with_visit("b", 2, 1), + revisit_test_started("a"), + revisit_test_completed_with_visit("a", 99, 2), + ]) + .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(), 2, "two distinct node_ids → two rows"); + assert_eq!( + stages[0]["stage"]["id"], "a", + "A appeared first → A's row first" + ); + assert_eq!(stages[1]["stage"]["id"], "b"); + let a_runtime = stages[0]["runtime_secs"].as_f64().unwrap(); + assert!( + (a_runtime - 0.1).abs() < f64::EPSILON, + "A should sum both visit durations (1ms + 99ms), got {a_runtime}" + ); + let b_runtime = stages[1]["runtime_secs"].as_f64().unwrap(); + assert!( + (b_runtime - 0.002).abs() < f64::EPSILON, + "B should carry its single visit's duration (2ms), got {b_runtime}" + ); +} + async fn append_raw_run_event( state: &Arc, run_id: RunId, diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index af467f557..ae23ff9fd 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -10,7 +10,7 @@ use fabro_types::{ BilledModelUsage, Checkpoint, Conclusion, EventBody, FailureSignature, InterviewQuestionRecord, Outcome, PendingInterviewRecord, PullRequestRecord, RunControlAction, RunEvent, RunId, RunProjection, RunSpec, RunStatus, RunSummary, SandboxRecord, StageCompletion, StageId, - StageOutcome, StageProjection, StartRecord, TerminalStatus, first_event_seq, + StageOutcome, StageProjection, StageState, StartRecord, TerminalStatus, first_event_seq, }; use fabro_util::error::render_with_causes; use serde_json::Value; @@ -290,11 +290,18 @@ impl RunProjectionReducer for RunProjection { let Some(stage_id) = stored.stage_id.as_ref() else { return Ok(()); }; - self.stage_entry( + let stage = self.stage_entry( stage_id.node_id(), stage_id.visit(), first_event_seq(event.seq), ); + stage.begin_attempt(ts); + } + EventBody::StageRetrying(_) => { + let Some(stage) = stage_at_stored_or_current_visit(self, stored, event.seq) else { + return Ok(()); + }; + stage.state = Some(StageState::Retrying); } EventBody::StagePrompt(props) => { let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) @@ -323,22 +330,25 @@ impl RunProjectionReducer for RunProjection { stage.completion = Some(completion); stage.duration_ms = Some(props.duration_ms); stage.usage.clone_from(&props.billing); + stage.state = Some(StageState::from(outcome.status)); } EventBody::StageFailed(props) => { let failure_reason = props.failure.as_ref().map(|detail| detail.message.clone()); let Some(stage) = stage_at_stored_or_current_visit(self, stored, event.seq) else { return Ok(()); }; + let outcome = StageOutcome::Failed { + retry_requested: props.will_retry, + }; stage.completion = Some(StageCompletion { - outcome: StageOutcome::Failed { - retry_requested: props.will_retry, - }, + outcome, notes: None, failure_reason, timestamp: ts, }); stage.duration_ms = Some(props.duration_ms); stage.usage.clone_from(&props.billing); + stage.state = Some(StageState::from(outcome)); } EventBody::AgentSessionStarted(props) => { let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) @@ -670,12 +680,13 @@ mod tests { use fabro_types::run_event::{ CheckpointCompletedProps, InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunControlEffectProps, StageCompletedProps, StageFailedProps, StagePromptProps, - StageStartedProps, + StageRetryingProps, StageStartedProps, }; use fabro_types::{ - BilledModelUsage, BlockedReason, Checkpoint, EventBody, FailureReason, Outcome, - QuestionType, RunBlobId, RunControlAction, RunEvent, RunStatus, StageOutcome, - SuccessReason, TerminalStatus, WorkflowSettings, first_event_seq, fixtures, + BilledModelUsage, BlockedReason, Checkpoint, EventBody, FailureCategory, FailureDetail, + FailureReason, Outcome, QuestionType, RunBlobId, RunControlAction, RunEvent, RunStatus, + StageOutcome, StageState, SuccessReason, TerminalStatus, WorkflowSettings, first_event_seq, + fixtures, }; use serde_json::json; @@ -1775,4 +1786,224 @@ mod tests { ); assert_eq!(state.status_updated_at, updated_at); } + + fn started_props() -> StageStartedProps { + StageStartedProps { + index: 0, + handler_type: "agent".to_string(), + attempt: 1, + max_attempts: 3, + } + } + + fn failed_props(duration_ms: u64, will_retry: bool) -> StageFailedProps { + StageFailedProps { + index: 0, + failure: Some(FailureDetail::new("boom", FailureCategory::TransientInfra)), + will_retry, + duration_ms, + billing: None, + } + } + + fn retrying_props() -> StageRetryingProps { + StageRetryingProps { + index: 0, + attempt: 2, + max_attempts: 3, + delay_ms: 0, + } + } + + fn completed_props(duration_ms: u64, status: StageOutcome) -> StageCompletedProps { + StageCompletedProps { + index: 0, + duration_ms, + status, + preferred_label: None, + suggested_next_ids: Vec::new(), + billing: None, + failure: None, + notes: None, + files_touched: Vec::new(), + context_updates: None, + jump_to_node: None, + context_values: None, + node_visits: None, + loop_failure_signatures: None, + restart_failure_signatures: None, + response: None, + attempt: 1, + max_attempts: 3, + } + } + + fn billed_usage() -> BilledModelUsage { + serde_json::from_value(json!({ + "input": { + "usage": { + "model": { + "provider": "openai", + "model_id": "gpt-test" + }, + "tokens": { + "input_tokens": 10, + "output_tokens": 5, + "reasoning_tokens": 2, + "cache_read_tokens": 3, + "cache_write_tokens": 4 + } + }, + "facts": { "provider": "open_ai" } + }, + "total_usd_micros": 123 + })) + .expect("billing fixture should deserialize") + } + + #[test] + fn stage_started_records_started_at_and_running_state() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("build", 1); + + state + .apply_event(&test_stage_event( + 3, + EventBody::StageStarted(started_props()), + stage_id.clone(), + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.state, Some(StageState::Running)); + assert!(stage.started_at.is_some()); + assert_eq!(stage.effective_state(), StageState::Running); + } + + #[test] + fn stage_completed_records_duration_usage_and_terminal_state() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("build", 1); + let usage = billed_usage(); + + state + .apply_event(&test_stage_event( + 1, + EventBody::StageStarted(started_props()), + stage_id.clone(), + )) + .unwrap(); + let mut props = completed_props(42, StageOutcome::Succeeded); + props.billing = Some(usage.clone()); + state + .apply_event(&test_event( + 2, + EventBody::StageCompleted(props), + Some("build"), + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.duration_ms, Some(42)); + assert_eq!(stage.usage.as_ref(), Some(&usage)); + assert_eq!(stage.state, Some(StageState::Succeeded)); + assert_eq!(stage.effective_state(), StageState::Succeeded); + } + + #[test] + fn stage_failed_records_duration_and_failed_state() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("build", 1); + + state + .apply_event(&test_stage_event( + 1, + EventBody::StageStarted(started_props()), + stage_id.clone(), + )) + .unwrap(); + state + .apply_event(&test_event( + 2, + EventBody::StageFailed(failed_props(10, false)), + Some("build"), + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.duration_ms, Some(10)); + assert_eq!(stage.state, Some(StageState::Failed)); + } + + #[test] + fn stage_retrying_sets_retrying_state() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("build", 1); + + state + .apply_event(&test_stage_event( + 1, + EventBody::StageStarted(started_props()), + stage_id.clone(), + )) + .unwrap(); + state + .apply_event(&test_event( + 2, + EventBody::StageFailed(failed_props(10, true)), + Some("build"), + )) + .unwrap(); + state + .apply_event(&test_event( + 3, + EventBody::StageRetrying(retrying_props()), + Some("build"), + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.state, Some(StageState::Retrying)); + } + + #[test] + fn stage_started_after_retrying_returns_to_running_and_resets_attempt_data() { + let mut state = RunProjection::default(); + let stage_id = StageId::new("build", 1); + + state + .apply_event(&test_stage_event( + 1, + EventBody::StageStarted(started_props()), + stage_id.clone(), + )) + .unwrap(); + state + .apply_event(&test_event( + 2, + EventBody::StageFailed(failed_props(10, true)), + Some("build"), + )) + .unwrap(); + state + .apply_event(&test_event( + 3, + EventBody::StageRetrying(retrying_props()), + Some("build"), + )) + .unwrap(); + state + .apply_event(&test_stage_event( + 4, + EventBody::StageStarted(started_props()), + stage_id.clone(), + )) + .unwrap(); + + let stage = state.stage(&stage_id).unwrap(); + assert_eq!(stage.state, Some(StageState::Running)); + // Prior attempt's terminal data must not leak into the new attempt. + assert!(stage.completion.is_none()); + assert_eq!(stage.duration_ms, None); + } } diff --git a/lib/crates/fabro-store/tests/serializable_projection.rs b/lib/crates/fabro-store/tests/serializable_projection.rs index 2b24cd54e..8809acda2 100644 --- a/lib/crates/fabro-store/tests/serializable_projection.rs +++ b/lib/crates/fabro-store/tests/serializable_projection.rs @@ -123,6 +123,10 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { let serialized = serde_json::to_value(SerializableProjection(&projection)) .expect("projection should serialize"); + assert!( + serialized["stages"]["build@2"].get("usage").is_none(), + "stage usage is server-internal and should not be serialized" + ); let round_tripped: RunProjection = serde_json::from_value(serialized).expect("serialized projection should deserialize"); let node = round_tripped.stage(&stage_id).expect("node should remain"); @@ -163,7 +167,7 @@ fn serializable_projection_round_trips_and_trims_bulky_node_fields() { Some(json!([{ "stage": "fanout@1" }])) ); assert_eq!(node.duration_ms, Some(1234)); - assert_eq!(node.usage, Some(sample_usage())); + assert_eq!(node.usage, None); } #[test] diff --git a/lib/crates/fabro-types/src/run_projection.rs b/lib/crates/fabro-types/src/run_projection.rs index 87321d348..90da9b59b 100644 --- a/lib/crates/fabro-types/src/run_projection.rs +++ b/lib/crates/fabro-types/src/run_projection.rs @@ -6,7 +6,7 @@ use chrono::{DateTime, Utc}; use crate::{ BilledModelUsage, Checkpoint, Conclusion, InterviewQuestionRecord, InvalidTransition, PullRequestRecord, Retro, RunControlAction, RunId, RunSpec, RunStatus, SandboxRecord, - StageCompletion, StageId, StartRecord, + StageCompletion, StageId, StageState, StartRecord, }; #[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] @@ -44,10 +44,6 @@ pub struct StageProjection { pub prompt: Option, pub response: Option, pub completion: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub duration_ms: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub usage: Option, pub provider_used: Option, pub diff: Option, pub script_invocation: Option, @@ -65,6 +61,17 @@ pub struct StageProjection { pub live_streaming: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub termination: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub started_at: Option>, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub duration_ms: Option, + /// Server-internal billing usage for the latest attempt; not part of the + /// wire contract because `BilledModelUsage` is not modeled in OpenAPI. + /// Read only in-process by the billing handler. + #[serde(skip)] + pub usage: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub state: Option, } /// Convert a 1-based event sequence number into the `NonZeroU32` form used for @@ -96,8 +103,55 @@ impl StageProjection { streams_separated: None, live_streaming: None, termination: None, + started_at: None, + state: None, } } + + /// Effective lifecycle state derived from stored event data. + /// + /// Falls back to deriving from `completion` for projections that predate + /// the stored `state` field, so old serialized projections still work + /// without a backfill. + #[must_use] + pub fn effective_state(&self) -> StageState { + self.state.unwrap_or_else(|| match &self.completion { + Some(completion) => StageState::from(completion.outcome), + None => StageState::Running, + }) + } + + /// Live wall-clock runtime in seconds. + /// + /// While the stage is non-terminal (`Pending`, `Running`, or `Retrying`), + /// this returns the elapsed time since `started_at` so the UI can tick + /// client-side. Once terminal, the stored `duration_ms` is returned. This + /// also handles retries safely: a new `StageStarted` resets the state + /// back to `Running` and keeps the live computation correct even if a + /// previous attempt left a stale `duration_ms`. + #[must_use] + pub fn runtime_secs(&self, now: DateTime) -> Option { + let state = self.effective_state(); + if matches!( + state, + StageState::Running | StageState::Retrying | StageState::Pending + ) { + return self.started_at.map(|started| { + now.signed_duration_since(started).num_milliseconds().max(0) as f64 / 1000.0 + }); + } + self.duration_ms.map(|ms| ms as f64 / 1000.0) + } + + /// Begin a new attempt (or visit) for this stage: clear every + /// per-attempt field so prior-attempt data does not leak, then record + /// `started_at` and `state = Running`. Preserves `first_event_seq` + /// (identity / sort key). + pub fn begin_attempt(&mut self, started_at: DateTime) { + *self = Self::new(self.first_event_seq); + self.started_at = Some(started_at); + self.state = Some(StageState::Running); + } } impl RunProjection { diff --git a/lib/packages/fabro-api-client/src/models/board-column-definition.ts b/lib/packages/fabro-api-client/src/models/board-column-definition.ts index 314cbcde1..20deba1f5 100644 --- a/lib/packages/fabro-api-client/src/models/board-column-definition.ts +++ b/lib/packages/fabro-api-client/src/models/board-column-definition.ts @@ -21,6 +21,3 @@ export interface BoardColumnDefinition { 'id': BoardColumn; 'name': string; } - - - diff --git a/lib/packages/fabro-api-client/src/models/run-billing-stage.ts b/lib/packages/fabro-api-client/src/models/run-billing-stage.ts index 5d376b7cc..acdf57c97 100644 --- a/lib/packages/fabro-api-client/src/models/run-billing-stage.ts +++ b/lib/packages/fabro-api-client/src/models/run-billing-stage.ts @@ -22,6 +22,9 @@ import type { BillingStageRef } from './billing-stage-ref'; // May contain unused imports in some cases // @ts-ignore import type { ModelReference } from './model-reference'; +// May contain unused imports in some cases +// @ts-ignore +import type { StageState } from './stage-state'; /** * Token counts and billed totals for one workflow node within a run. Rows are grouped by node; billing and runtime sum every visit of that node. @@ -34,5 +37,9 @@ export interface RunBillingStage { * Wall-clock runtime in seconds, summed across every visit of this node. */ 'runtime_secs': number; + /** + * Wall-clock time the latest attempt of this stage started, if known. + */ + 'started_at'?: string | null; + 'state'?: StageState | null; } - diff --git a/lib/packages/fabro-api-client/src/models/run-stage.ts b/lib/packages/fabro-api-client/src/models/run-stage.ts index ad01b9791..060a6166e 100644 --- a/lib/packages/fabro-api-client/src/models/run-stage.ts +++ b/lib/packages/fabro-api-client/src/models/run-stage.ts @@ -42,7 +42,8 @@ export interface RunStage { * 1-based visit count; bumped each time the workflow re-enters this node. */ 'visit': number; + /** + * Wall-clock time the latest attempt of this stage started, if known. + */ + 'started_at'?: string | null; } - - - 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 41321ec78..139fdd437 100644 --- a/lib/packages/fabro-api-client/src/models/stage-projection.ts +++ b/lib/packages/fabro-api-client/src/models/stage-projection.ts @@ -19,6 +19,9 @@ import type { CommandTermination } from './command-termination'; // May contain unused imports in some cases // @ts-ignore import type { StageCompletion } from './stage-completion'; +// May contain unused imports in some cases +// @ts-ignore +import type { StageState } from './stage-state'; /** * Observable projection data for one workflow stage execution. @@ -52,7 +55,13 @@ export interface StageProjection { 'streams_separated'?: boolean | null; 'live_streaming'?: boolean | null; 'termination'?: CommandTermination | null; + /** + * Wall-clock time the latest attempt of this stage started, if known. + */ + 'started_at'?: string | null; + /** + * Wall-clock duration of the stage\'s latest terminal attempt, if known. + */ + 'duration_ms'?: number | null; + 'state'?: StageState | null; } - - -