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) <noreply@anthropic.com>
This commit is contained in:
Yaowei Zheng
2026-07-29 23:30:13 +08:00
committed by GitHub
parent ed1ea509d3
commit b9de3da879
7 changed files with 289 additions and 78 deletions
+27 -3
View File
@@ -42,8 +42,29 @@ export interface ApiFetchOptions {
query?: Record<string, string | number | undefined>;
}
/** 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<T>(path: string, options: ApiFetchOptions = {}): Promise<T> {
return (await apiFetchWithMeta<T>(path, options)).data;
}
/** {@link apiFetch} plus the response metadata in {@link ApiFetchMeta}; identical in every other respect. */
export async function apiFetchWithMeta<T>(
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<T>(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 };
}
+11 -2
View File
@@ -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<void>(`/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<MessagesResponse>(`/api/sessions/${encodeURIComponent(sessionId)}/messages`);
apiFetchWithMeta<MessagesResponse>(
`/api/sessions/${encodeURIComponent(sessionId)}/messages`,
).then(({ data, serverNowMs }) => ({ ...data, serverNowMs }));
// Task execution, approval, abort, compaction ------------------------------------------------------
+14 -5
View File
@@ -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<void> => {
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.
+83 -35
View File
@@ -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);
+13 -5
View File
@@ -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,
+46 -5
View File
@@ -41,13 +41,21 @@ interface Harness {
errors: Array<string | null>;
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();
+95 -23
View File
@@ -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");