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 | })(); |