Read billing and stages from RunProjection with live runtimes (#213)

### Summary
Billing and stage lists now use the event-sourced `RunProjection` as
their source of truth, so running and retrying stages appear immediately
and runtimes keep advancing in the UI. This removes the checkpoint
completed-node bypass that hid in-flight work and froze totals until the
next server response.

### Plan Summary
- Store stage `started_at`, terminal `duration_ms`, server-internal
`usage`, and lifecycle `state` on `StageProjection`.
- Populate those fields from stage lifecycle events, including retry
transitions and per-attempt reset on new starts.
- Render `/runs/{id}/stages` and `/runs/{id}/billing` from
`RunProjection.iter_stages()`.
- Expose the new API/client fields and tick in-flight billing runtimes
on the web UI.

```mermaid
flowchart TB
  Events["Stage lifecycle events"] --> Projection["RunProjection StageProjection"]
  Projection --> StagesAPI["GET /runs/{id}/stages"]
  Projection --> BillingAPI["GET /runs/{id}/billing"]
  StagesAPI --> StageUI["Stage sidebar/stages view"]
  BillingAPI --> BillingUI["Billing tab live totals"]
```

### Key decisions
Retry and revisit handling stays one row per node id: latest visit data
wins, while first-seen event sequence keeps ordering stable with
finalize output. `state` is stored rather than derived so `Retrying` is
representable, and old serialized projections still work through the
`effective_state()` fallback. Billing `usage` remains server-internal
and is skipped on the wire; public schemas only expose the fields needed
by `/stages`, `/billing`, and the frontend live timer.

Added focused reducer, server retry/revisit, API round-trip, billing UI,
and event invalidation coverage.

⚒️ Generated with [Fabro](https://fabro.sh)

---------

Co-authored-by: Fabro <noreply@fabro.sh>
Co-authored-by: Bryan Helmkamp <bryan@brynary.com>
This commit is contained in:
fabro-sh-0530[bot] 2026-05-05 09:32:33 -04:00 • committed by GitHub
parent e901cd3a81
commit 786a6f67e1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
23 changed files with 954 additions and 213 deletions

View file

@ -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<Map<string, number>>(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;

View file

@ -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([]);
});
});
});

View file

@ -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"),
]);
});
});

View file

@ -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"),

View file

@ -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<StageState> = new Set(["running", "retrying"]);
export const ACTIVE_STAGE_STATES: ReadonlySet<StageState> = new Set([
StageState.RUNNING,
StageState.RETRYING,
]);
export const IN_FLIGHT_STAGE_STATES: ReadonlySet<StageState> = new Set([
StageState.PENDING,
StageState.RUNNING,
StageState.RETRYING,
]);
export const SUCCEEDED_STAGE_STATES: ReadonlySet<StageState> = new Set([
"succeeded",
"partially_succeeded",
StageState.SUCCEEDED,
StageState.PARTIALLY_SUCCEEDED,
]);
/**

View file

@ -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);
}

View file

@ -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;
}
});
});

View file

@ -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<MappedStageRow[]>(() => {
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<MappedStageRow[]>(() => {
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 (
<div className="py-12">
<EmptyState
title="No completed stages yet"
description="Stages will appear once the run produces completed nodes."
title="No stages yet"
description="Stages will appear as soon as the run starts executing."
/>
</div>
);
@ -105,7 +134,7 @@ export default function RunBilling({ params }: { params: { id: string } }) {
</tr>
</thead>
<tbody>
{stages.map((row) => (
{rows.map((row) => (
<tr key={row.stage} className="border-b border-line last:border-b-0">
<td className="px-4 py-3 text-fg-2">{row.stage}</td>
<td className="px-4 py-3 font-mono text-xs text-fg-3">
@ -115,7 +144,9 @@ export default function RunBilling({ params }: { params: { id: string } }) {
{formatTokens(row.inputTokens)} <span className="text-fg-muted">/</span>{" "}
{formatTokens(row.outputTokens)}
</td>
<td className="px-4 py-3 text-right font-mono text-xs text-fg-3">{row.runtime}</td>
<td className="px-4 py-3 text-right font-mono text-xs text-fg-3">
{formatDurationSecs(row.runtimeSecs)}
</td>
<td className="px-4 py-3 text-right font-mono text-xs text-fg-3">
{formatUsdMicros(row.totalUsdMicros)}
</td>
@ -131,7 +162,7 @@ export default function RunBilling({ params }: { params: { id: string } }) {
{formatTokens(totalOutput)}
</td>
<td className="px-4 py-3 text-right font-mono text-xs font-medium text-fg">
{totalRuntime}
{formatDurationSecs(totalRuntimeSecs)}
</td>
<td className="px-4 py-3 text-right font-mono text-xs font-medium text-fg">
{formatUsdMicros(totalUsdMicros)}

View file

@ -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<number | null>(() =>
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;
}

View file

@ -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.

View file

@ -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);
}

View file

@ -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();

View file

@ -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 {

View file

@ -570,6 +570,7 @@ pub(crate) fn run_stage_from_stage_id(
name: impl Into<String>,
status: StageState,
duration_secs: Option<f64>,
started_at: Option<chrono::DateTime<chrono::Utc>>,
) -> 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,
}
}

View file

@ -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<Arc<AppState>> {
@ -16,41 +17,6 @@ pub(super) fn routes() -> Router<Arc<AppState>> {
.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<StageState> {
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<StageId, StageState> {
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<Arc<AppState>>,
@ -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::<Vec<_>>();
(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::<Vec<_>>();
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::<HashMap<_, _>>();
let live_rows = live_billing_rows(&projection, Utc::now());
let runtime_secs = live_rows.iter().map(|row| row.runtime_secs).sum::<f64>();
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::<Vec<_>>();
@ -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<DateTime<Utc>>,
state: Option<StageState>,
latest_visit: u32,
}
fn live_billing_rows(projection: &RunProjection, now: DateTime<Utc>) -> Vec<LiveBillingRow> {
let mut row_indices = HashMap::<String, usize>::new();
let mut rows = Vec::<LiveBillingRow>::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<Utc>) -> Option<f64> {
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"))
}

View file

@ -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<AppState>,
run_id: RunId,

View file

@ -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);
}
}

View file

@ -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]

View file

@ -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<String>,
pub response: Option<String>,
pub completion: Option<StageCompletion>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub duration_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub usage: Option<BilledModelUsage>,
pub provider_used: Option<serde_json::Value>,
pub diff: Option<String>,
pub script_invocation: Option<serde_json::Value>,
@ -65,6 +61,17 @@ pub struct StageProjection {
pub live_streaming: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub termination: Option<crate::CommandTermination>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub started_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub duration_ms: Option<u64>,
/// 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<BilledModelUsage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub state: Option<StageState>,
}
/// 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<Utc>) -> Option<f64> {
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<Utc>) {
*self = Self::new(self.first_event_seq);
self.started_at = Some(started_at);
self.state = Some(StageState::Running);
}
}
impl RunProjection {

View file

@ -21,6 +21,3 @@ export interface BoardColumnDefinition {
'id': BoardColumn;
'name': string;
}

View file

@ -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;
}

View file

@ -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;
}

View file

@ -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;
}