| 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 |