返回 DeepSeek-Reasonix
sessionComposerPersistence.ts
根目录 / desktop / frontend / src / lib / sessionComposerPersistence.ts
1 import { useEffect, useMemo, useSyncExternalStore } from "react";
2 import { app } from "./bridge";
3 import type { SessionRef } from "./sessionRef";
4 import type { PersistentComposerDraft, PersistentComposerTarget } from "./composerDraftTypes";
5 import type { SessionComposerState } from "../generated/desktopContract.generated";
6 import { followupNotSubmitted } from "./pendingFollowup";
7 import { submissionOutcome, modelApplicationError, type ModelApplicationDetails } from "./modelApplication";
8
9 export const emptySessionInput = (): PersistentComposerDraft => ({ text: "", invocations: [], attachments: [], workspaceRefs: [], pastedBlocks: [], openPastedLabels: [], sessionRefs: [], selectedTextRefs: [] });
10 type Entry = { ref: SessionRef; state?: SessionComposerState; content: PersistentComposerDraft; generation: number; version: number; saved: number; users: number; touched: number;
11 application?: ModelApplicationDetails; error?: string; unreadable?: boolean; registering?: boolean; recovering?: Promise<void>; unsubmittedId?: string; settledId?: string; tasks: Set<Promise<unknown>>; loading?: Promise<void>; saving?: Promise<void>; timer?: ReturnType<typeof setTimeout> };
12 const entries = new Map<string, Entry>();
13 const tabs = new Map<string, string>();
14 const tabOwners = new Map<string, symbol>();
15
16 export function sendPersistedComposer<T>(tabId: string, display: string, submit: string, id: string, send: () => T | Promise<T>, expectedEdit?: number, sourceKey = tabs.get(tabId), kind = "turn") {
17 const entry = entries.get(sourceKey || "");
18 const owner = tabOwners.get(tabId);
19 const task = (async () => {
20 const complete = await beginPersistedComposerSubmission(tabId,id,JSON.stringify({kind,display,submit}),expectedEdit,sourceKey,owner);
21 try {
22 if (tabs.get(tabId) !== sourceKey || tabOwners.get(tabId) !== owner) throw new Error("reasonix_error:inbox_not_submitted");
23 const result = await send();
24 await complete?.("accepted");
25 return result;
26 } catch (error) {
27 const message=String(error);
28 const outcome = submissionOutcome(error);
29 const rejected = outcome ? outcome === "not_accepted" : followupNotSubmitted(error) || /^(?:Error:\s*)?submission not accepted(?:\n|$)/.test(message);
30 await complete?.(rejected ? "not_accepted" : "unknown");
31 throw error;
32 }
33 })().catch(async error => {
34 if (entry) await recoverInput(entry).catch(() => {});
35 throw error;
36 });
37 if (entry) { entry.tasks.add(task); void task.finally(()=>{entry.tasks.delete(task);prune();}).catch(()=>{}); }
38 return task;
39 }
40 const listeners = new Set<() => void>();
41 let serial = 0, revision = 0, exiting = false;
42 const keyFor = (ref: SessionRef) => `${ref.hostId}:${ref.sessionId}`;
43 const notify = () => { revision++; listeners.forEach(fn => fn()); };
44 const subscribe = (fn: () => void) => { listeners.add(fn); return () => { listeners.delete(fn); }; };
45 const snapshot = () => revision;
46 function prune() {
47 const idle=[...entries.entries()].filter(([,entry])=>entry.users===0 && entry.state && !entry.error && !entry.recovering && !entry.loading && !entry.saving && !entry.tasks.size && entry.version===entry.saved).sort((a,b)=>b[1].touched-a[1].touched);
48 for (const [key,entry] of idle.slice(20)) {
49 if (entry.timer) clearTimeout(entry.timer);
50 entries.delete(key);
51 for (const [tab,mapped] of tabs) if(mapped===key) tabs.delete(tab);
52 }
53 }
54 function parse(raw: string): PersistentComposerDraft {
55 const value = JSON.parse(raw || "{}");
56 if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Invalid saved input");
57 if (value.text !== undefined && typeof value.text !== "string") throw new Error("Invalid saved input text");
58 for (const field of ["invocations","attachments","workspaceRefs","pastedBlocks","openPastedLabels","sessionRefs","selectedTextRefs"]) {
59 if (value[field] !== undefined && !Array.isArray(value[field])) throw new Error(`Invalid saved input ${field}`);
60 }
61 const content = { ...emptySessionInput(), ...value } as PersistentComposerDraft;
62 // Preview URLs are temporary renderer resources and never durable data.
63 content.attachments = (content.attachments || []).map(attachment => ({ ...attachment, previewUrl: undefined, file: undefined }));
64 return content;
65 }
66 function encoded(content: PersistentComposerDraft) {
67 return JSON.stringify({ ...content, attachments: content.attachments.map(attachment => ({ ...attachment,
68 path: attachment.recoveryPath || attachment.path, draftId: attachment.recoveryPath ? undefined : attachment.draftId,
69 previewUrl: undefined, file: undefined })) });
70 }
71 function entryFor(ref: SessionRef) {
72 const key = keyFor(ref);
73 let entry = entries.get(key);
74 if (!entry) { entry = { ref, content: emptySessionInput(), generation: ++serial, version: 0, saved: 0, users:0, touched:Date.now(), tasks: new Set() }; entries.set(key, entry); }
75 return entry;
76 }
77 const locked = (entry: Entry) => !entry.state || Boolean(entry.registering || entry.recovering || entry.unreadable || entry.state.submissionId);
78 async function load(entry: Entry) {
79 if (entry.loading) return entry.loading;
80 entry.loading = (async () => {
81 try {
82 const state = await app.GetSessionComposerState(entry.ref);
83 if (entry.version !== entry.saved) return;
84 const content=parse(state.contentJson);
85 entry.state = state; entry.content = content; entry.error = undefined; entry.unreadable = false;
86 const tabId = [...tabs].find(([, key]) => key === keyFor(entry.ref))?.[0];
87 if (tabId && content.workspaceRefs.some(ref => ref.isDir && ref.path.startsWith("__reasonix_external_folder/"))) {
88 const { restoreExternalFolderReferences } = await import("./attachmentSubmit");
89 await restoreExternalFolderReferences(app, {kind:"session",session:entry.ref,tabId}, content.workspaceRefs);
90 }
91 const attachments = entry.content.attachments;
92 void Promise.all(attachments.map(async attachment => {
93 try {
94 const previewUrl = await app.AttachmentDataURLForComposerTarget({kind:"session",session:entry.ref},attachment.recoveryPath || attachment.path);
95 if (entries.get(keyFor(entry.ref))!==entry || !entry.content.attachments.includes(attachment)) return;
96 entry.content = {...entry.content,attachments:entry.content.attachments.map(item=>item===attachment ? {...item,previewUrl} : item)};
97 notify();
98 } catch { /* Missing files remain visible as repairable references. */ }
99 }));
100 } catch (error) { entry.error = String(error); entry.unreadable = true; }
101 finally { entry.loading = undefined; notify(); prune(); }
102 })();
103 return entry.loading;
104 }
105 async function flush(entry: Entry) {
106 if (entry.recovering) await entry.recovering;
107 if (entry.loading) await entry.loading;
108 if (entry.saving) { await entry.saving; if (entry.version > entry.saved) return flush(entry); return; }
109 if (entry.error || !entry.state) throw new Error(entry.error || "Input is not loaded");
110 if (entry.version === entry.saved && !entry.state.historyChanged) return;
111 entry.saving = (async () => {
112 let rebased = false;
113 while (entry.saved < entry.version || entry.state?.historyChanged) {
114 if (!entry.state || entry.state.submissionId) throw new Error("Check the previous submission before editing");
115 const version = entry.version;
116 const result = await app.SaveSessionComposerState({ ref: entry.ref, expectedRevision: entry.state.revision,
117 contentJson: encoded(entry.content), contentVersion: 1, acknowledgeHistory: true });
118 if (result.conflict) {
119 entry.state = result;
120 if (rebased || result.submissionId) throw new Error("Input changed during saving; retry to save the retained input");
121 // Rebase once without replacing the input currently being edited.
122 rebased = true; continue;
123 }
124 if (result.historyChanged) { entry.state = result; throw new Error("The conversation changed. Review the saved input."); }
125 entry.state = result; entry.saved = version;
126 }
127 })();
128 try { await entry.saving; } catch (error) { entry.error = String(error); throw error; }
129 finally { entry.saving = undefined; notify(); prune(); }
130 }
131 function edit(entry: Entry, content: PersistentComposerDraft) {
132 if (locked(entry)) return;
133 entry.content = content; entry.version++; entry.error = undefined; entry.application = undefined;
134 if (entry.timer) clearTimeout(entry.timer);
135 entry.timer = setTimeout(() => {
136 void flush(entry).catch(async () => { await recoverInput(entry); if (!locked(entry)) await flush(entry); }).catch(() => {});
137 }, 250);
138 notify();
139 }
140
141 export async function flushAllSessionComposers() {
142 exiting = true;
143 notify();
144 try {
145 for (const entry of entries.values()) {
146 while (entry.tasks.size) await Promise.allSettled([...entry.tasks]);
147 if (entry.version > entry.saved || entry.error) await flush(entry);
148 }
149 } catch (error) { exiting = false; throw error; }
150 }
151 export function resumeSessionComposerEditing() { exiting = false; notify(); }
152
153 export async function beginPersistedComposerSubmission(tabId: string, submissionId: string, requestJSON: string, expectedEdit?: number, key = tabs.get(tabId), owner = tabOwners.get(tabId)) {
154 const entry = key ? entries.get(key) : undefined;
155 if (!entry) return undefined;
156 await flush(entry);
157 if (!owner || tabs.get(tabId) !== key || tabOwners.get(tabId) !== owner) throw new Error("reasonix_error:inbox_not_submitted");
158 if (expectedEdit !== undefined && entry.version !== expectedEdit) throw new Error("Input changed while preparing submission; review and send again");
159 if (exiting || !entry.state || locked(entry)) throw new Error("Review the saved input before sending");
160 const submittedVersion = entry.version;
161 entry.registering = true;
162 notify();
163 try {
164 entry.state = await app.BeginSessionComposerSubmission(entry.ref, entry.state.revision, submissionId, requestJSON);
165 } catch (error) {
166 // The host may have committed despite transport failure. Reconcile before
167 // allowing edits or another send; never infer rejection from a lost reply.
168 entry.unsubmittedId = submissionId;
169 entry.unreadable = true; entry.error = String(error);
170 throw error;
171 } finally { entry.registering = false; notify(); }
172 if (entry.state.conflict) { notify(); throw new Error("The input changed before submission"); }
173 notify();
174 return (outcome: "accepted" | "not_accepted" | "unknown") => settleSubmission(entry, submissionId, outcome, submittedVersion);
175 }
176
177 async function settleSubmission(entry: Entry, submissionId: string, outcome: "accepted" | "not_accepted" | "unknown", submittedVersion: number) {
178 const result = await app.CompleteSessionComposerSubmission(entry.ref, submissionId, outcome);
179 if (result.conflict) { notify(); throw new Error("The input changed during submission"); }
180 if (!entry.state || BigInt(result.revision) >= BigInt(entry.state.revision)) {
181 entry.state = result;
182 // A recovery read may already have settled this send and enabled new
183 // edits. Only replace the captured version, using the host's actual input.
184 if (!result.submissionId && entry.version === submittedVersion) {
185 entry.content = parse(result.contentJson); entry.version++; entry.saved = entry.version;
186 }
187 if (!result.submissionId) entry.settledId = submissionId;
188 }
189 notify();
190 }
191
192 async function recoverInput(entry: Entry) {
193 if (entry.recovering) return entry.recovering;
194 if (entry.registering) return;
195 entry.recovering = (async () => {
196 if (entry.loading) await entry.loading;
197 if (entry.saving) await entry.saving.catch(() => {});
198 const previous = entry.state;
199 const pending = previous?.submissionId || entry.unsubmittedId;
200 try {
201 let state = await app.GetSessionComposerState(entry.ref);
202 // This renderer never invoked send after a failed registration. Only its
203 // own registration can be released; an absent receipt alone is not proof.
204 if (entry.unsubmittedId && state.submissionId === entry.unsubmittedId) {
205 state = await app.CompleteSessionComposerSubmission(entry.ref, entry.unsubmittedId, "not_accepted");
206 }
207 const content = parse(state.contentJson);
208 if (entry.state && BigInt(state.revision) < BigInt(entry.state.revision)) return;
209 entry.error = undefined; entry.unreadable = false;
210 entry.state = state;
211 if (state.submissionId) return;
212 if (pending) entry.settledId = pending;
213 if (previous?.submissionId && encoded(entry.content) === encoded(parse(previous.contentJson))) {
214 entry.content = content; entry.version++; entry.saved = entry.version;
215 } else if (encoded(content) === encoded(entry.content)) {
216 entry.saved = entry.version;
217 } else if (previous && encoded(content) !== encoded(parse(previous.contentJson))) {
218 // Continue with this window's input on its next save or explicit send.
219 if (entry.version === entry.saved) entry.version++;
220 } else if (!previous) {
221 entry.content = content; entry.saved = entry.version;
222 }
223 entry.unsubmittedId = undefined;
224 } catch (error) { entry.error = String(error); entry.unreadable = true; throw error; }
225 })();
226 notify();
227 try { await entry.recovering; } finally { entry.recovering = undefined; notify(); }
228 }
229
230 export function useSessionComposerPersistence(ref: SessionRef | undefined, tabId: string | undefined) {
231 const key = ref?.hostId === "local" ? keyFor(ref) : "";
232 const entry = useMemo(() => key && ref ? entryFor(ref) : undefined, [key]);
233 useSyncExternalStore(subscribe, snapshot);
234 useEffect(() => {
235 if (!entry) return;
236 entries.set(key,entry);
237 entry.users++; entry.touched=Date.now();
238 const owner = Symbol();
239 if (tabId) { tabs.set(tabId,key); tabOwners.set(tabId,owner); }
240 if (!entry.state && !entry.loading) void load(entry);
241 return () => {
242 if (tabId && tabOwners.get(tabId) === owner) { tabs.delete(tabId); tabOwners.delete(tabId); }
243 entry.users--; if (entry.version > entry.saved) void flush(entry).catch(() => {}); prune();
244 };
245 }, [entry,key,tabId]);
246 const target: PersistentComposerTarget | undefined = entry ? {
247 draftId:key, generation:entry.generation, initial:entry.content, revision:entry.version,
248 isCurrent:(id,generation) => id===key && generation===entry.generation,
249 canEdit:() => !locked(entry),
250 onChange:(_id,generation,content) => { if (generation===entry.generation) edit(entry,content); },
251 onPatch:(_id,generation,patch) => { if (generation===entry.generation) edit(entry,{ ...entry.content,...(typeof patch==="function" ? patch(entry.content) : patch) }); },
252 trackTask:(_id,_generation,promise) => { entry.tasks.add(promise); void promise.finally(()=>{entry.tasks.delete(promise);prune();}).catch(()=>{}); return promise; },
253 onTaskError:(_id,_generation,error) => { entry.error=error; notify(); },
254 } : undefined;
255 return { target, blocked:entry ? exiting || locked(entry) : false, error:entry?.error,
256 application:entry?.application,
257 setApplication:(value:ModelApplicationDetails|undefined)=>{if(entry){entry.application=value;notify();}},
258 reportSubmissionError:(error:unknown) => { if (entry) {entry.application=modelApplicationError(error);notify();} },
259 // Editable retained text does not authorize restoring a mode from older history.
260 goalDraft:!entry?.state?.historyChanged && entry?.content.goalDraft,
261 setGoalDraft:(enabled:boolean) => { if (entry) edit(entry,{...entry.content,goalDraft:enabled}); },
262 settleSubmission:async (id:string) => { if (entry?.state?.submissionId === id) await settleSubmission(entry,id,"accepted",entry.version); },
263 attention: Boolean(entry?.state?.submissionId),
264 recovering: Boolean(entry?.recovering),
265 settledId: entry?.settledId,
266 retry:async () => { if (!entry) return; await recoverInput(entry); if (!locked(entry)) await flush(entry); },
267 };
268 }
269
269 lines TYPESCRIPT