| 1 | import assert from "node:assert/strict"; |
| 2 | import test from "node:test"; |
| 3 | import type { FollowRequest, TranscriptFollowResponse, Change } 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 | installDesktopHostStub(commands); |
| 9 | const [{ TranscriptSessionFollower }, { initialState, reducer }, { getTranscriptStore }, { ChatSource }] = await Promise.all([ |
| 10 | import("../lib/transcriptSessionFollower"), import("../lib/useController"), import("../lib/transcriptStore"), import("../lib/chatViewSource"), |
| 11 | ]); |
| 12 | const text = "对比源码信息应该更准一些"; |
| 13 | const prefix = "[Mid-turn steer queued by the user. Do not treat this as a new task; use it only as additional guidance for the current task after completing the current step.]"; |
| 14 | const row = (id: string) => ({ recordId: `m:${id}`, messageId: id, role: "notice", content: `↪ ${text}`, turnId: "turn" }); |
| 15 | const initial = (tab: string): TranscriptFollowResponse => ({ |
| 16 | protocolVersion: 2, subscription: tab, changes: [], resetRequired: false, |
| 17 | snapshot: { |
| 18 | protocolVersion: 1, snapshotId: "cut", identity: { sessionId: tab, runtimeEpoch: "epoch", rewriteEpoch: 0, headId: "" }, |
| 19 | projectionRevision: 10, coveredThroughSeq: 4, durableSeq: 4, records: [], activeRecords: [], activeAttempts: [], |
| 20 | runtime: { status: "in_progress", turnId: "turn", pendingEvents: [], samplingCount: 1, toolCount: 0 }, |
| 21 | before: 0, hasOlder: false, totalRecords: 0, totalTurns: 1, stale: false, |
| 22 | }, |
| 23 | history: { status: "ready", snapshotSequence: 4, coverageSequence: 4, generation: "generation", totalTurns: 1, |
| 24 | hasOlder: false, hasNewer: false, messages: [] }, |
| 25 | }); |
| 26 | async function microtasks() { for (let i = 0; i < 60; i++) await Promise.resolve(); } |
| 27 | |
| 28 | for (const remote of [false, true]) for (const ordering of ["record-first", "event-first", "legacy-event"] as const) { |
| 29 | test(`${remote ? "remote" : "local"} ${ordering} steer keeps one row and acknowledges the exact inbox item`, async () => { |
| 30 | const tab = `steer-${remote}-${ordering}`, path = `/session/${tab}`; |
| 31 | let deliver!: (value: TranscriptFollowResponse) => void; |
| 32 | const pending = new Promise<TranscriptFollowResponse>(resolve => { deliver = resolve; }); |
| 33 | let polls = 0; |
| 34 | const response = initial(tab); |
| 35 | const wireRow = (id: string) => ordering === "legacy-event" |
| 36 | ? { ...row(id), messageId: undefined, recordId: `m:${id}:notice:0` } : row(id); |
| 37 | commands[remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"] = (_tab: string, request: FollowRequest) => request.close |
| 38 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 39 | : request.subscription ? polls++ === 0 ? pending : new Promise(() => {}) : Promise.resolve(response); |
| 40 | let state = initialState; |
| 41 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 42 | const source = new ChatSource(tab); |
| 43 | try { |
| 44 | await follower.start(); |
| 45 | const changes: Change[] = []; |
| 46 | let revision = 10, coverage = 4; |
| 47 | const record = (id: string) => changes.push({ revision: ++revision, firstSeq: ++coverage, commitSeq: coverage, |
| 48 | durableSeq: coverage, index: 0, records: [wireRow(id)] }); |
| 49 | const event = (id: string) => changes.push({ revision: ++revision, commitSeq: coverage, durableSeq: coverage, index: 0, |
| 50 | event: { kind: "steer", text, messageId: ordering === "legacy-event" ? undefined : id, itemId: `inbox-${id}` } }); |
| 51 | // Equal text belongs to two distinct user actions; each receipt is also |
| 52 | // redelivered to verify identity-based idempotency rather than text dedup. |
| 53 | for (const id of ["first", "second"]) { |
| 54 | if (ordering === "event-first") { event(id); record(id); } else { record(id); event(id); } |
| 55 | event(id); |
| 56 | } |
| 57 | deliver({ protocolVersion: 2, subscription: tab, resetRequired: false, changes }); |
| 58 | await microtasks(); |
| 59 | assert.deepEqual(state.items.map(item => item.id), ["he:m:first", "he:m:second"]); |
| 60 | assert.deepEqual(state.guidanceConsumed, { key: "inbox-second", itemId: "inbox-second", text }); |
| 61 | source.update({ items: state.items, running: state.running, hydrating: false, hasOlder: false, loadingOlder: false }); |
| 62 | assert.deepEqual(source.getOrderSnapshot().filter(id => id.startsWith("he:")), ["he:m:first", "he:m:second"]); |
| 63 | const projection = getTranscriptStore().peek(tab, path)!; |
| 64 | state = reducer(state, { type: "transcript_records", projection: { ...projection, removeIds: [] }, confirmedUsers: [] }); |
| 65 | assert.equal(state.items.length, 2, "subsequent refresh must not retain an extra live row"); |
| 66 | assert.equal(state.guidanceConsumed?.itemId, "inbox-second", "refresh preserves the independent consumption receipt"); |
| 67 | |
| 68 | // Reconnect joins canonical history with the runtime snapshot. Both must |
| 69 | // address the same message even though the provider role is user. |
| 70 | follower.stop(); |
| 71 | response.snapshot!.projectionRevision = revision; |
| 72 | response.snapshot!.coveredThroughSeq = coverage; |
| 73 | response.snapshot!.durableSeq = coverage; |
| 74 | response.snapshot!.totalRecords = 2; |
| 75 | response.snapshot!.records = ["first", "second"].map((id, order) => ({ id: wireRow(id).recordId, order, message: wireRow(id), refs: [] })); |
| 76 | response.history!.messages = ["first", "second"].map((id, position) => ({ messageId: id, position, version: 1, |
| 77 | role: "user", eventSequence: coverage, visibleTurn: 1, |
| 78 | inline: { id, role: "user", content: `${prefix}\n\n${text}`, raw_content: text, origin: "user" } })); |
| 79 | await follower.start(); |
| 80 | assert.deepEqual(state.items.map(item => item.id), ["he:m:first", "he:m:second"]); |
| 81 | } finally { source.dispose(); follower.stop(); getTranscriptStore().evictTab(tab); } |
| 82 | }); |
| 83 | } |
| 84 | |
| 85 | test("legacy standalone events render while v2 uncorrelated events only acknowledge", () => { |
| 86 | for (const protocol of [undefined, 2] as const) { |
| 87 | const state = reducer({ ...initialState, transcriptProtocol: protocol }, { type: "event", e: { kind: "steer", text, itemId: "queued" } }); |
| 88 | assert.equal(state.items.length, protocol === 2 ? 0 : 1); |
| 89 | assert.equal(state.guidanceConsumed?.itemId, "queued"); |
| 90 | assert.equal(reducer(state, { type: "reset" }).guidanceConsumed, undefined); |
| 91 | } |
| 92 | }); |
| 93 | |
| 94 | test("legacy native pages share steer identity while content refs retain their original owner", async () => { |
| 95 | const { nativeSnapshotWindow } = await import("../lib/nativeTranscriptHistory"); |
| 96 | const snapshot = initial("native-legacy").snapshot!; |
| 97 | snapshot.totalRecords = 1; |
| 98 | const recordId = "m:old-message:notice:0"; |
| 99 | snapshot.records = [{ id: recordId, order: 0, message: { role: "notice", content: `↪ ${text}` }, |
| 100 | refs: [{ snapshotId: "cut", recordId, path: ["content"], bytes: 100 }] }]; |
| 101 | const entry = nativeSnapshotWindow(snapshot as unknown as import("../lib/transcriptProtocol").TranscriptSnapshot).entries[0]; |
| 102 | assert.equal(entry.entryId, "m:old-message"); |
| 103 | assert.equal(entry.message.messageId, "old-message"); |
| 104 | assert.equal(entry.refs[0].transcriptRef?.recordId, recordId, "normalizing display must not invalidate the server's content ref"); |
| 105 | }); |
| 106 | |
| 107 | test("unapplied queue warnings join canonical records in either order", async () => { |
| 108 | const [{ historyMessagesToItems }, { canonicalMessage }, { nativeSnapshotWindow }] = await Promise.all([ |
| 109 | import("../lib/historyItems"), import("../lib/canonicalTranscriptBackend"), import("../lib/nativeTranscriptHistory"), |
| 110 | ]); |
| 111 | const warning = "Guidance was not applied because the turn ended before it could be processed. Send it again if it is still needed:\nSame guidance"; |
| 112 | for (const ordering of ["event-first", "record-first"] as const) { |
| 113 | let state: import("../lib/useController").State = { ...initialState, transcriptProtocol: 2 }; |
| 114 | const formal: import("../lib/useController").Item[] = []; |
| 115 | for (const id of ["first", "second"]) { |
| 116 | const event = { type: "event" as const, e: { kind: "notice" as const, code: "unapplied_steer", level: "warn" as const, |
| 117 | messageId: id, itemId: `inbox-${id}`, text: warning } }; |
| 118 | if (ordering === "event-first") state = reducer(state, event); |
| 119 | formal.push(...historyMessagesToItems([{ role: "notice", code: "unapplied_steer", level: "warn", content: warning, |
| 120 | messageId: id, recordId: `m:${id}:notice:0` }], "snapshot:").items); |
| 121 | state = reducer(state, { type: "transcript_records", confirmedUsers: [], projection: { |
| 122 | items: [...formal], removeIds: [], startTurn: 0, endTurn: 0, totalTurns: 0, hasOlder: false, hasNewer: false, |
| 123 | revision: formal.length, revisionKnown: true, digest: "unapplied", |
| 124 | } }); |
| 125 | state = reducer(state, event); |
| 126 | state = reducer(state, event); |
| 127 | } |
| 128 | assert.deepEqual(state.items.map(item => item.id), ["he:m:first", "he:m:second"], ordering); |
| 129 | } |
| 130 | |
| 131 | const oldRecordId = "m:old-warning:notice:0"; |
| 132 | const legacy = { role: "notice", code: "unapplied_steer", level: "warn" as const, content: warning, recordId: oldRecordId }; |
| 133 | assert.equal(historyMessagesToItems([legacy], "snapshot:").items[0]?.id, "he:m:old-warning"); |
| 134 | const snapshot = initial("old-warning").snapshot!; |
| 135 | snapshot.totalRecords = 1; |
| 136 | snapshot.records = [{ id: oldRecordId, order: 0, message: legacy, refs: [] }]; |
| 137 | assert.equal(nativeSnapshotWindow(snapshot as unknown as import("../lib/transcriptProtocol").TranscriptSnapshot).entries[0]?.entryId, "m:old-warning"); |
| 138 | |
| 139 | const persisted = canonicalMessage({ messageId: "queued", role: "tool", preview: "" } as import("../generated/desktopContract.generated").PersistentMessage, |
| 140 | { id: "queued", role: "tool", local_only: true, content: `${prefix}\nSame guidance` }); |
| 141 | assert.deepEqual({ role: persisted.role, messageId: persisted.messageId, code: persisted.code, content: persisted.content }, |
| 142 | { role: "notice", messageId: "queued", code: "unapplied_steer", content: "\nSame guidance" }); |
| 143 | const multiline = canonicalMessage({ messageId: "multiline", role: "tool", preview: "" } as import("../generated/desktopContract.generated").PersistentMessage, |
| 144 | { id: "multiline", role: "tool", local_only: true, content: `${prefix}\nFirst line\nSecond line` }); |
| 145 | const multilineItem = historyMessagesToItems([multiline], "canonical:").items[0]; |
| 146 | if (multilineItem?.kind !== "notice") throw new Error("multiline guidance must render as a notice"); |
| 147 | assert.ok(multilineItem.text.endsWith("First line\nSecond line")); |
| 148 | }); |
| 149 | |
| 150 | for (const remote of [false, true]) for (const ordering of ["record-first", "event-first"] as const) { |
| 151 | test(`${remote ? "remote" : "local"} ${ordering} unapplied warning survives follower reconnect as one row`, async () => { |
| 152 | const tab = `unapplied-${remote}-${ordering}`, path = `/session/${tab}`, id = "queued"; |
| 153 | const warning = "Guidance was not applied because the turn ended before it could be processed. Send it again if it is still needed:\nSame guidance"; |
| 154 | const wireRow = { role: "notice", messageId: id, recordId: `m:${id}:notice:0`, code: "unapplied_steer", level: "warn", content: warning }; |
| 155 | let deliver!: (value: TranscriptFollowResponse) => void; |
| 156 | const pending = new Promise<TranscriptFollowResponse>(resolve => { deliver = resolve; }); |
| 157 | let polls = 0; |
| 158 | const response = initial(tab); |
| 159 | commands[remote ? "RemoteTranscriptFollowForTab" : "TranscriptFollowForTab"] = (_tab: string, request: FollowRequest) => request.close |
| 160 | ? Promise.resolve({ protocolVersion: 2, subscription: tab, changes: [], resetRequired: false }) |
| 161 | : request.subscription ? polls++ === 0 ? pending : new Promise(() => {}) : Promise.resolve(response); |
| 162 | let state = initialState; |
| 163 | const follower = new TranscriptSessionFollower(tab, path, remote, action => { state = reducer(state, action); }); |
| 164 | try { |
| 165 | await follower.start(); |
| 166 | const record: Change = { revision: 11, firstSeq: 5, commitSeq: 5, durableSeq: 5, index: 0, records: [wireRow] }; |
| 167 | const event: Change = { revision: 12, commitSeq: ordering === "event-first" ? 4 : 5, durableSeq: 5, index: 0, |
| 168 | event: { kind: "notice", code: "unapplied_steer", level: "warn", messageId: id, itemId: "inbox-queued", text: warning } }; |
| 169 | deliver({ protocolVersion: 2, subscription: tab, resetRequired: false, |
| 170 | changes: ordering === "record-first" ? [record, event] : [{ ...event, revision: 11 }, { ...record, revision: 12 }] }); |
| 171 | await microtasks(); |
| 172 | assert.deepEqual(state.items.map(item => item.id), ["he:m:queued"]); |
| 173 | follower.stop(); |
| 174 | response.snapshot!.projectionRevision = 12; |
| 175 | response.snapshot!.coveredThroughSeq = 5; |
| 176 | response.snapshot!.durableSeq = 5; |
| 177 | response.snapshot!.totalRecords = 1; |
| 178 | response.snapshot!.records = [{ id: wireRow.recordId, order: 0, message: wireRow, refs: [] }]; |
| 179 | response.history!.messages = [{ messageId: id, position: 0, version: 1, role: "tool", eventSequence: 5, visibleTurn: 1, |
| 180 | inline: { id, role: "tool", content: `${prefix}\nSame guidance`, local_only: true, |
| 181 | tool_call_id: "__reasonix_local_only__", name: "__reasonix_local_only__" } }]; |
| 182 | await follower.start(); |
| 183 | assert.deepEqual(state.items.map(item => item.id), ["he:m:queued"]); |
| 184 | } finally { follower.stop(); getTranscriptStore().evictTab(tab); } |
| 185 | }); |
| 186 | } |
| 187 |