From b9de3da8797e54abfad230770dd5bb7ba3399139 Mon Sep 17 00:00:00 2001 From: Yaowei Zheng Date: Wed, 29 Jul 2026 23:30:13 +0800 Subject: [PATCH] fix(web): take the Session's elapsed time from Trace timestamps, so a reload stops restarting it (#124) Co-authored-by: Claude Opus 5 (1M context) --- packages/web/src/api/client.ts | 30 ++++- packages/web/src/api/endpoints.ts | 13 +- .../web/src/lib/omni/stream-controller.ts | 19 ++- packages/web/src/lib/omni/stream-model.ts | 118 ++++++++++++------ packages/web/src/lib/omni/task-stats.ts | 18 ++- packages/web/test/stream-controller.test.ts | 51 +++++++- packages/web/test/stream-model.test.ts | 118 ++++++++++++++---- 7 files changed, 289 insertions(+), 78 deletions(-) diff --git a/packages/web/src/api/client.ts b/packages/web/src/api/client.ts index a7a5c3c..b406702 100644 --- a/packages/web/src/api/client.ts +++ b/packages/web/src/api/client.ts @@ -42,8 +42,29 @@ export interface ApiFetchOptions { query?: Record; } +/** Response metadata a caller may need alongside the parsed body. */ +export interface ApiFetchMeta { + /** + * The server's own clock at the moment it produced the response, read from the HTTP `Date` + * header; null when absent or unparseable. Lets a caller measure a server-side interval + * entirely in server time — differencing it against a server-supplied timestamp cancels any + * client/server clock offset, which a local `Date.now()` cannot do. Whole-second precision + * (RFC 9110 fixes the header's format), so treat it as ±1s. `Date` is CORS-safelisted, so it + * is readable cross-origin too. + */ + serverNowMs: number | null; +} + /** Makes an API request; non-2xx responses uniformly throw ApiError; 204/empty body returns undefined. */ export async function apiFetch(path: string, options: ApiFetchOptions = {}): Promise { + return (await apiFetchWithMeta(path, options)).data; +} + +/** {@link apiFetch} plus the response metadata in {@link ApiFetchMeta}; identical in every other respect. */ +export async function apiFetchWithMeta( + path: string, + options: ApiFetchOptions = {}, +): Promise<{ data: T } & ApiFetchMeta> { let url = path; if (options.query) { const params = new URLSearchParams(); @@ -84,8 +105,11 @@ export async function apiFetch(path: string, options: ApiFetchOptions = {}): throw new ApiError(response.status, code, message); } - if (response.status === 204) return undefined as T; + const headerDate = Date.parse(response.headers.get("date") ?? ""); + const serverNowMs = Number.isFinite(headerDate) ? headerDate : null; + + if (response.status === 204) return { data: undefined as T, serverNowMs }; const text = await response.text(); - if (!text) return undefined as T; - return JSON.parse(text) as T; + if (!text) return { data: undefined as T, serverNowMs }; + return { data: JSON.parse(text) as T, serverNowMs }; } diff --git a/packages/web/src/api/endpoints.ts b/packages/web/src/api/endpoints.ts index 979e949..3d94aeb 100644 --- a/packages/web/src/api/endpoints.ts +++ b/packages/web/src/api/endpoints.ts @@ -71,7 +71,7 @@ import type { VersionResponse, WorkspaceFilesResponse, } from "@prismshadow/penguin-server/api"; -import { apiFetch } from "./client"; +import { apiFetch, apiFetchWithMeta } from "./client"; // Auth & user ----------------------------------------------------------------- @@ -246,8 +246,17 @@ export const patchSession = (sessionId: string, body: SessionPatchRequest) => export const deleteSession = (sessionId: string) => apiFetch(`/api/sessions/${encodeURIComponent(sessionId)}`, { method: "DELETE" }); +/** + * History rebuild. Carries the server's clock at read time (see ApiFetchMeta.serverNowMs) + * alongside the messages: a Task still running has no Trace entry for the event currently in + * flight, so its elapsed can only be measured by differencing this against the Task's first + * message timestamp — both server-side values, so no client clock offset enters the result + * (see pushMessages). + */ export const getMessages = (sessionId: string) => - apiFetch(`/api/sessions/${encodeURIComponent(sessionId)}/messages`); + apiFetchWithMeta( + `/api/sessions/${encodeURIComponent(sessionId)}/messages`, + ).then(({ data, serverNowMs }) => ({ ...data, serverNowMs })); // Task execution, approval, abort, compaction ------------------------------------------------------ diff --git a/packages/web/src/lib/omni/stream-controller.ts b/packages/web/src/lib/omni/stream-controller.ts index 0d66211..dd6963e 100644 --- a/packages/web/src/lib/omni/stream-controller.ts +++ b/packages/web/src/lib/omni/stream-controller.ts @@ -73,8 +73,17 @@ function parseEventId(id: string): { epoch: string; seq: number } | null { } export interface StreamControllerDeps { - /** Fetch history messages (GET /api/sessions/:id/messages), including the live tail while running. */ - loadMessages: () => Promise<{ messages: OmniMessage[]; live?: MessagesLiveTail }>; + /** + * Fetch history messages (GET /api/sessions/:id/messages), including the live tail while + * running. `serverNowMs` is the server's clock at read time (the response's `Date` header); + * omitted/null just costs a running Task's header the time its in-flight event has taken so + * far, which falls back to the Trace's own span (see pushMessages). + */ + loadMessages: () => Promise<{ + messages: OmniMessage[]; + live?: MessagesLiveTail; + serverNowMs?: number | null; + }>; /** Authoritative running state from the stream (covers both the subscription snapshot and transition events). */ onTaskState: (state: SessionStatus) => void; /** Queued follow-up count carried on task_state events (absent on old servers -> 0). */ @@ -183,7 +192,7 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro deps.onQueuedFollowUps?.(ev.queued ?? 0); if (ev.state === "idle") { // Task ended (or the snapshot confirms idle): finalize the current Task's stats; pending approvals have already converged server-side. - notifyTaskIdle(model, now()); + notifyTaskIdle(model); clearPending(); deps.onModelChange(); } @@ -263,14 +272,14 @@ export function createStreamController(deps: StreamControllerDeps): StreamContro const load = async (currentEpoch: number, freshModel?: StreamModel): Promise => { try { - const { messages, live } = await deps.loadMessages(); + const { messages, live, serverNowMs } = await deps.loadMessages(); if (disposed || currentEpoch !== epoch) return; // Rebuild path: make the freshly-built model visible only now, atomically — the old model // stayed on screen throughout the refetch above (see rebuild). Initial load / retry pass no // freshModel and keep operating on the current model. if (freshModel) model = freshModel; const target = model; - pushMessages(target, messages, now()); + pushMessages(target, messages, now(), serverNowMs ?? null); const dedup = buildDedupIndex(messages, 100); // Replay the buffer (events that arrived while fetching history), with dedup; while a // Task runs, the live tail is woven in so the in-progress message is seeded too. diff --git a/packages/web/src/lib/omni/stream-model.ts b/packages/web/src/lib/omni/stream-model.ts index 6622595..f8a8013 100644 --- a/packages/web/src/lib/omni/stream-model.ts +++ b/packages/web/src/lib/omni/stream-model.ts @@ -39,9 +39,9 @@ * session's user side starts a new Task; a Task ends when the live * stream receives task_state:idle (notifyTaskIdle), or — during history * rebuild — when the next Task starts / the stream ends - * (finalizeHistory). Live finalization takes the duration as the larger - * of the local-clock delta and the message-timestamp span (the local - * clock only covers the time since a mid-stream join). A stats row is only added if token_usage occurred during the Task. + * (finalizeHistory). Either way the duration comes from Trace timestamps + * alone and never the local clock, so a round settles to the same figure + * watched live and replayed after a reload. A stats row is only added if token_usage occurred during the Task. * - Overlap-dedup helpers: buildDedupIndex/isDuplicate * judge duplicates by exact match of the envelope JSON; when a complete * message hits the dedup check, discardFragmentFor also discards the corresponding in-flight fragment. @@ -378,13 +378,30 @@ export interface StreamModel { lastAuthFailureMs: number | null; /** Task segmentation state. */ taskOpen: boolean; + /** + * Local-clock instant the running Task's header elapsed ticks from — display + * only; no settled duration is ever derived from it (see finalizeOpenTask, + * which reads Trace timestamps alone). A live stream sets it to the real + * start; a history rebuild would otherwise stamp the page-load instant and + * restart the ticking value from zero on every reload, so pushMessages + * back-dates it by the elapsed already behind the Task — measured entirely + * in server time, so no client/server clock offset reaches it, and taken + * from the server's own clock rather than the Trace's tail so that an event + * still in flight is counted too. + */ taskStartLocalMs: number; + /** The Task's first message timestamp, in SERVER time: both the settled duration and the back-dated live anchor measure from it. */ taskFirstTsMs: number; /** - * The latest timestamp seen among this round's messages, used only as a - * **fallback for the round's end**: only used for a degenerate round that - * has no request_end at all (interrupted before its first Request even - * ran). The normal round-end is taken from taskLastReqEndMs. + * The latest timestamp seen among this round's messages. Two readers: + * - the **fallback for the round's end**, used only for a degenerate + * round that has no request_end at all (interrupted before its first + * Request even ran) — the normal round-end is taken from taskLastReqEndMs; + * - the floor under the anchor pushMessages back-dates taskStartLocalMs + * to, for a round still open when a history rebuild ends. That reader + * fires for every such round, degenerate or not, but only decides the + * anchor when the server's own clock did not come back with the + * response — it cannot see an event still in flight. */ taskLastTsMs: number; /** @@ -566,22 +583,55 @@ function advanceLastTs(model: StreamModel, timestamp: string): void { if (Number.isFinite(ms)) model.lastTsMs = ms; } +/** + * Replay a history rebuild. `serverNowMs` is the server's clock when it produced the + * response (see StreamControllerDeps.loadMessages); null when unavailable. + */ export function pushMessages( model: StreamModel, messages: OmniMessage[], nowMs: number = Date.now(), + serverNowMs: number | null = null, ): void { for (const msg of messages) pushMessage(model, msg, nowMs); + // Re-anchor a Task still open at the end of the replay. Every message in a + // rebuild is fed the same `nowMs`, so startTask stamped taskStartLocalMs + // with the instant the page loaded — and the header's live elapsed, which + // ticks over `now − taskStartLocalMs`, would restart from zero on every + // reload of a running Session. Back-date the anchor by the elapsed already + // behind this Task, so the ticking value resumes where it left off: + // + // now − anchor == (now − loadInstant) + elapsedSoFar + // + // elapsedSoFar is measured in SERVER time and applied to the local clock, so + // a client/server clock offset cancels out and never enters the result. Two + // readings of it, the larger winning: + // - serverNowMs − taskFirstTsMs: the true elapsed, and the only one that + // covers an event still in flight — a tool executing, a Request + // streaming, a compaction running — where nothing has been appended to + // the Trace since it began. Whole-second precision (the `Date` header's + // format), which a chip ticking in whole seconds cannot show. + // - taskLastTsMs − taskFirstTsMs: the span the Trace itself proves. The + // fallback when no `Date` header came back, and a floor under a stale + // one: a cached or intermediary-rewritten reading can only be older + // than the true now, so it can under-report but never overshoot. + // A live stream pushes one message at a time with the real current clock, + // where both readings are still zero at startTask and this is a no-op. + if (model.taskOpen) { + const tracedSpan = model.taskLastTsMs - model.taskFirstTsMs; + const serverSpan = serverNowMs === null ? 0 : serverNowMs - model.taskFirstTsMs; + model.taskStartLocalMs = nowMs - Math.max(0, tracedSpan, serverSpan); + } } -/** The live stream received task_state:idle: finalize the current Task using the local clock. */ -export function notifyTaskIdle(model: StreamModel, nowMs: number = Date.now()): void { - finalizeOpenTask(model, "live", nowMs); +/** The live stream received task_state:idle: finalize the current Task from its Trace timestamps. */ +export function notifyTaskIdle(model: StreamModel): void { + finalizeOpenTask(model); } /** History rebuild is complete (end of stream): finalize the last Task using message timestamps. */ export function finalizeHistory(model: StreamModel): void { - finalizeOpenTask(model, "history"); + finalizeOpenTask(model); } /** Register an approval clicked on this end (so the subsequent approval_decision event is labeled "manual"). */ @@ -620,12 +670,15 @@ export function findToolCard( // --------------------------------------------------------------------------- /** - * Advance this round's "latest timestamp seen among its messages" — used - * only as a **fallback for the round's end** (see taskLastReqEndMs). The - * normal round-end is set by request_end; this only guarantees a usable - * upper bound for a degenerate round with no request_end at all - * (interrupted before its first Request even ran). Compaction forms its - * own round, and messages within its range don't belong to this round, so this isn't advanced for them. + * Advance this round's "latest timestamp seen among its messages" — the + * **fallback for the round's end** (see taskLastReqEndMs: the normal + * round-end is set by request_end, and this only guarantees a usable + * upper bound for a degenerate round with no request_end at all, + * interrupted before its first Request even ran), and the floor under the + * live anchor back-dated on a history rebuild when the server's own clock + * did not come back with the response (see taskLastTsMs, pushMessages). + * Compaction forms its own round, and messages within its range don't + * belong to this round, so this isn't advanced for them. */ function touchTask(model: StreamModel, timestamp: string): void { if (!model.taskOpen) return; @@ -637,7 +690,7 @@ function touchTask(model: StreamModel, timestamp: string): void { function startTask(model: StreamModel, timestamp: string, nowMs: number): void { // The previous Task is finalized by "the next Task starting" (the history-rebuild convention). - finalizeOpenTask(model, "history"); + finalizeOpenTask(model); // Finalize any retry state left over from the previous Task: when the // server dies during a backoff window, the Trace's tail is // request_end(timeout) with no abort, and history rebuild would leave a @@ -659,7 +712,7 @@ function startTask(model: StreamModel, timestamp: string, nowMs: number): void { resetTaskCounters(model.stats); } -function finalizeOpenTask(model: StreamModel, mode: "history" | "live", nowMs?: number): void { +function finalizeOpenTask(model: StreamModel): void { if (!model.taskOpen) return; model.taskOpen = false; // No more tool output will arrive once a Task is finalized: close cards still "executing" and stop their LiveDuration. @@ -672,23 +725,18 @@ function finalizeOpenTask(model: StreamModel, mode: "history" | "live", nowMs?: // and is naturally excluded — no compaction wall-clock addition/subtraction // is needed at all, and history rebuild and live share the same // convention, consistent before and after a refresh (see taskLastReqEndMs). + // + // Trace timestamps are the ONLY source here — the local clock is never + // consulted, so a round settles to the same number whether it was watched + // live or replayed from the Trace after a reload. A degenerate round (no + // request_end at all, e.g. interrupted before its first Request even ran) + // falls back to its message span, which the abort event's own timestamp + // still bounds; that span is what a later reload would compute, so taking + // the local clock instead — as this did before — only bought a number that + // silently changed on refresh, along with idle-detection and mid-join + // latency folded into it. const endMs = model.taskLastReqEndMs ?? model.taskLastTsMs; - const tsElapsed = Math.max(0, endMs - model.taskFirstTsMs); - // Only a degenerate round (no request_end at all for the whole round, - // e.g. interrupted before its first Request even ran) falls back to the - // local clock: during a mid-stream join / resync rebuild, the local clock - // only covers the time since joining, so it's compared against the - // message span and the larger is taken to avoid underestimating. When - // there is a request_end, it's the accurate round-end and is used - // directly (not compared against the local clock, to avoid folding in noise like idle-detection latency). - let elapsed: number; - if (model.taskLastReqEndMs !== null) { - elapsed = tsElapsed; - } else if (mode === "live") { - elapsed = Math.max((nowMs ?? Date.now()) - model.taskStartLocalMs, tsElapsed); - } else { - elapsed = tsElapsed; - } + const elapsed = Math.max(0, endMs - model.taskFirstTsMs); const stats = endTask(model.stats, elapsed); if (model.nested) return; const reply = collectTaskAssistant(model); diff --git a/packages/web/src/lib/omni/task-stats.ts b/packages/web/src/lib/omni/task-stats.ts index 07efa93..8a5ecbc 100644 --- a/packages/web/src/lib/omni/task-stats.ts +++ b/packages/web/src/lib/omni/task-stats.ts @@ -316,11 +316,19 @@ export function endTask(t: TaskStatsTracker, elapsedMs: number): TaskStats | nul /** * Session elapsed time for live display (the header chip): the settled cross-Task cumulative * plus the currently running Task's wall clock so far. While no Task is open this is exactly - * `sessionElapsedMs`, so the idle rendering equals the settled value; the running addition - * counts from the model's task-start local clock (`taskStartLocalMs`), clamped so a clock - * anomaly never makes the sum go backwards. At task end {@link endTask} folds the Task's - * elapsed into `sessionElapsedMs` in the same model update that flips `taskOpen` off, so the - * live addition never double-counts across the boundary. + * `sessionElapsedMs`, so the idle rendering equals the settled value. At task end {@link endTask} + * folds the Task's elapsed into `sessionElapsedMs` in the same model update that flips `taskOpen` + * off, so the live addition never double-counts across the boundary. + * + * The running addition counts from the model's task-start local clock (`taskStartLocalMs`), + * clamped so a clock anomaly never makes the sum go backwards. + * + * This is the only place a local clock reaches the elapsed figure at all, and only to animate the + * Task in flight. `taskStartLocalMs` is not the instant this client loaded the page: on a history + * rebuild it is back-dated by the elapsed already behind the Task, measured entirely in server + * time (see pushMessages), so reloading mid-run resumes the ticking value instead of restarting + * it, covers an event still in flight, and carries no client/server clock offset. The moment the + * Task settles the number is replaced by one computed from Trace timestamps alone. */ export function liveSessionElapsedMs( t: TaskStatsTracker, diff --git a/packages/web/test/stream-controller.test.ts b/packages/web/test/stream-controller.test.ts index 8d30211..6c06f77 100644 --- a/packages/web/test/stream-controller.test.ts +++ b/packages/web/test/stream-controller.test.ts @@ -41,13 +41,21 @@ interface Harness { errors: Array; loadings: boolean[]; loadCalls: () => number; - resolveLoad: (messages: OmniMessage[], live?: MessagesLiveTail) => void; + resolveLoad: ( + messages: OmniMessage[], + live?: MessagesLiveTail, + serverNowMs?: number | null, + ) => void; rejectLoad: (err: Error) => void; } function createHarness(): Harness { const pendingLoads: Array<{ - resolve: (m: { messages: OmniMessage[]; live?: MessagesLiveTail }) => void; + resolve: (m: { + messages: OmniMessage[]; + live?: MessagesLiveTail; + serverNowMs?: number | null; + }) => void; reject: (e: unknown) => void; }> = []; const states: SessionStatus[] = []; @@ -56,7 +64,11 @@ function createHarness(): Harness { let calls = 0; const controller = createStreamController({ loadMessages: () => - new Promise<{ messages: OmniMessage[]; live?: MessagesLiveTail }>((resolve, reject) => { + new Promise<{ + messages: OmniMessage[]; + live?: MessagesLiveTail; + serverNowMs?: number | null; + }>((resolve, reject) => { calls += 1; pendingLoads.push({ resolve, reject }); }), @@ -73,8 +85,12 @@ function createHarness(): Harness { errors, loadings, loadCalls: () => calls, - resolveLoad: (messages, live) => - pendingLoads.shift()!.resolve({ messages, ...(live !== undefined ? { live } : {}) }), + resolveLoad: (messages, live, serverNowMs) => + pendingLoads.shift()!.resolve({ + messages, + ...(live !== undefined ? { live } : {}), + ...(serverNowMs !== undefined ? { serverNowMs } : {}), + }), rejectLoad: (err) => pendingLoads.shift()!.reject(err), }; } @@ -85,6 +101,31 @@ const HISTORY_TASK: OmniMessage[] = [ at(tokenUsage(counts(1000), counts(1000)), "2026-07-05T00:00:05.000Z"), ]; +describe("the running Task's live anchor is back-dated from the server's clock at read time", () => { + it("counts an event still in flight, which the Trace's own span cannot see", async () => { + const h = createHarness(); + const p = h.controller.load(); + // Still running, so the Task stays open. Its last Trace entry is 5s in, but the server's + // clock says 300s have passed — a tool has been executing throughout, appending nothing. + h.controller.handleServer({ type: "task_state", state: "running" }); + h.resolveLoad(HISTORY_TASK, undefined, Date.parse("2026-07-05T00:05:00.000Z")); + await p; + expect(h.controller.model.taskOpen).toBe(true); + // The harness's local clock is 1_000_000, nothing like the server's: the anchor is + // back-dated by the full 300s the server reports, not by the 5s the Trace shows. + expect(1_000_000 - h.controller.model.taskStartLocalMs).toBe(300_000); + }); + + it("falls back to the Trace's span when no Date header came back", async () => { + const h = createHarness(); + const p = h.controller.load(); + h.controller.handleServer({ type: "task_state", state: "running" }); + h.resolveLoad(HISTORY_TASK, undefined, null); + await p; + expect(1_000_000 - h.controller.model.taskStartLocalMs).toBe(5_000); + }); +}); + describe("in-stream task_state is the authoritative running state (history-closing decision)", () => { it("subscription snapshot idle: history closes, producing the last Task's stats row", async () => { const h = createHarness(); diff --git a/packages/web/test/stream-model.test.ts b/packages/web/test/stream-model.test.ts index 23927d4..14acfbf 100644 --- a/packages/web/test/stream-model.test.ts +++ b/packages/web/test/stream-model.test.ts @@ -44,6 +44,7 @@ import { pushMessages, registerLocalDecision, } from "../src/lib/omni/stream-model"; +import { liveSessionElapsedMs } from "../src/lib/omni/task-stats"; import type { AssistantTextItem, CompactionItem, @@ -731,7 +732,7 @@ describe("origin nested routing", () => { m, at(withOrigin(tokenUsage(counts(400), counts(400)), "c1"), "2026-07-05T00:00:02.000Z"), ); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; expect(stats.stats!.tokens).toBe(1400); expect(stats.stats!.tokensDelta).toBe(1400); @@ -814,13 +815,15 @@ describe("Task segmentation and stats triggering", () => { expect(items(m).filter((i) => i.kind === "task_stats")).toHaveLength(0); }); - it("live streams close at task_state:idle, measuring the delta with the local clock", () => { + it("live streams close at task_state:idle, measuring the delta from Trace timestamps", () => { const m = createStreamModel(); - pushMessage(m, userText("live question"), 10_000); - pushMessage(m, tokenUsage(counts(800), counts(800)), 11_000); - notifyTaskIdle(m, 15_100); + // The local clock advances 5.1s across this round, but only the message timestamps decide + // the settled figure — the same span a reload would replay out of the Trace. + pushMessage(m, at(userText("live question"), "2026-07-05T00:00:00.000Z"), 10_000); + pushMessage(m, at(tokenUsage(counts(800), counts(800)), "2026-07-05T00:00:01.000Z"), 11_000); + notifyTaskIdle(m); const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; - expect(stats.stats!.elapsedDeltaMs).toBe(5100); + expect(stats.stats!.elapsedDeltaMs).toBe(1_000); expect(stats.stats!.tokens).toBe(800); }); @@ -839,7 +842,7 @@ describe("Task segmentation and stats triggering", () => { const m = createStreamModel(); pushMessage(m, at(userText("one"), "2026-07-05T00:00:00.000Z")); pushMessage(m, at(tokenUsage(counts(1000), counts(1000)), "2026-07-05T00:00:01.000Z")); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); // Manual /compact (outside the Task boundary). pushMessage( m, @@ -850,7 +853,7 @@ describe("Task segmentation and stats triggering", () => { // Next Task. pushMessage(m, at(userText("two"), "2026-07-05T00:02:00.000Z")); pushMessage(m, at(tokenUsage(counts(1800), counts(500)), "2026-07-05T00:02:01.000Z")); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const statsItems = items(m).filter((i) => i.kind === "task_stats") as TaskStatsItem[]; const last = statsItems[statsItems.length - 1]!; expect(last.stats!.tokensDelta).toBe(500); // excludes the compaction's 300 @@ -888,7 +891,7 @@ describe("output TPS (request event pair timing)", () => { // happening between the two requests isn't counted). pushMessage(m, at(tokenUsage(out(900), out(900)), "2026-07-05T00:00:03.500Z")); pushMessage(m, at(requestEnd("completed"), "2026-07-05T00:00:04.000Z")); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; expect(stats.stats!.outputTps).toBe(300); // 900 / 3s expect(stats.stats!.tokensByBucket).toEqual({ cacheRead: 0, cacheWrite: 0, output: 900 }); @@ -899,7 +902,7 @@ describe("output TPS (request event pair timing)", () => { // No request events -> no LLM timing -> TPS is null. pushMessage(m, at(userText("q1"), "2026-07-05T00:00:00.000Z")); pushMessage(m, at(tokenUsage(out(100), out(100)), "2026-07-05T00:00:01.000Z")); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const s1 = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; expect(s1.stats!.outputTps).toBeNull(); // Next Task: two request rounds of 2s each, 400 output tokens each -> 800 / 4s = 200 tok/s. @@ -910,7 +913,7 @@ describe("output TPS (request event pair timing)", () => { pushMessage(m, at(requestBegin(), "2026-07-05T00:01:05.000Z")); pushMessage(m, at(tokenUsage(out(400), out(400)), "2026-07-05T00:01:06.500Z")); pushMessage(m, at(requestEnd("completed"), "2026-07-05T00:01:07.000Z")); // 2s - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const stats = items(m).filter((i) => i.kind === "task_stats") as TaskStatsItem[]; expect(stats[stats.length - 1]!.stats!.outputTps).toBe(200); }); @@ -932,7 +935,7 @@ describe("output TPS (request event pair timing)", () => { pushMessage(m, at(approvalDecision("allow", "t1"), "2026-07-05T00:00:32.000Z")); pushMessage(m, at(tokenUsage(out(1000), out(1000)), "2026-07-05T00:00:32.500Z")); pushMessage(m, at(requestEnd("completed"), "2026-07-05T00:00:33.000Z")); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; // Wall clock 32s, minus 30s approval -> 2s of generation: 1000 / 2s = 500 tok/s // (without subtracting it, it would be only 31 tok/s). @@ -963,7 +966,7 @@ describe("output TPS (request event pair timing)", () => { "2026-07-05T00:00:21.000Z", ), ); - notifyTaskIdle(m, Date.now()); + notifyTaskIdle(m); const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; expect(stats.stats!.outputTps).toBe(300); // 600 / 2s (the compaction request's 999 output and 15.5s are excluded) }); @@ -1242,7 +1245,64 @@ describe("compaction-internal messages (#17: history rebuild aligned with the li }); }); -describe("live close-out elapsed (#5/#20: mid-join takes the message-timestamp lower bound)", () => { +describe("elapsed comes from Trace timestamps (#5/#20: settled spans, reload-stable live anchor)", () => { + it("reloading mid-run resumes the header's live elapsed instead of restarting it", () => { + const m = createStreamModel(); + const loadNow = 1_000_000; + // The Task started 60s ago and is STILL running — nothing finalizes it, so the header + // renders sessionElapsedMs + (now − taskStartLocalMs). Every message in a rebuild is fed + // the same `nowMs`, so without the re-anchor that addend would be 0 and the chip would + // drop back to the settled total and climb from zero on every reload. + pushMessages( + m, + [ + at(userText("long-running task"), "2026-07-05T00:00:00.000Z"), + at(assistantText("working"), "2026-07-05T00:01:00.000Z"), + ], + loadNow, + ); + expect(m.taskOpen).toBe(true); + expect(loadNow - m.taskStartLocalMs).toBe(60_000); + // No `Date` header came back, so the Trace's own span decides the anchor. Note the local + // clock here is nowhere near the server timestamps, and the figure is unaffected: only + // differences between server-side values ever reach it. + expect(liveSessionElapsedMs(m.stats, m.taskOpen, m.taskStartLocalMs, loadNow + 5000)).toBe( + 65_000, + ); + }); + + it("a reload while an event is still in flight counts it, from the server's own clock", () => { + // A tool started executing 10s into the Task and is STILL running 300s later. Nothing has + // been appended to the Trace since it began, so its span reaches only those first 10s — + // the server's clock at read time is the only thing that sees the other 290s. + const replay = [ + at(userText("run the build"), "2026-07-05T00:00:00.000Z"), + at(toolCall({ name: "bash", arguments: "{}", toolCallId: "t1" }), "2026-07-05T00:00:10.000Z"), + ]; + const serverNow = Date.parse("2026-07-05T00:05:00.000Z"); + // The client's clock is deliberately nothing like the server's: a 90-minute offset that must + // not reach the figure, since both ends of the measured interval are server-side values. + const loadNow = serverNow + 90 * 60_000; + const m = createStreamModel(); + pushMessages(m, replay, loadNow, serverNow); + expect(m.taskOpen).toBe(true); + expect(liveSessionElapsedMs(m.stats, m.taskOpen, m.taskStartLocalMs, loadNow)).toBe(300_000); + // Without the header the Trace's span is the floor: short, but never an overshoot. + const noHeader = createStreamModel(); + pushMessages(noHeader, replay, loadNow, null); + expect( + liveSessionElapsedMs(noHeader.stats, noHeader.taskOpen, noHeader.taskStartLocalMs, loadNow), + ).toBe(10_000); + }); + + it("a live stream is unaffected: the anchor stays the real Task start", () => { + const m = createStreamModel(); + // One message at a time with the real current clock — the Trace span is still 0 when the + // Task opens, so the re-anchor is a no-op and must not shift the origin. + pushMessage(m, at(userText("live question"), "2026-07-05T00:00:00.000Z"), 10_000); + expect(m.taskStartLocalMs).toBe(10_000); + }); + it("a Task ending right after a refresh: elapsed takes the message-timestamp span, not the local-clock delta", () => { const m = createStreamModel(); const loadNow = 1_000_000; @@ -1257,19 +1317,31 @@ describe("live close-out elapsed (#5/#20: mid-join takes the message-timestamp l loadNow, ); // task_state:idle arrives 2s after joining. - notifyTaskIdle(m, loadNow + 2000); + notifyTaskIdle(m); const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; expect(stats.stats!.elapsedDeltaMs).toBe(60_000); expect(stats.stats!.elapsedMs).toBe(60_000); // sessionElapsedMs is corrected in sync }); - it("when the local-clock delta is larger (normal live stream), the local clock still wins", () => { - const m = createStreamModel(); - pushMessage(m, at(userText("live question"), "2026-07-05T00:00:00.000Z"), 10_000); - pushMessage(m, at(tokenUsage(counts(800), counts(800)), "2026-07-05T00:00:01.000Z"), 11_000); - notifyTaskIdle(m, 15_100); - const stats = items(m).find((i) => i.kind === "task_stats") as TaskStatsItem; - expect(stats.stats!.elapsedDeltaMs).toBe(5100); + it("a degenerate round settles to the same figure live and replayed — the local clock never leaks in", () => { + // No request_end anywhere (interrupted before its first Request ran), which used to be the + // one case that fell back to the local clock. Watching it live and replaying it out of the + // Trace must now agree, or the header would silently change on reload. + const msgs = [ + at(userText("live question"), "2026-07-05T00:00:00.000Z"), + at(tokenUsage(counts(800), counts(800)), "2026-07-05T00:00:01.000Z"), + ]; + const live = createStreamModel(); + pushMessage(live, msgs[0]!, 10_000); + pushMessage(live, msgs[1]!, 11_000); + notifyTaskIdle(live); // idle detected 4.1s later by the local clock — irrelevant now + const replayed = createStreamModel(); + pushMessages(replayed, msgs, 9_000_000); // reloaded much later, different clock entirely + finalizeHistory(replayed); + const of = (m: StreamModel) => + (items(m).find((i) => i.kind === "task_stats") as TaskStatsItem).stats!.elapsedDeltaMs; + expect(of(live)).toBe(1_000); + expect(of(replayed)).toBe(of(live)); }); }); @@ -1435,7 +1507,7 @@ describe("thinking/tool durations (collapsed-row display data)", () => { const m = createStreamModel(); pushMessage(m, at(userText("Q"), T0)); pushMessage(m, at(toolCall({ name: "x", arguments: "{}", toolCallId: "tb" }), T1)); - notifyTaskIdle(m, Date.parse(T2)); + notifyTaskIdle(m); const card = items(m).find((i) => i.kind === "tool_call") as ToolCallItem; expect(card.outputComplete).toBe(true); expect(card.outputStopReason).toBe("aborted");