返回 DeepSeek-Reasonix
transcriptSnapshotState.ts
根目录 / desktop / frontend / src / lib / transcriptSnapshotState.ts
1 import type { HistoryMessage, WireEvent } from "./types";
2 import type { Item, State } from "./useController";
3 import type { TranscriptRecord, TranscriptSnapshot } from "./transcriptProtocol";
4 import { canonicalUserConfirmations, settleLocalSubmissions, settleRebasedSubmissions } from "./localSubmissionState";
5 import { mergeSessionOperationItem, reconcileSessionOperationItems } from "./sessionMaintenanceOperation";
6 import { resetTurnTiming } from "./submissionReducer";
7
8 export function snapshotRecords(snapshot: TranscriptSnapshot): TranscriptRecord[] {
9 if (!Number.isSafeInteger(snapshot.totalRecords) || snapshot.totalRecords < 0) {
10 throw new Error("invalid transcript snapshot record count");
11 }
12 const records: TranscriptRecord[] = [];
13 const indexes = new Map<string, number>();
14 const normalize = (record: TranscriptRecord): TranscriptRecord => {
15 if (!Number.isSafeInteger(record.order) || record.order < 0 || record.order >= snapshot.totalRecords) {
16 throw new Error("invalid transcript snapshot record identity");
17 }
18 const derived = record.message.messageId ? `m:${record.message.messageId}`
19 : record.message.role === "tool" && record.message.toolCallId ? `tool:${record.message.toolCallId}` : undefined;
20 const referencedID = record.refs.find(ref => ref.recordId)?.recordId;
21 let id = record.id || record.message.recordId || derived || referencedID;
22 if (!id) {
23 if (record.refs.length > 0) throw new Error("transcript snapshot content identity missing");
24 id = `view:legacy:${snapshot.snapshotId || `${snapshot.identity.runtimeEpoch}:${snapshot.projectionRevision}`}:${record.order}`;
25 }
26 if (record.message.recordId && record.message.recordId !== id) {
27 throw new Error("transcript snapshot record identity mismatch");
28 }
29 if (record.refs.some(ref => ref.recordId !== id)) throw new Error("transcript snapshot content identity mismatch");
30 return { ...record, id, message: { ...record.message, recordId: id } };
31 };
32 for (const raw of snapshot.records ?? []) {
33 const record = normalize(raw);
34 if (indexes.has(record.id)) throw new Error("duplicate transcript snapshot record identity");
35 indexes.set(record.id, records.length);
36 records.push(record);
37 }
38 const activeIDs = new Set<string>();
39 for (const raw of snapshot.activeRecords ?? []) {
40 const record = normalize(raw);
41 if (activeIDs.has(record.id)) throw new Error("duplicate active transcript snapshot record identity");
42 activeIDs.add(record.id);
43 // A mutable owner may also fall inside the requested page. The backend
44 // omits that duplicate, but tolerate older/remote implementations that
45 // return it in both arrays and keep the active copy authoritative.
46 const index = indexes.get(record.id);
47 if (index !== undefined) records[index] = record;
48 else { indexes.set(record.id, records.length); records.push(record); }
49 }
50 return records.sort((a, b) => a.order - b.order);
51 }
52
53 type Convert = (messages: HistoryMessage[], prefix: string) => { items: Item[]; seq: number };
54 type ApplyEvent = (state: State, event: WireEvent) => State;
55
56 // Durable rows only match durable identities. Optimistic submission identity is
57 // owned by localSubmissionState and must never replace a canonical item id.
58 export function matchingSnapshotItem(items: Item[], item: Item): Item | undefined {
59 if (item.kind === "user") {
60 const users = items.filter((candidate): candidate is Extract<Item, { kind: "user" }> => candidate.kind === "user");
61 const message = item.messageId && users.find(candidate => candidate.messageId === item.messageId);
62 if (message) return message;
63 }
64 return items.find(candidate => candidate.id === item.id);
65 }
66
67 function recordItemOrder(records: TranscriptRecord[], convert: Convert): Record<string, number> {
68 const order: Record<string, number> = {};
69 for (const record of records) {
70 const items = convert([{ ...record.message, recordId: record.id }], "snapshot:").items;
71 items.forEach((item, index) => { order[item.id] = Math.min(order[item.id] ?? Infinity, record.order + index / (items.length + 1)); });
72 }
73 return order;
74 }
75
76 export function transcriptPageState(state: State, page: TranscriptSnapshot, convert: Convert): State {
77 const records = snapshotRecords({ ...page, activeRecords: [] });
78 const converted = convert(records.map((record) => ({ ...record.message, recordId: record.id })), "snapshot:");
79 const existing = new Map(state.items.map((item) => [item.id, item]));
80 const order = { ...state.transcriptItemOrder };
81 for (const [id, position] of Object.entries(recordItemOrder(records, convert))) order[id] = Math.min(order[id] ?? Infinity, position);
82 const prefix: Item[] = [];
83 for (const item of converted.items) {
84 const prior = existing.get(item.id) ?? matchingSnapshotItem(state.items, item);
85 if (!prior) { prefix.push(item); continue; }
86 if (prior.kind === "tool" && item.kind === "tool") {
87 prefix.push({ ...prior, args: prior.args || item.args, messageId: prior.messageId || item.messageId,
88 name: prior.name === "tool" ? item.name : prior.name, subject: prior.subject ?? item.subject,
89 summary: prior.summary ?? item.summary, fileDiff: prior.fileDiff ?? item.fileDiff });
90 } else if (prior.kind === "compaction" && item.kind === "compaction") {
91 prefix.push(mergeSessionOperationItem(prior, item));
92 } else prefix.push(prior.id === item.id ? prior : { ...prior, ...item, id: item.id } as Item);
93 }
94 const prefixIDs = new Set(prefix.map((item) => item.id));
95 const added = prefix.filter((item) => !existing.has(item.id)).length;
96 const users = records.map((record) => record.message.historyTurn).filter((turn): turn is number => typeof turn === "number" && turn > 0);
97 const items = [...prefix, ...state.items.filter((item) => !prefixIDs.has(item.id))];
98 items.sort((a, b) => (order[a.id] ?? Infinity) - (order[b.id] ?? Infinity));
99 return settleLocalSubmissions({ ...state, items, transcriptItemOrder: order,
100 seq: Math.max(state.seq, converted.seq), historyPrefixCount: state.historyPrefixCount + added,
101 historyStartTurn: Math.min(state.historyStartTurn, ...users.map((turn) => turn - 1)),
102 historyHasOlder: page.hasOlder, historyOlderLoading: false, historyOlderError: undefined,
103 historyHasNewer: false, historyNewerLoading: false, historyNewerError: undefined,
104 historyMutation: { seq: state.historyMutation.seq + 1, kind: "prepend" } }, items, canonicalUserConfirmations(converted.items));
105 }
106
107 /** One reducer transaction installs rows, runtime and the active attempt.
108 * The event projector advances coverage only after this function commits. */
109 export function transcriptSnapshotState(state: State, snapshot: TranscriptSnapshot, convert: Convert, applyEvent: ApplyEvent, clock: number, projectedItems?: Item[]): State {
110 const sessionId = snapshot.identity.sessionId;
111 const sameTurn = state.transcriptSessionId === sessionId && Boolean(snapshot.runtime.turnId)
112 && state.activeTurnId === snapshot.runtime.turnId && state.runtimeStatusEpoch === snapshot.identity.runtimeEpoch;
113 if (state.transcriptSessionId && state.transcriptSessionId !== sessionId) {
114 state = { ...state, localSubmissions: {}, localSubmissionOrder: [], visibleSubmissionHandoffs: {},
115 guidanceConsumed: undefined,
116 pendingSubmissionId: undefined, pendingUser: undefined, sessionGen: state.sessionGen + 1 };
117 }
118 const records = snapshotRecords(snapshot);
119 const messages = records.map((record) => ({ ...record.message, recordId: record.id }));
120 const convertedRaw = projectedItems ? { items: projectedItems, seq: state.seq } : convert(messages, "snapshot:");
121 const converted = {
122 ...convertedRaw,
123 items: reconcileSessionOperationItems(convertedRaw.items, state.items),
124 };
125 state = settleLocalSubmissions(state, converted.items, canonicalUserConfirmations(messages.map(message => ({
126 kind: message.role, messageId: message.messageId, submissionId: message.submissionId, turnId: message.turnId,
127 }))));
128 const active = snapshot.runtime.status === "queued" || snapshot.runtime.status === "in_progress" ||
129 snapshot.runtime.status === "waiting_user" || snapshot.runtime.status === "cancelling";
130 // Older serves and history-rebased projections omit message ids, so the
131 // id-keyed settlement above cannot retire an echo the server already owns.
132 state = settleRebasedSubmissions(state, state.items,
133 converted.items.filter((item): item is Extract<Item, { kind: "user" }> => item.kind === "user"), !active);
134 const order = projectedItems ? Object.fromEntries(projectedItems.map((item, index) => [item.id, index])) : recordItemOrder(records, convert);
135 const users = state.items.filter((item): item is Extract<Item, { kind: "user" }> => item.kind === "user");
136 const items = converted.items.map((item) => {
137 if (item.kind === "tool" && item.resultMissing && snapshot.runtime.status &&
138 ["queued", "in_progress", "waiting_user", "cancelling"].includes(snapshot.runtime.status) &&
139 messages.some(message => message.turnId === snapshot.runtime.turnId && message.toolCalls?.some(call => call.id === item.id))) {
140 return { ...item, status: "running" as const };
141 }
142 if (item.kind !== "user") return item;
143 const mounted = matchingSnapshotItem(users, item);
144 if (!mounted) return item;
145 const next = { ...mounted, ...item, id: item.id };
146 return Object.entries(next).every(([key, value]) => (mounted as unknown as Record<string, unknown>)[key] === value) ? mounted : next;
147 });
148 const hasLocalSubmission = Boolean(state.pendingSubmissionId && state.localSubmissions[state.pendingSubmissionId]
149 && state.localSubmissions[state.pendingSubmissionId].status !== "failed");
150 let next: State = {
151 ...state,
152 ...(!sameTurn ? resetTurnTiming(snapshot.runtime.startedAt ?? 0) : {}),
153 // Snapshot content predates this observer. Its provider-output intervals
154 // are unavailable, so start a fresh matched numerator/denominator window.
155 turnModelActiveAt: undefined,
156 turnModelActiveMs: 0,
157 pendingRequestModelMs: undefined,
158 lastRequestTps: null,
159 turnRateSample: active ? { outputQuarters: 0, requestStartQuarters: 0, requestStartModelMs: 0 } : undefined,
160 transcriptSessionId: sessionId,
161 transcriptProtocol: 1,
162 transcriptItemOrder: order,
163 discardTurn: false,
164 assistantSegmentOrdinal: active ? 1 : 0,
165 turnStartAt: snapshot.runtime.startedAt ?? (state.activeTurnId === snapshot.runtime.turnId ? state.turnStartAt : 0),
166 resolvedPromptId: undefined,
167 items: [...items, ...users.filter(user => user.failed && !items.some(item => item.id === user.id)),
168 ...state.items.filter(item => item.kind === "notice" && item.local)],
169 offscreenItems: undefined,
170 seq: Math.max(state.seq, converted.seq),
171 running: active || hasLocalSubmission,
172 turnActive: active,
173 pendingPrompt: false,
174 cancelRequested: snapshot.runtime.status === "cancelling",
175 cancellable: active || hasLocalSubmission,
176 activeTurnId: active ? snapshot.runtime.turnId : undefined,
177 turnPhase: active ? snapshot.runtime.phase : undefined,
178 completionSummary: snapshot.runtime.completionSummary,
179 runtimeStatusEpoch: snapshot.identity.runtimeEpoch,
180 runtimeStatusSeq: snapshot.coveredThroughSeq,
181 runtimeStatusSnapshotAt: clock,
182 turnLifecycleObservedAt: clock,
183 pendingUser: hasLocalSubmission ? state.pendingUser : undefined,
184 pendingSubmissionId: hasLocalSubmission ? state.pendingSubmissionId : undefined,
185 live: undefined,
186 currentAssistant: undefined,
187 streamAttemptJournal: undefined,
188 approval: undefined,
189 ask: undefined,
190 mcpInteraction: undefined,
191 retry: undefined,
192 promptArrivedAt: undefined,
193 promptArrivedId: undefined,
194 promptEpoch: state.promptEpoch + 1,
195 promptWaitStartedAt: undefined,
196 turnWaitAccumMs: 0,
197 hydrateHistoryLoaded: true,
198 hydratePlaceholderItems: undefined,
199 historyPrefixCount: items.length,
200 historyStartTurn: Math.max(0, Math.min(...messages.filter((m) => m.role === "user" && m.historyTurn).map((m) => m.historyTurn!), snapshot.totalTurns) - 1),
201 historyEndTurn: snapshot.totalTurns,
202 historyTotalTurns: snapshot.totalTurns,
203 historyHasOlder: snapshot.hasOlder,
204 historyHasNewer: false,
205 historyOlderLoading: false,
206 historyOlderError: undefined,
207 historyNewerLoading: false,
208 historyNewerError: undefined,
209 historyRevision: undefined,
210 historyDigest: undefined,
211 historyMutation: { seq: state.historyMutation.seq + 1, kind: "replace" },
212 };
213 for (const event of snapshot.runtime.pendingEvents ?? []) {
214 next = applyEvent(next, { ...event, runtimeEpoch: snapshot.identity.runtimeEpoch });
215 }
216 const attempts = snapshot.activeAttempts ?? [];
217 const attempt = attempts[attempts.length - 1];
218 if (active && attempt) {
219 const id = `m:${attempt.messageId}`;
220 const message = messages.find((message) => message.messageId === attempt.messageId);
221 const live = { id, text: message?.content ?? "", reasoning: message?.reasoning ?? "", reasoningComplete: false };
222 next = { ...next, currentAssistant: id, live, streamAttemptJournal: {
223 id: attempt.id,
224 baselineLive: { ...live, text: "", reasoning: "" },
225 baselineTurnArgChars: 0,
226 createdToolIds: message?.toolCalls?.map((tool) => tool.id).filter(Boolean) ?? [],
227 priorTools: {},
228 } };
229 }
230 return next;
231 }
232
232 lines TYPESCRIPT