| 1 | import type { |
| 2 | HistoryOutlineRequest, HistoryOutlinePage, |
| 3 | HistoryWindowPage, |
| 4 | HistoryWindowRequest, |
| 5 | MessageFieldPage, |
| 6 | MessageHistoryPage, |
| 7 | MessageLocation, |
| 8 | PersistentMessage, |
| 9 | Ref as SessionContentRef, |
| 10 | SearchHistoryPage, |
| 11 | SessionHistoryContentChunk, |
| 12 | SessionOpenView, |
| 13 | SessionHistoryReadHandle, |
| 14 | FollowRequest, TranscriptFollowResponse, Change, |
| 15 | } from "../generated/desktopContract.generated"; |
| 16 | import type { HistoryContentChunk, HistoryContentRef, HistoryMessage, HistorySlice, HistorySliceRequest, WireEvent } from "./types"; |
| 17 | |
| 18 | export interface SessionReaderBindings { |
| 19 | BeginSessionHistoryReadForTab(tabID: string): Promise<SessionHistoryReadHandle>; |
| 20 | ReleaseSessionHistoryRead(id: string): Promise<void>; |
| 21 | ReadSessionHistoryWindow(id: string, req: HistoryWindowRequest): Promise<HistoryWindowPage>; |
| 22 | ReadSessionHistorySlice(id: string, req: HistorySliceRequest): Promise<{ status: string; page: HistorySlice }>; |
| 23 | ReadSessionHistoryOutline(id: string, req: HistoryOutlineRequest): Promise<HistoryOutlinePage>; |
| 24 | LocateSessionHistoryMessage(id: string, messageID: string, snapshot: number): Promise<MessageLocation>; |
| 25 | SearchSessionHistoryRead(id: string, text: string, cursor: string, limit: number): Promise<SearchHistoryPage>; |
| 26 | SessionHistoryOutlineForTab?(tabID: string, req: HistoryOutlineRequest): Promise<HistoryOutlinePage>; |
| 27 | RemoteSessionHistoryOutlineForTab?(tabID: string, req: HistoryOutlineRequest): Promise<HistoryOutlinePage>; |
| 28 | TranscriptFollowForTab(tabID: string, request: FollowRequest): Promise<TranscriptFollowResponse>; |
| 29 | RemoteTranscriptFollowForTab(tabID: string, request: FollowRequest): Promise<TranscriptFollowResponse>; |
| 30 | SessionHistoryPageForTab(tabID: string, cursor: string, limit: number): Promise<MessageHistoryPage>; |
| 31 | SessionOpenForTab(tabID: string): Promise<SessionOpenView>; |
| 32 | SessionHistoryContentForTab(tabID: string, ref: SessionContentRef, offset: number): Promise<SessionHistoryContentChunk>; |
| 33 | RemoteSessionOpenForTab(tabID: string): Promise<SessionOpenView>; |
| 34 | RemoteSessionHistoryPageForTab(tabID: string, cursor: string, limit: number): Promise<MessageHistoryPage>; |
| 35 | RemoteSessionHistoryContentForTab(tabID: string, ref: SessionContentRef, offset: number): Promise<SessionHistoryContentChunk>; |
| 36 | SearchSessionHistoryForTab(tabID: string, textQuery: string, cursor: string, limit: number): Promise<SearchHistoryPage>; |
| 37 | RemoteSearchSessionHistoryForTab(tabID: string, textQuery: string, cursor: string, limit: number): Promise<SearchHistoryPage>; |
| 38 | LocateSessionMessageForTab(tabID: string, messageID: string, snapshot: number): Promise<MessageLocation>; |
| 39 | RemoteLocateSessionMessageForTab(tabID: string, messageID: string, snapshot: number): Promise<MessageLocation>; |
| 40 | // history-window-v1. A binding that does not implement these (or a remote |
| 41 | // service that does not advertise the capability) is read through the |
| 42 | // protocol-7 adapter in canonicalTranscriptBackend instead. |
| 43 | SessionHistoryWindowForTab(tabID: string, req: HistoryWindowRequest): Promise<HistoryWindowPage>; |
| 44 | SessionMessageFieldForTab(tabID: string, messageID: string, version: number, field: string, offset: number, length: number): Promise<MessageFieldPage>; |
| 45 | RemoteSessionHistoryWindowForTab(tabID: string, req: HistoryWindowRequest): Promise<HistoryWindowPage>; |
| 46 | RemoteSessionMessageFieldForTab(tabID: string, messageID: string, version: number, field: string, offset: number, length: number): Promise<MessageFieldPage>; |
| 47 | } |
| 48 | |
| 49 | interface MockSessionReaderHost { |
| 50 | MetaForTab?(tabID: string): Promise<{ running?: boolean; runtime?: { epoch?: string } }>; |
| 51 | HistoryForTab?(tabID: string): Promise<HistoryMessage[]>; |
| 52 | HistorySliceForTab(tabID: string, req: HistorySliceRequest): Promise<HistorySlice>; |
| 53 | HistoryContentForTab(tabID: string, ref: HistoryContentRef, chunkIndex: number): Promise<HistoryContentChunk>; |
| 54 | SessionHistoryPageForTab(tabID: string, cursor: string, limit: number): Promise<MessageHistoryPage>; |
| 55 | SessionOpenForTab(tabID: string): Promise<SessionOpenView>; |
| 56 | SessionHistoryContentForTab(tabID: string, ref: SessionContentRef, offset: number): Promise<SessionHistoryContentChunk>; |
| 57 | SessionHistoryWindowForTab(tabID: string, req: HistoryWindowRequest): Promise<HistoryWindowPage>; |
| 58 | } |
| 59 | |
| 60 | async function mockReaderSlice(host: MockSessionReaderHost, tab: string, request: HistorySliceRequest): Promise<HistorySlice> { |
| 61 | if (typeof host.HistorySliceForTab === "function") return host.HistorySliceForTab(tab, request); |
| 62 | const messages = await host.HistoryForTab?.(tab) ?? []; |
| 63 | let turn = 0; |
| 64 | const entries = messages.map((message, order) => { |
| 65 | if (message.role === "user") turn++; |
| 66 | return { entryId: message.messageId ?? `mock:${tab}:${order}`, order, turn, message, refs: [] }; |
| 67 | }).slice(-(request.entries ?? 32)); |
| 68 | return { entries, totalTurns: turn, startTurn: entries[0]?.turn ?? 0, endTurn: turn, nextCursor: "", |
| 69 | hasOlder: false, hasNewer: false, revision: 1, revisionKnown: true, digest: `mock:${tab}`, stale: false }; |
| 70 | } |
| 71 | |
| 72 | // Test/browser host implementation of the mandatory protocol. Production |
| 73 | // transports never route an unsupported Follow request through this adapter. |
| 74 | const mockFollowers = new Map<string, { tab: string; epoch?: string; revision: number; coverage: number; queue: Change[]; wake?: () => void }>(); |
| 75 | const mockPrompts = new Map<string, Map<string, WireEvent>>(); |
| 76 | const mockActive = new Map<string, { id: string; text: string; reasoning: string }>(); |
| 77 | const mockMetadata = new Map<string, { turnId?: string; running?: boolean; runtime?: { epoch?: string } }>(); |
| 78 | export function setMockTranscriptMetadata(tab: string, meta: { turnId?: string; running?: boolean; runtime?: { epoch?: string } }): void { |
| 79 | const prior = mockMetadata.get(tab)?.runtime?.epoch; |
| 80 | if (prior && meta.runtime?.epoch && prior !== meta.runtime.epoch) { mockActive.delete(tab); mockPrompts.delete(tab); } |
| 81 | mockMetadata.set(tab, { ...mockMetadata.get(tab), ...meta }); |
| 82 | } |
| 83 | let mockSubscription = 0; |
| 84 | let mockMessage = 0; |
| 85 | export function publishMockTranscriptEvent(event: WireEvent): void { |
| 86 | if (event.tabId) { |
| 87 | const tab = event.tabId; |
| 88 | const epoch = mockMetadata.get(tab)?.runtime?.epoch; |
| 89 | if (event.runtimeEpoch && epoch && event.runtimeEpoch !== epoch) return; |
| 90 | if (event.kind === "turn_started") setMockTranscriptMetadata(tab, { running: true }); |
| 91 | if (event.kind === "stream_attempt" && event.streamAttempt?.action === "begin" && event.messageId) { |
| 92 | mockActive.set(tab, { id: event.messageId, text: "", reasoning: "" }); |
| 93 | } |
| 94 | if (event.kind === "text" || event.kind === "reasoning") { |
| 95 | const active = mockActive.get(tab) ?? { id: event.messageId ?? `mock-active:${tab}:${++mockMessage}`, text: "", reasoning: "" }; |
| 96 | if (event.kind === "text") active.text += event.text ?? ""; else active.reasoning += event.text ?? ""; |
| 97 | mockActive.set(tab, active); |
| 98 | event = { ...event, messageId: active.id, attemptId: active.id }; |
| 99 | setMockTranscriptMetadata(tab, { running: true }); |
| 100 | } |
| 101 | if (event.kind === "turn_done") { mockActive.delete(tab); setMockTranscriptMetadata(tab, { running: false }); } |
| 102 | const prompts = mockPrompts.get(tab) ?? new Map<string, WireEvent>(); |
| 103 | const id = event.approval?.id ?? event.ask?.id; |
| 104 | if ((event.kind === "approval_request" || event.kind === "ask_request") && id) prompts.set(id, event); |
| 105 | if (event.kind === "prompt_answered" && event.itemId) prompts.delete(event.itemId); |
| 106 | if (event.kind === "turn_done") prompts.clear(); |
| 107 | mockPrompts.set(tab, prompts); |
| 108 | } |
| 109 | for (const follower of mockFollowers.values()) { |
| 110 | if (event.tabId && event.tabId !== follower.tab) continue; |
| 111 | if (event.runtimeEpoch && follower.epoch && event.runtimeEpoch !== follower.epoch) continue; |
| 112 | follower.queue.push({ revision: ++follower.revision, commitSeq: follower.coverage, durableSeq: follower.coverage, |
| 113 | index: 0, event: { ...event, seq: undefined } as unknown as Change["event"] }); |
| 114 | follower.wake?.(); |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | function messageIndex(entryId: string): number { |
| 119 | return Number(/:m(\d+):o\d+$/.exec(entryId)?.[1] ?? -1); |
| 120 | } |
| 121 | |
| 122 | function canonicalBody(message: HistoryMessage | undefined, messageId: string): Uint8Array { |
| 123 | if (!message) return new Uint8Array(); |
| 124 | const body = { |
| 125 | ...message, |
| 126 | id: message.messageId ?? messageId, |
| 127 | reasoning_content: message.reasoning, |
| 128 | tool_calls: message.toolCalls?.map(call => ({ |
| 129 | ...call, |
| 130 | resolved_name: call.resolvedName, |
| 131 | capability_id: call.capabilityId, |
| 132 | resolved_read_only: call.resolvedReadOnly, |
| 133 | })), |
| 134 | tool_call_id: message.toolCallId, |
| 135 | name: message.toolName, |
| 136 | tool_execution: message.execution, |
| 137 | presented_files: message.presentedFiles ? { files: message.presentedFiles } : undefined, |
| 138 | server_search: message.serverSearch, |
| 139 | }; |
| 140 | return new TextEncoder().encode(JSON.stringify(body)); |
| 141 | } |
| 142 | |
| 143 | function persistentMessages(slice: HistorySlice, history: HistoryMessage[]): PersistentMessage[] { |
| 144 | return slice.entries.map(entry => { |
| 145 | const index = messageIndex(entry.entryId); |
| 146 | const messageId = entry.message.messageId ?? entry.entryId; |
| 147 | const body = canonicalBody(history[index], messageId); |
| 148 | return { |
| 149 | messageId, position: entry.order, |
| 150 | version: 1, role: entry.message.role, preview: entry.message.content ?? "", |
| 151 | eventSequence: 0, visibleTurn: entry.turn, inline: JSON.parse(new TextDecoder().decode(canonicalBody(entry.message, messageId))), |
| 152 | contentRef: (entry.refs ?? []).length > 0 ? { digest: `mock-canonical:${index}`, bytes: body.length, mediaType: "application/json" } : undefined, |
| 153 | }; |
| 154 | }); |
| 155 | } |
| 156 | |
| 157 | function encodedChunk(bytes: Uint8Array): string { |
| 158 | let binary = ""; |
| 159 | for (let start = 0; start < bytes.length; start += 0x8000) binary += String.fromCharCode(...bytes.subarray(start, start + 0x8000)); |
| 160 | return btoa(binary); |
| 161 | } |
| 162 | |
| 163 | export function makeMockSessionReaderBindings(): SessionReaderBindings { |
| 164 | const search = (): SearchHistoryPage => ({ hits: [], snapshotSequence: 0, coverageSequence: 0, status: "preparing", hasMore: false }); |
| 165 | async function follow(this: MockSessionReaderHost, tab: string, request: FollowRequest): Promise<TranscriptFollowResponse> { |
| 166 | const response: TranscriptFollowResponse = { protocolVersion: 2, subscription: request.subscription ?? "", changes: [], resetRequired: false }; |
| 167 | if (request.close) { mockFollowers.get(response.subscription)?.wake?.(); mockFollowers.delete(response.subscription); return response; } |
| 168 | if (!request.subscription) { |
| 169 | const id = `mock-follow-${++mockSubscription}`; |
| 170 | const follower = { tab, epoch: undefined as string | undefined, revision: 1, coverage: 0, queue: [] as Change[] }; |
| 171 | mockFollowers.set(id, follower); |
| 172 | const meta = mockMetadata.get(tab); |
| 173 | const active = mockActive.get(tab) ? { ...mockActive.get(tab)! } : undefined; |
| 174 | const pendingEvents = [...(mockPrompts.get(tab)?.values() ?? [])]; |
| 175 | follower.epoch = meta?.runtime?.epoch; |
| 176 | const page = await (this.SessionHistoryWindowForTab ?? bindings.SessionHistoryWindowForTab).call(this, tab, { anchor: "newest", limit: 32 }); |
| 177 | follower.coverage = page.snapshotSequence; |
| 178 | follower.queue = follower.queue.map(change => ({ ...change, commitSeq: page.snapshotSequence, durableSeq: page.snapshotSequence })); |
| 179 | return { ...response, subscription: id, history: page, snapshot: { |
| 180 | protocolVersion: 2, snapshotId: id, identity: { sessionId: tab, runtimeEpoch: meta?.runtime?.epoch ?? "mock-runtime", headId: "", rewriteEpoch: 0 }, |
| 181 | projectionRevision: 1, coveredThroughSeq: page.snapshotSequence, durableSeq: page.snapshotSequence, |
| 182 | records: [], activeRecords: active ? [{ id: `m:${active.id}`, order: page.messages.length, message: { role: "assistant", messageId: active.id, content: active.text, reasoning: active.reasoning }, refs: [] }] : [], |
| 183 | activeAttempts: active ? [{ id: active.id, messageId: active.id, turnId: "mock-turn", nextIndex: 0 }] : [], totalRecords: page.messages.length + (active ? 1 : 0), totalTurns: page.totalTurns, |
| 184 | before: 0, hasOlder: page.hasOlder, stale: false, |
| 185 | runtime: { turnId: meta?.turnId, status: pendingEvents.length ? "waiting_user" : meta?.running ? "in_progress" : "completed", pendingEvents: pendingEvents as unknown as NonNullable<TranscriptFollowResponse["snapshot"]>["runtime"]["pendingEvents"], samplingCount: 0, toolCount: 0 }, |
| 186 | } }; |
| 187 | } |
| 188 | const follower = mockFollowers.get(request.subscription); |
| 189 | if (!follower) return { ...response, resetRequired: true }; |
| 190 | follower.queue = follower.queue.filter(change => change.revision > (request.afterRevision ?? 0)); |
| 191 | if (!follower.queue.length) await new Promise<void>(resolve => { follower.wake = resolve; }); |
| 192 | follower.wake = undefined; |
| 193 | return { ...response, changes: [...follower.queue] }; |
| 194 | } |
| 195 | const bindings: SessionReaderBindings = { |
| 196 | async BeginSessionHistoryReadForTab(tabID) { return { id: tabID, storageBackend: "canonical", sessionGeneration: 0, capabilities: [] }; }, |
| 197 | async ReleaseSessionHistoryRead() {}, |
| 198 | async ReadSessionHistoryWindow(this: MockSessionReaderHost, id, req) { return this.SessionHistoryWindowForTab(id, req); }, |
| 199 | async ReadSessionHistorySlice(this: MockSessionReaderHost, id, req) { return { status: "ready", page: await this.HistorySliceForTab(id, req) }; }, |
| 200 | async ReadSessionHistoryOutline() { return { entries: [], status: "unsupported", totalTurns: 0, nextTurn: 0, done: true, snapshotSequence: 0, coverageSequence: 0, generation: "" }; }, |
| 201 | async LocateSessionHistoryMessage(_id, messageID) { return { status: "not_found", messageId: messageID, snapshotSequence: 0, coverageSequence: 0 }; }, |
| 202 | async SearchSessionHistoryRead() { return search(); }, |
| 203 | TranscriptFollowForTab: follow, |
| 204 | RemoteTranscriptFollowForTab: follow, |
| 205 | async SessionHistoryPageForTab(this: MockSessionReaderHost, tabID, cursor, limit) { |
| 206 | const slice = await mockReaderSlice(this, tabID, { cursor, entries: limit, turns: limit }); |
| 207 | const history = slice.entries.some(entry => entry.refs?.length) ? await this.HistoryForTab?.(tabID) ?? [] : []; |
| 208 | return { messages: persistentMessages(slice, history), snapshotSequence: slice.revision, coverageSequence: slice.revision, status: "ready", totalTurns: slice.totalTurns, generation: slice.digest ?? "", nextCursor: slice.nextCursor, hasMore: slice.hasOlder }; |
| 209 | }, |
| 210 | async SessionOpenForTab(this: MockSessionReaderHost, tabID) { |
| 211 | const slice = await mockReaderSlice(this, tabID, { cursor: "", entries: 100, turns: 100 }); |
| 212 | const history = slice.entries.some(entry => entry.refs?.length) ? await this.HistoryForTab?.(tabID) ?? [] : []; |
| 213 | const entries = persistentMessages(slice, history); |
| 214 | return { session: { hostId: "local", sessionId: tabID }, storageGeneration: slice.digest, snapshotSequence: slice.revision, acceptedSequence: slice.revision, durableSequence: slice.revision, recent: { version: 1, sessionId: tabID, storageGeneration: slice.digest ?? "", durableSequence: slice.revision, totalTurns: slice.totalTurns, entries }, recovery: "ready", history: "ready", search: "preparing", canExecute: true }; |
| 215 | }, |
| 216 | async SessionHistoryContentForTab(this: MockSessionReaderHost, tabID, ref, offset) { |
| 217 | const index = Number(ref.digest.replace("mock-canonical:", "")); |
| 218 | const history = await this.HistoryForTab?.(tabID) ?? []; |
| 219 | const message = history[index]; |
| 220 | const entryId = `smock-${tabID}:r0:m${index}:o0`; |
| 221 | await this.HistoryContentForTab(tabID, { entryId, field: "content", size: message?.content?.length ?? 0, chunks: 1, revision: 0, digest: "mock" }, 0); |
| 222 | const body = canonicalBody(message, entryId); |
| 223 | const nextOffset = Math.min(body.length, offset + (1 << 20)); |
| 224 | return { data: encodedChunk(body.subarray(offset, nextOffset)), nextOffset, done: nextOffset >= body.length }; |
| 225 | }, |
| 226 | async RemoteSessionOpenForTab(this: MockSessionReaderHost, tabID) { return this.SessionOpenForTab(tabID); }, |
| 227 | async RemoteSessionHistoryPageForTab(this: MockSessionReaderHost, tabID, cursor, limit) { return this.SessionHistoryPageForTab(tabID, cursor, limit); }, |
| 228 | async RemoteSessionHistoryContentForTab(this: MockSessionReaderHost, tabID, ref, offset) { return this.SessionHistoryContentForTab(tabID, ref, offset); }, |
| 229 | async SearchSessionHistoryForTab() { return search(); }, |
| 230 | async RemoteSearchSessionHistoryForTab() { return search(); }, |
| 231 | async LocateSessionMessageForTab(_tabID, messageID) { return { status: "not_found", messageId: messageID, snapshotSequence: 0, coverageSequence: 0 }; }, |
| 232 | async RemoteLocateSessionMessageForTab(_tabID, messageID) { return { status: "not_found", messageId: messageID, snapshotSequence: 0, coverageSequence: 0 }; }, |
| 233 | // The in-memory mock has one direction of history: a newest page and its |
| 234 | // older cursors. It answers a window request without inventing a newer |
| 235 | // cursor, which is exactly how a protocol-7 service behaves. |
| 236 | async SessionHistoryWindowForTab(this: MockSessionReaderHost, tabID, req) { |
| 237 | const slice = await mockReaderSlice(this, tabID, { |
| 238 | cursor: req.anchor === "cursor" ? req.cursor ?? "" : "", |
| 239 | entries: req.limit, |
| 240 | turns: req.limit, |
| 241 | }); |
| 242 | const history = slice.entries.some(entry => entry.refs?.length) ? await this.HistoryForTab?.(tabID) ?? [] : []; |
| 243 | return { |
| 244 | messages: persistentMessages(slice, history), |
| 245 | status: slice.stale ? "stale_cursor" : "ready", |
| 246 | snapshotSequence: slice.revision, |
| 247 | coverageSequence: slice.revision, |
| 248 | generation: slice.digest, |
| 249 | totalTurns: slice.totalTurns, |
| 250 | hasOlder: slice.hasOlder, |
| 251 | // The in-memory mock has one direction of history, exactly like a |
| 252 | // protocol-7 service: it never invents a newer cursor. |
| 253 | hasNewer: false, |
| 254 | olderCursor: slice.nextCursor, |
| 255 | newerCursor: "", |
| 256 | }; |
| 257 | }, |
| 258 | async SessionMessageFieldForTab() { return { status: "not_found", messageId: "", version: 0, field: "", totalBytes: 0, offset: 0, encoding: "utf-8" }; }, |
| 259 | async RemoteSessionHistoryWindowForTab(this: MockSessionReaderHost, tabID, req) { return this.SessionHistoryWindowForTab(tabID, req); }, |
| 260 | async RemoteSessionMessageFieldForTab() { return { status: "not_found", messageId: "", version: 0, field: "", totalBytes: 0, offset: 0, encoding: "utf-8" }; }, |
| 261 | }; |
| 262 | return bindings; |
| 263 | } |
| 264 |