| 1 | import assert from "node:assert/strict"; |
| 2 | import { initialState, reducer, type State } from "../lib/useController"; |
| 3 | import { turnMetrics } from "../lib/turnMetrics"; |
| 4 | import { historyMessagesToItems } from "../lib/historyItems"; |
| 5 | import type { TranscriptSnapshot } from "../lib/transcriptProtocol"; |
| 6 | import type { WireEvent } from "../lib/types"; |
| 7 | |
| 8 | const originalNow = Date.now; |
| 9 | let now = 426_000; |
| 10 | Date.now = () => now; |
| 11 | const snapshot = (sessionId = "A", content = "x".repeat(52_000)): TranscriptSnapshot => ({ |
| 12 | protocolVersion: 1, snapshotId: `snapshot-${sessionId}`, |
| 13 | identity: { sessionId, headId: "head", rewriteEpoch: 0, runtimeEpoch: "epoch" }, |
| 14 | projectionRevision: 1, coveredThroughSeq: 10, |
| 15 | records: [{ id: "m:answer", order: 0, refs: [], message: { |
| 16 | role: "assistant", messageId: "answer", turnId: `turn-${sessionId}`, content, |
| 17 | } }], activeRecords: [], |
| 18 | runtime: { status: "in_progress", turnId: `turn-${sessionId}`, startedAt: 1_000, pendingEvents: [] }, |
| 19 | activeAttempts: [{ id: "attempt", messageId: "answer" }], |
| 20 | before: 0, hasOlder: false, totalRecords: 1, totalTurns: 1, stale: false, |
| 21 | }); |
| 22 | const event = (s: State, e: WireEvent) => reducer(s, { type: "event", e }); |
| 23 | const metrics = (s: State) => turnMetrics({ ...s, now, waitAccumMs: 0, |
| 24 | turnRateOutputQuarters: s.turnRateSample?.outputQuarters })!; |
| 25 | const usage = (completionTokens: number): WireEvent => ({ kind: "usage", usage: { |
| 26 | promptTokens: 0, completionTokens, totalTokens: completionTokens, cacheHitTokens: 0, cacheMissTokens: 0, |
| 27 | sessionCacheHitTokens: 0, sessionCacheMissTokens: 0, |
| 28 | } }); |
| 29 | |
| 30 | try { |
| 31 | // Both local and remote followers use the same snapshot reducer. Exercise |
| 32 | // both protocol entry points, including the v2 projected history path. |
| 33 | for (const remote of [false, true]) for (const v2 of [false, true]) { |
| 34 | const install = (state: State, snap = snapshot()) => v2 |
| 35 | ? reducer(state, { type: "transcript_v2_snapshot", snapshot: snap, remote, |
| 36 | projection: { items: historyMessagesToItems(snap.records.map(r => r.message), "snapshot:").items, |
| 37 | startTurn: 0, endTurn: 1, totalTurns: 1, hasOlder: false, hasNewer: false, revision: 10, revisionKnown: true, digest: "" } }) |
| 38 | : reducer(state, { type: "transcript_snapshot", snapshot: snap, remote }); |
| 39 | now = 426_000; |
| 40 | let s = install(initialState); |
| 41 | assert.equal(metrics(s).tokens, 13_000, "restored output remains in the token total"); |
| 42 | assert.equal(metrics(s).tps, null, "a snapshot alone cannot provide throughput"); |
| 43 | s = event(s, { kind: "text", messageId: "answer", text: "x".repeat(40) }); |
| 44 | now += 500; |
| 45 | assert.equal(metrics(s).tps, 20, "only 10 newly observed tokens belong to the 500ms window"); |
| 46 | assert.equal(metrics(s).elapsedMs, 425_500, "the turn clock retains its original start"); |
| 47 | assert.equal(metrics(s).tokens, 13_010); |
| 48 | now += 500; |
| 49 | s = event(s, usage(13_010)); |
| 50 | assert.equal(s.lastRequestTps, 10, "full-request usage cannot inflate the partial observation"); |
| 51 | assert.equal(metrics(s).tps, 10, "usage settlement keeps the matched sample"); |
| 52 | s = event(s, { kind: "message", messageId: "answer", text: "x".repeat(52_040) }); |
| 53 | |
| 54 | now += 20_000; // tool execution is outside the provider-output clock |
| 55 | s = event(s, { kind: "reasoning", messageId: "answer-2", text: "中".repeat(40) }); |
| 56 | now += 1_000; |
| 57 | assert.equal(metrics(s).tps, 20, "CJK deltas and tool gaps share the same rate window"); |
| 58 | s = event(s, usage(30)); |
| 59 | assert.equal(s.lastRequestTps, 30, "the next request measures only its own output"); |
| 60 | s = event(s, { kind: "turn_done", turnId: "turn-A" }); |
| 61 | assert.equal(metrics(s).tps, 20, "completion does not restore the inflated numerator"); |
| 62 | const answer = s.items.find(item => item.id === "m:answer-2"); |
| 63 | if (!v2) assert.equal(answer?.kind === "assistant" && answer.tokensPerSecond, 20); |
| 64 | |
| 65 | s = event(s, { kind: "turn_started", turnId: "next-turn" }); |
| 66 | assert.equal(s.turnRateSample, undefined, "a new turn returns to ordinary billed telemetry"); |
| 67 | s = install(s, snapshot("B")); |
| 68 | assert.equal(s.turnDoneAt, 0, "another session does not inherit completion or timing"); |
| 69 | assert.equal(s.turnOutputTokens, 0, "another session does not inherit billed output"); |
| 70 | s = event(s, { kind: "text", messageId: "answer", text: "x".repeat(40) }); |
| 71 | now += 1_000; |
| 72 | s = install(s, snapshot("B", "x".repeat(100_000))); |
| 73 | assert.equal(metrics(s).tps, null, "reconnection replaces both sides of the observation window"); |
| 74 | s = event(s, { kind: "text", messageId: "answer", text: "x".repeat(40) }); |
| 75 | now += 1_000; |
| 76 | assert.equal(metrics(s).tps, 10, "reconnection backlog never becomes newly sampled output"); |
| 77 | |
| 78 | // Cumulative tool argument progress needs its own first-observation baseline. |
| 79 | s = install(initialState); |
| 80 | s = event(s, { kind: "tool_dispatch", tool: { id: "tool", name: "write_file", readOnly: false, partial: true, argChars: 52_000 } }); |
| 81 | now += 1_000; |
| 82 | s = event(s, { kind: "tool_dispatch", tool: { id: "tool", name: "write_file", readOnly: false, partial: true, argChars: 52_040 } }); |
| 83 | assert.equal(metrics(s).tps, 10, "restored cumulative arguments count only their increment"); |
| 84 | s = event(s, { kind: "tool_dispatch", tool: { id: "child", parentId: "parent", name: "write_file", readOnly: false, partial: true, argChars: 1_000_000 } }); |
| 85 | assert.equal(metrics(s).tps, 10, "child tool output cannot enter the executor sample"); |
| 86 | s = event(s, { kind: "tool_dispatch", tool: { id: "tool", name: "write_file", readOnly: false, partial: true } }); |
| 87 | assert.equal(metrics(s).tps, 10, "a name-only partial cannot sample a child's cumulative counter"); |
| 88 | s = event(s, { kind: "tool_dispatch", tool: { id: "tool", name: "write_file", readOnly: false, partial: true, argChars: 52_080 } }); |
| 89 | assert.equal(metrics(s).tps, 20, "parent argument progress retains its own baseline"); |
| 90 | |
| 91 | s = install(initialState); |
| 92 | s = event(s, { kind: "text", messageId: "answer", text: "x".repeat(40) }); |
| 93 | now += 1_000; |
| 94 | s = event(s, { kind: "stream_attempt", messageId: "answer", streamAttempt: { id: "attempt", action: "discard" } }); |
| 95 | now += 20_000; |
| 96 | s = event(s, { kind: "stream_attempt", messageId: "retry", streamAttempt: { id: "retry", action: "begin" } }); |
| 97 | s = event(s, { kind: "text", messageId: "retry", text: "x".repeat(40) }); |
| 98 | now += 1_000; |
| 99 | s = event(s, usage(10)); |
| 100 | assert.equal(metrics(s).tps, 10, "retry backoff is outside the sampled provider time"); |
| 101 | assert.equal(s.lastRequestTps, 10, "request settlement pairs all sampled attempts with their intervals"); |
| 102 | } |
| 103 | } finally { |
| 104 | Date.now = originalNow; |
| 105 | } |
| 106 | console.log("restored turn rate: local/remote snapshots, reconnect, usage, tools and completion passed"); |
| 107 |