返回 DeepSeek-Reasonix
sessionReaderBridge.ts
根目录 / desktop / frontend / src / lib / sessionReaderBridge.ts
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
264 lines TYPESCRIPT