返回 DeepSeek-Reasonix
localSubmissionState.ts
根目录 / desktop / frontend / src / lib / localSubmissionState.ts
1 export type LocalSubmissionStatus = "sending" | "accepted" | "unknown" | "failed";
2
3 export interface LocalSubmission {
4 submissionId: string;
5 localId: string;
6 text: string;
7 submitText?: string;
8 createdAt: number;
9 sequence: number;
10 anchorItemId?: string;
11 placement?: "start" | "after" | "latest";
12 status: LocalSubmissionStatus;
13 messageId?: string;
14 turnId?: string;
15 checkpointTurn?: number;
16 /** The first terminal event consumed this submission's checkpoint authority. */
17 settled?: boolean;
18 }
19
20 export interface LocalSubmissionFields {
21 localSubmissions: Record<string, LocalSubmission>;
22 localSubmissionOrder: string[];
23 localSubmissionSendRevision: number;
24 visibleSubmissionHandoffs: Record<string, { submissionId: string }>;
25 }
26
27 export type CanonicalUserConfirmation = { messageId: string; submissionId?: string; turnId?: string };
28 type CanonicalUserIdentity = {
29 kind: string;
30 messageId?: string;
31 submissionId?: string;
32 turnId?: string;
33 };
34
35 export function canonicalUserConfirmations(items: readonly CanonicalUserIdentity[]): CanonicalUserConfirmation[] {
36 return items.flatMap(item => item.kind === "user" && item.messageId ? [{ messageId: item.messageId, submissionId: item.submissionId, turnId: item.turnId }] : []);
37 }
38
39 /** Identity matching is shared by state settlement and defensive presentation. */
40 export function matchLocalSubmissions(submissions: readonly LocalSubmission[], confirmations: readonly CanonicalUserConfirmation[]) {
41 const consumed = new Set<string>();
42 const matches: Array<{ submissionId: string; messageId: string }> = [];
43 for (const confirmation of confirmations) {
44 const available = (local: LocalSubmission) => !consumed.has(local.submissionId);
45 const local = submissions.find(local => available(local) && local.messageId === confirmation.messageId)
46 ?? submissions.find(local => available(local) && !local.messageId && Boolean(confirmation.submissionId) && local.submissionId === confirmation.submissionId);
47 if (!local) continue;
48 consumed.add(local.submissionId);
49 matches.push({ submissionId: local.submissionId, messageId: confirmation.messageId });
50 }
51 return matches;
52 }
53
54 export function pruneSubmissionHandoffs<T extends LocalSubmissionFields>(state: T, items: readonly CanonicalUserIdentity[]): T {
55 const visible = new Set(canonicalUserConfirmations(items).map(item => item.messageId));
56 const entries = Object.entries(state.visibleSubmissionHandoffs).filter(([id]) => visible.has(id));
57 return entries.length === Object.keys(state.visibleSubmissionHandoffs).length ? state
58 : { ...state, visibleSubmissionHandoffs: Object.fromEntries(entries) };
59 }
60
61 export function isUnknownSubmissionError(error: unknown): boolean {
62 const outcome=(error as {data?:{submissionOutcome?:string}} | undefined)?.data?.submissionOutcome;
63 if(outcome) return outcome === "unknown";
64 return /timeout|timed out|network|connection|socket|channel.*closed|fetch failed|failed to fetch|\beof\b/i.test(error instanceof Error ? error.message : String(error));
65 }
66
67 export function orderedLocalSubmissions(state: LocalSubmissionFields): LocalSubmission[] {
68 return state.localSubmissionOrder.flatMap((submissionId) => {
69 const submission = state.localSubmissions[submissionId];
70 return submission ? [submission] : [];
71 });
72 }
73
74 export function beginLocalSubmission<T extends LocalSubmissionFields>(
75 state: T,
76 submission: Omit<LocalSubmission, "status">,
77 ): T {
78 return {
79 ...state,
80 localSubmissions: { ...state.localSubmissions, [submission.submissionId]: { ...submission, status: "sending" } },
81 localSubmissionOrder: state.localSubmissionOrder.includes(submission.submissionId)
82 ? state.localSubmissionOrder
83 : [...state.localSubmissionOrder, submission.submissionId],
84 localSubmissionSendRevision: state.localSubmissionSendRevision + 1,
85 };
86 }
87
88 export function updateLocalSubmission<T extends LocalSubmissionFields>(
89 state: T,
90 submissionId: string | undefined,
91 patch: Partial<LocalSubmission>,
92 ): T {
93 if (!submissionId) return state;
94 const current = state.localSubmissions[submissionId];
95 if (!current) return state;
96 if (current.messageId && patch.messageId && current.messageId !== patch.messageId) return state;
97 return {
98 ...state,
99 localSubmissions: { ...state.localSubmissions, [submissionId]: { ...current, ...patch, submissionId } },
100 };
101 }
102
103 export function removeLocalSubmission<T extends LocalSubmissionFields>(state: T, submissionId: string | undefined): T {
104 if (!submissionId || !state.localSubmissions[submissionId]) return state;
105 const localSubmissions = { ...state.localSubmissions };
106 delete localSubmissions[submissionId];
107 return {
108 ...state,
109 localSubmissions,
110 localSubmissionOrder: state.localSubmissionOrder.filter((candidate) => candidate !== submissionId),
111 };
112 }
113
114 /** Retire each local echo once its durable user message is installed. */
115 export function settleLocalSubmissions<T extends LocalSubmissionFields>(
116 state: T,
117 items: readonly CanonicalUserIdentity[],
118 confirmations: readonly CanonicalUserConfirmation[] = canonicalUserConfirmations(items),
119 ): T {
120 state = pruneSubmissionHandoffs(state, items);
121 if (state.localSubmissionOrder.length === 0) return state;
122 const remaining = { ...state.localSubmissions };
123 const matches = matchLocalSubmissions(orderedLocalSubmissions(state), confirmations);
124 const consumed = new Set(matches.map(match => match.submissionId));
125 const visible = new Set(canonicalUserConfirmations(items).map(item => item.messageId));
126 const handoffs = { ...state.visibleSubmissionHandoffs };
127 for (const match of matches) {
128 delete remaining[match.submissionId];
129 if (visible.has(match.messageId) && !handoffs[match.messageId]) handoffs[match.messageId] = { submissionId: match.submissionId };
130 }
131 if (consumed.size === 0) return state;
132 return {
133 ...state,
134 localSubmissions: remaining,
135 visibleSubmissionHandoffs: handoffs,
136 localSubmissionOrder: state.localSubmissionOrder.filter((submissionId) => !consumed.has(submissionId)),
137 };
138 }
139
140 export type RebasedUserRecord = { submissionId?: string; text: string; submitText?: string };
141 type EchoLike = { kind: string; id: string; text?: string; submitText?: string };
142
143 function submittedText(entry: { text: string; submitText?: string }): string {
144 return entry.submitText ?? entry.text;
145 }
146
147 /** User echoes that were already on screen when this submission was created. */
148 function preexistingEchoes(previousItems: readonly EchoLike[], local: LocalSubmission): EchoLike[] {
149 if (!local.anchorItemId) return local.placement === "start" ? [] : [...previousItems];
150 const anchor = previousItems.findIndex(item => item.id === local.anchorItemId);
151 return anchor < 0 ? [...previousItems] : previousItems.slice(0, anchor + 1);
152 }
153
154 /**
155 * A snapshot rebase (runtime rebuild, serve restart, model switch) may carry
156 * user records without the message ids settleLocalSubmissions keys on, so an
157 * echo the server already journaled would otherwise survive as a duplicate
158 * turn stuck at "processing". Retire only echoes the rebased projection
159 * provably owns: one whose submission id a durable record repeats, or — with
160 * no turn in flight — the oldest echo whose text the newest durable user
161 * record repeats. The text match is bounded by the echoes' own anchor: a
162 * record that already existed when the echo was created can never absorb it,
163 * so a lost re-send of an identical message stays visible. An active runtime
164 * never absorbs by text because the trailing record may be an older sibling
165 * of the in-flight submission.
166 */
167 export function settleRebasedSubmissions<T extends LocalSubmissionFields>(
168 state: T,
169 previousItems: readonly EchoLike[],
170 durable: readonly RebasedUserRecord[],
171 idle: boolean,
172 ): T {
173 const pending = orderedLocalSubmissions(state).filter(local => local.status !== "failed");
174 if (pending.length === 0) return state;
175 const consumed = new Set<string>();
176 const durableSubmissionIds = new Set(durable.map(record => record.submissionId).filter(Boolean));
177 for (const local of pending) if (durableSubmissionIds.has(local.submissionId)) consumed.add(local.submissionId);
178 const latest = durable[durable.length - 1];
179 if (idle && latest) {
180 const text = submittedText(latest);
181 const durableCopies = durable.filter(record => submittedText(record) === text).length;
182 const echo = pending.find(local => !consumed.has(local.submissionId) && submittedText(local) === text);
183 if (echo) {
184 const priorCopies = preexistingEchoes(previousItems, echo)
185 .filter(item => item.kind === "user" && item.text !== undefined && submittedText({ text: item.text, submitText: item.submitText }) === text).length;
186 if (durableCopies > priorCopies) consumed.add(echo.submissionId);
187 }
188 }
189 if (consumed.size === 0) return state;
190 const localSubmissions = { ...state.localSubmissions };
191 for (const submissionId of consumed) delete localSubmissions[submissionId];
192 return {
193 ...state,
194 localSubmissions,
195 localSubmissionOrder: state.localSubmissionOrder.filter(submissionId => !consumed.has(submissionId)),
196 };
197 }
198
199 export function checkpointLocalSubmission<T extends LocalSubmissionFields>(
200 state: T,
201 submissionId: string | undefined,
202 checkpointTurn: number | undefined,
203 ): T {
204 if (!submissionId) return state;
205 const current = state.localSubmissions[submissionId];
206 if (!current || current.settled) return state;
207 const validTurn = Number.isInteger(checkpointTurn) && checkpointTurn! >= 0;
208 return updateLocalSubmission(state, submissionId, {
209 checkpointTurn: validTurn ? checkpointTurn : current.checkpointTurn,
210 settled: true,
211 });
212 }
213
213 lines TYPESCRIPT