返回 DeepSeek-Reasonix
transcriptSessionFollowerRuntime.ts
根目录 / desktop / frontend / src / lib / transcriptSessionFollowerRuntime.ts
1 import { noteSessionObservation } from "./sessionObservationDiagnostics";
2 import { app } from "./bridge";
3 import { bindNativeTranscriptHistory, nativeSnapshotWindow } from "./nativeTranscriptHistory";
4 import { entriesFor, registerTranscriptContentRecovery } from "./canonicalTranscriptBackend";
5 import { TranscriptFollowClient } from "./transcriptFollowClient";
6 import { getTranscriptStore } from "./transcriptStore";
7 import type { Action, State } from "./useController";
8 import type { HistoryEntry, HistoryMessage, WireEvent } from "./types";
9 import type { TranscriptSnapshot } from "./transcriptProtocol";
10 import type { Message, TranscriptFollowResponse } from "../generated/desktopContract.generated";
11 import { canonicalUserConfirmations } from "./localSubmissionState";
12 import { snapshotRecords } from "./transcriptSnapshotState";
13 import { signalOutline } from "./transcriptOutlineSignals";
14 import { historyMessageIdentity } from "./historyItemIds";
15
16 /** Canonical positions and projection indexes count different record sets.
17 * Join the two ordered windows by shared message identity before assigning
18 * positions to projection-only rows. */
19 function mergeFollowEntries(canonical: HistoryEntry[], snapshot: HistoryEntry[], merged: Map<string, HistoryEntry>): HistoryEntry[] {
20 const canonicalIds = new Set(canonical.map(entry => entry.entryId));
21 const ordered = [...canonical].sort((a, b) => a.order - b.order).map(entry => merged.get(entry.entryId)!);
22 const present = new Set(canonicalIds);
23 for (let snapshotIndex = 0; snapshotIndex < snapshot.length; snapshotIndex++) {
24 // Multiple projection aliases can normalize to one durable identity. The
25 // merge map owns its body/ref precedence; this pass only places that owner.
26 const entry = merged.get(snapshot[snapshotIndex].entryId)!;
27 if (present.has(entry.entryId)) continue;
28 let at = -1;
29 for (let index = snapshotIndex - 1; index >= 0; index--) {
30 const anchor = snapshot[index].entryId;
31 if (!present.has(anchor)) continue;
32 at = ordered.findIndex(candidate => candidate.entryId === anchor) + 1;
33 break;
34 }
35 if (at < 0) {
36 for (let index = snapshotIndex + 1; index < snapshot.length; index++) {
37 const anchor = snapshot[index].entryId;
38 if (!present.has(anchor)) continue;
39 at = ordered.findIndex(candidate => candidate.entryId === anchor);
40 break;
41 }
42 }
43 if (at < 0) {
44 at = ordered.findIndex(candidate => candidate.turn > entry.turn);
45 if (at < 0) at = ordered.length;
46 }
47 ordered.splice(at, 0, entry);
48 present.add(entry.entryId);
49 }
50 for (let start = 0; start < ordered.length;) {
51 if (canonicalIds.has(ordered[start].entryId)) { start++; continue; }
52 let end = start;
53 while (end < ordered.length && !canonicalIds.has(ordered[end].entryId)) end++;
54 const left = start > 0 ? ordered[start - 1].order : undefined;
55 const right = end < ordered.length ? ordered[end].order : undefined;
56 for (let index = start; index < end; index++) {
57 const fraction = (index - start + 1) / (end - start + 1);
58 let order = ordered[index].order;
59 if (left !== undefined && right !== undefined) order = left + (right - left) * fraction;
60 else if (left !== undefined) order = left + fraction;
61 else if (right !== undefined) order = right - 1 + fraction;
62 ordered[index] = { ...ordered[index], order };
63 }
64 start = end;
65 }
66 return ordered;
67 }
68
69 export class TranscriptSessionFollowerRuntime {
70 private readonly client: TranscriptFollowClient;
71 private readonly orders = new Map<string, number>();
72 private nextOrder = 0;
73 private turn = 0;
74 private generation = 0;
75 private releaseContentRecovery?: () => void;
76 private releaseNativeHistory?: () => void;
77 private recoveringContent = false;
78 private coverage = 0;
79 private recoveryScheduled = false;
80 private readonly confirmationReads = new Map<string, { coverage: number; pending: boolean }>();
81 private readonly submissionCoverage = new Map<string, number>();
82 metrics = { entries: 0, inlineBytes: 0 };
83
84 constructor(private readonly tabId: string, private readonly path: string, private readonly remote: boolean,
85 private readonly dispatch: (action: Action) => void,
86 private readonly state: () => State | undefined = () => getTranscriptStore().states.get(tabId)) {
87 this.client = new TranscriptFollowClient(request => {
88 const read = remote ? app.RemoteTranscriptFollowForTab : app.TranscriptFollowForTab;
89 if (!read) return Promise.reject(new Error("Transcript v2 is required. Upgrade Desktop and Serve together."));
90 return read(tabId, request);
91 }, { transport: remote ? "remote" : "local" });
92 }
93
94 async start(): Promise<void> {
95 this.generation++;
96 this.releaseNativeHistory?.();
97 this.releaseNativeHistory = undefined;
98 noteSessionObservation(this.path, { action: "subscribe", tabId: this.tabId, generation: this.generation, sequence: this.coverage });
99 this.recoveryScheduled = false;
100 this.coverage = 0;
101 this.confirmationReads.clear();
102 this.submissionCoverage.clear();
103 this.observeSubmissions();
104 this.releaseContentRecovery?.();
105 this.releaseContentRecovery = registerTranscriptContentRecovery(this.tabId, () => {
106 if (this.recoveringContent) return;
107 this.recoveringContent = true;
108 void this.start().catch(() => undefined).finally(() => { this.recoveringContent = false; });
109 });
110 await this.client.start({
111 install: response => this.install(response),
112 changes: changes => {
113 if (changes.some(change => change.records?.some(record => record.role === "user" || record.role === "assistant") || change.durableSeq > this.coverage)) signalOutline(this.tabId);
114 this.observeSubmissions();
115 for (const change of changes) {
116 if (change.records?.length) {
117 const entries = change.records.map(message => this.entry(message));
118 const projection = getTranscriptStore().upsertEntries(this.tabId, this.path, entries, change.commitSeq);
119 if (!projection) throw new Error("Transcript window is unavailable; synchronize again");
120 this.dispatch({ type: "transcript_records", projection,
121 confirmedUsers: canonicalUserConfirmations(entries.map(entry => ({ ...entry.message, kind: entry.message.role }))) });
122 this.coverage = Math.max(this.coverage, change.commitSeq);
123 }
124 if (change.event) this.dispatch({ type: "event", e: { ...change.event, tabId: this.tabId } as unknown as WireEvent, remote: this.remote });
125 if (change.runtime) this.dispatch({ type: "transcript_runtime", runtime: change.runtime });
126 }
127 this.scheduleConfirmationRecovery();
128 },
129 connection: (status, error) => { noteSessionObservation(this.path, { action: "connection", tabId: this.tabId, generation: this.generation, sequence: this.coverage, status }); this.dispatch({ type: "transcript_connection", status, error }); },
130 });
131 }
132
133 stop(closeSubscription = true, reason?: "service_stopping"): void { noteSessionObservation(this.path, { action: "unsubscribe", tabId: this.tabId, generation: this.generation, sequence: this.coverage }); this.generation++; this.confirmationReads.clear(); this.submissionCoverage.clear(); this.releaseContentRecovery?.(); this.releaseContentRecovery = undefined; this.releaseNativeHistory?.(); this.releaseNativeHistory = undefined; this.client.stop(closeSubscription, reason); }
134
135 private observeSubmissions(): void {
136 const pending = new Set(this.state()?.localSubmissionOrder ?? []);
137 for (const id of this.submissionCoverage.keys()) if (!pending.has(id)) this.submissionCoverage.delete(id);
138 for (const id of pending) if (!this.submissionCoverage.has(id)) this.submissionCoverage.set(id, this.coverage);
139 }
140
141 private scheduleConfirmationRecovery(): void {
142 if (this.recoveryScheduled) return;
143 this.recoveryScheduled = true;
144 const generation = this.generation;
145 queueMicrotask(() => {
146 this.recoveryScheduled = false;
147 if (generation !== this.generation) return;
148 const state = this.state();
149 if (!state) return;
150 this.observeSubmissions();
151 const unresolved = new Set(state.localSubmissionOrder);
152 for (const key of this.confirmationReads.keys()) if (!unresolved.has(key)) this.confirmationReads.delete(key);
153 for (const local of Object.values(state.localSubmissions)) {
154 // Identity events alone do not justify a history read. Only recover
155 // after a committed cut could have hidden an earlier formal record.
156 if (!local.messageId || this.coverage <= (this.submissionCoverage.get(local.submissionId) ?? this.coverage)) continue;
157 const previous = this.confirmationReads.get(local.submissionId);
158 if (previous?.pending || previous?.coverage === this.coverage) continue;
159 const read = this.remote ? app.RemoteSessionHistoryWindowForTab : app.SessionHistoryWindowForTab;
160 if (!read) continue;
161 const ticket = { coverage: this.coverage, pending: true };
162 this.confirmationReads.set(local.submissionId, ticket);
163 const current = () => generation === this.generation && this.state()?.sessionGen === state.sessionGen
164 && this.state()?.localSubmissions[local.submissionId]?.messageId === local.messageId;
165 void read(this.tabId, { anchor: "message", messageId: local.messageId, limit: 1 }).then(page => {
166 if (!current() || page.status !== "ready") return;
167 const message = page.messages.find(message => message.messageId === local.messageId && message.role === "user");
168 if (message) this.dispatch({ type: "submission_verified", submissionId: local.submissionId, messageId: local.messageId! });
169 }).catch(() => { /* An inconclusive read retains the echo until the next committed cut. */ }).finally(() => {
170 if (this.confirmationReads.get(local.submissionId) !== ticket) return;
171 ticket.pending = false;
172 if (!current()) this.confirmationReads.delete(local.submissionId);
173 else if (ticket.coverage !== this.coverage) this.scheduleConfirmationRecovery();
174 });
175 }
176 });
177 }
178
179 private entry(message: Message | HistoryMessage, outerRecordId?: string, snapshotOrder?: number): HistoryEntry {
180 // Canonical history addresses every persisted message as m:<messageId>,
181 // including tool results. A snapshot may instead carry its projection
182 // identity (for example tool:<toolCallId>); that explicit identity is
183 // valid metadata, while messageId remains the merge key used by history.
184 const messageId = historyMessageIdentity(message, outerRecordId);
185 const derived = messageId ? `m:${messageId}`
186 : message.role === "tool" && message.toolCallId ? `tool:${message.toolCallId}` : undefined;
187 if (outerRecordId && message.recordId && outerRecordId !== message.recordId) {
188 throw new Error("transcript snapshot record identity mismatch");
189 }
190 const entryId = derived ?? message.recordId ?? outerRecordId;
191 if (!entryId) throw new Error("invalid transcript snapshot record identity");
192 const normalized = message.recordId === entryId && message.messageId === messageId ? message : { ...message, messageId, recordId: entryId };
193 let order = this.orders.get(entryId);
194 if (order === undefined) {
195 order = snapshotOrder ?? this.nextOrder;
196 this.nextOrder = Math.max(this.nextOrder, order + 1);
197 this.orders.set(entryId, order);
198 // A snapshot's total already includes users outside its history page.
199 // Only newly followed user records advance the turn counter.
200 if (message.role === "user" && snapshotOrder === undefined) this.turn++;
201 // Only the resident tail needs an order index. Older pages carry their
202 // canonical positions and are owned by the bounded transcript store.
203 while (this.orders.size > 192) this.orders.delete(this.orders.keys().next().value!);
204 }
205 return { entryId, order, turn: message.historyTurn || this.turn, message: normalized as unknown as HistoryMessage, refs: [] };
206 }
207
208 private async install(response: TranscriptFollowResponse): Promise<void> {
209 const generation = this.generation;
210 const snapshot = response.snapshot!;
211 // Validate the untrusted bridge cut before resolving refs, merging rows or
212 // changing the resident store. The outer record id is authoritative for
213 // older peers that did not repeat it inside the message.
214 const records = snapshotRecords(snapshot as unknown as TranscriptSnapshot);
215 // A suffix must never be applied to a truncated prefix. Resolve active
216 // snapshot references before publishing any part of this recovery cut.
217 const activeIds = new Set(snapshot.activeAttempts.map(attempt => attempt.messageId));
218 for (const record of records) {
219 if (!record.message.messageId || !activeIds.has(record.message.messageId)) continue;
220 for (const ref of record.refs) {
221 const read = this.remote ? app.RemoteTranscriptContentForTab : app.TranscriptContentForTab;
222 if (!read) throw new Error("Transcript v2 content is unavailable");
223 let offset = 0, text = "";
224 while (true) {
225 const chunk = await read(this.tabId, { ...ref, offset });
226 if (generation !== this.generation) return;
227 if (chunk.stale) throw new Error("Active transcript snapshot expired; synchronize again");
228 text += chunk.data;
229 if (chunk.done) break;
230 if (chunk.nextOffset <= offset) throw new Error("Active transcript content did not advance");
231 offset = chunk.nextOffset;
232 }
233 let target = record.message as unknown as Record<string, unknown>;
234 for (const key of ref.path.slice(0, -1)) target = target[key] as Record<string, unknown>;
235 target[ref.path[ref.path.length - 1]] = text;
236 }
237 record.refs = [];
238 }
239 if (generation !== this.generation) return;
240 this.observeSubmissions();
241 const page = response.history;
242 this.releaseNativeHistory?.();
243 this.releaseNativeHistory = response.storageBackend === "legacy"
244 ? bindNativeTranscriptHistory(this.tabId, snapshot as unknown as TranscriptSnapshot, this.remote) : undefined;
245 const nativeWindow = response.storageBackend === "legacy" ? nativeSnapshotWindow(snapshot as unknown as TranscriptSnapshot) : undefined;
246 this.turn = page?.totalTurns ?? snapshot.totalTurns;
247 this.orders.clear();
248 const entries = nativeWindow?.entries ?? entriesFor(page?.messages ?? [], page?.snapshotSequence ?? snapshot.coveredThroughSeq);
249 for (const entry of entries) this.orders.set(entry.entryId, entry.order);
250 this.nextOrder = Math.max(0, ...entries.map(entry => entry.order + 1));
251 const merged = new Map(entries.map(entry => [entry.entryId, entry]));
252 const snapshotEntries: HistoryEntry[] = [];
253 for (const record of records) {
254 const entry = this.entry(record.message, record.id, record.order);
255 snapshotEntries.push(entry);
256 // Durable canonical refs remain loadable after a view snapshot expires.
257 const canonical = merged.get(entry.entryId);
258 if (canonical && !snapshot.activeAttempts.some(attempt => attempt.messageId === record.message.messageId)) continue;
259 entry.refs = record.refs.map(ref => ({ entryId: entry.entryId, field: ref.path[0], size: ref.bytes, chunks: 1,
260 revision: snapshot.coveredThroughSeq, revKnown: true, digest: snapshot.snapshotId, transcriptRef: ref }));
261 merged.set(entry.entryId, entry);
262 }
263 // Native pages and activeRecords share the frozen snapshot's order space.
264 // Keep those positions so pages fetched into an active-prefix gap can sort
265 // between its original neighbors.
266 const all = nativeWindow ? [...merged.values()].sort((a, b) => a.order - b.order)
267 : mergeFollowEntries(entries, snapshotEntries, merged);
268 this.orders.clear();
269 for (const entry of all.slice(-192)) this.orders.set(entry.entryId, entry.order);
270 this.nextOrder = Math.max(this.nextOrder, ...all.map(entry => entry.order + 1));
271 this.metrics = { entries: all.length, inlineBytes: all.reduce((bytes, entry) => bytes + entry.message.content.length + (entry.message.reasoning?.length ?? 0), 0) };
272 const prepared = getTranscriptStore().prepareInstallSlice(this.tabId, this.path, {
273 entries: all, nextCursor: page?.olderCursor ?? nativeWindow?.olderCursor ?? "", newerCursor: page?.newerCursor ?? nativeWindow?.newerCursor ?? "",
274 hasOlder: Boolean(page?.hasOlder ?? nativeWindow?.hasOlder), hasNewer: Boolean(page?.hasNewer ?? nativeWindow?.hasNewer), totalTurns: this.turn,
275 startTurn: Math.min(this.turn, ...all.map(entry => entry.turn)), endTurn: this.turn,
276 revision: page?.snapshotSequence ?? snapshot.coveredThroughSeq, revisionKnown: true, digest: page?.generation ?? nativeWindow?.digest ?? "", stale: false,
277 });
278 const combined: TranscriptSnapshot = {
279 ...snapshot as unknown as TranscriptSnapshot,
280 records: all.map((entry, order) => ({ id: entry.entryId, order, message: entry.message, refs: [] })),
281 activeRecords: [], totalRecords: all.length, totalTurns: this.turn,
282 };
283 this.dispatch({ type: "transcript_v2_snapshot", snapshot: combined, projection: prepared.projection, remote: this.remote });
284 prepared.commit();
285 signalOutline(this.tabId);
286 this.dispatch({ type: "transcript_runtime", runtime: snapshot.runtime });
287 this.coverage = snapshot.coveredThroughSeq;
288 noteSessionObservation(this.path, { action: "snapshot_installed", tabId: this.tabId, generation: this.generation, sequence: this.coverage, status: snapshot.runtime.status });
289 this.scheduleConfirmationRecovery();
290 }
291 }
292
292 lines TYPESCRIPT