Oregami
Repositories/oxedyne/daimond

oxedyne/daimond/www/js/journal.js

14.4 KiB, 1 run

created by r2519314175:1387, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1/* journal.js — a write-ahead log, so work in flight survives the tab dying.
2 *
3 * Daimond persists a COMPLETED turn to localStorage (the snapshot). Everything between the moment
4 * you press Send and the moment the turn finishes lived only in memory: a crash, a shut browser,
5 * a discarded tab took the whole turn — the prompt included — with it. This is the other half.
6 *
7 * The model is crash-only (Candea & Fox): assume the tab can vanish between any two instructions,
8 * keep the durable log ALWAYS CURRENT, and make recovery the only startup path. There is no
9 * clean-shutdown step to rely on, because on the web there is no clean shutdown to rely on.
10 *
11 * Division of labour:
12 * - The SNAPSHOT (localStorage `daimond-chats`) is the durable record of COMPLETED turns.
13 * - The JOURNAL (this, in IndexedDB) is the write-ahead log of the IN-FLIGHT turn: the prompt,
14 * the reply as it streams, each tool as it is called and as it returns. When a turn finishes
15 * it is folded into the snapshot and its journal events are pruned, so the journal only ever
16 * holds what has not yet been made durable elsewhere.
17 *
18 * On the next boot, whatever is still in the journal was interrupted. Recovery reconstructs it —
19 * the prompt, the partial reply, the tools that ran — hands it back to the app to show as
20 * interrupted, and offers to continue it.
21 *
22 * IndexedDB, not localStorage: it is roomy, structured, appends without rewriting the whole value,
23 * and its writes commit to disk. It is async, but that is fine because we journal CONTINUOUSLY as
24 * events happen, never in a last-gasp unload handler (async writes do not finish there, and the
25 * events would already be safe). Per account, via the same namespace the rest of storage uses.
26 */
27(function () {
28 'use strict';
29
30 var DB_BASE = 'daimond-journal';
31 var STORE = 'events';
32 var VERSION = 1;
33
34 var db = null;
35 var dbOpen = null; // the name `db` is open on, to notice an account switch
36 var pending = []; // events buffered since the last flush
37 var timer = null;
38 var FLUSH_MS = 200; // deltas coalesce inside this window; markers flush at once
39 var _seq = 0;
40
41 function dbName() {
42 var ns = (window.DaimondAccounts && DaimondAccounts.opfsNs()) || '';
43 return ns ? DB_BASE + '-' + ns : DB_BASE;
44 }
45
46 function open(name) {
47 return new Promise(function (resolve, reject) {
48 var req = indexedDB.open(name || dbName(), VERSION);
49 req.onupgradeneeded = function () {
50 var d = req.result;
51 if (!d.objectStoreNames.contains(STORE)) {
52 var os = d.createObjectStore(STORE, { keyPath: 'k', autoIncrement: true });
53 os.createIndex('by_stream', 'stream', { unique: false });
54 }
55 };
56 req.onsuccess = function () { resolve(req.result); };
57 req.onerror = function () { reject(req.error); };
58 });
59 }
60
61 async function init() {
62 var want = dbName();
63 // The account can change under us (a switch points storage at a different namespace). The
64 // cached handle is for the OLD account's store, so reopen on the new name — and drop any
65 // buffered events, which belonged to the account we are leaving (a real switch reloads the
66 // page, so this only bites a same-session namespace change, but it must still isolate).
67 if (db && dbOpen === want) return db;
68 if (db && dbOpen !== want) { try { db.close(); } catch (e) { /* already */ } db = null; pending = []; }
69 try { db = await open(want); dbOpen = want; } catch (e) { db = null; dbOpen = null; }
70 // A tab that dies while a transaction is open leaves the connection unusable on the next
71 // event; drop the handle so the next call reopens.
72 if (db) {
73 db.onclose = function () { db = null; dbOpen = null; };
74 // Let a delete (forgetIdentity's deleteDatabase) proceed instead of blocking on this
75 // open connection: close on the version-change request.
76 db.onversionchange = function () { try { db.close(); } catch (e) { /* already */ } db = null; dbOpen = null; };
77 }
78 return db;
79 }
80
81 function tx(mode) {
82 var t = db.transaction(STORE, mode);
83 return { store: t.objectStore(STORE), done: new Promise(function (res, rej) {
84 t.oncomplete = res; t.onerror = function () { rej(t.error); }; t.onabort = function () { rej(t.error); };
85 }) };
86 }
87
88 /// Push an event onto the buffer. Markers (anything but a delta) flush immediately, so the
89 /// write-ahead record of a tool call or a turn boundary is on disk without waiting; deltas
90 /// coalesce, so a fast stream is not one disk write per token.
91 function push(ev, immediate) {
92 ev.stream = ev.stream || ev.chatId || (ev.runId ? 'agent:' + ev.runId : 'x');
93 ev.t = ev.t || 0; // ts is stamped by the caller (Date.now is fine in JS)
94 pending.push(ev);
95 if (immediate) return flush();
96 if (!timer) timer = setTimeout(function () { timer = null; flush(); }, FLUSH_MS);
97 return Promise.resolve();
98 }
99
100 /// Coalesce consecutive text deltas for the same stream (one turn has one assistant stream) into
101 /// one event, so a burst of tokens costs one row, not hundreds.
102 function coalesce(evs) {
103 var out = [];
104 for (var i = 0; i < evs.length; i++) {
105 var e = evs[i], last = out[out.length - 1];
106 if (e.type === 'delta' && last && last.type === 'delta' && last.stream === e.stream) {
107 last.text += e.text;
108 } else {
109 out.push(e);
110 }
111 }
112 return out;
113 }
114
115 async function flush() {
116 if (!pending.length) return;
117 // Always through init(), which reopens if the account changed — so buffered events can
118 // never be written into the wrong account's store.
119 await init();
120 if (!db) return; // storage unavailable: stay in memory
121 var batch = coalesce(pending);
122 pending = [];
123 try {
124 var t = tx('readwrite');
125 // The key `k` is auto-generated: it must be ABSENT from the value, not present-and-
126 // undefined (which IndexedDB rejects as an invalid key rather than auto-filling).
127 batch.forEach(function (e) { if ('k' in e) delete e.k; t.store.add(e); });
128 await t.done;
129 } catch (e) {
130 // A failed flush must not lose the events; put them back to try again next tick.
131 pending = batch.concat(pending);
132 }
133 }
134
135 function stamp(ev) { ev.t = nowMs(); return ev; }
136 // Date.now() is available in the browser (this is not a workflow script); wrapped for one place.
137 function nowMs() { try { return Date.now(); } catch (e) { return 0; } }
138
139 // ── The turn ────────────────────────────────────────────────────
140 //
141 // Every event of ONE turn shares a `stream` = its turn id (the user message's id), NOT the
142 // chat id. Keying by chat conflated successive turns of the same chat and let one turn's prune
143 // wipe the next turn's events; keying by turn means each turn is opened, closed and pruned in
144 // isolation. `chatId` rides on every event too, so a turn can always be placed back in its
145 // chat even if its opening event is the one that was lost.
146
147 function turnOpen(turnId, chatId, text, meta) {
148 return push(stamp({ type: 'turn_open', stream: turnId, chatId: chatId, text: text, meta: meta || null }), true);
149 }
150 function delta(turnId, chatId, text) {
151 if (!text) return Promise.resolve();
152 return push(stamp({ type: 'delta', stream: turnId, chatId: chatId, text: text }), false);
153 }
154 function toolOpen(turnId, chatId, callId, name, args) {
155 return push(stamp({ type: 'tool_open', stream: turnId, chatId: chatId, callId: callId, name: name, args: args }), true);
156 }
157 /// A tool call ended, and HOW it ended travels as the engine's own word.
158 ///
159 /// `outcome` is `AgentEvent::ToolResult`'s third field -- exactly `done`, `refused` or
160 /// `failed` -- and is stored, not flattened. It used to be a BOOLEAN, and the crash-recovery
161 /// path in daimond.js then rebuilt the missing third value by reading the result TEXT, through
162 /// the one reader dev/CONTRACT_OUTCOME.md §4 quarantines for conversations written before the
163 /// field existed. §4 says the outcome is stored with the tool log so that exception stops
164 /// growing; this is the write that makes that true. A refusal recovered from an interrupted
165 /// turn came back as a FAILURE, which is what the user saw and what the Optimiser was told.
166 ///
167 /// Stored verbatim rather than normalised: a value this build does not recognise is the
168 /// engine's, and inventing one here would be the second source of truth the contract exists to
169 /// remove. An empty string means the field was not sent, which the reader treats as §4's case.
170 function toolDone(turnId, chatId, callId, result, outcome) {
171 return push(stamp({ type: 'tool_done', stream: turnId, chatId: chatId, callId: callId, result: result, outcome: String(outcome || '') }), true);
172 }
173 function turnError(turnId, chatId, message) {
174 return push(stamp({ type: 'turn_error', stream: turnId, chatId: chatId, message: message }), true);
175 }
176 /// Close a turn: flush the last of it, then prune only THIS turn's events. Scoped to the turn,
177 /// so a newer turn already opened on the same chat is never in range.
178 async function turnClose(turnId, chatId, pTok, cTok) {
179 await push(stamp({ type: 'turn_close', stream: turnId, chatId: chatId, pTok: pTok || 0, cTok: cTok || 0 }), true);
180 await clearStream(turnId);
181 }
182 function clearTurn(turnId) { return clearStream(turnId); }
183
184 // ── Agents (the conductor's dispatched workers) ─────────────────
185
186 function agentOpen(runId, rec) {
187 return push(stamp({ type: 'agent_open', runId: runId, rec: rec }), true);
188 }
189 function agentDelta(runId, text) {
190 if (!text) return Promise.resolve();
191 return push(stamp({ type: 'agent_delta', runId: runId, text: text }), false);
192 }
193 function agentClose(runId, status, pTok, cTok) {
194 var p = push(stamp({ type: 'agent_close', runId: runId, status: status, pTok: pTok || 0, cTok: cTok || 0 }), true);
195 return p.then(function () { return clearStream('agent:' + runId); });
196 }
197
198 // ── Pruning ─────────────────────────────────────────────────────
199
200 async function clearStream(stream) {
201 await flush(); // land anything buffered first
202 if (!db) return;
203 try {
204 var t = tx('readwrite');
205 var idx = t.store.index('by_stream');
206 var range = IDBKeyRange.only(stream);
207 await new Promise(function (res) {
208 var cur = idx.openCursor(range);
209 cur.onsuccess = function () { var c = cur.result; if (c) { c.delete(); c.continue(); } else res(); };
210 cur.onerror = function () { res(); };
211 });
212 await t.done;
213 } catch (e) { /* best effort */ }
214 }
215
216 function clearAgent(runId) { return clearStream('agent:' + runId); }
217
218 /// Wipe the whole journal — used when an account is forgotten.
219 async function clearAll() {
220 await init();
221 if (!db) return;
222 try { var t = tx('readwrite'); t.store.clear(); await t.done; } catch (e) { /* ignore */ }
223 }
224
225 // ── Recovery ────────────────────────────────────────────────────
226
227 /// Read the whole journal back, grouped per TURN into what was in flight when the tab died.
228 /// Returns { turns: [ {turnId, chatId, userText, text, tools, closed} ], agents: [ {runId, rec, text} ] }.
229 /// A turn is returned only if it never closed — i.e. it was interrupted — and can be placed in a
230 /// chat (its chatId is known). Because each turn is its own stream, successive turns of one chat
231 /// never conflate, and one turn's failed prune can never hide another's interruption.
232 async function recover() {
233 await init();
234 await flush(); // land anything still buffered before we read
235 var empty = { turns: [], agents: [] };
236 if (!db) return empty;
237 var rows = [];
238 try {
239 var t = tx('readonly');
240 await new Promise(function (res) {
241 var cur = t.store.openCursor();
242 cur.onsuccess = function () { var c = cur.result; if (c) { rows.push(c.value); c.continue(); } else res(); };
243 cur.onerror = function () { res(); };
244 });
245 await t.done;
246 } catch (e) { return empty; }
247
248 rows.sort(function (a, b) { return (a.k || 0) - (b.k || 0); });
249
250 var turns = {}, agents = {};
251 rows.forEach(function (e) {
252 if (e.type && e.type.indexOf('agent_') === 0 && e.runId) {
253 var a = agents[e.runId] || (agents[e.runId] = { runId: e.runId, rec: null, text: '', closed: false });
254 if (e.type === 'agent_open') a.rec = e.rec;
255 else if (e.type === 'agent_delta') a.text += (e.text || '');
256 else if (e.type === 'agent_close') a.closed = true;
257 return;
258 }
259 var tid = e.stream;
260 if (!tid) return;
261 var c = turns[tid] || (turns[tid] = { turnId: tid, chatId: '', userText: '', text: '', tools: [], closed: false });
262 if (e.chatId) c.chatId = e.chatId; // present on every turn event, so placement survives a lost open
263 if (e.type === 'turn_open') { c.userText = e.text || ''; }
264 else if (e.type === 'delta') { c.text += (e.text || ''); }
265 else if (e.type === 'tool_open') { c.tools.push({ callId: e.callId, name: e.name, args: e.args, result: null, outcome: '', done: false }); }
266 // `e.outcome` is the engine's word; `e.failed` is what a record written before this
267 // field carried, and is passed through UNREAD so the caller can tell "not stored" from
268 // "stored as done" -- which is the difference between §4's named exception and a guess.
269 else if (e.type === 'tool_done') { for (var i = c.tools.length - 1; i >= 0; i--) { if (c.tools[i].callId === e.callId) { c.tools[i].result = e.result; c.tools[i].outcome = String(e.outcome || ''); c.tools[i].failed = !!e.failed; c.tools[i].done = true; break; } } }
270 else if (e.type === 'turn_close') { c.closed = true; }
271 else if (e.type === 'turn_error') { c.closed = true; } // errored is terminal, not interrupted
272 });
273
274 var interruptedTurns = [];
275 Object.keys(turns).forEach(function (id) { var t = turns[id]; if (!t.closed && t.chatId) interruptedTurns.push(t); });
276 var interruptedAgents = [];
277 Object.keys(agents).forEach(function (id) { if (!agents[id].closed && agents[id].rec) interruptedAgents.push(agents[id]); });
278
279 return { turns: interruptedTurns, agents: interruptedAgents };
280 }
281
282 window.DaimondJournal = {
283 init: init,
284 flush: flush,
285 turnOpen: turnOpen,
286 delta: delta,
287 toolOpen: toolOpen,
288 toolDone: toolDone,
289 turnError: turnError,
290 turnClose: turnClose,
291 clearTurn: clearTurn,
292 agentOpen: agentOpen,
293 agentDelta: agentDelta,
294 agentClose: agentClose,
295 clearAgent: clearAgent,
296 clearAll: clearAll,
297 recover: recover,
298 available: function () { return !!window.indexedDB; },
299 };
300})();