summaryrefslogtreecommitdiffhomepage
path: root/src/core/metrics
diff options
context:
space:
mode:
Diffstat (limited to 'src/core/metrics')
-rw-r--r--src/core/metrics/format.test.ts375
-rw-r--r--src/core/metrics/format.ts179
-rw-r--r--src/core/metrics/index.ts31
-rw-r--r--src/core/metrics/place.test.ts621
-rw-r--r--src/core/metrics/place.ts298
-rw-r--r--src/core/metrics/reducer.test.ts689
-rw-r--r--src/core/metrics/reducer.ts345
-rw-r--r--src/core/metrics/types.ts97
8 files changed, 2635 insertions, 0 deletions
diff --git a/src/core/metrics/format.test.ts b/src/core/metrics/format.test.ts
new file mode 100644
index 0000000..97170d0
--- /dev/null
+++ b/src/core/metrics/format.test.ts
@@ -0,0 +1,375 @@
+import type { StepId, StepMetrics, TurnMetrics } from "@dispatch/wire";
+import { describe, expect, it } from "vitest";
+import {
+ computeCachePct,
+ computeContextUsage,
+ computeExpectedCachePct,
+ computeTps,
+ formatCompactTokens,
+ formatContextSize,
+ viewCacheRate,
+ viewExpectedCache,
+ viewStepMetrics,
+ viewTurnMetrics,
+} from "./format";
+
+describe("computeTps", () => {
+ it("null when elapsed missing", () => {
+ expect(computeTps(100, undefined)).toBeNull();
+ });
+
+ it("null when elapsed is zero", () => {
+ expect(computeTps(100, 0)).toBeNull();
+ });
+
+ it("null when elapsed is negative", () => {
+ expect(computeTps(100, -100)).toBeNull();
+ });
+
+ it("computes tokens per second", () => {
+ expect(computeTps(1000, 2000)).toBe(500);
+ });
+
+ it("computes fractional tps", () => {
+ expect(computeTps(100, 3000)).toBeCloseTo(33.33, 1);
+ });
+});
+
+describe("viewStepMetrics", () => {
+ it("formats tokens with thousands separator, tps, and durations", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 1234, outputTokens: 567 },
+ ttftMs: 820,
+ decodeMs: 1200,
+ genTotalMs: 2020,
+ };
+ const view = viewStepMetrics(step, 0);
+ expect(view.label).toBe("step 1");
+ expect(view.tokensLabel).toBe("1,801 tok");
+ expect(view.tps).toBe("473 tok/s");
+ expect(view.ttft).toBe("820ms");
+ expect(view.decode).toBe("1.2s");
+ expect(view.genTotal).toBe("2.0s");
+ });
+
+ it("handles missing timing fields", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50 },
+ };
+ const view = viewStepMetrics(step, 0);
+ expect(view.tps).toBeNull();
+ expect(view.ttft).toBeNull();
+ expect(view.decode).toBeNull();
+ expect(view.genTotal).toBeNull();
+ });
+
+ it("formats duration < 1s as ms", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 10, outputTokens: 5 },
+ ttftMs: 42,
+ };
+ const view = viewStepMetrics(step, 0);
+ expect(view.ttft).toBe("42ms");
+ });
+
+ it("formats duration >= 1s as seconds", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 10, outputTokens: 5 },
+ genTotalMs: 3200,
+ };
+ const view = viewStepMetrics(step, 0);
+ expect(view.genTotal).toBe("3.2s");
+ });
+
+ it("uses step index for label", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 10, outputTokens: 5 },
+ };
+ expect(viewStepMetrics(step, 2).label).toBe("step 3");
+ });
+
+ it("tps uses decodeMs (not genTotalMs)", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50 },
+ decodeMs: 500,
+ genTotalMs: 800,
+ };
+ const view = viewStepMetrics(step, 0);
+ // 50 / (500/1000) = 100 tok/s, NOT 50/(800/1000)=62.5
+ expect(view.tps).toBe("100 tok/s");
+ });
+
+ it("tps falls back to genTotalMs when decodeMs absent", () => {
+ const step: StepMetrics = {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50 },
+ genTotalMs: 800,
+ };
+ const view = viewStepMetrics(step, 0);
+ // 50 / (800/1000) = 62.5 → rounds to 63
+ expect(view.tps).toBe("63 tok/s");
+ });
+});
+
+describe("viewTurnMetrics", () => {
+ it("formats total tokens and breakdown", () => {
+ const turn: TurnMetrics = {
+ turnId: "t1",
+ usage: { inputTokens: 1000, outputTokens: 234 },
+ durationMs: 5000,
+ steps: [
+ {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 1000, outputTokens: 234 },
+ decodeMs: 3000,
+ genTotalMs: 4000,
+ },
+ ],
+ };
+ const view = viewTurnMetrics(turn);
+ expect(view.tokensLabel).toBe("1,234 tok");
+ expect(view.breakdown).toBe("1,000 in / 234 out");
+ expect(view.tps).toBe("78 tok/s");
+ expect(view.duration).toBe("5.0s");
+ });
+
+ it("breakdown includes cache only when present", () => {
+ const turn: TurnMetrics = {
+ turnId: "t1",
+ usage: { inputTokens: 1000, outputTokens: 234, cacheReadTokens: 500 },
+ steps: [],
+ };
+ const view = viewTurnMetrics(turn);
+ expect(view.breakdown).toBe("1,000 in / 234 out / 500 cache");
+ });
+
+ it("breakdown omits cache when not present", () => {
+ const turn: TurnMetrics = {
+ turnId: "t1",
+ usage: { inputTokens: 100, outputTokens: 50 },
+ steps: [],
+ };
+ const view = viewTurnMetrics(turn);
+ expect(view.breakdown).toBe("100 in / 50 out");
+ });
+
+ it("tps is null when no step has decodeMs or genTotalMs", () => {
+ const turn: TurnMetrics = {
+ turnId: "t1",
+ usage: { inputTokens: 100, outputTokens: 50 },
+ steps: [
+ {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50 },
+ },
+ ],
+ };
+ const view = viewTurnMetrics(turn);
+ expect(view.tps).toBeNull();
+ });
+
+ it("duration is null when durationMs absent", () => {
+ const turn: TurnMetrics = {
+ turnId: "t1",
+ usage: { inputTokens: 100, outputTokens: 50 },
+ steps: [],
+ };
+ const view = viewTurnMetrics(turn);
+ expect(view.duration).toBeNull();
+ });
+
+ it("sums decodeMs across steps (fallback genTotalMs per step) for tps", () => {
+ const turn: TurnMetrics = {
+ turnId: "t1",
+ usage: { inputTokens: 300, outputTokens: 150 },
+ steps: [
+ {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50 },
+ decodeMs: 800,
+ genTotalMs: 1000,
+ },
+ {
+ stepId: "s2" as StepId,
+ usage: { inputTokens: 200, outputTokens: 100 },
+ genTotalMs: 2000,
+ },
+ ],
+ };
+ const view = viewTurnMetrics(turn);
+ // step1 uses decodeMs=800, step2 falls back to genTotalMs=2000 → total=2800ms
+ // 150 / (2800/1000) = 53.57 → rounds to 54
+ expect(view.tps).toBe("54 tok/s");
+ });
+});
+
+describe("computeCachePct", () => {
+ it("is cacheReadTokens / inputTokens as a rounded percentage", () => {
+ expect(computeCachePct({ inputTokens: 2737, outputTokens: 10, cacheReadTokens: 2560 })).toBe(
+ 94,
+ );
+ expect(computeCachePct({ inputTokens: 2669, outputTokens: 10, cacheReadTokens: 384 })).toBe(14);
+ });
+
+ it("is 0 when cacheReadTokens absent (legitimate miss, not missing data)", () => {
+ expect(computeCachePct({ inputTokens: 1000, outputTokens: 50 })).toBe(0);
+ });
+
+ it("is 0 when there are no input tokens (guard divide-by-zero)", () => {
+ expect(computeCachePct({ inputTokens: 0, outputTokens: 0, cacheReadTokens: 5 })).toBe(0);
+ });
+
+ it("clamps to 100 if read somehow exceeds input", () => {
+ expect(computeCachePct({ inputTokens: 100, outputTokens: 0, cacheReadTokens: 250 })).toBe(100);
+ });
+});
+
+describe("viewCacheRate", () => {
+ it("success level for a high hit rate (>= 66)", () => {
+ const v = viewCacheRate({ inputTokens: 100, outputTokens: 0, cacheReadTokens: 93 });
+ expect(v.pct).toBe(93);
+ expect(v.level).toBe("success");
+ expect(v.isHit).toBe(true);
+ });
+
+ it("warning level for a mid hit rate (33..65)", () => {
+ const v = viewCacheRate({ inputTokens: 100, outputTokens: 0, cacheReadTokens: 54 });
+ expect(v.pct).toBe(54);
+ expect(v.level).toBe("warning");
+ });
+
+ it("error level for a low hit rate (< 33), including a legitimate 0%", () => {
+ expect(viewCacheRate({ inputTokens: 100, outputTokens: 0, cacheReadTokens: 14 }).level).toBe(
+ "error",
+ );
+ const miss = viewCacheRate({ inputTokens: 1000, outputTokens: 50 });
+ expect(miss.pct).toBe(0);
+ expect(miss.level).toBe("error");
+ expect(miss.isHit).toBe(false);
+ });
+});
+
+describe("computeExpectedCachePct", () => {
+ it("null when there is no prior turn (first turn has no baseline)", () => {
+ expect(computeExpectedCachePct({ inputTokens: 100, outputTokens: 0 }, null)).toBeNull();
+ });
+
+ it("null when the prior turn cached nothing (denominator 0)", () => {
+ const prev = { inputTokens: 100, outputTokens: 0 };
+ const current = { inputTokens: 200, outputTokens: 0, cacheReadTokens: 50 };
+ expect(computeExpectedCachePct(current, prev)).toBeNull();
+ });
+
+ it("100% when the whole prior cached prefix was read back (backend worked example)", () => {
+ // turn 1: cacheRead 0, cacheWrite 5146 → prefix 5146; turn 2 reads 5146 back.
+ const prev = { inputTokens: 5149, outputTokens: 0, cacheReadTokens: 0, cacheWriteTokens: 5146 };
+ const current = {
+ inputTokens: 8462,
+ outputTokens: 0,
+ cacheReadTokens: 5146,
+ cacheWriteTokens: 3313,
+ };
+ expect(computeExpectedCachePct(current, prev)).toBe(100);
+ });
+
+ it("drops below 100% when the cache busted (read < prior prefix)", () => {
+ const prev = {
+ inputTokens: 1000,
+ outputTokens: 0,
+ cacheReadTokens: 100,
+ cacheWriteTokens: 900,
+ };
+ const current = { inputTokens: 1000, outputTokens: 0, cacheReadTokens: 500 };
+ // 500 / (100 + 900) = 50%
+ expect(computeExpectedCachePct(current, prev)).toBe(50);
+ });
+
+ it("clamps to 100 if read somehow exceeds the prior prefix", () => {
+ const prev = { inputTokens: 100, outputTokens: 0, cacheWriteTokens: 100 };
+ const current = { inputTokens: 100, outputTokens: 0, cacheReadTokens: 250 };
+ expect(computeExpectedCachePct(current, prev)).toBe(100);
+ });
+});
+
+describe("viewExpectedCache", () => {
+ it("null view when it cannot be derived (no prior turn)", () => {
+ expect(viewExpectedCache({ inputTokens: 100, outputTokens: 0 }, null)).toBeNull();
+ });
+
+ it("success level + hit flag for full retention", () => {
+ const prev = { inputTokens: 5149, outputTokens: 0, cacheWriteTokens: 5146 };
+ const current = { inputTokens: 8462, outputTokens: 0, cacheReadTokens: 5146 };
+ const v = viewExpectedCache(current, prev);
+ expect(v?.pct).toBe(100);
+ expect(v?.level).toBe("success");
+ expect(v?.isHit).toBe(true);
+ });
+});
+
+describe("formatContextSize", () => {
+ it("formats a defined count with thousands separators", () => {
+ expect(formatContextSize(34102)).toBe("34,102 tokens in context");
+ });
+
+ it("renders a placeholder for undefined (never 0)", () => {
+ expect(formatContextSize(undefined)).toBe("context size unknown");
+ });
+
+ it("renders an explicit 0 as zero tokens (a real reported value)", () => {
+ expect(formatContextSize(0)).toBe("0 tokens in context");
+ });
+});
+
+describe("formatCompactTokens", () => {
+ it("renders sub-1k counts as-is", () => {
+ expect(formatCompactTokens(0)).toBe("0");
+ expect(formatCompactTokens(812)).toBe("812");
+ });
+
+ it("renders thousands with one decimal (rounded ≥100k)", () => {
+ expect(formatCompactTokens(12300)).toBe("12.3k");
+ expect(formatCompactTokens(150000)).toBe("150k");
+ });
+
+ it("renders millions with one decimal", () => {
+ expect(formatCompactTokens(1_200_000)).toBe("1.2M");
+ expect(formatCompactTokens(1_000_000)).toBe("1.0M");
+ });
+});
+
+describe("computeContextUsage", () => {
+ it("computes an unrounded clamped percent against the limit", () => {
+ const u = computeContextUsage(34102, 1_000_000);
+ expect(u.current).toBe(34102);
+ expect(u.max).toBe(1_000_000);
+ expect(u.percent).toBeCloseTo(3.4102, 4);
+ });
+
+ it("treats unknown contextSize as current null (never 0)", () => {
+ const u = computeContextUsage(undefined, 1_000_000);
+ expect(u.current).toBeNull();
+ expect(u.percent).toBeNull();
+ });
+
+ it("an explicit 0 context size is a real reported value (current 0)", () => {
+ const u = computeContextUsage(0, 1_000_000);
+ expect(u.current).toBe(0);
+ expect(u.percent).toBe(0);
+ });
+
+ it("clamps percent to [0,100] and over-limit reads 100", () => {
+ expect(computeContextUsage(2_000_000, 1_000_000).percent).toBe(100);
+ });
+
+ it("max null (no/zero limit) ⇒ percent null", () => {
+ expect(computeContextUsage(5000, null).percent).toBeNull();
+ expect(computeContextUsage(5000, 0).percent).toBeNull();
+ expect(computeContextUsage(5000, null).max).toBeNull();
+ });
+});
diff --git a/src/core/metrics/format.ts b/src/core/metrics/format.ts
new file mode 100644
index 0000000..894bd54
--- /dev/null
+++ b/src/core/metrics/format.ts
@@ -0,0 +1,179 @@
+import type { StepMetrics, TurnMetrics, Usage } from "@dispatch/wire";
+import type { CacheRateView, StepMetricsView, TurnMetricsView } from "./types";
+
+function formatTokens(n: number): string {
+ return n.toLocaleString("en-US");
+}
+
+function formatDuration(ms: number | undefined): string | null {
+ if (ms === undefined || ms <= 0) return null;
+ if (ms < 1000) return `${Math.round(ms)}ms`;
+ return `${(ms / 1000).toFixed(1)}s`;
+}
+
+function formatTps(tps: number | null): string | null {
+ if (tps === null) return null;
+ if (tps < 10) return `${tps.toFixed(1)} tok/s`;
+ return `${Math.round(tps)} tok/s`;
+}
+
+/**
+ * Format the current context size for display. A defined count renders as
+ * `"<n> tokens in context"` (thousands-separated); `undefined` ("unknown" — no
+ * per-step usage reported yet) renders the placeholder `"context size unknown"`.
+ * Never renders `0` for the unknown case.
+ */
+export function formatContextSize(n: number | undefined): string {
+ if (n === undefined) return "context size unknown";
+ return `${formatTokens(n)} tokens in context`;
+}
+
+/**
+ * Compact token count for a slim status bar: `812`, `12.3k`, `1.2M`. Full
+ * thousands-separated numbers live elsewhere; this trades precision for width.
+ */
+export function formatCompactTokens(n: number): string {
+ if (n < 1000) return `${n}`;
+ if (n < 1_000_000) {
+ const k = n / 1000;
+ return `${k >= 100 ? Math.round(k) : k.toFixed(1)}k`;
+ }
+ const m = n / 1_000_000;
+ return `${m >= 100 ? Math.round(m) : m.toFixed(1)}M`;
+}
+
+/**
+ * Context-window occupancy: the current size against a max window limit.
+ *
+ * `current` is the latest turn's context size, or `null` when unknown (no
+ * per-step usage reported yet) — NEVER coerced to `0`, so a consumer cannot
+ * silently render "0 tokens / 1M"; it must branch on `current === null` and show
+ * a placeholder instead. `max` is the model's window limit (or `null` when
+ * unknown). `percent` is `current / max * 100` clamped to [0, 100], UNROUNDED
+ * (the UI picks the precision) — so a few-thousand-token context against a
+ * 1,000,000 window still reads non-zero. `percent` is `null` when `current` OR
+ * `max` is unknown (no bar/denominator).
+ */
+export interface ContextUsage {
+ readonly current: number | null;
+ readonly max: number | null;
+ readonly percent: number | null;
+}
+
+export function computeContextUsage(
+ contextSize: number | undefined,
+ contextLimit: number | null | undefined,
+): ContextUsage {
+ const current = contextSize ?? null;
+ const max = typeof contextLimit === "number" && contextLimit > 0 ? contextLimit : null;
+ const percent =
+ current === null || max === null ? null : Math.max(0, Math.min(100, (current / max) * 100));
+ return { current, max, percent };
+}
+
+/** Compute tokens-per-second. Returns null when elapsed time is absent or zero. */
+export function computeTps(outputTokens: number, elapsedMs: number | undefined): number | null {
+ if (elapsedMs === undefined || elapsedMs <= 0) return null;
+ return outputTokens / (elapsedMs / 1000);
+}
+
+function totalTokens(u: Usage): number {
+ return u.inputTokens + u.outputTokens;
+}
+
+function formatBreakdown(u: Usage): string {
+ let s = `${formatTokens(u.inputTokens)} in / ${formatTokens(u.outputTokens)} out`;
+ if (u.cacheReadTokens !== undefined && u.cacheReadTokens > 0) {
+ s += ` / ${formatTokens(u.cacheReadTokens)} cache`;
+ }
+ return s;
+}
+
+/** Build a formatted view of a single step's metrics. */
+export function viewStepMetrics(step: StepMetrics, index: number): StepMetricsView {
+ const total = totalTokens(step.usage);
+ const tps = computeTps(step.usage.outputTokens, step.decodeMs ?? step.genTotalMs);
+ return {
+ label: `step ${index + 1}`,
+ tokensLabel: `${formatTokens(total)} tok`,
+ tps: formatTps(tps),
+ ttft: formatDuration(step.ttftMs),
+ decode: formatDuration(step.decodeMs),
+ genTotal: formatDuration(step.genTotalMs),
+ };
+}
+
+/**
+ * Cache hit rate as a 0..100 integer percentage: `cacheReadTokens / inputTokens`,
+ * clamped to [0,1]. Absent cache field counts as 0; a 0% rate is legitimate (not
+ * missing data). Returns 0 when there are no input tokens.
+ */
+export function computeCachePct(u: Usage): number {
+ const read = u.cacheReadTokens ?? 0;
+ if (u.inputTokens <= 0) return 0;
+ const rate = read / u.inputTokens;
+ const clamped = rate < 0 ? 0 : rate > 1 ? 1 : rate;
+ return Math.round(clamped * 100);
+}
+
+/** Colour severity for a cache hit percentage (badge colour). */
+function cacheLevel(pct: number): "success" | "warning" | "error" {
+ if (pct >= 66) return "success";
+ if (pct >= 33) return "warning";
+ return "error";
+}
+
+/** Build a view of a cache hit rate (percentage + colour level + hit flag). */
+export function viewCacheRate(u: Usage): CacheRateView {
+ const pct = computeCachePct(u);
+ return { pct, level: cacheLevel(pct), isHit: (u.cacheReadTokens ?? 0) > 0 };
+}
+
+/**
+ * Expected cache (retention): of the cache that existed going INTO this turn, how
+ * much was read back — `clamp01(cacheRead_N / (cacheRead_{N-1} + cacheWrite_{N-1}))`.
+ * The denominator is the PRIOR turn's cached prefix (what it read + what it wrote).
+ * Ideally ~100% on every turn after the first; <100% = the cache busted/expired.
+ *
+ * Returns `null` when it cannot be derived: no prior turn (`prev === null`) or the
+ * prior turn cached nothing (denominator <= 0) — distinct from a real 0%.
+ */
+export function computeExpectedCachePct(current: Usage, prev: Usage | null): number | null {
+ if (prev === null) return null;
+ const denom = (prev.cacheReadTokens ?? 0) + (prev.cacheWriteTokens ?? 0);
+ if (denom <= 0) return null;
+ const read = current.cacheReadTokens ?? 0;
+ const rate = read / denom;
+ const clamped = rate < 0 ? 0 : rate > 1 ? 1 : rate;
+ return Math.round(clamped * 100);
+}
+
+/**
+ * Build a view of the cross-turn retention (percentage + colour level + hit flag),
+ * or `null` when it can't be derived (see `computeExpectedCachePct`).
+ */
+export function viewExpectedCache(current: Usage, prev: Usage | null): CacheRateView | null {
+ const pct = computeExpectedCachePct(current, prev);
+ if (pct === null) return null;
+ return { pct, level: cacheLevel(pct), isHit: (current.cacheReadTokens ?? 0) > 0 };
+}
+
+/** Build a formatted view of a turn's aggregate metrics. */
+export function viewTurnMetrics(turn: TurnMetrics, turnNumber?: number): TurnMetricsView {
+ const total = totalTokens(turn.usage);
+ let totalGenMs: number | undefined;
+ for (const step of turn.steps) {
+ const stepMs = step.decodeMs ?? step.genTotalMs;
+ if (stepMs !== undefined) {
+ totalGenMs = (totalGenMs ?? 0) + stepMs;
+ }
+ }
+ const tps = computeTps(turn.usage.outputTokens, totalGenMs);
+ return {
+ label: turnNumber !== undefined ? `turn ${turnNumber}` : "turn",
+ tokensLabel: `${formatTokens(total)} tok`,
+ breakdown: formatBreakdown(turn.usage),
+ tps: formatTps(tps),
+ duration: formatDuration(turn.durationMs),
+ };
+}
diff --git a/src/core/metrics/index.ts b/src/core/metrics/index.ts
new file mode 100644
index 0000000..d3c9669
--- /dev/null
+++ b/src/core/metrics/index.ts
@@ -0,0 +1,31 @@
+export {
+ type ContextUsage,
+ computeCachePct,
+ computeContextUsage,
+ computeExpectedCachePct,
+ computeTps,
+ formatCompactTokens,
+ formatContextSize,
+ viewCacheRate,
+ viewExpectedCache,
+ viewStepMetrics,
+ viewTurnMetrics,
+} from "./format";
+export { interleaveTurnMetrics } from "./place";
+export {
+ applyDurableMetrics,
+ foldMetricsEvent,
+ initialMetricsState,
+ selectCurrentContextSize,
+ selectOrderedTurnMetrics,
+} from "./reducer";
+export type {
+ CacheRateView,
+ MetricsRow,
+ MetricsState,
+ StepMetrics,
+ StepMetricsView,
+ TurnMetrics,
+ TurnMetricsEntry,
+ TurnMetricsView,
+} from "./types";
diff --git a/src/core/metrics/place.test.ts b/src/core/metrics/place.test.ts
new file mode 100644
index 0000000..9c925a3
--- /dev/null
+++ b/src/core/metrics/place.test.ts
@@ -0,0 +1,621 @@
+import type { StepId, StepMetrics, TurnMetrics } from "@dispatch/wire";
+import { describe, expect, it } from "vitest";
+import type { RenderGroup } from "../chunks";
+import { interleaveTurnMetrics } from "./place";
+import type { MetricsRow, TurnMetricsEntry } from "./types";
+
+function userGroup(seq: number, text: string): RenderGroup {
+ return {
+ kind: "single",
+ chunk: {
+ seq,
+ role: "user",
+ chunk: { type: "text", text },
+ provisional: false,
+ },
+ };
+}
+
+function assistantGroup(seq: number, text: string): RenderGroup {
+ return {
+ kind: "single",
+ chunk: {
+ seq,
+ role: "assistant",
+ chunk: { type: "text", text },
+ provisional: false,
+ },
+ };
+}
+
+function toolCallGroup(seq: number, stepId: string, toolCallId: string): RenderGroup {
+ return {
+ kind: "single",
+ chunk: {
+ seq,
+ role: "assistant",
+ chunk: {
+ type: "tool-call",
+ toolCallId,
+ toolName: "test",
+ input: {},
+ stepId: stepId as StepId,
+ },
+ provisional: false,
+ },
+ };
+}
+
+function toolResultGroup(seq: number, stepId: string, toolCallId: string): RenderGroup {
+ return {
+ kind: "single",
+ chunk: {
+ seq,
+ role: "tool",
+ chunk: {
+ type: "tool-result",
+ toolCallId,
+ toolName: "test",
+ content: "",
+ isError: false,
+ stepId: stepId as StepId,
+ },
+ provisional: false,
+ },
+ };
+}
+
+function toolBatchGroup(stepId: string, toolCallIds: string[]): RenderGroup {
+ return {
+ kind: "tool-batch",
+ stepId,
+ entries: toolCallIds.map((id) => ({
+ call: {
+ type: "tool-call" as const,
+ toolCallId: id,
+ toolName: "test",
+ input: {},
+ stepId: stepId as StepId,
+ },
+ result: null,
+ })),
+ provisional: false,
+ };
+}
+
+function makeStep(stepId: string, inputTokens: number, outputTokens: number): StepMetrics {
+ return {
+ stepId: stepId as StepId,
+ usage: { inputTokens, outputTokens },
+ };
+}
+
+function makeTurn(
+ turnId: string,
+ inputTokens: number,
+ outputTokens: number,
+ steps: StepMetrics[] = [],
+): TurnMetrics {
+ return {
+ turnId,
+ usage: { inputTokens, outputTokens },
+ steps,
+ };
+}
+
+function makeEntry(
+ turnId: string,
+ inputTokens: number,
+ outputTokens: number,
+ steps: StepMetrics[] = [],
+): TurnMetricsEntry {
+ return {
+ turnId,
+ steps,
+ total: makeTurn(turnId, inputTokens, outputTokens, steps),
+ };
+}
+
+function makeProgressiveEntry(turnId: string, steps: StepMetrics[]): TurnMetricsEntry {
+ return {
+ turnId,
+ steps,
+ total: null,
+ };
+}
+
+function expectGroupAt(
+ rows: readonly { readonly kind: string }[],
+ index: number,
+ expected: RenderGroup,
+): void {
+ const row = rows[index];
+ expect(row?.kind).toBe("group");
+ expect((row as { readonly group: RenderGroup } | undefined)?.group).toBe(expected);
+}
+
+function expectStepMetricsAt(
+ rows: readonly { readonly kind: string }[],
+ index: number,
+ expectedStepId: string,
+ expectedIndex: number,
+): void {
+ const row = rows[index];
+ expect(row?.kind).toBe("step-metrics");
+ const sm = row as { readonly step: StepMetrics; readonly index: number } | undefined;
+ expect(sm?.step.stepId).toBe(expectedStepId);
+ expect(sm?.index).toBe(expectedIndex);
+}
+
+function expectTurnMetricsAt(
+ rows: readonly { readonly kind: string }[],
+ index: number,
+ expectedTurnId: string,
+): void {
+ const row = rows[index];
+ expect(row?.kind).toBe("turn-metrics");
+ expect((row as { readonly turn: TurnMetrics } | undefined)?.turn.turnId).toBe(expectedTurnId);
+}
+
+describe("interleaveTurnMetrics", () => {
+ it("no metrics: rows are all groups, unchanged order", () => {
+ const g1 = userGroup(1, "q");
+ const g2 = assistantGroup(2, "a");
+ const rows = interleaveTurnMetrics([g1, g2], []);
+ expect(rows).toHaveLength(2);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ });
+
+ it("head-aligned: segment i gets entries[i]", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolCallGroup(2, "s1", "c1");
+ const g3 = userGroup(3, "q2");
+ const g4 = toolCallGroup(4, "s2", "c2");
+ const step1 = makeStep("s1", 100, 50);
+ const step2 = makeStep("s2", 200, 80);
+ const rows = interleaveTurnMetrics(
+ [g1, g2, g3, g4],
+ [makeEntry("t1", 100, 50, [step1]), makeEntry("t2", 200, 80, [step2])],
+ );
+
+ expect(rows).toHaveLength(8);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectTurnMetricsAt(rows, 3, "t1");
+ expectGroupAt(rows, 4, g3);
+ expectGroupAt(rows, 5, g4);
+ expectStepMetricsAt(rows, 6, "s2", 0);
+ expectTurnMetricsAt(rows, 7, "t2");
+ });
+
+ it("a trailing segment with no entry (in-flight turn) renders no metrics", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolCallGroup(2, "s1", "c1");
+ const g3 = userGroup(3, "q2");
+ const g4 = assistantGroup(4, "a2");
+ const step = makeStep("s1", 100, 50);
+ const rows = interleaveTurnMetrics([g1, g2, g3, g4], [makeEntry("t1", 100, 50, [step])]);
+
+ expect(rows).toHaveLength(6);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectTurnMetricsAt(rows, 3, "t1");
+ expectGroupAt(rows, 4, g3);
+ expectGroupAt(rows, 5, g4);
+ });
+
+ it("single text-only turn: no step row (unanchored), turn-metrics at tail", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = assistantGroup(2, "a1");
+ const step = makeStep("s1", 100, 50);
+ const turn = makeEntry("t1", 100, 50, [step]);
+ const rows = interleaveTurnMetrics([g1, g2], [turn]);
+
+ expect(rows).toHaveLength(3);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectTurnMetricsAt(rows, 2, "t1");
+ });
+
+ it("tool step anchors inline after its tool-batch group", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolBatchGroup("t#0", ["c1", "c2"]);
+ const g3 = assistantGroup(3, "a1");
+ const step0 = makeStep("t#0", 100, 50);
+ const step1 = makeStep("t#1", 200, 80);
+ const turn = makeEntry("t1", 300, 130, [step0, step1]);
+ const rows = interleaveTurnMetrics([g1, g2, g3], [turn]);
+
+ expect(rows).toHaveLength(5);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "t#0", 0);
+ expectGroupAt(rows, 3, g3);
+ expectTurnMetricsAt(rows, 4, "t1");
+ });
+
+ it("single tool-call group anchors its step", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolCallGroup(2, "s1", "c1");
+ const g3 = assistantGroup(3, "a1");
+ const step = makeStep("s1", 100, 50);
+ const turn = makeEntry("t1", 100, 50, [step]);
+ const rows = interleaveTurnMetrics([g1, g2, g3], [turn]);
+
+ expect(rows).toHaveLength(5);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectGroupAt(rows, 3, g3);
+ expectTurnMetricsAt(rows, 4, "t1");
+ });
+
+ it("single tool-result group anchors its step", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolResultGroup(2, "s1", "c1");
+ const g3 = assistantGroup(3, "a1");
+ const step = makeStep("s1", 100, 50);
+ const turn = makeEntry("t1", 100, 50, [step]);
+ const rows = interleaveTurnMetrics([g1, g2, g3], [turn]);
+
+ expect(rows).toHaveLength(5);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectGroupAt(rows, 3, g3);
+ expectTurnMetricsAt(rows, 4, "t1");
+ });
+
+ it("multi-step: each tool step inline, unanchored text step skipped", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolBatchGroup("t#0", ["c1"]);
+ const g3 = assistantGroup(2, "thinking");
+ const g4 = toolBatchGroup("t#1", ["c2", "c3"]);
+ const g5 = assistantGroup(3, "a1");
+ const step0 = makeStep("t#0", 100, 50);
+ const step1 = makeStep("t#1", 200, 80);
+ const step2 = makeStep("t#2", 50, 20);
+ const turn = makeEntry("t1", 350, 150, [step0, step1, step2]);
+ const rows = interleaveTurnMetrics([g1, g2, g3, g4, g5], [turn]);
+
+ expect(rows).toHaveLength(8);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "t#0", 0);
+ expectGroupAt(rows, 3, g3);
+ expectGroupAt(rows, 4, g4);
+ expectStepMetricsAt(rows, 5, "t#1", 1);
+ expectGroupAt(rows, 6, g5);
+ expectTurnMetricsAt(rows, 7, "t1");
+ });
+
+ it("multiple turns head-aligned with inline steps", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolBatchGroup("s1", ["c1"]);
+ const g3 = assistantGroup(2, "a1");
+ const g4 = userGroup(3, "q2");
+ const g5 = toolCallGroup(4, "s2", "c2");
+ const step1 = makeStep("s1", 100, 50);
+ const step2 = makeStep("s2", 200, 80);
+ const rows = interleaveTurnMetrics(
+ [g1, g2, g3, g4, g5],
+ [makeEntry("t1", 100, 50, [step1]), makeEntry("t2", 200, 80, [step2])],
+ );
+
+ expect(rows).toHaveLength(9);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectGroupAt(rows, 3, g3);
+ expectTurnMetricsAt(rows, 4, "t1");
+ expectGroupAt(rows, 5, g4);
+ expectGroupAt(rows, 6, g5);
+ expectStepMetricsAt(rows, 7, "s2", 0);
+ expectTurnMetricsAt(rows, 8, "t2");
+ });
+
+ it("unanchored step (stepId not in groups) is skipped — only turn-metrics", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = assistantGroup(2, "a1");
+ const step0 = makeStep("orphan", 100, 50);
+ const turn = makeEntry("t1", 100, 50, [step0]);
+ const rows = interleaveTurnMetrics([g1, g2], [turn]);
+
+ expect(rows).toHaveLength(3);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectTurnMetricsAt(rows, 2, "t1");
+ });
+
+ it("fewer metrics than segments: trailing segments are bare", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolCallGroup(2, "s1", "c1");
+ const g3 = userGroup(3, "q2");
+ const g4 = assistantGroup(4, "a2");
+ const g5 = userGroup(5, "q3");
+ const g6 = assistantGroup(6, "a3");
+ const step = makeStep("s1", 300, 120);
+ const rows = interleaveTurnMetrics(
+ [g1, g2, g3, g4, g5, g6],
+ [makeEntry("t1", 300, 120, [step])],
+ );
+
+ expect(rows).toHaveLength(8);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectTurnMetricsAt(rows, 3, "t1");
+ expectGroupAt(rows, 4, g3);
+ expectGroupAt(rows, 5, g4);
+ expectGroupAt(rows, 6, g5);
+ expectGroupAt(rows, 7, g6);
+ });
+
+ it("trimmed leading turns: a mixed tool+text transcript tail-aligns text-only turns to their OWN (newest) entries, not stale trimmed ones", () => {
+ // A long conversation where the chat limit unloaded the oldest turn (t1).
+ // Metrics still hold all three turns; the loaded transcript is turns 2-3.
+ // Turn 2 is a tool turn (matched by stepId); turn 3 is text-only (no
+ // stepId groups) — the failure case. The text-only turn MUST get its OWN
+ // entry (t3), NOT the trimmed t1's stale metrics.
+ const g3 = userGroup(3, "q2");
+ const g4 = toolBatchGroup("s2", ["c2"]);
+ const g5 = assistantGroup(4, "tool-reply");
+ const g6 = userGroup(5, "q3");
+ const g7 = assistantGroup(6, "text-reply");
+ const step1 = makeStep("s1", 11, 1); // t1 (trimmed)
+ const step2 = makeStep("s2", 22, 2); // t2 (loaded, tool)
+ const step3 = makeStep("s3", 33, 3); // t3 (loaded, text-only — unanchored)
+ const entries = [
+ makeEntry("t1", 11, 1, [step1]),
+ makeEntry("t2", 22, 2, [step2]),
+ makeEntry("t3", 33, 3, [step3]),
+ ];
+ const rows = interleaveTurnMetrics([g3, g4, g5, g6, g7], entries);
+
+ const tmRows = rows.filter(
+ (r): r is Extract<MetricsRow, { kind: "turn-metrics" }> => r.kind === "turn-metrics",
+ );
+ // Two loaded turns → two turn-metrics rows. The trimmed t1 does NOT render.
+ expect(tmRows).toHaveLength(2);
+ // CRITICAL: the text-only turn (segment 1) got t3 (its own newest entry),
+ // not t1 (the stale trimmed one). A misaligned head-align would show t1.
+ expect(tmRows[1]?.turn.turnId).toBe("t3");
+ expect(tmRows[0]?.turn.turnId).toBe("t2");
+ // And t1 never appears as a rendered row.
+ expect(tmRows.some((r) => r.turn.turnId === "t1")).toBe(false);
+ });
+
+ it("trimmed turn still counts toward the cumulative 'chat total' on the first visible turn", () => {
+ // t1 is trimmed (no segment) but finalized; t2 is the loaded visible turn.
+ // t2's "Chat Total" cumulative must INCLUDE t1's usage (the whole chat),
+ // even though t1 renders no row of its own.
+ const g1 = userGroup(2, "q2");
+ const g2 = assistantGroup(3, "a2");
+ const entries = [
+ {
+ turnId: "t1",
+ steps: [],
+ total: {
+ turnId: "t1",
+ usage: { inputTokens: 1000, outputTokens: 10, cacheReadTokens: 500 },
+ steps: [],
+ },
+ },
+ {
+ turnId: "t2",
+ steps: [],
+ total: {
+ turnId: "t2",
+ usage: { inputTokens: 2000, outputTokens: 20, cacheReadTokens: 1600 },
+ steps: [],
+ },
+ },
+ ];
+ const rows = interleaveTurnMetrics([g1, g2], entries);
+ const tmRows = rows.filter(
+ (r): r is Extract<MetricsRow, { kind: "turn-metrics" }> => r.kind === "turn-metrics",
+ );
+ // Only the loaded turn renders a row; the trimmed t1 does not.
+ expect(tmRows).toHaveLength(1);
+ expect(tmRows[0]?.turn.turnId).toBe("t2");
+ // Cumulative includes BOTH turns (t1 + t2): input 3000, cacheRead 2100.
+ expect(tmRows[0]?.cumulativeUsage.inputTokens).toBe(3000);
+ expect(tmRows[0]?.cumulativeUsage.cacheReadTokens).toBe(2100);
+ // Retention baseline is the prior finalized turn (t1, even though trimmed).
+ expect(tmRows[0]?.prevTurnUsage?.inputTokens).toBe(1000);
+ expect(tmRows[0]?.prevTurnUsage?.cacheReadTokens).toBe(500);
+ });
+
+ it("in-flight turn (no durationMs) still produces turn row", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolCallGroup(2, "s1", "c1");
+ const step = makeStep("s1", 100, 50);
+ const turn: TurnMetricsEntry = {
+ turnId: "t1",
+ steps: [step],
+ total: {
+ turnId: "t1",
+ usage: { inputTokens: 100, outputTokens: 50 },
+ steps: [step],
+ },
+ };
+ const rows = interleaveTurnMetrics([g1, g2], [turn]);
+
+ expect(rows).toHaveLength(4);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectTurnMetricsAt(rows, 3, "t1");
+ const metricsRow = rows[3] as { readonly turn: TurnMetrics } | undefined;
+ expect(metricsRow?.turn.durationMs).toBeUndefined();
+ });
+
+ it("leading non-turn groups emit as plain group rows", () => {
+ const g0 = assistantGroup(1, "system msg");
+ const g1 = userGroup(2, "q1");
+ const g2 = toolCallGroup(3, "s1", "c1");
+ const step = makeStep("s1", 100, 50);
+ const rows = interleaveTurnMetrics([g0, g1, g2], [makeEntry("t1", 100, 50, [step])]);
+
+ expect(rows).toHaveLength(5);
+ expectGroupAt(rows, 0, g0);
+ expect(rows[1]?.kind).toBe("group");
+ expect(rows[2]?.kind).toBe("group");
+ expectStepMetricsAt(rows, 3, "s1", 0);
+ expectTurnMetricsAt(rows, 4, "t1");
+ });
+
+ it("trimmed turn (more metrics than segments) does NOT emit a standalone row at the top", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolCallGroup(2, "s1", "c1");
+ const step1 = makeStep("s1", 100, 50);
+ const step2 = makeStep("s2", 200, 80);
+ const rows = interleaveTurnMetrics(
+ [g1, g2],
+ [makeEntry("t1", 100, 50, [step1]), makeEntry("t2", 200, 80, [step2])],
+ );
+
+ // t2's content was unloaded by the chat limit (no segment for it); its
+ // metrics must NOT render a standalone row piled at the top. Only the
+ // loaded turn's content + its matched metrics appear. (t2 still counts
+ // toward the cumulative "chat total" — see the cache-total tests.)
+ expect(rows).toHaveLength(4);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectTurnMetricsAt(rows, 3, "t1");
+ // No standalone turn-metrics row for t2 anywhere.
+ const tmRows = rows.filter((r) => r.kind === "turn-metrics");
+ expect(tmRows).toHaveLength(1);
+ expect((tmRows[0] as { readonly turn: TurnMetrics }).turn.turnId).toBe("t1");
+ });
+
+ it("turn with no steps emits only turn-metrics (no step-metrics)", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = assistantGroup(2, "a1");
+ const rows = interleaveTurnMetrics([g1, g2], [makeEntry("t1", 100, 50)]);
+
+ expect(rows).toHaveLength(3);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectTurnMetricsAt(rows, 2, "t1");
+ });
+
+ it("progressive: entry with steps but total=null emits step rows and NO turn-metrics row", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolBatchGroup("s1", ["c1"]);
+ const g3 = assistantGroup(2, "a1");
+ const step1 = makeStep("s1", 100, 50);
+ const entry = makeProgressiveEntry("t1", [step1]);
+ const rows = interleaveTurnMetrics([g1, g2, g3], [entry]);
+
+ expect(rows).toHaveLength(4);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectGroupAt(rows, 3, g3);
+ });
+
+ it("entry with total emits step rows + a turn-metrics row", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = toolBatchGroup("s1", ["c1"]);
+ const g3 = assistantGroup(2, "a1");
+ const step1 = makeStep("s1", 100, 50);
+ const entry = makeEntry("t1", 100, 50, [step1]);
+ const rows = interleaveTurnMetrics([g1, g2, g3], [entry]);
+
+ expect(rows).toHaveLength(5);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ expectStepMetricsAt(rows, 2, "s1", 0);
+ expectGroupAt(rows, 3, g3);
+ expectTurnMetricsAt(rows, 4, "t1");
+ });
+
+ it("progressive multi-step: unanchored steps skipped, no turn-metrics", () => {
+ const g1 = userGroup(1, "q1");
+ const g2 = assistantGroup(2, "a1");
+ const step0 = makeStep("s1", 100, 50);
+ const step1 = makeStep("s2", 200, 80);
+ const entry = makeProgressiveEntry("t1", [step0, step1]);
+ const rows = interleaveTurnMetrics([g1, g2], [entry]);
+
+ expect(rows).toHaveLength(2);
+ expectGroupAt(rows, 0, g1);
+ expectGroupAt(rows, 1, g2);
+ });
+});
+
+describe("interleaveTurnMetrics — cumulative usage (cache total)", () => {
+ function turnMetricsRows(rows: readonly MetricsRow[]) {
+ return rows.filter((r): r is Extract<MetricsRow, { kind: "turn-metrics" }> => {
+ return r.kind === "turn-metrics";
+ });
+ }
+
+ function cacheEntry(
+ turnId: string,
+ inputTokens: number,
+ outputTokens: number,
+ cacheReadTokens: number,
+ ): TurnMetricsEntry {
+ const total: TurnMetrics = {
+ turnId,
+ usage: { inputTokens, outputTokens, cacheReadTokens },
+ steps: [],
+ };
+ return { turnId, steps: [], total };
+ }
+
+ it("turn-metrics row carries this turn's usage and the running cumulative", () => {
+ const rows = interleaveTurnMetrics(
+ [userGroup(1, "q1"), assistantGroup(2, "a1")],
+ [makeEntry("t1", 1000, 100)],
+ );
+ const tm = turnMetricsRows(rows);
+ expect(tm).toHaveLength(1);
+ expect(tm[0]?.turn.turnId).toBe("t1");
+ expect(tm[0]?.cumulativeUsage).toEqual({ inputTokens: 1000, outputTokens: 100 });
+ });
+
+ it("accumulates cache read + input across turns (chat total)", () => {
+ const rows = interleaveTurnMetrics(
+ [userGroup(1, "q1"), assistantGroup(2, "a1"), userGroup(3, "q2"), assistantGroup(4, "a2")],
+ [cacheEntry("t1", 2669, 10, 384), cacheEntry("t2", 2737, 10, 2560)],
+ );
+ const tm = turnMetricsRows(rows);
+ expect(tm).toHaveLength(2);
+ // turn 1: only its own usage
+ expect(tm[0]?.cumulativeUsage.inputTokens).toBe(2669);
+ expect(tm[0]?.cumulativeUsage.cacheReadTokens).toBe(384);
+ // turn 2: sum of both (input 5406, cacheRead 2944 → matches the backend's 54% example)
+ expect(tm[1]?.cumulativeUsage.inputTokens).toBe(5406);
+ expect(tm[1]?.cumulativeUsage.cacheReadTokens).toBe(2944);
+ });
+
+ it("an in-flight (total=null) turn does not contribute to the cumulative", () => {
+ const rows = interleaveTurnMetrics(
+ [userGroup(1, "q1"), assistantGroup(2, "a1"), userGroup(3, "q2"), assistantGroup(4, "a2")],
+ [cacheEntry("t1", 1000, 10, 500), makeProgressiveEntry("t2", [makeStep("s1", 200, 5)])],
+ );
+ const tm = turnMetricsRows(rows);
+ // only the finalized turn emits a turn-metrics row; its cumulative is just itself
+ expect(tm).toHaveLength(1);
+ expect(tm[0]?.cumulativeUsage.inputTokens).toBe(1000);
+ expect(tm[0]?.cumulativeUsage.cacheReadTokens).toBe(500);
+ });
+
+ it("carries the prior finalized turn's usage as the retention baseline", () => {
+ const rows = interleaveTurnMetrics(
+ [userGroup(1, "q1"), assistantGroup(2, "a1"), userGroup(3, "q2"), assistantGroup(4, "a2")],
+ [cacheEntry("t1", 2669, 10, 384), cacheEntry("t2", 2737, 10, 2560)],
+ );
+ const tm = turnMetricsRows(rows);
+ // first finalized turn has no earlier baseline
+ expect(tm[0]?.prevTurnUsage).toBeNull();
+ // second turn's baseline is the first turn's usage
+ expect(tm[1]?.prevTurnUsage?.inputTokens).toBe(2669);
+ expect(tm[1]?.prevTurnUsage?.cacheReadTokens).toBe(384);
+ });
+});
diff --git a/src/core/metrics/place.ts b/src/core/metrics/place.ts
new file mode 100644
index 0000000..7122b09
--- /dev/null
+++ b/src/core/metrics/place.ts
@@ -0,0 +1,298 @@
+import type { Usage } from "@dispatch/wire";
+import type { RenderGroup } from "../chunks";
+import type { MetricsRow, TurnMetricsEntry } from "./types";
+
+function groupStepId(g: RenderGroup): string | undefined {
+ if (g.kind === "tool-batch") return g.stepId;
+ const c = g.chunk.chunk;
+ return c.type === "tool-call" || c.type === "tool-result" ? c.stepId : undefined;
+}
+
+/** Element-wise sum of two token usages (cache fields included only when nonzero). */
+function addUsage(a: Usage, b: Usage): Usage {
+ const out: Usage = {
+ inputTokens: a.inputTokens + b.inputTokens,
+ outputTokens: a.outputTokens + b.outputTokens,
+ };
+ const read = (a.cacheReadTokens ?? 0) + (b.cacheReadTokens ?? 0);
+ const write = (a.cacheWriteTokens ?? 0) + (b.cacheWriteTokens ?? 0);
+ if (read > 0) (out as { cacheReadTokens?: number }).cacheReadTokens = read;
+ if (write > 0) (out as { cacheWriteTokens?: number }).cacheWriteTokens = write;
+ return out;
+}
+
+/**
+ * Interleave turn metrics into the rendered transcript.
+ *
+ * Splits groups into per-turn segments: a new segment begins at each `single`
+ * group with `group.chunk.role === "user"`. Segments are matched to entries
+ * by `stepId` presence when possible (robust against chat-limit trimming: when
+ * a turn's user message is trimmed, positional alignment would be off, but
+ * stepId matching still finds the right entry). Segments with no stepId-bearing
+ * groups (text-only turns) fall back to POSITIONAL tail-alignment: since the
+ * loaded transcript is always a SUFFIX of the full turn history (the chat limit
+ * keeps the newest and unloads the oldest), segment `seg` ↔ entry `K - T + seg`.
+ *
+ * Within a segment that has a matched entry, each completed step's metrics
+ * are placed INLINE right after the last group bearing that step's `stepId`.
+ * Steps whose `stepId` does not appear in any group ("unanchored"):
+ * - If the segment HAS stepId-bearing groups (tool chunks exist but this step's
+ * were trimmed): SKIPPED (no blank "step N · 0 tok" bubbles).
+ * - If the segment has NO stepId-bearing groups (text-only turn): placed at the
+ * segment tail before the turn-metrics row (the original behavior).
+ *
+ * A `turn-metrics` row is emitted ONLY when `entry.total !== null` (i.e. the turn
+ * is finalized via `done` or durable data). A still-generating turn emits no
+ * turn-total row.
+ *
+ * Fully trimmed turns (entries whose content was unloaded by the chat limit and
+ * which match no segment) are NOT rendered as standalone rows — that previously
+ * piled a wall of stale cache badges at the top of a long, trimmed transcript.
+ * Their usage still counts toward the per-turn "chat total" cumulative (computed
+ * across ALL finalized turns in entry-array order), so the running cache rate
+ * stays correct regardless of which turns were trimmed; paging earlier history
+ * back in ("Show earlier messages") re-matches them and re-renders their rows.
+ */
+export function interleaveTurnMetrics(
+ groups: readonly RenderGroup[],
+ entries: readonly TurnMetricsEntry[],
+): readonly MetricsRow[] {
+ if (entries.length === 0) {
+ return groups.map((g) => ({ kind: "group" as const, group: g }));
+ }
+
+ const segmentStarts: number[] = [];
+ for (let i = 0; i < groups.length; i++) {
+ const g = groups[i];
+ if (g !== undefined && g.kind === "single" && g.chunk.role === "user") {
+ segmentStarts.push(i);
+ }
+ }
+
+ let T = segmentStarts.length;
+
+ // No user messages — e.g. a compacted conversation whose history starts
+ // with a system summary. Treat the entire transcript as one segment so
+ // turn/step metrics can still be placed.
+ if (T === 0 && entries.length > 0) {
+ segmentStarts.push(0);
+ T = 1;
+ }
+
+ if (T === 0) {
+ return groups.map((g) => ({ kind: "group" as const, group: g }));
+ }
+
+ const K = entries.length;
+
+ // Build stepId → entry-index lookup for matching.
+ const entryStepIds: Set<string>[] = entries.map((e) => new Set(e.steps.map((s) => s.stepId)));
+
+ // Match segments to entries. Pass 1: match by stepId overlap (handles
+ // trimming where positional alignment alone could be ambiguous). Pass 2:
+ // positional tail-alignment fallback for unmatched segments (text-only turns
+ // with no stepId-bearing groups).
+ const usedEntries = new Set<number>();
+ const segmentEntry = new Map<number, TurnMetricsEntry>();
+ const segmentEntryIndex = new Map<number, number>();
+
+ // Pass 1: stepId matching.
+ for (let seg = 0; seg < T; seg++) {
+ const start = segmentStarts[seg] ?? 0;
+ const end = seg + 1 < T ? (segmentStarts[seg + 1] ?? groups.length) : groups.length;
+
+ const segStepIds = new Set<string>();
+ for (let i = start; i < end; i++) {
+ const g = groups[i];
+ if (g === undefined) continue;
+ const sid = groupStepId(g);
+ if (sid !== undefined) segStepIds.add(sid);
+ }
+ if (segStepIds.size === 0) continue; // text-only — defer to pass 2
+
+ let bestEntry = -1;
+ let bestMatch = 0;
+ for (let i = 0; i < K; i++) {
+ if (usedEntries.has(i)) continue;
+ let match = 0;
+ for (const sid of segStepIds) {
+ if (entryStepIds[i]?.has(sid)) match++;
+ }
+ if (match > bestMatch) {
+ bestMatch = match;
+ bestEntry = i;
+ }
+ }
+ if (bestEntry >= 0) {
+ usedEntries.add(bestEntry);
+ const e = entries[bestEntry];
+ if (e !== undefined) {
+ segmentEntry.set(seg, e);
+ segmentEntryIndex.set(seg, bestEntry);
+ }
+ }
+ }
+
+ // Pass 2: positional fallback for segments pass 1 left unmatched
+ // (text-only turns with no stepId-bearing groups to anchor on).
+ //
+ // The loaded transcript is always a SUFFIX of the full turn history —
+ // chat-limit/windowing keeps the NEWEST chunks and unloads the OLDEST — so
+ // the T loaded segments correspond to the LAST T entries. TAIL-ALIGNMENT
+ // (segment `seg` ↔ entry `K - T + seg`) is therefore correct whenever the
+ // metrics hold at least as many turns as there are loaded segments
+ // (`K >= T`): the leading `K - T` entries are TRIMMED turns (their content
+ // was unloaded) and must be skipped, never matched to a newer segment.
+ //
+ // This MUST run even when pass 1 matched SOME segments (tool turns). The
+ // earlier code only tail-aligned when pass 1 matched NONE, falling back to
+ // HEAD-alignment otherwise — which, with leading trimmed entries, matched a
+ // brand-new text-only turn to an old (trimmed) entry's STALE metrics (the
+ // "new steps show no / wrong cache" failure). Tail-aligning by position is
+ // safe alongside pass 1: stepIds are unique per turn, so pass 1 already
+ // grabbed each tool turn's positionally-correct entry, leaving the right
+ // entry free for each text-only turn.
+ //
+ // Only when `K < T` (fewer entries than segments — some loaded turns have no
+ // metrics yet, e.g. a metrics sync still pending or a freshly loaded
+ // transcript) do we head-align, assigning the first K entries to the first K
+ // unmatched segments (the turns that DO have metrics sit at the front).
+ if (K >= T) {
+ // Tail-align: skip the first K-T entries (trimmed turns).
+ for (let seg = 0; seg < T; seg++) {
+ if (segmentEntry.has(seg)) continue;
+ const entryIdx = K - T + seg;
+ if (entryIdx >= 0 && entryIdx < K && !usedEntries.has(entryIdx)) {
+ usedEntries.add(entryIdx);
+ const e = entries[entryIdx];
+ if (e !== undefined) {
+ segmentEntry.set(seg, e);
+ segmentEntryIndex.set(seg, entryIdx);
+ }
+ }
+ }
+ } else {
+ // Head-align fallback (K < T): first K entries to first K unmatched segments.
+ let nextUnused = 0;
+ for (let seg = 0; seg < T; seg++) {
+ if (segmentEntry.has(seg)) continue;
+ while (nextUnused < K && usedEntries.has(nextUnused)) nextUnused++;
+ if (nextUnused < K) {
+ usedEntries.add(nextUnused);
+ const e = entries[nextUnused];
+ if (e !== undefined) {
+ segmentEntry.set(seg, e);
+ segmentEntryIndex.set(seg, nextUnused);
+ }
+ nextUnused++;
+ }
+ }
+ }
+
+ // Running cumulative usage across ALL finalized turns (in entry order), for
+ // the per-turn "chat total" cache rate. Alongside it, the previous finalized
+ // turn's usage at each index — the baseline for cross-turn retention.
+ const cumulativeByEntry: Usage[] = [];
+ const prevUsageByEntry: (Usage | null)[] = [];
+ let runningUsage: Usage = { inputTokens: 0, outputTokens: 0 };
+ let lastFinalizedUsage: Usage | null = null;
+ for (const e of entries) {
+ prevUsageByEntry.push(lastFinalizedUsage);
+ if (e.total !== null) {
+ runningUsage = addUsage(runningUsage, e.total.usage);
+ lastFinalizedUsage = e.total.usage;
+ }
+ cumulativeByEntry.push(runningUsage);
+ }
+
+ const rows: MetricsRow[] = [];
+
+ const firstUserIdx = segmentStarts[0] ?? 0;
+
+ for (let i = 0; i < firstUserIdx; i++) {
+ const g = groups[i];
+ if (g !== undefined) {
+ rows.push({ kind: "group", group: g });
+ }
+ }
+
+ for (let seg = 0; seg < T; seg++) {
+ const start = segmentStarts[seg] ?? 0;
+ const end = seg + 1 < T ? (segmentStarts[seg + 1] ?? groups.length) : groups.length;
+
+ const entry = segmentEntry.get(seg);
+
+ if (entry === undefined) {
+ for (let i = start; i < end; i++) {
+ const g = groups[i];
+ if (g !== undefined) {
+ rows.push({ kind: "group", group: g });
+ }
+ }
+ continue;
+ }
+
+ const entryIdx = segmentEntryIndex.get(seg) ?? 0;
+
+ // Build anchor map: for each stepId, the LAST group index in this segment.
+ const anchorByStepId = new Map<string, number>();
+ for (let i = start; i < end; i++) {
+ const g = groups[i];
+ if (g === undefined) continue;
+ const sid = groupStepId(g);
+ if (sid !== undefined) {
+ anchorByStepId.set(sid, i);
+ }
+ }
+
+ // Classify each step as anchored or unanchored. Unanchored steps
+ // (content trimmed, or text-only steps with no tool chunks) are SKIPPED —
+ // step-metrics are only shown inline next to the content they describe.
+ const anchored: Map<number, { stepIndex: number; step: (typeof entry.steps)[number] }[]> =
+ new Map();
+
+ for (let i = 0; i < entry.steps.length; i++) {
+ const step = entry.steps[i];
+ if (step === undefined) continue;
+ const anchorGroupIdx = anchorByStepId.get(step.stepId);
+ if (anchorGroupIdx !== undefined) {
+ let arr = anchored.get(anchorGroupIdx);
+ if (arr === undefined) {
+ arr = [];
+ anchored.set(anchorGroupIdx, arr);
+ }
+ arr.push({ stepIndex: i, step });
+ }
+ // Unanchored steps (no matching group) are skipped — no tail bubbles.
+ }
+
+ // Emit groups; after each anchored group, emit its step-metrics rows.
+ for (let i = start; i < end; i++) {
+ const g = groups[i];
+ if (g !== undefined) {
+ rows.push({ kind: "group", group: g });
+ }
+ const stepsHere = anchored.get(i);
+ if (stepsHere !== undefined) {
+ stepsHere.sort((a, b) => a.stepIndex - b.stepIndex);
+ for (const { step, stepIndex } of stepsHere) {
+ rows.push({ kind: "step-metrics", step, index: stepIndex });
+ }
+ }
+ }
+
+ // Turn-metrics row (only when the turn is finalized). Unanchored steps
+ // are skipped — no tail bubbles.
+ if (entry.total !== null) {
+ rows.push({
+ kind: "turn-metrics",
+ turn: entry.total,
+ turnNumber: entryIdx + 1,
+ cumulativeUsage: cumulativeByEntry[entryIdx] ?? entry.total.usage,
+ prevTurnUsage: prevUsageByEntry[entryIdx] ?? null,
+ });
+ }
+ }
+
+ return rows;
+}
diff --git a/src/core/metrics/reducer.test.ts b/src/core/metrics/reducer.test.ts
new file mode 100644
index 0000000..581a8b7
--- /dev/null
+++ b/src/core/metrics/reducer.test.ts
@@ -0,0 +1,689 @@
+import type { StepId, TurnDoneEvent, TurnStepCompleteEvent, TurnUsageEvent } from "@dispatch/wire";
+import { describe, expect, it } from "vitest";
+import {
+ applyDurableMetrics,
+ foldMetricsEvent,
+ initialMetricsState,
+ selectCurrentContextSize,
+ selectOrderedTurnMetrics,
+} from "./reducer";
+
+const usageEvent = (
+ turnId: string,
+ inputTokens: number,
+ outputTokens: number,
+ stepId?: string,
+): TurnUsageEvent => {
+ const base = {
+ type: "usage" as const,
+ conversationId: "c1",
+ turnId,
+ usage: { inputTokens, outputTokens },
+ };
+ if (stepId !== undefined) {
+ return { ...base, stepId: stepId as StepId };
+ }
+ return base;
+};
+
+const stepCompleteEvent = (
+ turnId: string,
+ stepId: string,
+ timing: { ttftMs?: number; decodeMs?: number; genTotalMs?: number } = {},
+): TurnStepCompleteEvent => ({
+ type: "step-complete",
+ conversationId: "c1",
+ turnId,
+ stepId: stepId as StepId,
+ ...timing,
+});
+
+const doneEvent = (
+ turnId: string,
+ extra: {
+ durationMs?: number;
+ usage?: { inputTokens: number; outputTokens: number };
+ contextSize?: number;
+ } = {},
+): TurnDoneEvent => ({
+ type: "done",
+ conversationId: "c1",
+ turnId,
+ reason: "stop",
+ ...extra,
+});
+
+describe("initialMetricsState", () => {
+ it("starts empty", () => {
+ const s = initialMetricsState();
+ expect(s.live.size).toBe(0);
+ expect(s.liveOrder).toEqual([]);
+ expect(s.durable.size).toBe(0);
+ expect(s.durableOrder).toEqual([]);
+ });
+});
+
+describe("foldMetricsEvent", () => {
+ it("folds per-step usage by stepId into a turn", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(1);
+ expect(ordered[0]?.turnId).toBe("t1");
+ expect(ordered[0]?.steps).toHaveLength(2);
+ expect(ordered[0]?.steps[0]?.stepId).toBe("s1");
+ expect(ordered[0]?.steps[0]?.usage).toEqual({ inputTokens: 100, outputTokens: 50 });
+ expect(ordered[0]?.steps[1]?.stepId).toBe("s2");
+ expect(ordered[0]?.steps[1]?.usage).toEqual({ inputTokens: 200, outputTokens: 80 });
+ });
+
+ it("folds step-complete timing and merges with same-step usage", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(
+ s,
+ stepCompleteEvent("t1", "s1", { ttftMs: 200, decodeMs: 800, genTotalMs: 1000 }),
+ );
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(1);
+ const step = ordered[0]?.steps[0];
+ expect(step?.usage).toEqual({ inputTokens: 100, outputTokens: 50 });
+ expect(step?.ttftMs).toBe(200);
+ expect(step?.decodeMs).toBe(800);
+ expect(step?.genTotalMs).toBe(1000);
+ });
+
+ it("step-complete before usage defaults usage to zeros", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1", { genTotalMs: 500 }));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ const step = ordered[0]?.steps[0];
+ expect(step?.usage).toEqual({ inputTokens: 0, outputTokens: 0 });
+ expect(step?.genTotalMs).toBe(500);
+ });
+
+ it("done sets durationMs and aggregate usage", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(
+ s,
+ doneEvent("t1", {
+ durationMs: 5000,
+ usage: { inputTokens: 300, outputTokens: 150 },
+ }),
+ );
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.total?.durationMs).toBe(5000);
+ expect(ordered[0]?.total?.usage).toEqual({ inputTokens: 300, outputTokens: 150 });
+ });
+
+ it("aggregate usage sums steps when done.usage absent", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.total?.usage).toEqual({ inputTokens: 300, outputTokens: 130 });
+ });
+
+ it("aggregate usage includes cache only when a step had cache", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, {
+ type: "usage",
+ conversationId: "c1",
+ turnId: "t1",
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50, cacheReadTokens: 30 },
+ });
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.total?.usage.cacheReadTokens).toBe(30);
+ expect(ordered[0]?.total?.usage.cacheWriteTokens).toBeUndefined();
+ });
+
+ it("tolerates missing clock (no genTotalMs/ttft/decode)", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ const step = ordered[0]?.steps[0];
+ expect(step?.ttftMs).toBeUndefined();
+ expect(step?.decodeMs).toBeUndefined();
+ expect(step?.genTotalMs).toBeUndefined();
+ expect(ordered[0]?.total?.durationMs).toBeUndefined();
+ });
+
+ it("usage without stepId does not create a turn", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(0);
+ });
+
+ it("ignores non-metrics events", () => {
+ const s = initialMetricsState();
+ const next = foldMetricsEvent(s, {
+ type: "status",
+ conversationId: "c1",
+ status: "running",
+ });
+ expect(next).toBe(s);
+ });
+
+ it("preserves first-seen order of steps", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 10, 5, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ s = foldMetricsEvent(s, usageEvent("t1", 20, 8, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.steps[0]?.stepId).toBe("s2");
+ expect(ordered[0]?.steps[1]?.stepId).toBe("s1");
+ });
+
+ it("preserves first-seen order of turns", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t2", 10, 5, "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 20, 8, "s1"));
+ s = foldMetricsEvent(s, doneEvent("t2"));
+ s = foldMetricsEvent(s, doneEvent("t1"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.turnId).toBe("t2");
+ expect(ordered[1]?.turnId).toBe("t1");
+ });
+});
+
+describe("selectOrderedTurnMetrics", () => {
+ it("durable wins over live by turnId, live-done appended last", () => {
+ let s = initialMetricsState();
+
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, usageEvent("t2", 200, 80, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t2"));
+
+ s = applyDurableMetrics(s, [
+ {
+ turnId: "t1",
+ usage: { inputTokens: 999, outputTokens: 999 },
+ durationMs: 3000,
+ steps: [
+ {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 999, outputTokens: 999 },
+ genTotalMs: 3000,
+ },
+ ],
+ },
+ ]);
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(2);
+ expect(ordered[0]?.turnId).toBe("t1");
+ expect(ordered[0]?.total?.usage.inputTokens).toBe(999);
+ expect(ordered[0]?.total?.durationMs).toBe(3000);
+ expect(ordered[1]?.turnId).toBe("t2");
+ expect(ordered[1]?.total?.durationMs).toBeUndefined();
+ });
+
+ it("empty state returns empty", () => {
+ const s = initialMetricsState();
+ expect(selectOrderedTurnMetrics(s)).toEqual([]);
+ });
+
+ it("selectOrderedTurnMetrics: in-flight turn exposes only completed steps and total=null", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1", { genTotalMs: 1000 }));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(1);
+ expect(ordered[0]?.turnId).toBe("t1");
+ expect(ordered[0]?.steps).toHaveLength(1);
+ expect(ordered[0]?.steps[0]?.stepId).toBe("s1");
+ expect(ordered[0]?.total).toBeNull();
+ });
+
+ it("selectOrderedTurnMetrics: a turn with no complete step and not done is omitted", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(0);
+ });
+
+ it("selectOrderedTurnMetrics: after done, total is present", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1", { genTotalMs: 1000 }));
+ s = foldMetricsEvent(s, doneEvent("t1", { durationMs: 2000 }));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(1);
+ expect(ordered[0]?.turnId).toBe("t1");
+ expect(ordered[0]?.total?.durationMs).toBe(2000);
+ expect(ordered[0]?.steps).toHaveLength(1);
+ });
+
+ it("step-complete marks the step complete", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1", { genTotalMs: 500 }));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(1);
+ expect(ordered[0]?.steps).toHaveLength(1);
+ expect(ordered[0]?.steps[0]?.stepId).toBe("s1");
+ expect(ordered[0]?.steps[0]?.genTotalMs).toBe(500);
+ });
+
+ it("selectOrderedTurnMetrics: durable turn → steps + total present", () => {
+ let s = initialMetricsState();
+ s = applyDurableMetrics(s, [
+ {
+ turnId: "t1",
+ usage: { inputTokens: 300, outputTokens: 150 },
+ durationMs: 5000,
+ steps: [
+ {
+ stepId: "s1" as StepId,
+ usage: { inputTokens: 100, outputTokens: 50 },
+ genTotalMs: 1000,
+ },
+ {
+ stepId: "s2" as StepId,
+ usage: { inputTokens: 200, outputTokens: 100 },
+ genTotalMs: 2000,
+ },
+ ],
+ },
+ ]);
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered).toHaveLength(1);
+ expect(ordered[0]?.turnId).toBe("t1");
+ expect(ordered[0]?.steps).toHaveLength(2);
+ expect(ordered[0]?.steps[0]?.stepId).toBe("s1");
+ expect(ordered[0]?.steps[1]?.stepId).toBe("s2");
+ expect(ordered[0]?.total?.usage.inputTokens).toBe(300);
+ expect(ordered[0]?.total?.durationMs).toBe(5000);
+ });
+});
+
+describe("applyDurableMetrics", () => {
+ it("stores durable turns in order", () => {
+ let s = initialMetricsState();
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 10, outputTokens: 5 }, steps: [] },
+ { turnId: "t2", usage: { inputTokens: 20, outputTokens: 8 }, steps: [] },
+ ]);
+ expect(s.durableOrder).toEqual(["t1", "t2"]);
+ expect(s.durable.size).toBe(2);
+ });
+
+ it("is idempotent for same turnId", () => {
+ let s = initialMetricsState();
+ const turn = {
+ turnId: "t1",
+ usage: { inputTokens: 10, outputTokens: 5 },
+ steps: [],
+ };
+ s = applyDurableMetrics(s, [turn]);
+ s = applyDurableMetrics(s, [turn]);
+ expect(s.durableOrder).toEqual(["t1"]);
+ expect(s.durable.size).toBe(1);
+ });
+
+ it("overwrites durable turn data for same turnId", () => {
+ let s = initialMetricsState();
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 10, outputTokens: 5 }, steps: [] },
+ ]);
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 99, outputTokens: 99 }, steps: [] },
+ ]);
+ expect(s.durable.get("t1")?.usage.inputTokens).toBe(99);
+ });
+});
+
+describe("contextSize / selectCurrentContextSize", () => {
+ it("live done carries contextSize onto the turn total", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 1234 }));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.total?.contextSize).toBe(1234);
+ expect(selectCurrentContextSize(s)).toBe(1234);
+ });
+
+ it("contextSize is NOT the aggregate usage sum (multi-step turn)", () => {
+ let s = initialMetricsState();
+ // Two steps: usage sums to 300 in / 130 out = 430, but contextSize is the
+ // backend-stamped final-step occupancy, independent of the sum.
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 250 }));
+
+ const ordered = selectOrderedTurnMetrics(s);
+ expect(ordered[0]?.total?.usage).toEqual({ inputTokens: 300, outputTokens: 130 });
+ expect(ordered[0]?.total?.contextSize).toBe(250);
+ expect(selectCurrentContextSize(s)).toBe(250);
+ });
+
+ it("persisted (durable) contextSize is preserved and selected", () => {
+ let s = initialMetricsState();
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 10, outputTokens: 5 }, steps: [], contextSize: 4096 },
+ ]);
+ expect(s.durable.get("t1")?.contextSize).toBe(4096);
+ expect(selectCurrentContextSize(s)).toBe(4096);
+ });
+
+ it("selectCurrentContextSize returns the LATEST turn's value", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 100 }));
+ s = foldMetricsEvent(s, doneEvent("t2", { contextSize: 900 }));
+ expect(selectCurrentContextSize(s)).toBe(900);
+ });
+
+ it("selectCurrentContextSize skips a later turn that lacks contextSize", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 700 }));
+ // t2 finishes but the provider reported no per-step usage → no contextSize.
+ s = foldMetricsEvent(s, doneEvent("t2"));
+ expect(selectCurrentContextSize(s)).toBe(700);
+ });
+
+ it("selectCurrentContextSize is undefined (not 0) when nothing reported", () => {
+ let s = initialMetricsState();
+ expect(selectCurrentContextSize(s)).toBeUndefined();
+ s = foldMetricsEvent(s, doneEvent("t1"));
+ expect(selectCurrentContextSize(s)).toBeUndefined();
+ });
+
+ it("durable contextSize wins over live for a shared turnId", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 111 }));
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 1, outputTokens: 1 }, steps: [], contextSize: 222 },
+ ]);
+ expect(selectCurrentContextSize(s)).toBe(222);
+ });
+
+ it("in-flight turn updates context size after the first step completes", () => {
+ // Before the requirement: an in-flight turn had total=null so its step usage
+ // was ignored until `done`. Now the latest step's input+output is used.
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 5000, 200, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+
+ // Still generating (no done) — context = step 1 input+output = 5200.
+ expect(selectCurrentContextSize(s)).toBe(5200);
+ });
+
+ it("in-flight turn updates progressively as each step reports usage", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 5000, 200, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ expect(selectCurrentContextSize(s)).toBe(5200);
+
+ // Step 2 reports usage mid-stream (before its step-complete): each step's
+ // input already includes all prior context, so the last step's input+output
+ // is the current occupancy.
+ s = foldMetricsEvent(s, usageEvent("t1", 5200, 150, "s2"));
+ expect(selectCurrentContextSize(s)).toBe(5350);
+
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ expect(selectCurrentContextSize(s)).toBe(5350);
+ });
+
+ it("in-flight context size is the latest step with usage, NOT the aggregate sum", () => {
+ // Mirrors the finalized-turn test: contextSize is the FINAL step's
+ // input+output, not the sum across steps (which would overcount a
+ // multi-step turn because every step re-prefills the growing prompt).
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 200, 80, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ // Aggregate would be 300+130=430; the latest step is 200+80=280.
+ expect(selectCurrentContextSize(s)).toBe(280);
+ });
+
+ it("in-flight turn with a step-complete but no usage falls back to older turn", () => {
+ // step-complete before usage → the step has no usage yet, so the in-flight
+ // turn exposes no context size and the display falls back to the prior
+ // finalized turn's value (never 0).
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 700 }));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1", { genTotalMs: 500 }));
+
+ expect(selectCurrentContextSize(s)).toBe(700);
+ });
+
+ it("in-flight turn with no steps/usage returns undefined (falls back)", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 700 }));
+ // t2 just started — no usage, no complete step — omitted entirely.
+ s = foldMetricsEvent(s, { type: "turn-start", conversationId: "c1", turnId: "t2" });
+ expect(selectCurrentContextSize(s)).toBe(700);
+
+ // t2's first step reports usage → the display jumps to t2's live value.
+ s = foldMetricsEvent(s, usageEvent("t2", 800, 10, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1"));
+ expect(selectCurrentContextSize(s)).toBe(810);
+ });
+
+ it("done finalizes the in-flight progressive value with the authoritative contextSize", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 5000, 200, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ expect(selectCurrentContextSize(s)).toBe(5200);
+
+ s = foldMetricsEvent(s, usageEvent("t1", 5200, 150, "s2"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ expect(selectCurrentContextSize(s)).toBe(5350);
+
+ // done stamps the authoritative contextSize (the final step's input+output).
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 5350 }));
+ expect(selectCurrentContextSize(s)).toBe(5350);
+ });
+
+ it("in-flight context size excludes cache tokens (they are a subset of inputTokens)", () => {
+ // cacheReadTokens / cacheWriteTokens are portions of inputTokens already
+ // counted — adding them would double-count. Only input+output is occupancy.
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, {
+ type: "usage",
+ conversationId: "c1",
+ turnId: "t1",
+ stepId: "s1" as StepId,
+ usage: {
+ inputTokens: 5000,
+ outputTokens: 200,
+ cacheReadTokens: 4000,
+ cacheWriteTokens: 1000,
+ },
+ });
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ // 5000+200=5200, NOT 9200 (with cacheRead) or 10200 (with both).
+ expect(selectCurrentContextSize(s)).toBe(5200);
+ });
+
+ it("multiple in-flight turns: the newest turn's live value wins", () => {
+ let s = initialMetricsState();
+ // t1 (older) in-flight with one completed step → 5200.
+ s = foldMetricsEvent(s, usageEvent("t1", 5000, 200, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ // t2 (newer, seen later → last in liveOrder) in-flight → 8000.
+ s = foldMetricsEvent(s, usageEvent("t2", 7800, 200, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1"));
+ expect(selectCurrentContextSize(s)).toBe(8000);
+ });
+
+ it("out-of-order step IDs: usage for step 2 before step 1's step-complete still scans newest-first", () => {
+ // stepOrder is FIRST-SEEN: s1 (its usage arrived first), then s2. So s2 is
+ // the newest step regardless of when each step's step-complete arrives.
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 5000, 200, "s1"));
+ s = foldMetricsEvent(s, usageEvent("t1", 5200, 150, "s2"));
+ // Neither step complete yet → the turn is omitted (no complete step), so the
+ // display can't update until the first step completes.
+ expect(selectCurrentContextSize(s)).toBeUndefined();
+
+ // s1 completes AFTER s2's usage was reported. The turn is now visible; the
+ // newest-first scan picks s2 (the later step), not s1 (the just-completed one).
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ expect(selectCurrentContextSize(s)).toBe(5350);
+
+ // s2 completes — still s2, unchanged.
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ expect(selectCurrentContextSize(s)).toBe(5350);
+ });
+
+ it("done turn without contextSize falls back to an older turn (even with step usage)", () => {
+ // Contract lock-in: a done turn's step usage is NOT consulted for the
+ // context display — only its authoritative total.contextSize is. When that
+ // is absent, the display falls back to the next older finalized turn rather
+ // than synthesizing a value from the step usage.
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 700 }));
+ // t2 done WITH step usage but NO done.contextSize (edge case: the done event
+ // omitted contextSize despite per-step usage).
+ s = foldMetricsEvent(s, usageEvent("t2", 800, 10, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t2"));
+ expect(selectCurrentContextSize(s)).toBe(700);
+ });
+
+ it("in-flight context size skips a step with unsafe usage (NaN / negative)", () => {
+ // A corrupt provider report must never reach the status bar. The newest
+ // step with invalid counters is skipped, falling back to the prior valid one.
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 5000, 200, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ // s2 reports NaN input (e.g. a non-numeric provider field coerced).
+ s = foldMetricsEvent(s, {
+ type: "usage",
+ conversationId: "c1",
+ turnId: "t1",
+ stepId: "s2" as StepId,
+ usage: { inputTokens: Number.NaN, outputTokens: 150 },
+ });
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s2"));
+ // s2 skipped (NaN) → falls back to s1's 5200, NOT NaN.
+ expect(selectCurrentContextSize(s)).toBe(5200);
+
+ // Negative tokens are likewise skipped.
+ s = foldMetricsEvent(s, {
+ type: "usage",
+ conversationId: "c1",
+ turnId: "t1",
+ stepId: "s3" as StepId,
+ usage: { inputTokens: -10, outputTokens: 5 },
+ });
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s3"));
+ expect(selectCurrentContextSize(s)).toBe(5200);
+ });
+});
+
+describe("applyDurableMetrics pruning", () => {
+ it("prunes a live turn once durable data covers it (no unbounded growth)", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 150 }));
+ expect(s.live.has("t1")).toBe(true);
+ expect(s.liveOrder).toContain("t1");
+
+ s = applyDurableMetrics(s, [
+ {
+ turnId: "t1",
+ usage: { inputTokens: 100, outputTokens: 50 },
+ steps: [{ stepId: "s1" as StepId, usage: { inputTokens: 100, outputTokens: 50 } }],
+ contextSize: 150,
+ },
+ ]);
+ // The live copy is gone; the durable (authoritative) entry replaces it.
+ expect(s.live.has("t1")).toBe(false);
+ expect(s.liveOrder).not.toContain("t1");
+ expect(s.durable.has("t1")).toBe(true);
+ // The display still reads the durable value atomically (no gap).
+ expect(selectCurrentContextSize(s)).toBe(150);
+ });
+
+ it("prunes only the turns present in the durable batch (leaves other live turns)", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t1", 100, 50, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t1", "s1"));
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 150 }));
+ // t2 still in flight — must NOT be pruned when only t1 seals.
+ s = foldMetricsEvent(s, usageEvent("t2", 800, 10, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1"));
+
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 100, outputTokens: 50 }, steps: [], contextSize: 150 },
+ ]);
+ expect(s.live.has("t1")).toBe(false);
+ expect(s.live.has("t2")).toBe(true);
+ expect(s.liveOrder).toEqual(["t2"]);
+ // The newest (in-flight) turn's live value still wins.
+ expect(selectCurrentContextSize(s)).toBe(810);
+ });
+
+ it("is a no-op when no incoming turn is live (no live mutation)", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, usageEvent("t2", 800, 10, "s1"));
+ s = foldMetricsEvent(s, stepCompleteEvent("t2", "s1"));
+ const before = s;
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 1, outputTokens: 1 }, steps: [] },
+ ]);
+ // t1 was never live → the live map/order are unchanged (same reference).
+ expect(s.live).toBe(before.live);
+ expect(s.liveOrder).toBe(before.liveOrder);
+ // t1 (durable) is older; the in-flight t2 still wins.
+ expect(selectCurrentContextSize(s)).toBe(810);
+ });
+
+ it("durable wins over live for a shared turnId (pruned live no longer consulted)", () => {
+ let s = initialMetricsState();
+ s = foldMetricsEvent(s, doneEvent("t1", { contextSize: 111 }));
+ s = applyDurableMetrics(s, [
+ { turnId: "t1", usage: { inputTokens: 1, outputTokens: 1 }, steps: [], contextSize: 222 },
+ ]);
+ // The live (111) copy is pruned; only durable (222) remains.
+ expect(s.live.has("t1")).toBe(false);
+ expect(selectCurrentContextSize(s)).toBe(222);
+ });
+});
diff --git a/src/core/metrics/reducer.ts b/src/core/metrics/reducer.ts
new file mode 100644
index 0000000..39fc5ee
--- /dev/null
+++ b/src/core/metrics/reducer.ts
@@ -0,0 +1,345 @@
+import type { AgentEvent, StepId, StepMetrics, TurnMetrics, Usage } from "@dispatch/wire";
+import type { BuildingStep, LiveTurn, MetricsState, TurnMetricsEntry } from "./types";
+
+function sumStepUsages(steps: readonly BuildingStep[]): Usage {
+ let inputTokens = 0;
+ let outputTokens = 0;
+ let hasCacheRead = false;
+ let hasCacheWrite = false;
+ let cacheReadTokens = 0;
+ let cacheWriteTokens = 0;
+
+ for (const step of steps) {
+ if (step.usage === undefined) continue;
+ inputTokens += step.usage.inputTokens;
+ outputTokens += step.usage.outputTokens;
+ if (step.usage.cacheReadTokens !== undefined && step.usage.cacheReadTokens > 0) {
+ hasCacheRead = true;
+ cacheReadTokens += step.usage.cacheReadTokens;
+ }
+ if (step.usage.cacheWriteTokens !== undefined && step.usage.cacheWriteTokens > 0) {
+ hasCacheWrite = true;
+ cacheWriteTokens += step.usage.cacheWriteTokens;
+ }
+ }
+
+ const base: Usage = { inputTokens, outputTokens };
+ if (hasCacheRead) {
+ (base as { cacheReadTokens?: number }).cacheReadTokens = cacheReadTokens;
+ }
+ if (hasCacheWrite) {
+ (base as { cacheWriteTokens?: number }).cacheWriteTokens = cacheWriteTokens;
+ }
+ return base;
+}
+
+function buildingStepToMetrics(bs: BuildingStep): StepMetrics {
+ const usage: Usage = bs.usage ?? { inputTokens: 0, outputTokens: 0 };
+ const base: StepMetrics = { stepId: bs.stepId as StepId, usage };
+ if (bs.ttftMs !== undefined) {
+ (base as { ttftMs?: number }).ttftMs = bs.ttftMs;
+ }
+ if (bs.decodeMs !== undefined) {
+ (base as { decodeMs?: number }).decodeMs = bs.decodeMs;
+ }
+ if (bs.genTotalMs !== undefined) {
+ (base as { genTotalMs?: number }).genTotalMs = bs.genTotalMs;
+ }
+ return base;
+}
+
+function getStep(lt: LiveTurn, id: string): BuildingStep {
+ const step = lt.stepMap.get(id);
+ if (step === undefined) throw new Error(`Missing step ${id} in live turn`);
+ return step;
+}
+
+function liveTurnToMetrics(lt: LiveTurn): TurnMetrics {
+ const buildingSteps = lt.stepOrder.map((id) => getStep(lt, id));
+ const steps = buildingSteps.map((bs) => buildingStepToMetrics(bs));
+ const usage = lt.doneUsage ?? sumStepUsages(buildingSteps);
+ const base: TurnMetrics = { turnId: lt.turnId, usage, steps };
+ if (lt.durationMs !== undefined) {
+ (base as { durationMs?: number }).durationMs = lt.durationMs;
+ }
+ if (lt.doneContextSize !== undefined) {
+ (base as { contextSize?: number }).contextSize = lt.doneContextSize;
+ }
+ return base;
+}
+
+/**
+ * A step's contribution to the live context size: `inputTokens + outputTokens`,
+ * or `undefined` when the step has no usage yet OR its counters are not safe to
+ * sum (non-finite / negative — defensive: a corrupt provider report must never
+ * reach the status bar as NaN/Infinity). Cache tokens are deliberately NOT
+ * included: `cacheReadTokens` / `cacheWriteTokens` are a SUBSET of
+ * `inputTokens`, so adding them would double-count.
+ */
+function stepContextSize(usage: Usage | undefined): number | undefined {
+ if (usage === undefined) return undefined;
+ const { inputTokens, outputTokens } = usage;
+ if (!Number.isFinite(inputTokens) || !Number.isFinite(outputTokens)) return undefined;
+ if (inputTokens < 0 || outputTokens < 0) return undefined;
+ return inputTokens + outputTokens;
+}
+
+/**
+ * The context size an IN-FLIGHT (not-done) turn occupies right now — for
+ * progressive display DURING a turn (before it seals), so the indicator updates
+ * after each step instead of waiting for `done`.
+ *
+ * CONTRACT: only call this on a turn whose `done` event has NOT arrived (the
+ * caller, `selectCurrentContextSize`, reaches it solely for entries with
+ * `total === null`, i.e. `lt.done === false`). Finalized turns use their
+ * authoritative `contextSize` instead; `doneContextSize` is read on the
+ * `total` path, never here.
+ *
+ * Returns the most recent step WITH USABLE USAGE's `inputTokens + outputTokens`
+ * (scanning newest → oldest by first-seen step order): each step's input
+ * already includes all prior context (the prompt is re-prefilled every step), so
+ * the last step's input+output is the true occupancy — the same definition
+ * `TurnDoneEvent.contextSize` stamps at turn end. A just-reported step's usage
+ * wins immediately, even mid-stream. Steps with no usage or unsafe usage are
+ * skipped, falling back to the next older usable step. `undefined` when no step
+ * has reported usable usage yet.
+ */
+function liveTurnContextSize(lt: LiveTurn): number | undefined {
+ for (let i = lt.stepOrder.length - 1; i >= 0; i--) {
+ const step = lt.stepMap.get(lt.stepOrder[i] ?? "");
+ const ctx = stepContextSize(step?.usage);
+ if (ctx !== undefined) return ctx;
+ }
+ return undefined;
+}
+
+function ensureLiveTurn(state: MetricsState, turnId: string): [MetricsState, LiveTurn] {
+ const existing = state.live.get(turnId);
+ if (existing !== undefined) return [state, existing];
+
+ const newTurn: LiveTurn = {
+ turnId,
+ done: false,
+ durationMs: undefined,
+ doneUsage: undefined,
+ doneContextSize: undefined,
+ stepMap: new Map(),
+ stepOrder: [],
+ };
+ const newLive = new Map(state.live);
+ newLive.set(turnId, newTurn);
+ return [{ ...state, live: newLive, liveOrder: [...state.liveOrder, turnId] }, newTurn];
+}
+
+function upsertStep(lt: LiveTurn, stepId: string, update: Partial<BuildingStep>): LiveTurn {
+ const existing = lt.stepMap.get(stepId);
+ if (existing !== undefined) {
+ const merged: BuildingStep = {
+ stepId,
+ usage: update.usage ?? existing.usage,
+ ttftMs: update.ttftMs ?? existing.ttftMs,
+ decodeMs: update.decodeMs ?? existing.decodeMs,
+ genTotalMs: update.genTotalMs ?? existing.genTotalMs,
+ complete: update.complete ?? existing.complete,
+ };
+ const newMap = new Map(lt.stepMap);
+ newMap.set(stepId, merged);
+ return { ...lt, stepMap: newMap };
+ }
+
+ const fresh: BuildingStep = {
+ stepId,
+ usage: update.usage,
+ ttftMs: update.ttftMs,
+ decodeMs: update.decodeMs,
+ genTotalMs: update.genTotalMs,
+ complete: update.complete ?? false,
+ };
+ const newMap = new Map(lt.stepMap);
+ newMap.set(stepId, fresh);
+ return { ...lt, stepMap: newMap, stepOrder: [...lt.stepOrder, stepId] };
+}
+
+/** The initial empty metrics state. */
+export function initialMetricsState(): MetricsState {
+ return {
+ live: new Map(),
+ liveOrder: [],
+ durable: new Map(),
+ durableOrder: [],
+ };
+}
+
+/**
+ * Fold one live AgentEvent into the metrics state.
+ *
+ * - `usage` with `stepId`: upsert that step's usage.
+ * - `usage` without `stepId`: ignored.
+ * - `step-complete`: upsert that step's timing; default usage to zeros if absent.
+ * - `done`: set turn's `durationMs`, optional aggregate `usage`, and optional `contextSize`.
+ * - All other event types: return state unchanged.
+ */
+export function foldMetricsEvent(state: MetricsState, event: AgentEvent): MetricsState {
+ switch (event.type) {
+ case "usage": {
+ if (event.stepId === undefined) return state;
+ const [s1, lt] = ensureLiveTurn(state, event.turnId);
+ const updated = upsertStep(lt, event.stepId, { usage: event.usage });
+ const newLive = new Map(s1.live);
+ newLive.set(event.turnId, updated);
+ return { ...s1, live: newLive };
+ }
+
+ case "step-complete": {
+ const [s1, lt] = ensureLiveTurn(state, event.turnId);
+ const updated = upsertStep(lt, event.stepId, {
+ ttftMs: event.ttftMs,
+ decodeMs: event.decodeMs,
+ genTotalMs: event.genTotalMs,
+ complete: true,
+ });
+ const newLive = new Map(s1.live);
+ newLive.set(event.turnId, updated);
+ return { ...s1, live: newLive };
+ }
+
+ case "done": {
+ const [s1, lt] = ensureLiveTurn(state, event.turnId);
+ const updated: LiveTurn = {
+ ...lt,
+ done: true,
+ durationMs: event.durationMs ?? lt.durationMs,
+ doneUsage: event.usage ?? lt.doneUsage,
+ doneContextSize: event.contextSize ?? lt.doneContextSize,
+ };
+ const newLive = new Map(s1.live);
+ newLive.set(event.turnId, updated);
+ return { ...s1, live: newLive };
+ }
+
+ default:
+ return state;
+ }
+}
+
+/**
+ * Store durable (sealed) metrics from the backend. These win over live data
+ * for any shared `turnId`.
+ *
+ * Once durable (authoritative) data covers a turn, its live (in-memory) copy
+ * is REDUNDANT and is pruned from `state.live` / `liveOrder` so the live map
+ * doesn't grow unbounded over a long conversation. There is no display gap:
+ * the durable entry replaces the live one atomically in the same fold, and
+ * `selectOrderedTurnMetrics` / `selectCurrentContextSize` read durable for it.
+ */
+export function applyDurableMetrics(
+ state: MetricsState,
+ turns: readonly TurnMetrics[],
+): MetricsState {
+ const newDurable = new Map(state.durable);
+ const newDurableOrder = [...state.durableOrder];
+ const prunedIds = new Set<string>();
+ for (const turn of turns) {
+ if (!newDurable.has(turn.turnId)) {
+ newDurableOrder.push(turn.turnId);
+ }
+ newDurable.set(turn.turnId, turn);
+ if (state.live.has(turn.turnId)) prunedIds.add(turn.turnId);
+ }
+
+ if (prunedIds.size === 0) {
+ return { ...state, durable: newDurable, durableOrder: newDurableOrder };
+ }
+
+ const newLive = new Map(state.live);
+ for (const id of prunedIds) newLive.delete(id);
+ const newLiveOrder = state.liveOrder.filter((id) => !prunedIds.has(id));
+
+ return {
+ ...state,
+ live: newLive,
+ liveOrder: newLiveOrder,
+ durable: newDurable,
+ durableOrder: newDurableOrder,
+ };
+}
+
+/**
+ * Select the merged ordered list of turn metrics entries.
+ * Durable turns come first (in their order), then any live turns whose
+ * `turnId` is not in durable (in live first-seen order).
+ *
+ * Each entry contains the completed steps so far and an optional total
+ * (null until the turn is finalized via `done` or durable data).
+ * Live turns with no completed steps and not done are omitted.
+ */
+export function selectOrderedTurnMetrics(state: MetricsState): readonly TurnMetricsEntry[] {
+ const result: TurnMetricsEntry[] = [];
+ const seen = new Set<string>();
+
+ for (const turnId of state.durableOrder) {
+ const tm = state.durable.get(turnId);
+ if (tm !== undefined) {
+ result.push({ turnId, steps: tm.steps, total: tm });
+ seen.add(turnId);
+ }
+ }
+
+ for (const turnId of state.liveOrder) {
+ if (seen.has(turnId)) continue;
+ const lt = state.live.get(turnId);
+ if (lt === undefined) continue;
+
+ const completeSteps = lt.stepOrder
+ .map((id) => lt.stepMap.get(id))
+ .filter((s): s is BuildingStep => s?.complete === true)
+ .map((s) => buildingStepToMetrics(s));
+
+ if (completeSteps.length === 0 && !lt.done) continue;
+
+ result.push({
+ turnId,
+ steps: completeSteps,
+ total: lt.done ? liveTurnToMetrics(lt) : null,
+ });
+ }
+
+ return result;
+}
+
+/**
+ * Select the conversation's CURRENT context size — the tokens it occupies right
+ * now. Per the wire contract a client reads the LATEST turn's `contextSize`; we
+ * scan the merged ordered turns NEWEST → OLDEST and return the first DEFINED
+ * value.
+ *
+ * For a FINALIZED turn (`done` event or durable data) we use its authoritative
+ * `contextSize`. For an IN-FLIGHT (not-done) turn we compute it PROGRESSIVELY
+ * from the most recent step WITH USAGE — its `inputTokens + outputTokens` is the
+ * current occupancy (mirroring `TurnDoneEvent.contextSize`'s definition) — so
+ * the indicator updates after each step completes instead of waiting for the
+ * turn to seal. An in-flight turn with no step usage yet is skipped, falling
+ * back to the next older finalized turn.
+ *
+ * Returns `undefined` ("unknown") when no turn carries a context size — the
+ * caller renders a placeholder, NEVER `0`. Durable (sealed) data wins over
+ * live for a shared `turnId` (it is the persisted, authoritative value).
+ */
+export function selectCurrentContextSize(state: MetricsState): number | undefined {
+ const ordered = selectOrderedTurnMetrics(state);
+ for (let i = ordered.length - 1; i >= 0; i--) {
+ const entry = ordered[i];
+ if (entry === undefined) continue;
+ if (entry.total !== null) {
+ if (entry.total.contextSize !== undefined) return entry.total.contextSize;
+ continue;
+ }
+ // In-flight turn: progressive context size from the latest step with usage.
+ const lt = state.live.get(entry.turnId);
+ if (lt !== undefined) {
+ const live = liveTurnContextSize(lt);
+ if (live !== undefined) return live;
+ }
+ }
+ return undefined;
+}
diff --git a/src/core/metrics/types.ts b/src/core/metrics/types.ts
new file mode 100644
index 0000000..84d1904
--- /dev/null
+++ b/src/core/metrics/types.ts
@@ -0,0 +1,97 @@
+import type { StepMetrics, TurnMetrics, Usage } from "@dispatch/wire";
+import type { RenderGroup } from "../chunks";
+
+export type { StepMetrics, TurnMetrics };
+
+/** A step being built from live events (may be incomplete). */
+export interface BuildingStep {
+ readonly stepId: string;
+ readonly usage: Usage | undefined;
+ readonly ttftMs: number | undefined;
+ readonly decodeMs: number | undefined;
+ readonly genTotalMs: number | undefined;
+ readonly complete: boolean;
+}
+
+/** A turn being built from live events (in-flight). */
+export interface LiveTurn {
+ readonly turnId: string;
+ readonly done: boolean;
+ readonly durationMs: number | undefined;
+ readonly doneUsage: Usage | undefined;
+ /**
+ * Context size carried on the turn's `done` event (the turn's FINAL step
+ * `inputTokens + outputTokens` — current context occupancy). `undefined` when
+ * the provider reported no per-step usage; never coerced to `0`.
+ */
+ readonly doneContextSize: number | undefined;
+ readonly stepMap: ReadonlyMap<string, BuildingStep>;
+ readonly stepOrder: readonly string[];
+}
+
+/**
+ * Reducer state for per-turn / per-step token + timing metrics.
+ *
+ * - `live`: in-flight turns keyed by `turnId` in FIRST-SEEN order.
+ * - `durable`: sealed turns keyed by `turnId` in the order they arrived.
+ */
+export interface MetricsState {
+ readonly live: ReadonlyMap<string, LiveTurn>;
+ readonly liveOrder: readonly string[];
+ readonly durable: ReadonlyMap<string, TurnMetrics>;
+ readonly durableOrder: readonly string[];
+}
+
+/** Per-turn placement entry: completed steps so far + optional turn total. */
+export interface TurnMetricsEntry {
+ readonly turnId: string;
+ readonly steps: readonly StepMetrics[];
+ readonly total: TurnMetrics | null;
+}
+
+/** A row in the interleaved transcript: a render group, per-step metrics, or turn metrics. */
+export type MetricsRow =
+ | { readonly kind: "group"; readonly group: RenderGroup }
+ | { readonly kind: "step-metrics"; readonly step: StepMetrics; readonly index: number }
+ | {
+ readonly kind: "turn-metrics";
+ readonly turn: TurnMetrics;
+ /** 1-based turn number (the entry's position in the metrics array + 1). */
+ readonly turnNumber: number;
+ /** Cumulative usage across all finalized turns up to and including this one. */
+ readonly cumulativeUsage: Usage;
+ /**
+ * Usage of the most recent EARLIER finalized turn, or `null` when this is the
+ * first finalized turn. The baseline for cross-turn retention (expected cache).
+ */
+ readonly prevTurnUsage: Usage | null;
+ };
+
+/** Formatted cache hit-rate view: percentage + colour severity + hit flag. */
+export interface CacheRateView {
+ /** Cache hit rate as a 0..100 integer percentage (`cacheReadTokens / inputTokens`). */
+ readonly pct: number;
+ /** Colour severity for a badge (maps to DaisyUI `badge-{level}`). */
+ readonly level: "success" | "warning" | "error";
+ /** Whether any input tokens were served from cache. */
+ readonly isHit: boolean;
+}
+
+/** Formatted per-step view for display. */
+export interface StepMetricsView {
+ readonly label: string;
+ readonly tokensLabel: string;
+ readonly tps: string | null;
+ readonly ttft: string | null;
+ readonly decode: string | null;
+ readonly genTotal: string | null;
+}
+
+/** Formatted per-turn view for display. */
+export interface TurnMetricsView {
+ readonly label: string;
+ readonly tokensLabel: string;
+ readonly breakdown: string;
+ readonly tps: string | null;
+ readonly duration: string | null;
+}