| 1 | import assert from "node:assert/strict"; |
| 2 | import test from "node:test"; |
| 3 | import type { FollowRequest, TranscriptFollowResponse } from "../generated/desktopContract.generated"; |
| 4 | import { installDesktopHostStub } from "./desktopHostStub"; |
| 5 | |
| 6 | Object.defineProperty(globalThis, "window", { configurable: true, value: {} }); |
| 7 | const commands: Record<string, unknown> = {}; |
| 8 | const desktopStub = installDesktopHostStub(commands); |
| 9 | const [{ TranscriptSessionFollower }, { initialState, reducer }, { getTranscriptStore }] = await Promise.all([ |
| 10 | import("../lib/transcriptSessionFollower"), import("../lib/useController"), import("../lib/transcriptStore"), |
| 11 | ]); |
| 12 | const { ChatSource } = await import("../lib/chatViewSource"); |
| 13 | function deferred<T>() { |
| 14 | let resolve!: (value: T) => void; |
| 15 | let reject!: (error: unknown) => void; |
| 16 | const promise = new Promise<T>((done, fail) => { resolve = done; reject = fail; }); |
| 17 | return { promise, resolve, reject }; |
| 18 | } |
| 19 | |
| 20 | test("expired settled content requests the owning follower to resynchronize", async () => { |
| 21 | const { canonicalHistoryContent, registerTranscriptContentRecovery } = await import("../lib/canonicalTranscriptBackend"); |
| 22 | let recoveries = 0; |
| 23 | commands.TranscriptContentForTab = async () => ({ stale: true }); |
| 24 | const release = registerTranscriptContentRecovery("expired-content", () => { recoveries++; }); |
| 25 | const ref = { entryId: "m:answer", field: "content", size: 1, chunks: 1, revision: 1, digest: "expired", |
| 26 | transcriptRef: { snapshotId: "expired", recordId: "m:answer", path: ["content"], bytes: 1 } }; |
| 27 | await assert.rejects(canonicalHistoryContent("expired-content", ref, 0), /synchronizing/); |
| 28 | assert.equal(recoveries, 1); |
| 29 | release(); |
| 30 | await assert.rejects(canonicalHistoryContent("expired-content", ref, 0), /synchronizing/); |
| 31 | assert.equal(recoveries, 1, "released owner cannot be restarted by a late content result"); |
| 32 | const delayed = deferred<{ stale: boolean }>(); |
| 33 | commands.TranscriptContentForTab = () => delayed.promise; |
| 34 | const oldRelease = registerTranscriptContentRecovery("expired-content", () => { recoveries++; }); |
| 35 | const staleRead = canonicalHistoryContent("expired-content", ref, 0); |
| 36 | oldRelease(); |
| 37 | const newRelease = registerTranscriptContentRecovery("expired-content", () => { recoveries += 100; }); |
| 38 | delayed.resolve({ stale: true }); |
| 39 | await assert.rejects(staleRead, /synchronizing/); |
| 40 | assert.equal(recoveries, 1, "late stale reference cannot resynchronize a replacement session"); |
| 41 | newRelease(); |
| 42 | }); |
| 43 | |
| 44 | test("reading old pages isolates live output and rejoins the current stable node", () => { |
| 45 | let state: import("../lib/useController").State = { ...initialState, transcriptProtocol: 2, historyHasNewer: true, |
| 46 | items: [{ kind: "user" as const, id: "old-reader", text: "old page" }] }; |
| 47 | const resident = state.items; |
| 48 | state = reducer(state, { type: "event", e: { kind: "text", messageId: "active", text: "full prefix" } }); |
| 49 | assert.equal(state.items, resident); |
| 50 | assert.equal(state.offscreenItems?.find(item => item.id === "m:active")?.kind, "assistant"); |
| 51 | const active = state.offscreenItems!.find(item => item.id === "m:active")!; |
| 52 | state = reducer(state, { type: "history_append", items: [], startTurn: 0, endTurn: 1, totalTurns: 1, hasOlder: true, hasNewer: false }); |
| 53 | assert.equal(state.items.find(item => item.id === "m:active"), active); |
| 54 | assert.equal(state.offscreenItems, undefined); |
| 55 | }); |
| 56 | |
| 57 | test("business replacement preserves authoritative final turn annotations", () => { |
| 58 | const final = { kind: "assistant" as const, id: "m:final", text: "answer", reasoning: "", streaming: false, turnFinal: true, turnDurationMs: 933524, samplingCount: 72, toolCount: 72 }; |
| 59 | const state = reducer({ ...initialState, items: [final] }, { type: "transcript_records", confirmedUsers: [], projection: { |
| 60 | items: [{ ...final, turnFinal: undefined, turnDurationMs: undefined, samplingCount: undefined, toolCount: undefined }], |
| 61 | removeIds: [], startTurn: 0, endTurn: 1, totalTurns: 1, hasOlder: false, hasNewer: false, revision: 2, revisionKnown: true, digest: "cut", |
| 62 | } }); |
| 63 | assert.equal(state.items[0].kind === "assistant" && state.items[0].turnDurationMs, 933524); |
| 64 | }); |
| 65 | |
| 66 | test("business projection removes stale duplicate nodes and restores authoritative order", () => { |
| 67 | const user = { kind: "user" as const, id: "m:user", text: "build" }; |
| 68 | const tool = { kind: "tool" as const, id: "call", name: "edit_file", args: "{}", readOnly: false, status: "done" as const, output: "written" }; |
| 69 | const final = { kind: "assistant" as const, id: "m:final", text: "done", reasoning: "", streaming: false }; |
| 70 | const state = reducer({ ...initialState, items: [user, final, tool, { ...tool }] }, { type: "transcript_records", confirmedUsers: [], projection: { |
| 71 | items: [user, tool, final], removeIds: [], startTurn: 1, endTurn: 1, totalTurns: 1, hasOlder: false, hasNewer: false, |
| 72 | revision: 2, revisionKnown: true, digest: "cut", |
| 73 | } }); |
| 74 | assert.deepEqual(state.items.map(item => item.id), ["m:user", "call", "m:final"]); |
| 75 | }); |
| 76 | |
| 77 | test("content-driven projection preserves a completed event result", () => { |
| 78 | const live = { kind: "tool" as const, id: "call", name: "write_file", args: "{}", readOnly: false, |
| 79 | status: "done" as const, output: "written", execution: { state: "completed" as const, durationMs: 99 } }; |
| 80 | const projected = { ...live, status: "unknown" as const, output: "", execution: undefined, |
| 81 | resultMissing: true, resultEvidence: "missing" as const }; |
| 82 | const state = reducer({ ...initialState, items: [live], transcriptProjectedIds: ["call"] }, { |
| 83 | type: "transcript_records", confirmedUsers: [], projection: { |
| 84 | items: [projected], removeIds: [], mutation: "patch", startTurn: 1, endTurn: 1, totalTurns: 1, |
| 85 | hasOlder: false, hasNewer: false, revision: 2, revisionKnown: true, digest: "cut", |
| 86 | }, |
| 87 | }); |
| 88 | const tool = state.items.find((item): item is typeof live => item.kind === "tool"); |
| 89 | assert.equal(tool?.status, "done"); |
| 90 | assert.equal(tool?.output, "written"); |
| 91 | assert.equal(tool?.execution?.durationMs, 99); |
| 92 | }); |
| 93 | |
| 94 | test("formal empty result replaces an earlier event preview", () => { |
| 95 | const live = { kind: "tool" as const, id: "call", name: "write_file", args: "{}", readOnly: false, |
| 96 | status: "done" as const, output: "preview", error: "temporary", execution: { state: "completed" as const } }; |
| 97 | const formal = { ...live, output: "", error: undefined, resultEvidence: "formal" as const }; |
| 98 | const state = reducer({ ...initialState, items: [live], transcriptProjectedIds: ["call"] }, { |
| 99 | type: "transcript_records", confirmedUsers: [], projection: { |
| 100 | items: [formal], removeIds: [], mutation: "patch", startTurn: 1, endTurn: 1, totalTurns: 1, |
| 101 | hasOlder: false, hasNewer: false, revision: 2, revisionKnown: true, digest: "cut", |
| 102 | }, |
| 103 | }); |
| 104 | const tool = state.items[0]; |
| 105 | assert.equal(tool.kind === "tool" && tool.output, ""); |
| 106 | assert.equal(tool.kind === "tool" && tool.error, undefined); |
| 107 | }); |
| 108 | |
| 109 | test("duplicate authoritative projection keys are rejected without changing the mounted state", () => { |
| 110 | const tool = { kind: "tool" as const, id: "call", name: "write_file", args: "{}", readOnly: false, |
| 111 | status: "done" as const, output: "written" }; |
| 112 | const before = { ...initialState, items: [tool], transcriptProjectedIds: [tool.id] }; |
| 113 | const state = reducer(before, { type: "transcript_records", confirmedUsers: [], projection: { |
| 114 | items: [tool, { ...tool, output: "conflicting" }], removeIds: [], mutation: "patch", |
| 115 | startTurn: 1, endTurn: 1, totalTurns: 1, hasOlder: false, hasNewer: false, |
| 116 | revision: 2, revisionKnown: true, digest: "cut", |
| 117 | } }); |
| 118 | assert.equal(state, before); |
| 119 | assert.equal(state.items[0].kind === "tool" && state.items[0].output, "written"); |
| 120 | }); |
| 121 | |
| 122 | test("an id-only late event cannot overwrite an ambiguous formal result", () => { |
| 123 | const first = { kind: "tool" as const, id: "call", name: "bash", args: "", readOnly: false, |
| 124 | status: "done" as const, output: "first", identityConflict: true }; |
| 125 | const second = { ...first, id: "call:conflict:second", output: "second" }; |
| 126 | const before = { ...initialState, items: [first, second], transcriptProjectedIds: [first.id, second.id] }; |
| 127 | const state = reducer(before, { |
| 128 | type: "event", e: { kind: "tool_result", tool: { id: "call", name: "bash", output: "late", readOnly: false } }, |
| 129 | }); |
| 130 | assert.equal(state, before); |
| 131 | assert.deepEqual(state.items.map(item => item.kind === "tool" && item.output), ["first", "second"]); |
| 132 | }); |
| 133 | |
| 134 | test("reapplying an identical authoritative projection is a state no-op", () => { |
| 135 | const tool = { kind: "tool" as const, id: "call", name: "write_file", args: "{}", readOnly: false, |
| 136 | status: "done" as const, output: "written", resultEvidence: "formal" as const }; |
| 137 | const projection = { items: [tool], removeIds: [], mutation: "patch" as const, |
| 138 | startTurn: 1, endTurn: 1, totalTurns: 1, hasOlder: false, hasNewer: false, |
| 139 | revision: 2, revisionKnown: true, digest: "cut" }; |
| 140 | const first = reducer(initialState, { type: "transcript_records", confirmedUsers: [], projection }); |
| 141 | const second = reducer(first, { type: "transcript_records", confirmedUsers: [], projection }); |
| 142 | assert.equal(second, first); |
| 143 | assert.equal(second.historyMutation.seq, first.historyMutation.seq); |
| 144 | }); |
| 145 | |
| 146 | test("unrelated lazy body hydration cannot downgrade a completed tool event", async () => { |
| 147 | const tab = "completed-tool-hydration", path = "/session/completed-tool-hydration"; |
| 148 | const body = "loaded old text"; |
| 149 | commands.HistoryContentForTab = async (_tab: string, ref: { entryId: string; field: string }) => ({ |
| 150 | entryId: ref.entryId, field: ref.field, chunk: 0, chunks: 1, data: body, done: true, stale: false, |
| 151 | }); |
| 152 | const store = getTranscriptStore(); |
| 153 | const projection = store.installSlice(tab, path, { |
| 154 | entries: [ |
| 155 | { entryId: "m:old", turn: 1, order: 0, message: { role: "assistant", messageId: "old", content: "preview" }, |
| 156 | refs: [{ entryId: "m:old", field: "content", size: body.length, chunks: 1, revision: 1, digest: "body" }] }, |
| 157 | { entryId: "m:call", turn: 1, order: 1, message: { role: "assistant", messageId: "call", content: "", |
| 158 | toolCalls: [{ id: "write", name: "write_file", arguments: "{}" }] }, refs: [] }, |
| 159 | ], nextCursor: "", newerCursor: "", hasOlder: false, hasNewer: false, totalTurns: 1, startTurn: 1, endTurn: 1, |
| 160 | revision: 1, revisionKnown: true, digest: "cut", stale: false, |
| 161 | }); |
| 162 | let state = reducer({ ...initialState, transcriptProtocol: 2 }, { type: "transcript_records", |
| 163 | projection: { ...projection, removeIds: [], mutation: "replace" }, confirmedUsers: [] }); |
| 164 | state = reducer(state, { type: "event", remote: true, e: { kind: "tool_result", tool: { |
| 165 | id: "write", name: "write_file", output: "written", readOnly: false, |
| 166 | execution: { state: "completed", durationMs: 99 }, |
| 167 | } } }); |
| 168 | const unsubscribe = store.subscribe(tab, change => { |
| 169 | if (change.projection) state = reducer(state, { type: "transcript_records", projection: change.projection, confirmedUsers: [] }); |
| 170 | }); |
| 171 | try { |
| 172 | await store.requestFullContent(tab, "m:old", "content"); |
| 173 | const tool = state.items.find(item => item.kind === "tool"); |
| 174 | assert.equal(tool?.kind === "tool" && tool.status, "done"); |
| 175 | assert.equal(tool?.kind === "tool" && tool.output, "written"); |
| 176 | assert.equal(tool?.kind === "tool" && tool.execution?.durationMs, 99); |
| 177 | } finally { unsubscribe(); store.evictTab(tab); } |
| 178 | }); |
| 179 | |
| 180 | test("authoritative projection keeps local notices at their persisted anchors", () => { |
| 181 | const user = { kind: "user" as const, id: "m:user", text: "build" }; |
| 182 | const notice = { kind: "notice" as const, id: "local:notice", local: true, level: "info" as const, text: "checking" }; |
| 183 | const final = { kind: "assistant" as const, id: "m:final", text: "done", reasoning: "", streaming: false }; |
| 184 | const state = reducer({ ...initialState, items: [user, notice, final], transcriptProjectedIds: [user.id, final.id] }, { |
| 185 | type: "transcript_records", confirmedUsers: [], projection: { |
| 186 | items: [user, final], removeIds: [], mutation: "patch", startTurn: 1, endTurn: 1, totalTurns: 1, |
| 187 | hasOlder: false, hasNewer: false, revision: 2, revisionKnown: true, digest: "cut", |
| 188 | }, |
| 189 | }); |
| 190 | assert.deepEqual(state.items.map(item => item.id), ["m:user", "local:notice", "m:final"]); |
| 191 | }); |
| 192 | |
| 193 | test("durable user handoff keeps live rows after the user without remounting them", () => { |
| 194 | const previous = { kind: "assistant" as const, id: "m:previous", text: "before", reasoning: "", streaming: false }; |
| 195 | const process = { kind: "tool" as const, id: "tool:live", name: "shell", args: "{}", readOnly: true, |
| 196 | status: "running" as const, turnId: "turn-live" }; |
| 197 | const answer = { kind: "assistant" as const, id: "m:answer-live", text: "done", reasoning: "", streaming: false, turnId: "turn-live" }; |
| 198 | const durableUser = { kind: "user" as const, id: "m:user-live", messageId: "user-live", text: "build", turnId: "turn-live" }; |
| 199 | const state = reducer({ ...initialState, items: [previous, process, answer], transcriptProjectedIds: [previous.id], |
| 200 | localSubmissions: { submit: { submissionId: "submit", localId: "u1", text: "build", createdAt: 1, sequence: 1, |
| 201 | status: "accepted" as const, messageId: "user-live", turnId: "turn-live" } }, localSubmissionOrder: ["submit"] }, { |
| 202 | type: "transcript_records", confirmedUsers: [], projection: { |
| 203 | items: [previous, durableUser], removeIds: [], mutation: "patch", startTurn: 1, endTurn: 2, totalTurns: 2, |
| 204 | hasOlder: false, hasNewer: false, revision: 2, revisionKnown: true, digest: "cut", |
| 205 | }, |
| 206 | }); |
| 207 | assert.deepEqual(state.items.map(item => item.id), ["m:previous", "m:user-live", "tool:live", "m:answer-live"]); |
| 208 | assert.equal(state.items[2], process); |
| 209 | assert.equal(state.items[3], answer); |
| 210 | }); |
| 211 | async function microtasks() { for (let i = 0; i < 16; i++) await Promise.resolve(); } |
| 212 | |
| 213 | for (const remote of [false, true]) for (const scenario of ["direct", "event-first", "late", "stale", "missing"] as const) { |
| 214 | const lateBinding = scenario !== "direct" && scenario !== "event-first"; |
| 215 | test(`${remote ? "remote" : "local"} offscreen submission confirmation (${scenario})`, async () => { |
| 216 | const tab = `offscreen-${remote}-${scenario}`, path = `/session/${tab}`; |
| 217 | let state = { ...initialState }; |
| 218 | const polls: Array<ReturnType<typeof deferred<TranscriptFollowResponse>>> = []; |
| 219 | const readKey = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 220 | commands[readKey] = (_tab: string, request: FollowRequest) => { |
| 221 | if (request.close) return Promise.resolve({ protocolVersion: 2, changes: [], resetRequired: false, subscription: tab }); |
| 222 | if (!request.subscription) return Promise.resolve(initial(tab)); |
| 223 | const poll = deferred<TranscriptFollowResponse>(); polls.push(poll); return poll.promise; |
| 224 | }; |
| 225 | let lookups = 0; |
| 226 | const lookup = deferred<{ status: string; messages: Array<{ role: string; messageId: string }> }>(); |
| 227 | commands[remote ? "RemoteSessionHistoryWindowForTab" : "SessionHistoryWindowForTab"] = (_tab: string, request: { anchor: string; messageId: string }) => { |
| 228 | lookups++; assert.equal(request.anchor, "message"); assert.equal(request.messageId, "sent"); return lookup.promise; |
| 229 | }; |
| 230 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }, () => state); |
| 231 | await follower.start(); |
| 232 | try { |
| 233 | state = reducer(state, { type: "user", seq: 0, text: "question", submissionId: "submit" }); |
| 234 | const old = getTranscriptStore().installSlice(tab, path, { |
| 235 | entries: [{ entryId: "m:old", turn: 1, order: 0, message: { role: "user", content: "old page", messageId: "old" }, refs: [] }], |
| 236 | nextCursor: "", newerCursor: "next", hasOlder: false, hasNewer: true, startTurn: 1, endTurn: 1, totalTurns: 100, |
| 237 | revision: 4, digest: "generation", stale: false, |
| 238 | }); |
| 239 | state = reducer(state, { type: "history_replace", ...old }); |
| 240 | const resident = state.items; |
| 241 | if (scenario === "event-first") { |
| 242 | polls.shift()!.resolve({ protocolVersion: 2, subscription: tab, changes: [{ revision: 11, commitSeq: 4, durableSeq: 4, index: 0, |
| 243 | event: { kind: "user_message", messageId: "sent", submissionId: "submit", source: "executor" } }], resetRequired: false }); |
| 244 | await microtasks(); |
| 245 | assert.equal(lookups, 0, "an identity event alone must not initiate a history read"); |
| 246 | assert.ok(state.localSubmissions.submit, "identity alone keeps the echo until the formal record"); |
| 247 | } |
| 248 | polls.shift()!.resolve({ protocolVersion: 2, subscription: tab, changes: [{ revision: scenario === "event-first" ? 12 : 11, firstSeq: 5, commitSeq: 5, durableSeq: 5, index: 0, |
| 249 | records: [{ role: "user", messageId: "sent", submissionId: lateBinding ? undefined : "submit", content: "question" }] }], resetRequired: false }); |
| 250 | await microtasks(); |
| 251 | if (lateBinding) { |
| 252 | assert.ok(state.localSubmissions.submit); |
| 253 | polls.shift()!.resolve({ protocolVersion: 2, subscription: tab, changes: [{ revision: 12, commitSeq: 5, durableSeq: 5, index: 0, |
| 254 | event: { kind: "user_message", messageId: "sent", submissionId: "submit", source: "executor" } }], resetRequired: false }); |
| 255 | await microtasks(); |
| 256 | assert.equal(lookups, 1); |
| 257 | if (scenario === "missing") { |
| 258 | polls.shift()!.resolve({ protocolVersion: 2, subscription: tab, changes: [{ revision: 13, commitSeq: 5, durableSeq: 5, index: 0, |
| 259 | event: { kind: "user_message", messageId: "sent", submissionId: "submit", source: "executor" } }], resetRequired: false }); |
| 260 | await microtasks(); |
| 261 | assert.equal(lookups, 1, "duplicate binding does not create a second in-flight read"); |
| 262 | } |
| 263 | if (scenario === "stale") { |
| 264 | state = reducer(state, { type: "reset" }); |
| 265 | state = reducer(state, { type: "user", seq: 0, text: "replacement", submissionId: "submit" }); |
| 266 | state = { ...state, localSubmissions: { submit: { ...state.localSubmissions.submit, messageId: "sent" } } }; |
| 267 | } |
| 268 | lookup.resolve({ status: "ready", messages: scenario === "missing" ? [] : [{ role: "user", messageId: "sent" }] }); |
| 269 | await microtasks(); |
| 270 | if (scenario === "stale" || scenario === "missing") { |
| 271 | assert.ok(state.localSubmissions.submit, "stale or inconclusive reads must retain the current echo"); |
| 272 | assert.equal(state.localSubmissions.submit.status, scenario === "stale" ? "sending" : "accepted"); |
| 273 | if (scenario === "missing") { |
| 274 | assert.deepEqual(state.items, resident); |
| 275 | polls.shift()!.resolve({ protocolVersion: 2, subscription: tab, changes: [{ revision: 14, firstSeq: 6, commitSeq: 6, durableSeq: 6, index: 0, |
| 276 | records: [{ role: "user", messageId: "unrelated", content: "unrelated" }] }], resetRequired: false }); |
| 277 | await microtasks(); |
| 278 | assert.equal(lookups, 2, "new committed coverage permits one retry of an inconclusive read"); |
| 279 | assert.ok(state.localSubmissions.submit); |
| 280 | } |
| 281 | return; |
| 282 | } |
| 283 | } else assert.equal(lookups, 0, "ordinary formal handoff makes no extra history read"); |
| 284 | assert.equal(state.localSubmissionOrder.length, 0); |
| 285 | assert.deepEqual(state.items, resident); |
| 286 | assert.equal(Object.keys(state.visibleSubmissionHandoffs).length, 0); |
| 287 | state = reducer(state, { type: "history_replace", items: [{ kind: "assistant", id: "m:later", text: "later", reasoning: "", streaming: false }], |
| 288 | startTurn: 101, totalTurns: 101, hasOlder: true, hasNewer: false, revision: 6 }); |
| 289 | assert.equal(state.localSubmissionOrder.length, 0); |
| 290 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 291 | }); |
| 292 | } |
| 293 | function initial(subscription: string): TranscriptFollowResponse { |
| 294 | return { |
| 295 | protocolVersion: 2, subscription, changes: [], resetRequired: false, |
| 296 | snapshot: { |
| 297 | protocolVersion: 1, snapshotId: "cut", identity: { sessionId: subscription, runtimeEpoch: "epoch", rewriteEpoch: 0, headId: "" }, |
| 298 | projectionRevision: 10, coveredThroughSeq: 4, durableSeq: 4, records: [], activeRecords: [], activeAttempts: [], |
| 299 | runtime: { status: "completed", pendingEvents: [], samplingCount: 0, toolCount: 0 }, before: 0, hasOlder: false, totalRecords: 1, totalTurns: 1, stale: false, |
| 300 | }, |
| 301 | history: { |
| 302 | status: "ready", snapshotSequence: 4, coverageSequence: 4, generation: "generation", totalTurns: 1, hasOlder: false, hasNewer: false, |
| 303 | messages: [{ messageId: "answer", position: 0, version: 1, role: "assistant", eventSequence: 4, visibleTurn: 1, |
| 304 | preview: "", contentRef: { digest: "canonical-message", bytes: 8192 } }], |
| 305 | }, |
| 306 | }; |
| 307 | } |
| 308 | |
| 309 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} running tab reconnect keeps its user before live reasoning when history and snapshot orders differ`, async () => { |
| 310 | const tab = `running-order-${remote}`, path = `/session/${tab}`; |
| 311 | const response = initial(tab); |
| 312 | response.history!.totalTurns = 2; |
| 313 | response.history!.messages = [ |
| 314 | { messageId: "previous", position: 10, version: 1, role: "assistant", eventSequence: 1, visibleTurn: 1, |
| 315 | inline: { role: "assistant", content: "previous answer" } }, |
| 316 | { messageId: "user", position: 11, version: 1, role: "user", eventSequence: 2, visibleTurn: 2, |
| 317 | preview: "make a tiger fly", inline: { role: "user", content: "make a tiger fly" } }, |
| 318 | ]; |
| 319 | response.snapshot!.totalRecords = 3; |
| 320 | response.snapshot!.totalTurns = 2; |
| 321 | response.snapshot!.runtime = { status: "in_progress", turnId: "current", pendingEvents: [], samplingCount: 1, toolCount: 0 }; |
| 322 | response.snapshot!.activeAttempts = [{ id: "attempt", messageId: "live", turnId: "current", nextIndex: 1 }]; |
| 323 | response.snapshot!.records = [ |
| 324 | { id: "m:previous", order: 0, message: { role: "assistant", messageId: "previous", content: "previous answer", historyTurn: 1 }, refs: [] }, |
| 325 | { id: "m:user", order: 1, message: { role: "user", messageId: "user", content: "make a tiger fly", historyTurn: 2, turnId: "current" }, refs: [] }, |
| 326 | { id: "m:live", order: 2, message: { role: "assistant", messageId: "live", content: "", reasoning: "thinking", historyTurn: 2, turnId: "current" }, refs: [] }, |
| 327 | ]; |
| 328 | const key = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 329 | const pending = deferred<TranscriptFollowResponse>(); |
| 330 | let polls = 0; |
| 331 | commands[key] = (_tab: string, request: FollowRequest) => request.close |
| 332 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 333 | : request.subscription ? polls++ === 0 ? pending.promise : new Promise<TranscriptFollowResponse>(() => {}) : Promise.resolve(response); |
| 334 | let state = initialState; |
| 335 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 336 | const source = new ChatSource(tab); |
| 337 | const assertVisibleOrder = () => { |
| 338 | source.update({ items: state.items, running: state.running, hydrating: false, hasOlder: false, loadingOlder: false }); |
| 339 | const order = source.getOrderSnapshot(); |
| 340 | assert.ok(order.indexOf("m:user") >= 0 && order.indexOf("m:user") < order.indexOf("m:live:reasoning"), |
| 341 | "the rendered reasoning stays in the user's turn"); |
| 342 | }; |
| 343 | try { |
| 344 | await follower.start(); |
| 345 | assert.deepEqual(state.items.filter(item => item.kind === "user" || item.kind === "assistant").map(item => item.id), |
| 346 | ["m:previous", "m:user", "m:live"]); |
| 347 | assert.deepEqual(getTranscriptStore().peek(tab, path)?.items.map(item => item.id), |
| 348 | ["m:previous", "m:user", "m:live"]); |
| 349 | assertVisibleOrder(); |
| 350 | pending.resolve({ protocolVersion: 2, subscription: tab, resetRequired: false, changes: [{ |
| 351 | revision: 11, firstSeq: 5, commitSeq: 5, durableSeq: 5, index: 0, |
| 352 | records: [{ role: "assistant", messageId: "live", turnId: "current", historyTurn: 2, |
| 353 | content: "final answer", reasoning: "thinking" }], |
| 354 | runtime: { status: "completed", turnId: "current", finalMessageId: "live", durationMs: 1000, |
| 355 | pendingEvents: [], samplingCount: 1, toolCount: 0 }, |
| 356 | }] }); |
| 357 | await microtasks(); |
| 358 | assertVisibleOrder(); |
| 359 | assert.deepEqual(state.items.filter(item => item.kind === "user" || item.kind === "assistant").map(item => item.id), |
| 360 | ["m:previous", "m:user", "m:live"]); |
| 361 | assert.equal(state.running, false); |
| 362 | assert.ok(state.items.some(item => item.kind === "assistant" && item.id === "m:live" && item.turnFinal)); |
| 363 | } finally { source.dispose(); follower.stop(); getTranscriptStore().evictTab(tab); } |
| 364 | }); |
| 365 | |
| 366 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} native active prefix keeps its fixed-cut order when older pages fill the gap`, async () => { |
| 367 | const tab = `native-prefix-order-${remote}`, path = `/session/${tab}`; |
| 368 | const response = initial(tab); |
| 369 | response.storageBackend = "legacy"; |
| 370 | delete response.history; |
| 371 | const notice = (order: number) => ({ id: `notice:${order}`, order, |
| 372 | message: { role: "notice", content: `step ${order}`, historyTurn: 1, turnId: "current" }, refs: [] }); |
| 373 | Object.assign(response.snapshot!, { |
| 374 | before: 4, hasOlder: true, totalRecords: 5, totalTurns: 1, |
| 375 | records: [notice(4)], |
| 376 | activeRecords: [ |
| 377 | { id: "m:user", order: 0, message: { role: "user", messageId: "user", content: "build", historyTurn: 1, turnId: "current" }, refs: [] }, |
| 378 | { id: "m:live", order: 1, message: { role: "assistant", messageId: "live", content: "", reasoning: "thinking", historyTurn: 1, turnId: "current" }, refs: [] }, |
| 379 | ], |
| 380 | activeAttempts: [{ id: "attempt", messageId: "live", turnId: "current", nextIndex: 1 }], |
| 381 | runtime: { status: "in_progress", turnId: "current", pendingEvents: [], samplingCount: 1, toolCount: 0 }, |
| 382 | }); |
| 383 | const pending = deferred<TranscriptFollowResponse>(); |
| 384 | commands[remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"] = (_tab: string, request: FollowRequest) => request.close |
| 385 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 386 | : request.subscription ? pending.promise : Promise.resolve(response); |
| 387 | commands[remote ? "RemoteTranscriptPageForTab" : "TranscriptPageForTab"] = async () => ({ |
| 388 | ...response.snapshot!, records: [notice(2), notice(3)], activeRecords: [], before: 2, hasOlder: true, |
| 389 | }); |
| 390 | let state = initialState; |
| 391 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 392 | try { |
| 393 | await follower.start(); |
| 394 | const older = await getTranscriptStore().loadOlder(tab, path); |
| 395 | assert.deepEqual(older?.items.map(item => item.id), ["m:user", "m:live", "he:notice:2", "he:notice:3", "he:notice:4"], |
| 396 | "paging into the gap cannot move earlier output above the active user"); |
| 397 | assert.equal(state.historyTotalTurns, 1, "installing an existing active user does not invent a new turn"); |
| 398 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 399 | }); |
| 400 | |
| 401 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} snapshot aliases retain the authoritative active message body`, async () => { |
| 402 | const tab = `active-alias-${remote}`, path = `/session/${tab}`; |
| 403 | const response = initial(tab); |
| 404 | response.history!.messages = []; |
| 405 | response.snapshot!.runtime.status = "in_progress"; |
| 406 | response.snapshot!.activeAttempts = [{ id: "attempt", messageId: "answer", turnId: "current", nextIndex: 1 }]; |
| 407 | response.snapshot!.records = [{ id: "view:answer", order: 0, |
| 408 | message: { role: "assistant", messageId: "answer", content: "old preview" }, refs: [] }]; |
| 409 | response.snapshot!.activeRecords = [{ id: "m:answer", order: 0, |
| 410 | message: { role: "assistant", messageId: "answer", content: "current active body" }, refs: [] }]; |
| 411 | const pending = deferred<TranscriptFollowResponse>(); |
| 412 | commands[remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"] = (_tab: string, request: FollowRequest) => request.close |
| 413 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 414 | : request.subscription ? pending.promise : Promise.resolve(response); |
| 415 | let state = initialState; |
| 416 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 417 | try { |
| 418 | await follower.start(); |
| 419 | assert.equal(state.items.length, 1, "aliases share one stable message node"); |
| 420 | assert.equal(state.items[0].kind === "assistant" && state.items[0].text, "current active body"); |
| 421 | assert.equal(state.live?.text, "current active body", "later deltas append to the authoritative active prefix"); |
| 422 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 423 | }); |
| 424 | |
| 425 | for (const remote of [false, true]) for (const during of ["baseline", "baseline rejection", "delta", "retry", "load"] as const) { |
| 426 | test(`${remote ? "remote" : "local"} stopping during ${during} fences current and future followers`, async t => { |
| 427 | t.mock.timers.enable({ apis: ["setTimeout"] }); |
| 428 | desktopStub.emitServiceState({ phase: "ready", generation: "running" }); |
| 429 | const tab = `stopping-${remote}-${during}`; |
| 430 | const requests: FollowRequest[] = []; |
| 431 | const baseline = deferred<TranscriptFollowResponse>(); |
| 432 | const delta = deferred<TranscriptFollowResponse>(); |
| 433 | const readStarted = deferred<void>(); |
| 434 | commands[remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"] = (_tab: string, request: FollowRequest) => { |
| 435 | requests.push(request); |
| 436 | if (request.close) return Promise.resolve({ protocolVersion: 2, changes: [] }); |
| 437 | readStarted.resolve(); |
| 438 | if (!request.subscription) return baseline.promise; |
| 439 | return delta.promise; |
| 440 | }; |
| 441 | let state = { ...initialState }; |
| 442 | const follower = new TranscriptSessionFollower(tab, "", remote, action => { state = reducer(state, action); }); |
| 443 | const starting = follower.start(); |
| 444 | if (during !== "load") await readStarted.promise; |
| 445 | if (during === "delta" || during === "retry") { |
| 446 | baseline.resolve(initial(tab)); await starting; |
| 447 | if (during === "retry") { |
| 448 | delta.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: true }); |
| 449 | await microtasks(); |
| 450 | } |
| 451 | } |
| 452 | const before = requests.length; |
| 453 | desktopStub.emitServiceState({ phase: "stopping", generation: "running" }); |
| 454 | const visible = state; |
| 455 | if (during === "baseline rejection") baseline.reject(new Error("stopping service rejected pending baseline")); |
| 456 | else baseline.resolve(initial(tab)); |
| 457 | delta.resolve({ protocolVersion: 2, subscription: tab, changes: [{ revision: 11, commitSeq: 4, durableSeq: 4, index: 0, event: { kind: "text", messageId: "answer", text: "late" } }], resetRequired: false }); |
| 458 | await starting; await microtasks(); |
| 459 | t.mock.timers.tick(1000); await microtasks(); |
| 460 | const late = new TranscriptSessionFollower(`${tab}-late`, "", remote, () => { throw new Error("stopped service must not publish"); }); |
| 461 | await late.start(); await follower.start(); |
| 462 | assert.equal(requests.length, before, "no cleanup, retry or new baseline after stopping"); |
| 463 | assert.equal(state, visible, "a late response cannot update the retained transcript"); |
| 464 | follower.stop(); late.stop(); |
| 465 | desktopStub.emitServiceState({ phase: "ready", generation: "replacement" }); |
| 466 | getTranscriptStore().evictTab(tab); |
| 467 | }); |
| 468 | } |
| 469 | |
| 470 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} follower preserves an outer snapshot identity from an older peer`, async () => { |
| 471 | const tab = `outer-record-${remote}`, path = `/session/${tab}`; |
| 472 | const response = initial(tab); |
| 473 | response.snapshot!.records = [ |
| 474 | { id: "view:older:1", order: 0, message: { role: "notice", content: "outer identity" }, refs: [] }, |
| 475 | { id: "tool:older-call", order: 1, message: { role: "tool", messageId: "tool-message", toolCallId: "older-call", toolName: "read_file", content: "result" }, refs: [] }, |
| 476 | { id: "", order: 2, message: { role: "notice", content: "legacy empty identity" }, refs: [] }, |
| 477 | ]; |
| 478 | response.snapshot!.totalRecords = 3; |
| 479 | const key = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 480 | const pending = deferred<TranscriptFollowResponse>(); |
| 481 | commands[key] = (_tab: string, request: FollowRequest) => request.close |
| 482 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 483 | : request.subscription ? pending.promise : Promise.resolve(response); |
| 484 | let state = initialState; |
| 485 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 486 | try { |
| 487 | await follower.start(); |
| 488 | assert.ok(state.items.some(item => item.kind === "notice" && item.text === "outer identity")); |
| 489 | assert.ok(getTranscriptStore().peek(tab, path)?.items.some(item => item.id === "he:view:older:1")); |
| 490 | assert.ok(state.items.some(item => item.kind === "tool" && item.id === "older-call")); |
| 491 | assert.ok(state.items.some(item => item.kind === "notice" && item.text === "legacy empty identity")); |
| 492 | assert.ok(!getTranscriptStore().peek(tab, path)?.items.some(item => item.id === "he:undefined")); |
| 493 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 494 | }); |
| 495 | |
| 496 | test("canonical tool history keeps its message identity when a tool call id is also present", async () => { |
| 497 | const tab = "canonical-tool-identity", path = `/session/${tab}`; |
| 498 | const response = initial(tab); |
| 499 | response.history!.messages = [{ |
| 500 | messageId: "tool-message", position: 0, version: 1, role: "tool", eventSequence: 4, visibleTurn: 1, |
| 501 | preview: "result", inline: { id: "tool-message", role: "tool", tool_call_id: "older-call", name: "read_file", content: "result" }, |
| 502 | }]; |
| 503 | commands.TranscriptFollowForTab = (_tab: string, request: FollowRequest) => Promise.resolve(request.close |
| 504 | ? { protocolVersion: 2, subscription: tab, changes: [], resetRequired: false } |
| 505 | : response); |
| 506 | let state = initialState; |
| 507 | const follower = new TranscriptSessionFollower(tab, path, false, action => { state = reducer(state, action); }); |
| 508 | try { |
| 509 | await follower.start(); |
| 510 | assert.ok(getTranscriptStore().peek(tab, path)?.items.some(item => item.kind === "tool" && item.id === "older-call")); |
| 511 | assert.ok(state.items.some(item => item.kind === "tool" && item.id === "older-call")); |
| 512 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 513 | }); |
| 514 | |
| 515 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} follower coalesces a legacy tool alias before the final answer`, async () => { |
| 516 | const tab = `legacy-tool-alias-${remote}`, path = `/session/${tab}`; |
| 517 | const response = initial(tab); |
| 518 | response.history!.messages = [ |
| 519 | { messageId: "user-1", position: 0, version: 1, role: "user", eventSequence: 1, visibleTurn: 1, |
| 520 | preview: "make it", inline: { id: "user-1", role: "user", content: "make it" } }, |
| 521 | { messageId: "call-owner", position: 1, version: 1, role: "assistant", eventSequence: 2, visibleTurn: 1, |
| 522 | preview: "", inline: { id: "call-owner", role: "assistant", content: "", tool_calls: [{ id: "call-1", name: "edit_file", arguments: '{"path":"blackhole.html"}' }] } }, |
| 523 | { messageId: "result-1", position: 2, version: 1, role: "tool", eventSequence: 3, visibleTurn: 1, |
| 524 | preview: "written", inline: { id: "result-1", role: "tool", tool_call_id: "call-1", name: "edit_file", content: "written" } }, |
| 525 | { messageId: "final-1", position: 3, version: 1, role: "assistant", eventSequence: 4, visibleTurn: 1, turnFinal: true, |
| 526 | preview: "done", inline: { id: "final-1", role: "assistant", content: "done" } }, |
| 527 | ]; |
| 528 | response.snapshot!.records = [{ id: "tool:call-1", order: 4, |
| 529 | message: { recordId: "tool:call-1", role: "tool", toolCallId: "call-1", toolName: "edit_file", content: "written" }, refs: [] }]; |
| 530 | response.snapshot!.totalRecords = 5; |
| 531 | const key = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 532 | const pending = deferred<TranscriptFollowResponse>(); |
| 533 | commands[key] = (_tab: string, request: FollowRequest) => request.close |
| 534 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 535 | : request.subscription ? pending.promise : Promise.resolve(response); |
| 536 | let state = initialState; |
| 537 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 538 | try { |
| 539 | await follower.start(); |
| 540 | const tools = state.items.filter(item => item.kind === "tool" && item.id === "call-1"); |
| 541 | assert.equal(tools.length, 1, "formal result and tool:<callId> alias project one node"); |
| 542 | assert.equal(tools[0].kind === "tool" && tools[0].args, '{"path":"blackhole.html"}'); |
| 543 | assert.ok(state.items.findIndex(item => item.id === "call-1") < state.items.findIndex(item => item.id === "m:final-1")); |
| 544 | const source = new ChatSource(tab); |
| 545 | source.update({ items: state.items, running: false, hydrating: false, hasOlder: false, loadingOlder: false }); |
| 546 | await Promise.resolve(); |
| 547 | const process = source.getNodeSnapshot("m:user-1:process"); |
| 548 | assert.ok(process?.kind === "process" && process.foldable && process.collapsed, |
| 549 | "normal completion folds after the duplicate trailing alias is removed"); |
| 550 | source.dispose(); |
| 551 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 552 | }); |
| 553 | |
| 554 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} follower joins a lazy formal result through its observation`, async () => { |
| 555 | const tab = `lazy-formal-result-${remote}`, path = `/session/${tab}`; |
| 556 | const response = initial(tab); |
| 557 | response.history!.messages = [ |
| 558 | { messageId: "user", position: 0, version: 1, role: "user", eventSequence: 1, visibleTurn: 1, |
| 559 | preview: "build", inline: { id: "user", role: "user", content: "build" } }, |
| 560 | { messageId: "owner", position: 1, version: 1, role: "assistant", eventSequence: 2, visibleTurn: 1, |
| 561 | preview: "", inline: { id: "owner", role: "assistant", content: "", tool_calls: [{ id: "lazy", name: "write_file", arguments: '{"path":"lazy.html"}' }] }, |
| 562 | toolObservations: { lazy: { state: "completed", messageId: "result", version: 1, contentRef: { digest: "body", bytes: 8192 } } } }, |
| 563 | { messageId: "result", position: 2, version: 1, role: "tool", eventSequence: 3, visibleTurn: 1, |
| 564 | preview: "written", contentRef: { digest: "body", bytes: 8192 } }, |
| 565 | { messageId: "final", position: 3, version: 1, role: "assistant", eventSequence: 4, visibleTurn: 1, turnFinal: true, |
| 566 | preview: "done", inline: { id: "final", role: "assistant", content: "done" } }, |
| 567 | ]; |
| 568 | response.snapshot!.records = [{ id: "tool:lazy", order: 4, |
| 569 | message: { recordId: "tool:lazy", role: "tool", toolCallId: "lazy", toolName: "write_file", content: "written" }, refs: [] }]; |
| 570 | response.snapshot!.totalRecords = 5; |
| 571 | const key = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 572 | const pending = deferred<TranscriptFollowResponse>(); |
| 573 | commands[key] = (_tab: string, request: FollowRequest) => request.close |
| 574 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 575 | : request.subscription ? pending.promise : Promise.resolve(response); |
| 576 | let state = initialState; |
| 577 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 578 | try { |
| 579 | await follower.start(); |
| 580 | const tools = state.items.filter((item): item is Extract<import("../lib/useController").Item, { kind: "tool" }> => item.kind === "tool"); |
| 581 | assert.equal(tools.length, 1); |
| 582 | assert.equal(tools[0].id, "lazy"); |
| 583 | assert.equal(tools[0].args, '{"path":"lazy.html"}'); |
| 584 | assert.ok(state.items.findIndex(item => item.id === "lazy") < state.items.findIndex(item => item.id === "m:final")); |
| 585 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 586 | }); |
| 587 | |
| 588 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} malformed snapshot leaves the resident store unchanged`, async () => { |
| 589 | const tab = `invalid-record-${remote}`, path = `/session/${tab}`; |
| 590 | getTranscriptStore().installSlice(tab, path, { |
| 591 | entries: [{ entryId: "m:resident", turn: 1, order: 0, message: { role: "assistant", messageId: "resident", content: "resident" }, refs: [] }], |
| 592 | nextCursor: "", hasOlder: false, newerCursor: "", hasNewer: false, totalTurns: 1, startTurn: 1, endTurn: 1, |
| 593 | revision: 1, revisionKnown: true, digest: "resident", stale: false, |
| 594 | }); |
| 595 | const response = initial(tab); |
| 596 | response.snapshot!.records = [{ id: "", order: 0, message: { role: "notice", content: "invalid" }, |
| 597 | refs: [{ snapshotId: "cut", recordId: "", path: ["content"], bytes: 100 }] }]; |
| 598 | response.snapshot!.totalRecords = 1; |
| 599 | const key = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 600 | commands[key] = (_tab: string, request: FollowRequest) => Promise.resolve(request.close |
| 601 | ? { protocolVersion: 2, subscription: tab, changes: [], resetRequired: false } |
| 602 | : response); |
| 603 | const follower = new TranscriptSessionFollower(tab, path, remote, () => undefined); |
| 604 | try { |
| 605 | await assert.rejects(follower.start(), /transcript snapshot content identity missing/); |
| 606 | const resident = getTranscriptStore().peek(tab, path); |
| 607 | assert.equal(resident?.digest, "resident"); |
| 608 | assert.ok(resident?.items.some(item => item.id === "m:resident")); |
| 609 | assert.ok(!resident?.items.some(item => item.id === "he:undefined")); |
| 610 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 611 | }); |
| 612 | |
| 613 | test("reducer rejection does not commit a prepared transcript replacement", async () => { |
| 614 | const tab = "reducer-reject", path = `/session/${tab}`; |
| 615 | getTranscriptStore().installSlice(tab, path, { |
| 616 | entries: [{ entryId: "m:resident", turn: 1, order: 0, message: { role: "assistant", messageId: "resident", content: "resident" }, refs: [] }], |
| 617 | nextCursor: "", hasOlder: false, newerCursor: "", hasNewer: false, totalTurns: 1, startTurn: 1, endTurn: 1, |
| 618 | revision: 1, revisionKnown: true, digest: "resident", stale: false, |
| 619 | }); |
| 620 | commands.TranscriptFollowForTab = (_tab: string, request: FollowRequest) => Promise.resolve(request.close |
| 621 | ? { protocolVersion: 2, subscription: tab, changes: [], resetRequired: false } |
| 622 | : initial(tab)); |
| 623 | const follower = new TranscriptSessionFollower(tab, path, false, action => { |
| 624 | if (action.type === "transcript_v2_snapshot") throw new Error("reducer rejected snapshot"); |
| 625 | }); |
| 626 | try { |
| 627 | await assert.rejects(follower.start(), /reducer rejected snapshot/); |
| 628 | const resident = getTranscriptStore().peek(tab, path); |
| 629 | assert.equal(resident?.digest, "resident"); |
| 630 | assert.ok(resident?.items.some(item => item.id === "m:resident")); |
| 631 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 632 | }); |
| 633 | |
| 634 | test("snapshot rejects duplicate durable identities and mismatched content references", async () => { |
| 635 | const cases: Array<{ name: string; records: NonNullable<TranscriptFollowResponse["snapshot"]>["records"]; pattern: RegExp }> = [ |
| 636 | { name: "duplicate", records: [ |
| 637 | { id: "same", order: 0, message: { role: "notice", recordId: "same", content: "first" }, refs: [] }, |
| 638 | { id: "same", order: 1, message: { role: "notice", recordId: "same", content: "second" }, refs: [] }, |
| 639 | ], pattern: /duplicate transcript snapshot record identity/ }, |
| 640 | { name: "content-ref", records: [{ id: "owner", order: 0, message: { role: "notice", recordId: "owner", content: "preview" }, |
| 641 | refs: [{ snapshotId: "cut", recordId: "different", path: ["content"], bytes: 100 }] }], pattern: /content identity mismatch/ }, |
| 642 | { name: "embedded-record", records: [{ id: "owner", order: 0, message: { role: "notice", recordId: "different", content: "preview" }, refs: [] }], |
| 643 | pattern: /record identity mismatch/ }, |
| 644 | ]; |
| 645 | for (const fixture of cases) { |
| 646 | const tab = `invalid-${fixture.name}`; |
| 647 | const response = initial(tab); |
| 648 | response.snapshot!.records = fixture.records; |
| 649 | response.snapshot!.totalRecords = fixture.records.length; |
| 650 | commands.TranscriptFollowForTab = (_tab: string, request: FollowRequest) => Promise.resolve(request.close |
| 651 | ? { protocolVersion: 2, subscription: tab, changes: [], resetRequired: false } |
| 652 | : response); |
| 653 | const follower = new TranscriptSessionFollower(tab, "", false, () => undefined); |
| 654 | try { await assert.rejects(follower.start(), fixture.pattern); } |
| 655 | finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 656 | } |
| 657 | }); |
| 658 | |
| 659 | for (const remote of [false, true]) { |
| 660 | test(`${remote ? "remote" : "local"} follower retains an empty canonical body reference as a loadable assistant node`, async () => { |
| 661 | const tab = remote ? "follow-ref-remote" : "follow-ref-local"; |
| 662 | let state = initialState; |
| 663 | const requests: FollowRequest[] = []; |
| 664 | const poll = deferred<TranscriptFollowResponse>(); |
| 665 | const read = async (tabId: string, request: FollowRequest): Promise<TranscriptFollowResponse> => { |
| 666 | assert.equal(tabId, tab); requests.push(request); |
| 667 | if (request.close) return { protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }; |
| 668 | if (!request.subscription) return initial(tab); |
| 669 | return poll.promise; |
| 670 | }; |
| 671 | const key = remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"; |
| 672 | const wrong = remote ? "TranscriptFollowForTab" : "RemoteTranscriptFollowForTab"; |
| 673 | commands[key] = read; |
| 674 | commands[wrong] = () => { throw new Error("cross-host fallback is forbidden"); }; |
| 675 | commands.SendForTab = () => { throw new Error("history recovery must not invoke the model"); }; |
| 676 | commands.RemoteSendForTab = commands.SendForTab; |
| 677 | const follower = new TranscriptSessionFollower(tab, `/session/${tab}`, remote, action => { state = reducer(state, action); }); |
| 678 | try { |
| 679 | await follower.start(); |
| 680 | const assistant = state.items.find(item => item.kind === "assistant"); |
| 681 | assert.ok(assistant, "empty inline preview with a canonical ref must not disappear"); |
| 682 | assert.equal(assistant.id, "m:answer"); |
| 683 | assert.equal(getTranscriptStore().hasContentReference(tab, "m:answer", "content"), true); |
| 684 | assert.equal(state.running, false); |
| 685 | assert.equal(state.transcriptProtocol, 2); |
| 686 | assert.equal(requests.length, 2, "follow starts only after initial installation"); |
| 687 | } finally { follower.stop(); } |
| 688 | await microtasks(); |
| 689 | assert.ok(requests.some(request => request.close && request.subscription === tab)); |
| 690 | }); |
| 691 | } |
| 692 | |
| 693 | test("stopped session follower cannot install a delayed baseline into its replaced tab", async () => { |
| 694 | const delayed = deferred<TranscriptFollowResponse>(); |
| 695 | const started = deferred<void>(); |
| 696 | const requests: FollowRequest[] = []; |
| 697 | commands.TranscriptFollowForTab = (_tab: string, request: FollowRequest) => { |
| 698 | requests.push(request); |
| 699 | started.resolve(); |
| 700 | return request.close ? Promise.resolve({ protocolVersion: 2, subscription: "stale", changes: [], resetRequired: false }) : delayed.promise; |
| 701 | }; |
| 702 | let state = initialState; |
| 703 | const follower = new TranscriptSessionFollower("stale-tab", "/session/stale", false, action => { state = reducer(state, action); }); |
| 704 | const loading = follower.start(); await started.promise; follower.stop(); delayed.resolve(initial("stale")); await loading; await microtasks(); |
| 705 | assert.equal(state.items.length, 0); |
| 706 | assert.ok(requests.some(request => request.close && request.subscription === "stale")); |
| 707 | }); |
| 708 | |
| 709 | test("stopping before lazy follow startup prevents a backend subscription", async () => { |
| 710 | let requests = 0; |
| 711 | commands.TranscriptFollowForTab = async () => { requests++; return initial("cancelled-load"); }; |
| 712 | const follower = new TranscriptSessionFollower("cancelled-load", "/session/cancelled-load", false, () => { |
| 713 | assert.fail("cancelled module load cannot publish state"); |
| 714 | }); |
| 715 | const loading = follower.start(); |
| 716 | follower.stop(); |
| 717 | await loading; |
| 718 | assert.equal(requests, 0); |
| 719 | }); |
| 720 | |
| 721 | for (const remote of [false, true]) test(`${remote ? "remote" : "local"} active reference is complete before suffix polling`, async () => { |
| 722 | const response = initial(`prefix-${remote}`); |
| 723 | response.snapshot!.runtime.status = "in_progress"; |
| 724 | response.snapshot!.activeAttempts = [{ id: "attempt", messageId: "answer", turnId: "turn", nextIndex: 1 }]; |
| 725 | response.snapshot!.records = [{ id: "m:answer", order: 0, message: { role: "assistant", messageId: "answer", content: "truncated" }, |
| 726 | refs: [{ snapshotId: "cut", recordId: "m:answer", path: ["content"], bytes: 20 }] }]; |
| 727 | const content = deferred<{ data: string; nextOffset: number; done: boolean; stale: boolean }>(); |
| 728 | const poll = deferred<TranscriptFollowResponse>(); |
| 729 | let polls = 0; |
| 730 | commands[remote ? "RemoteTranscriptContentForTab" : "TranscriptContentForTab"] = () => content.promise; |
| 731 | commands[remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"] = (_tab: string, request: FollowRequest) => { |
| 732 | if (request.close) return Promise.resolve({ protocolVersion: 2, changes: [], resetRequired: false, subscription: response.subscription }); |
| 733 | if (!request.subscription) return Promise.resolve(response); |
| 734 | polls++; return poll.promise; |
| 735 | }; |
| 736 | let state = initialState; |
| 737 | const follower = new TranscriptSessionFollower(`prefix-${remote}`, "", remote, action => { state = reducer(state, action); }); |
| 738 | const starting = follower.start(); |
| 739 | await microtasks(); |
| 740 | assert.equal(polls, 0); |
| 741 | assert.equal(state.items.length, 0, "partial baseline is not published"); |
| 742 | content.resolve({ data: "complete prefix", nextOffset: 15, done: true, stale: false }); |
| 743 | await starting; |
| 744 | assert.equal(state.live?.text, "complete prefix"); |
| 745 | assert.equal(polls, 1); |
| 746 | follower.stop(); |
| 747 | }); |
| 748 | |
| 749 | test("v2 keeps deferred bodies through subsequent samples and attaches terminal time by backend message identity", () => { |
| 750 | let state: import("../lib/useController").State = { ...initialState, transcriptProtocol: 2, running: true, activeTurnId: "turn", |
| 751 | items: [ |
| 752 | { kind: "assistant" as const, id: "m:deferred", text: "", reasoning: "", streaming: false }, |
| 753 | { kind: "assistant" as const, id: "m:final", text: "answer", reasoning: "thinking", streaming: false }, |
| 754 | ] }; |
| 755 | state = reducer(state, { type: "event", e: { kind: "text", messageId: "next", text: "later sample" } }); |
| 756 | assert.ok(state.items.some(item => item.id === "m:deferred"), "text frames must not remove unloaded body owners"); |
| 757 | state = reducer(state, { type: "event", e: { kind: "turn_done", turnId: "turn" } }); |
| 758 | state = reducer(state, { type: "transcript_runtime", runtime: { turnId: "turn", status: "completed", |
| 759 | finalMessageId: "final", durationMs: 933524, samplingCount: 72, toolCount: 72, pendingEvents: [] } }); |
| 760 | const final = state.items.find(item => item.id === "m:final"); |
| 761 | assert.ok(final?.kind === "assistant"); |
| 762 | assert.equal(final.turnDurationMs, 933524); |
| 763 | assert.equal(final.turnFinal, true); |
| 764 | assert.ok(state.items.some(item => item.id === "m:deferred"), "terminal must retain deferred body owners"); |
| 765 | const later = state.items.find(item => item.id === "m:next"); |
| 766 | assert.ok(later?.kind === "assistant"); |
| 767 | assert.equal(later.turnDurationMs, undefined, "array-tail sample is not the final reply"); |
| 768 | }); |
| 769 |