oxedyne/daimond/www/js/peer.js
96.5 KiB, 151 runs
created by r2519314175:1413, 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 | /* ============================================================ |
| 2 | Daimond — the persistent desktop peer, seam layer (peer.js) |
| 3 | ------------------------------------------------------------ |
| 4 | Step 1 of dev/PEER_DESIGN.md: PROVE THE SEAM. One tab seals a |
| 5 | small errand envelope to its OWN account and drops it in the |
| 6 | post box; a second tab of the same account collects it, runs |
| 7 | the ordinary turn, folds the answer into the transcript and |
| 8 | pushes the parcel; the first tab collects the answer by the |
| 9 | ordinary sync merge. This module is the client-only glue for |
| 10 | that, and NOTHING here is a gateway change: an errand and a |
| 11 | report ride the same `/api/post` door a message does, sealed |
| 12 | with the same seal, read only by the account that holds the key. |
| 13 | |
| 14 | ── WHAT THIS LAYER OWNS, AND WHAT IT DOES NOT ────────────── |
| 15 | |
| 16 | It OWNS the envelope shapes (errand, report), the self-seal to |
| 17 | the account's own sealing key, the open that only the account |
| 18 | can perform, and the route-by-type a collector does. It also |
| 19 | owns the ONE fold the peer's result needs -- an assistant |
| 20 | message appended under the turn's id -- which §3.1 of the |
| 21 | design calls "an append, so it merges with nothing new". |
| 22 | |
| 23 | It does NOT own the transport (the raw put and the collect are |
| 24 | DaimondPost's, injected here so step 1 needs no post.js edit -- |
| 25 | see the recommendation for step 2 at the foot of the design), |
| 26 | the turn engine (`runTurn`, injected), or the parcel merge |
| 27 | (sync.js / daimond.js `mergeMessages`, which unions this fold |
| 28 | in unchanged). The lease is STEP 3 and is deliberately absent: |
| 29 | step 1 assumes a single peer, so a double run is not yet a |
| 30 | concern -- the point here is only that the errand travels, |
| 31 | seals, opens to the same account ALONE, runs, and merges back. |
| 32 | |
| 33 | ── THE SEAL, REUSED NOT REINVENTED ───────────────────────── |
| 34 | |
| 35 | The envelope is sealed with `DaimondPost.seal` (post.js:291) to |
| 36 | a single recipient: the account's OWN sealing key, |
| 37 | `DaimondIdentity.sealingKeyRaw()`. That is exactly the self-slot |
| 38 | `compose` already adds to every message so a Sent copy opens on |
| 39 | the account's other devices (post.js:539-546). The gateway names |
| 40 | no recipient in the clear and a reader trial-decrypts, so the |
| 41 | gateway learns only that an account posted to itself. |
| 42 | |
| 43 | Attaches one global, `window.DaimondPeer`. |
| 44 | ============================================================ */ |
| 45 | (function () { |
| 46 | 'use strict'; |
| 47 | |
| 48 | // The schema version the errand and report carry. Bumped when a field's |
| 49 | // meaning changes, so a peer never runs an envelope it half-understands. |
| 50 | var ENVELOPE_V = 1; |
| 51 | |
| 52 | var T_ERRAND = 'errand'; // a turn dispatched to a peer |
| 53 | var T_REPORT = 'report'; // a peer's account of how the turn went |
| 54 | // The two envelopes remote consent rides. A runner blocked on a genuinely |
| 55 | // per-turn question (a `web_type`, a `web_click`, an overlong address) that the |
| 56 | // synced account policy does NOT already cover cannot raise a dialog where nobody |
| 57 | // is, so it seals a `consent-ask` to the account and awaits the answer; an |
| 58 | // attended device raises the tile and seals back a `consent-grant`. They ride the |
| 59 | // same self-seal and the same signature as the errand -- no new transport, no new |
| 60 | // crypto -- and slot into the same route-by-type the collector already does. |
| 61 | var T_ASK = 'consent-ask'; // runner -> the account: a live question for a human |
| 62 | var T_GRANT = 'consent-grant'; // an attended device -> the runner: the answer |
| 63 | |
| 64 | /// Is `t` a peer envelope tag this layer owns? The collector routes exactly these |
| 65 | /// four and hands everything else (a message artefact) to the message path. |
| 66 | function peerType(t) { |
| 67 | return t === T_ERRAND || t === T_REPORT || t === T_ASK || t === T_GRANT; |
| 68 | } |
| 69 | |
| 70 | // ── Bytes and text ───────────────────────────────────────── |
| 71 | |
| 72 | function utf8(s) { return new TextEncoder().encode(String(s)); } |
| 73 | function fromUtf8(b) { return new TextDecoder().decode(b); } |
| 74 | |
| 75 | /// Standard base64 of some bytes, and back. The seal hands base64 across the |
| 76 | /// wire, so the envelope does too. |
| 77 | function b64enc(bytes) { |
| 78 | var b = (bytes instanceof Uint8Array) ? bytes : new Uint8Array(bytes); |
| 79 | var s = ''; |
| 80 | for (var i = 0; i < b.length; i++) s += String.fromCharCode(b[i]); |
| 81 | return btoa(s); |
| 82 | } |
| 83 | function b64dec(s) { |
| 84 | var raw = atob(String(s)); |
| 85 | var out = new Uint8Array(raw.length); |
| 86 | for (var i = 0; i < raw.length; i++) out[i] = raw.charCodeAt(i); |
| 87 | return out; |
| 88 | } |
| 89 | |
| 90 | // The peer envelope is sealed with the account's SHARED symmetric key (the |
| 91 | // passphrase+salt one every paired device derives), NOT the per-device X25519 |
| 92 | // sealing key. The old per-device seal was the "same account, different sealing |
| 93 | // key -> silent drop" bug that needed a manual re-pair: a device that lazily |
| 94 | // minted its own sealing key could not open a sibling's errand even though the |
| 95 | // account was one. The symmetric key travels whole in the pairing bundle (the |
| 96 | // salt does), so every device of the account opens it and the gateway -- which |
| 97 | // never holds it -- opens nothing. |
| 98 | // |
| 99 | // The current scheme is `DPY2`: AES-GCM under the account key with the purpose |
| 100 | // string bound in as additional data, so the ciphertext is cryptographically |
| 101 | // domain-separated from every other thing sealed under that one key (the parcel, |
| 102 | // the wrapped keys, the voice) and cannot be opened where any of them is |
| 103 | // expected. The AAD authenticates the PURPOSE, not the literal tag bytes -- the |
| 104 | // tag only routes the decrypt -- but that is enough: a body sealed for one |
| 105 | // purpose fails GCM under another's AAD, so a mis-routed or flipped tag fails to |
| 106 | // open rather than opening something. The open path also reads a legacy `DPY1` |
| 107 | // (the same key, no AAD, the first-hour form) and a legacy post.js `DPS1` X25519 |
| 108 | // envelope, so a hand-off in flight across a rollout is never dropped. |
| 109 | var PEER_AAD = 'daimond/peer/env/1'; // the envelope's GCM domain |
| 110 | var SYM_MAGIC = new Uint8Array([0x44, 0x50, 0x59, 0x32]); // "DPY2" -- AAD-bound |
| 111 | var SYM_MAGIC_1 = new Uint8Array([0x44, 0x50, 0x59, 0x31]); // "DPY1" -- legacy, no AAD |
| 112 | |
| 113 | /// Do the first four bytes match this scheme tag? |
| 114 | function tagged(bytes, tag) { |
| 115 | if (!bytes || bytes.length < tag.length) return false; |
| 116 | for (var i = 0; i < tag.length; i++) { |
| 117 | if (bytes[i] !== tag[i]) return false; |
| 118 | } |
| 119 | return true; |
| 120 | } |
| 121 | |
| 122 | /// Concatenate byte arrays into one Uint8Array. |
| 123 | function cat(parts) { |
| 124 | var n = 0, i; |
| 125 | for (i = 0; i < parts.length; i++) n += parts[i].length; |
| 126 | var out = new Uint8Array(n), off = 0; |
| 127 | for (i = 0; i < parts.length; i++) { out.set(parts[i], off); off += parts[i].length; } |
| 128 | return out; |
| 129 | } |
| 130 | |
| 131 | function hex(bytes) { |
| 132 | var b = (bytes instanceof Uint8Array) ? bytes : new Uint8Array(bytes); |
| 133 | var s = ''; |
| 134 | for (var i = 0; i < b.length; i++) s += ('0' + b[i].toString(16)).slice(-2); |
| 135 | return s; |
| 136 | } |
| 137 | |
| 138 | /// A random 128-bit id, hex. Names one dispatch (`eid`), distinct from the |
| 139 | /// turn id, so a re-dispatch of the same turn is still a different errand. |
| 140 | function newId() { |
| 141 | return hex(crypto.getRandomValues(new Uint8Array(16))); |
| 142 | } |
| 143 | |
| 144 | function unhex(s) { |
| 145 | var str = String(s), out = new Uint8Array(str.length / 2); |
| 146 | for (var i = 0; i < out.length; i++) out[i] = parseInt(str.substr(i * 2, 2), 16); |
| 147 | return out; |
| 148 | } |
| 149 | |
| 150 | /// The content address of the sealed bytes: a SHA-256, hex. This is the `addr` |
| 151 | /// the post body carries beside the envelope -- the relay addresses a row by it |
| 152 | /// and a re-post of the identical envelope collapses to one row, exactly as a |
| 153 | /// message's address does (post.js `compose` -> `address`). |
| 154 | async function addressOf(sealedBytes) { |
| 155 | var d = await crypto.subtle.digest('SHA-256', sealedBytes); |
| 156 | return hex(new Uint8Array(d)); |
| 157 | } |
| 158 | |
| 159 | // ── The signature, off the wasm message path ─────────────── |
| 160 | // |
| 161 | // The seal restricts WHO CAN OPEN the errand to the account (the one slot is |
| 162 | // the account's own sealing key). It does NOT restrict who can WRITE one: the |
| 163 | // account's PUBLIC sealing key is on its card and in its QR code, so anyone |
| 164 | // holding the card could seal an errand to the account and drop it in the box, |
| 165 | // and the account would open it. A peer that ran that would run a stranger's |
| 166 | // errand on the account's money. The signature closes exactly that hole: the |
| 167 | // envelope is signed with the account's PRIVATE signing key, which is on no |
| 168 | // card, so `verifyEnvelope` accepts only an envelope this account authored. |
| 169 | // |
| 170 | // It is a DETACHED signature over the canonical bytes, NOT the wasm |
| 171 | // `signingInput`/`assemble` path a message takes -- the errand stays off the |
| 172 | // bridge and out of the message renderer, which is the whole reason it is raw |
| 173 | // JSON and not a message artefact. |
| 174 | |
| 175 | /// Canonical JSON of a value: object keys sorted, arrays in order, so the |
| 176 | /// signer and the verifier serialise byte-for-byte the same thing however |
| 177 | /// their field insertion order happened to differ. |
| 178 | function canonical(v) { |
| 179 | if (v === null || typeof v !== 'object') return JSON.stringify(v); |
| 180 | if (Array.isArray(v)) return '[' + v.map(canonical).join(',') + ']'; |
| 181 | var keys = Object.keys(v).sort(); |
| 182 | return '{' + keys.map(function (k) { |
| 183 | return JSON.stringify(k) + ':' + canonical(v[k]); |
| 184 | }).join(',') + '}'; |
| 185 | } |
| 186 | |
| 187 | /// The bytes a signature is over: the envelope WITHOUT its own `author`/`sig`, |
| 188 | /// canonicalised. Both sides compute these identically -- the signer before it |
| 189 | /// adds the two fields, the verifier after it strips them. |
| 190 | function signedBytes(obj) { |
| 191 | var base = {}; |
| 192 | Object.keys(obj).forEach(function (k) { |
| 193 | if (k !== 'author' && k !== 'sig') base[k] = obj[k]; |
| 194 | }); |
| 195 | return utf8(canonical(base)); |
| 196 | } |
| 197 | |
| 198 | /// Sign an envelope with this account's signing key and answer a copy carrying |
| 199 | /// `author` (this account's public key, hex) and `sig` (base64, as |
| 200 | /// `DaimondIdentity.sign` answers). |
| 201 | async function signEnvelope(obj) { |
| 202 | if (!window.DaimondIdentity || !window.DaimondIdentity.sign) { |
| 203 | throw new Error('peer: no identity, so an errand cannot be signed.'); |
| 204 | } |
| 205 | var author = await window.DaimondIdentity.publicKeyRaw(); |
| 206 | if (!author) throw new Error('peer: this device has no signing key.'); |
| 207 | var sig = await window.DaimondIdentity.sign(signedBytes(obj)); // base64 |
| 208 | var out = {}; |
| 209 | Object.keys(obj).forEach(function (k) { out[k] = obj[k]; }); |
| 210 | out.author = hex(author); |
| 211 | out.sig = sig; |
| 212 | return out; |
| 213 | } |
| 214 | |
| 215 | /// Whether an envelope was authored by THIS account: the signature verifies |
| 216 | /// AND `author` is this account's own public key. Both halves matter -- a valid |
| 217 | /// signature by a stranger's key is a stranger's errand, and an `author` set to |
| 218 | /// our key with no matching signature is a forgery that named us. |
| 219 | async function verifyEnvelope(obj) { |
| 220 | if (!obj || !obj.author || !obj.sig) return false; |
| 221 | if (!window.DaimondIdentity || !window.DaimondIdentity.verifySig) return false; |
| 222 | var mine = await window.DaimondIdentity.publicKeyRaw(); |
| 223 | if (!mine || hex(mine) !== String(obj.author)) return false; // authored by us? |
| 224 | try { return await window.DaimondIdentity.verifySig(unhex(obj.author), obj.sig, signedBytes(obj)); } |
| 225 | catch (e) { return false; } |
| 226 | } |
| 227 | |
| 228 | // ── The envelopes ────────────────────────────────────────── |
| 229 | |
| 230 | /// Build an errand envelope object (not yet sealed). The full §1.1 schema; a |
| 231 | /// step-1 dispatcher fills only the core (`turnId`, `chatId`, `prompt`, |
| 232 | /// `model`) and leaves the lease/freshness fields for later steps. The tag `t` |
| 233 | /// is what the collector routes on, and it lives INSIDE the sealed plaintext, |
| 234 | /// so the gateway -- which cannot open the seal -- never sees it. |
| 235 | function makeErrand(f) { |
| 236 | var o = f || {}; |
| 237 | return { |
| 238 | t: T_ERRAND, |
| 239 | v: ENVELOPE_V, |
| 240 | eid: o.eid || newId(), |
| 241 | turnId: String(o.turnId || ''), |
| 242 | chatId: String(o.chatId || ''), |
| 243 | prompt: String(o.prompt == null ? '' : o.prompt), |
| 244 | model: o.model || null, // { provider, model, url } -- models.js:127-129 |
| 245 | scope: o.scope || null, // the workspace fence -- scopeChatTo, daimond.js:17553 |
| 246 | pause: o.pause || null, // pause-tree snapshot at dispatch -- §1.1 |
| 247 | parcelVersion: o.parcelVersion | 0, // the freshness anchor -- §1.3 |
| 248 | deadline: +o.deadline || 0, // epoch-ms after which no peer starts (NOT |0: ms overflows 32 bits) |
| 249 | dispatchedBy: String(o.dispatchedBy || ''), |
| 250 | // How many times THIS turn has already parked (below MAX_PARKS by |
| 251 | // construction -- a dispatch is refused at the bound). Rides the errand and |
| 252 | // is carried forward through the parked report and the synced dispatched |
| 253 | // placeholder, so the ≤MAX_PARKS bound is GLOBAL across devices rather than |
| 254 | // a per-device count that would multiply the spend cap by device count. |
| 255 | parkCount: o.parkCount | 0, |
| 256 | ts: o.ts || Date.now(), |
| 257 | }; |
| 258 | } |
| 259 | |
| 260 | /// Build a report envelope object (not yet sealed). §4.4: the report is the |
| 261 | /// NUDGE, never the answer -- the answer is already in the parcel. |
| 262 | function makeReport(f) { |
| 263 | var o = f || {}; |
| 264 | return { |
| 265 | t: T_REPORT, |
| 266 | v: ENVELOPE_V, |
| 267 | eid: String(o.eid || ''), |
| 268 | turnId: String(o.turnId || ''), |
| 269 | chatId: String(o.chatId || ''), |
| 270 | status: String(o.status || 'done'), // done | refused-spend | error | aborted | parked |
| 271 | parcelVersion: o.parcelVersion | 0, // which version already carries the answer |
| 272 | cost: o.cost || null, |
| 273 | why: o.why ? String(o.why) : '', // a human sentence for the failure states |
| 274 | // The GLOBAL park counter carried home so a re-dispatcher (any device that |
| 275 | // collects this report) bumps from the true total, never from a device-local |
| 276 | // zero. Only meaningful on a `parked` report; 0 elsewhere. |
| 277 | parkCount: o.parkCount | 0, |
| 278 | ts: o.ts || Date.now(), |
| 279 | }; |
| 280 | } |
| 281 | |
| 282 | /// Build a consent-ask (not yet sealed): a runner's live question for a human. |
| 283 | /// `cid` names THIS question and is minted FRESH on every ask (including a |
| 284 | /// re-raise), so a captured or replayed grant for a spent `cid` matches nothing. |
| 285 | /// `detail` is the EXACT uncut string the human must authorise -- never a summary; |
| 286 | /// `dispatchedBy` is where the answer routes back to (the runner's device id). |
| 287 | function makeAsk(f) { |
| 288 | var o = f || {}; |
| 289 | return { |
| 290 | t: T_ASK, |
| 291 | v: ENVELOPE_V, |
| 292 | cid: o.cid || newId(), // names THIS question; fresh per ask |
| 293 | eid: String(o.eid || ''), |
| 294 | turnId: String(o.turnId || ''), |
| 295 | chatId: String(o.chatId || ''), |
| 296 | tool: String(o.tool || ''), |
| 297 | host: String(o.host || ''), |
| 298 | detail: String(o.detail == null ? '' : o.detail), // the uncut string to authorise |
| 299 | deadline: +o.deadline || 0, // epoch-ms (NOT |0: ms overflows 32 bits) |
| 300 | dispatchedBy: String(o.dispatchedBy || ''), // the runner, so the grant routes home |
| 301 | ts: o.ts || Date.now(), |
| 302 | }; |
| 303 | } |
| 304 | |
| 305 | /// Build a consent-grant (not yet sealed): an attended device's answer. It signs |
| 306 | /// `cid`/`turnId`/`verdict` but NOT `tool`/`host`/`detail` (they are not fields |
| 307 | /// here, so the signature cannot cover them) -- so the runner must replay the EXACT |
| 308 | /// act it bound to `cid`, never re-derive it from the grant. `verdict` is exactly |
| 309 | /// what `egressAllowed` returns, so the runner needs no translation layer. |
| 310 | function makeGrant(f) { |
| 311 | var o = f || {}; |
| 312 | return { |
| 313 | t: T_GRANT, |
| 314 | v: ENVELOPE_V, |
| 315 | cid: String(o.cid || ''), |
| 316 | eid: String(o.eid || ''), |
| 317 | turnId: String(o.turnId || ''), |
| 318 | verdict: o.verdict === 'allow' ? 'allow' : 'deny', |
| 319 | by: String(o.by || ''), // the answering device, for the UI |
| 320 | ts: o.ts || Date.now(), |
| 321 | }; |
| 322 | } |
| 323 | |
| 324 | // ── Seal, and the open only the account can do ───────────── |
| 325 | |
| 326 | /// Seal one envelope object to the account's SHARED symmetric key and answer the |
| 327 | /// post body `{ to, addr, envelope }` -- the identical shape `send` posts |
| 328 | /// (post.js:1127). `to` is the account's own public address; `addr` is the |
| 329 | /// sealed artefact's address; `envelope` is the base64 sealed bytes. |
| 330 | /// |
| 331 | /// Every device of the account derives the SAME symmetric key from the passphrase |
| 332 | /// and the account salt (the salt travels whole in the pairing bundle, identity.js |
| 333 | /// `exportBundle`), so a peer of the SAME account opens it and the gateway -- which |
| 334 | /// never holds the key -- opens nothing. This is deliberately NOT the per-device |
| 335 | /// X25519 sealing key: two siblings of one account can hold different sealing keys |
| 336 | /// (one was minted lazily after the other paired), and the old per-device seal then |
| 337 | /// dropped the errand silently, which is why a re-pair was needed to hand off. |
| 338 | async function sealForSelf(obj) { |
| 339 | if (!window.DaimondIdentity || !window.DaimondIdentity.wrapBytes) { |
| 340 | throw new Error('peer: no identity, so there is no key to seal with.'); |
| 341 | } |
| 342 | if (window.DaimondIdentity.isUnlocked && !window.DaimondIdentity.isUnlocked()) { |
| 343 | throw new Error('peer: Daimond is locked, so nothing can be sealed for a peer.'); |
| 344 | } |
| 345 | // Signed BEFORE sealing, so the signature is inside the seal and the gateway |
| 346 | // -- which cannot open the seal -- never sees author or sig. An envelope |
| 347 | // already carrying a `sig` (a re-seal) is not signed twice. |
| 348 | var signed = obj.sig ? obj : await signEnvelope(obj); |
| 349 | var plain = utf8(JSON.stringify(signed)); |
| 350 | // The tag rides in front of the AES-GCM `IV || ciphertext` so the open path |
| 351 | // tells this scheme from a legacy one without a trial decrypt; the purpose |
| 352 | // (not the tag bytes) is bound in as AAD, which domain-separates this key's |
| 353 | // uses -- a body for another purpose fails to open here, and vice versa. |
| 354 | var sealed = cat([SYM_MAGIC, await window.DaimondIdentity.wrapBytesAad(plain, PEER_AAD)]); |
| 355 | // The delivery address is the account's public key in the BASE64URL form the |
| 356 | // gateway binds an account to (identity.js:publicKeyB64url). It is NOT the hex |
| 357 | // of the raw key: the gateway looks a delivery up by the b64url string, so a |
| 358 | // hex `to` matches no account and every post 404s ("No account holds that key"). |
| 359 | var to = window.DaimondIdentity.publicKeyB64url |
| 360 | ? window.DaimondIdentity.publicKeyB64url() : ''; |
| 361 | return { |
| 362 | to: to || '', |
| 363 | addr: await addressOf(sealed), |
| 364 | envelope: b64enc(sealed), |
| 365 | }; |
| 366 | } |
| 367 | |
| 368 | /// Open a sealed peer body to its plaintext bytes, whichever scheme sealed it: |
| 369 | /// the current symmetric account key (`DPY1`), or a legacy per-device X25519 |
| 370 | /// envelope (`DPS1`, post.js) still in flight across a rollout. Throws the same |
| 371 | /// way each underlying open does, so the callers' refusal handling is unchanged. |
| 372 | /// It DECRYPTS only; `openEnvelope` is where the signature is verified. |
| 373 | async function openSealed(bytes) { |
| 374 | // The three schemes are mutually exclusive at byte 0-3, so each tag routes |
| 375 | // exactly one decrypt and a mis-tagged body fails its own scheme rather than |
| 376 | // cross-opening another's. DPY2 (current, AAD-bound) and DPY1 (the first-hour |
| 377 | // legacy, same key, no AAD) both open under this account's symmetric key; a |
| 378 | // GCM failure is "sealed under a different account key" -- named, not left as |
| 379 | // a raw OperationError -- the symmetric analogue of post.js's "not for you". |
| 380 | var sym = tagged(bytes, SYM_MAGIC) ? SYM_MAGIC : (tagged(bytes, SYM_MAGIC_1) ? SYM_MAGIC_1 : null); |
| 381 | if (sym) { |
| 382 | if (!window.DaimondIdentity || !window.DaimondIdentity.unwrapBytesAad) { |
| 383 | throw new Error('peer: no identity, so a sealed peer body cannot be opened.'); |
| 384 | } |
| 385 | var body = bytes.subarray(sym.length); |
| 386 | try { |
| 387 | return sym === SYM_MAGIC |
| 388 | ? await window.DaimondIdentity.unwrapBytesAad(body, PEER_AAD) |
| 389 | : await window.DaimondIdentity.unwrapBytes(body); // DPY1 legacy, drop next release |
| 390 | } catch (e) { |
| 391 | throw new Error('peer: this errand was not sealed to this account, so it is not for this device.'); |
| 392 | } |
| 393 | } |
| 394 | if (!window.DaimondPost || !window.DaimondPost.unseal) { |
| 395 | throw new Error('peer: the post seal is not loaded, so nothing can be opened.'); |
| 396 | } |
| 397 | return await window.DaimondPost.unseal(bytes); // DPS1 legacy X25519 |
| 398 | } |
| 399 | |
| 400 | /// Open a sealed envelope, VERIFY its signature, and answer the parsed object -- |
| 401 | /// or THROW. Two refusals, both about authorship: |
| 402 | /// |
| 403 | /// - `openSealed` refuses an envelope not sealed under this account's key -- the |
| 404 | /// same-account-can-OPEN property; |
| 405 | /// - `verifyEnvelope` refuses one this account did not SIGN -- the |
| 406 | /// same-account-WROTE-it property, which is what stops a correspondent who |
| 407 | /// knows our public sealing key from forging an errand into the box. |
| 408 | /// |
| 409 | /// A caller that wants the object without acting on it -- to inspect a rejected |
| 410 | /// forgery -- catches the throw; the collector uses `peek`/`absorb` instead. |
| 411 | async function openEnvelope(b64) { |
| 412 | var plain = await openSealed(b64dec(b64)); |
| 413 | var obj; |
| 414 | try { obj = JSON.parse(fromUtf8(plain)); } |
| 415 | catch (e) { throw new Error('peer: an opened envelope was not an errand or report.'); } |
| 416 | if (!obj || !peerType(obj.t)) { |
| 417 | throw new Error('peer: an opened envelope carried no known type tag.'); |
| 418 | } |
| 419 | if (!(await verifyEnvelope(obj))) { |
| 420 | throw new Error('peer: an opened envelope was not signed by this account, so it is refused.'); |
| 421 | } |
| 422 | return obj; |
| 423 | } |
| 424 | |
| 425 | /// Classify a collected row's sealed body WITHOUT verifying or throwing: unseal |
| 426 | /// and parse, answer the object if it is a peer envelope (`t` in errand/report), |
| 427 | /// or null for everything else -- a message artefact (not JSON), a row this |
| 428 | /// device cannot open, or a shape with no peer tag. This is the cheap peek |
| 429 | /// `takeRow` does before the message read: a null falls straight through to the |
| 430 | /// message path unchanged. Verification is deferred to `absorb`, so the peek |
| 431 | /// stays a classify and nothing more. |
| 432 | async function peek(b64) { |
| 433 | var obj; |
| 434 | try { obj = JSON.parse(fromUtf8(await openSealed(b64dec(b64)))); } |
| 435 | catch (e) { return null; } |
| 436 | if (obj && peerType(obj.t)) return obj; |
| 437 | return null; |
| 438 | } |
| 439 | |
| 440 | // The registered runners, set by daimond.js (step 4/5). Absent here, `absorb` |
| 441 | // verifies and drops -- routing without a runner is a no-op, not a crash. |
| 442 | var _onErrand = null; |
| 443 | var _onReport = null; |
| 444 | var _onAsk = null; // a runner's live consent question, raised on an attended device |
| 445 | var _onGrant = null; // an attended device's answer, delivered to the awaiting runner |
| 446 | function onErrand(fn) { _onErrand = fn; } |
| 447 | function onReport(fn) { _onReport = fn; } |
| 448 | function onAsk(fn) { _onAsk = fn; } |
| 449 | function onGrant(fn) { _onGrant = fn; } |
| 450 | |
| 451 | /// Verify a peeked envelope and, if it was authored by this account, hand it to |
| 452 | /// the registered runner. An envelope that does not verify is DROPPED with a |
| 453 | /// note, never run -- that is the forged-errand defence, applied at the one door |
| 454 | /// the collector routes through. Answers `{ routed, verified }`. |
| 455 | async function absorb(obj, row) { |
| 456 | var verified = await verifyEnvelope(obj); |
| 457 | if (!verified) { |
| 458 | if (window.console) console.log('peer: a ' + obj.t + ' failed signature check; dropped.'); |
| 459 | return { routed: false, verified: false }; |
| 460 | } |
| 461 | // The errand runner's answer is propagated so takeRow can read a stand-down: |
| 462 | // a non-nominee that deferred to the awake nominee (`why:'nominee'`) must leave |
| 463 | // the errand on the relay (HOLD), not ack it away before the nominee collects. |
| 464 | var result = null; |
| 465 | if (obj.t === T_ERRAND && _onErrand) result = await _onErrand(obj, row); |
| 466 | else if (obj.t === T_REPORT && _onReport) await _onReport(obj, row); |
| 467 | else if (obj.t === T_ASK && _onAsk) await _onAsk(obj, row); |
| 468 | else if (obj.t === T_GRANT && _onGrant) await _onGrant(obj, row); |
| 469 | return { routed: true, verified: true, result: result }; |
| 470 | } |
| 471 | |
| 472 | /// Route the post box's rows the way a collector does: open each, dispatch by |
| 473 | /// the sealed `t` tag. A row that is not ours to open, or is an ordinary |
| 474 | /// message, is handed to `onOther` rather than dropped -- the collector still |
| 475 | /// owes it to the message list. Answers a small tally, so a caller adds rather |
| 476 | /// than branches. |
| 477 | /// |
| 478 | /// This mirrors `takeRow`'s routing (post.js:1322), and is retained as the |
| 479 | /// direct-drive door a test uses; the real collect path goes through |
| 480 | /// `takeRow` -> `peek` -> `absorb` (post.js, step 2). Unlike `absorb`, this |
| 481 | /// verifies via `openEnvelope` (which throws on a bad signature), so a forgery |
| 482 | /// lands in `onOther`. |
| 483 | async function routeRows(rows, handlers) { |
| 484 | var h = handlers || {}; |
| 485 | var tally = { errands: 0, reports: 0, other: 0, unopened: 0 }; |
| 486 | var list = rows || []; |
| 487 | for (var i = 0; i < list.length; i++) { |
| 488 | var row = list[i]; |
| 489 | var obj = null; |
| 490 | try { obj = await openEnvelope(row.envelope); } |
| 491 | catch (e) { |
| 492 | // Not ours, not JSON, or an ordinary message: hand it on untouched. |
| 493 | tally.other++; |
| 494 | if (h.onOther) await h.onOther(row, e); |
| 495 | continue; |
| 496 | } |
| 497 | if (obj.t === T_ERRAND) { |
| 498 | tally.errands++; |
| 499 | if (h.onErrand) await h.onErrand(obj, row); |
| 500 | } else if (obj.t === T_REPORT) { |
| 501 | tally.reports++; |
| 502 | if (h.onReport) await h.onReport(obj, row); |
| 503 | } else if (obj.t === T_ASK) { |
| 504 | tally.asks = (tally.asks | 0) + 1; |
| 505 | if (h.onAsk) await h.onAsk(obj, row); |
| 506 | } else if (obj.t === T_GRANT) { |
| 507 | tally.grants = (tally.grants | 0) + 1; |
| 508 | if (h.onGrant) await h.onGrant(obj, row); |
| 509 | } |
| 510 | } |
| 511 | return tally; |
| 512 | } |
| 513 | |
| 514 | // ── The one fold the result needs ────────────────────────── |
| 515 | |
| 516 | /// Fold a peer's answer into a chat's transcript as an APPEND. §3.1: an errand |
| 517 | /// result is a new assistant message on an existing chat -- a pure append -- so |
| 518 | /// the parcel's append-only union takes it with no new rule. The message |
| 519 | /// carries `iturn` = the turn id, which is what lets the phone's own tombstone |
| 520 | /// path (§2.6, daimond.js:11796,11808) displace its "dispatched" placeholder |
| 521 | /// rather than sit a duplicate beside the answer. |
| 522 | /// |
| 523 | /// Mints a fresh `mid` the same shape daimond.js mints (`newMid`, |
| 524 | /// daimond.js:968): time-in-base36 plus a random tail, so the union keys on it |
| 525 | /// and never duplicates the message across a re-pull. |
| 526 | function foldAssistant(chat, f) { |
| 527 | var o = f || {}; |
| 528 | if (!chat.messages) chat.messages = []; |
| 529 | var msg = { |
| 530 | mid: o.mid || (Date.now().toString(36) + '-' + Math.random().toString(36).slice(2, 9)), |
| 531 | role: 'assistant', |
| 532 | content: String(o.text == null ? '' : o.text), |
| 533 | iturn: String(o.turnId || ''), |
| 534 | model: o.model || null, |
| 535 | ts: o.ts || Date.now(), |
| 536 | }; |
| 537 | chat.messages.push(msg); |
| 538 | chat.updatedAt = msg.ts; |
| 539 | return msg; |
| 540 | } |
| 541 | |
| 542 | // ── The dispatcher (dev/PEER_DESIGN.md §4.1) ─────────────── |
| 543 | // |
| 544 | // STEP 4. The phone can only dispatch while awake -- sealing, signing and |
| 545 | // posting all need the JS context running. `buildDispatch` is the PURE core: it |
| 546 | // assembles the full errand and fixes the STRICT ORDER, and daimond.js does only |
| 547 | // the thin wiring that runs that order. The order is load-bearing (§4.1): the |
| 548 | // prompt parcel is pushed FIRST so a peer can never claim an errand whose prompt |
| 549 | // it cannot yet read, the local turn is marked peer-held SECOND, and the errand |
| 550 | // is posted LAST, carrying the version the prompt push returned. |
| 551 | |
| 552 | var DISPATCH_DEADLINE_MS = 15 * 60 * 1000; // no peer should start a turn older than this |
| 553 | var REASON_DISPATCHED = 'dispatched'; // the interrupted-reason, beside 'offline'/'unloading' |
| 554 | |
| 555 | // The ordered step tags. The SEQUENCE is the safety property, so it is data a |
| 556 | // test can assert, not just the shape of the wiring. |
| 557 | var STEP_PUSH_PROMPT = 'push-prompt'; // push the prompt parcel, capture parcelVersion |
| 558 | var STEP_MARK_DISPATCH = 'mark-dispatched'; // mark the local turn why:'dispatched' |
| 559 | var STEP_POST_ERRAND = 'post-errand'; // seal + post the errand carrying that version |
| 560 | |
| 561 | /// Assemble a dispatch. PURE: it reads `chat` and the raw materials in `opts` |
| 562 | /// (already gathered by daimond.js -- the turn id, the prompt, the scope from |
| 563 | /// `scopeChatTo`, a `DaimondPause.snapshot()` pinned to this moment, the chat's |
| 564 | /// model, this device's id) and answers the ordered plan plus a `errand(version)` |
| 565 | /// finaliser. It does NO I/O, so the order and the whole envelope are testable |
| 566 | /// without the app. `parcelVersion` is NOT known until the push returns, so the |
| 567 | /// errand is finalised by the caller once it has it. |
| 568 | function buildDispatch(chat, opts) { |
| 569 | var o = opts || {}, c = chat || {}; |
| 570 | var now = o.now || Date.now(); |
| 571 | var turnId = String(o.turnId || ''); |
| 572 | var chatId = String(o.chatId || c.id || ''); |
| 573 | var eid = o.eid || newId(); |
| 574 | var prompt = String(o.prompt == null ? '' : o.prompt); |
| 575 | // The peer must run the model the CHAT chose, not the peer's own default. |
| 576 | var model = o.model || { provider: c.provider || '', model: c.model || '', url: String(o.url || '') }; |
| 577 | var scope = o.scope || null; // the workspace fence -- scopeChatTo, daimond.js:17553 |
| 578 | var pause = o.pause || null; // pause-tree snapshot pinned to dispatch -- §1.1 |
| 579 | var deadline = leaseMs(o.deadline) || (now + DISPATCH_DEADLINE_MS); |
| 580 | var by = String(o.dispatchedBy || ''); |
| 581 | // The prior-park total this dispatch inherits (0 on a first dispatch). A |
| 582 | // re-dispatch of a parked turn carries the GLOBAL count read from the synced |
| 583 | // placeholder / parked report, so the ≤MAX_PARKS bound holds across devices. |
| 584 | var parkCount = o.parkCount | 0; |
| 585 | return { |
| 586 | order: [STEP_PUSH_PROMPT, STEP_MARK_DISPATCH, STEP_POST_ERRAND], |
| 587 | turnId: turnId, chatId: chatId, eid: eid, |
| 588 | // What daimond.js writes on the local turn BETWEEN the push and the post, |
| 589 | // so recoverInterrupted and Continue treat it as peer-held, not a local |
| 590 | // interruption (§3.3). The runner/guards consult it in step 6. |
| 591 | mark: { interrupted: true, why: REASON_DISPATCHED, iturn: turnId, itext: prompt, dispatchedBy: by, parkCount: parkCount }, |
| 592 | /// Finalise the errand once the prompt push has returned its version. |
| 593 | errand: function (parcelVersion) { |
| 594 | return makeErrand({ |
| 595 | eid: eid, turnId: turnId, chatId: chatId, prompt: prompt, model: model, |
| 596 | scope: scope, pause: pause, parcelVersion: parcelVersion, |
| 597 | deadline: deadline, dispatchedBy: by, parkCount: parkCount, ts: now, |
| 598 | }); |
| 599 | }, |
| 600 | // The fully-resolved fields (bar parcelVersion), exposed for inspection. |
| 601 | fields: { |
| 602 | turnId: turnId, chatId: chatId, prompt: prompt, model: model, scope: scope, |
| 603 | pause: pause, deadline: deadline, dispatchedBy: by, eid: eid, parkCount: parkCount, |
| 604 | }, |
| 605 | }; |
| 606 | } |
| 607 | |
| 608 | /// What a turn marked `why:'dispatched'` should be treated as, given the lease. |
| 609 | /// The pure decision recoverInterrupted and the Continue button consult (§3.3, |
| 610 | /// wired in step 6): |
| 611 | /// - `not-dispatched` -> an ordinary turn, ordinary recovery; |
| 612 | /// - `peer-held` -> a LIVE FOREIGN lease holds it: do NOT recover locally, |
| 613 | /// show "running on your other device"; |
| 614 | /// - `reclaimable` -> the lease is vacant/expired or ours: recover locally / |
| 615 | /// offer Continue. |
| 616 | function dispatchState(turn, leaseRec, selfId, now) { |
| 617 | if (!turn || turn.why !== REASON_DISPATCHED) return 'not-dispatched'; |
| 618 | var n = now == null ? Date.now() : now; |
| 619 | if (liveLease(leaseRec, n) && leaseRec.holder !== String(selfId)) return 'peer-held'; |
| 620 | return 'reclaimable'; |
| 621 | } |
| 622 | |
| 623 | /// The §5 DISPLAY state of a dispatched turn, for the phone's UI. Pure, so the |
| 624 | /// renderer only draws what this decides. `turn` is the local dispatched turn, |
| 625 | /// `lease` its lease record (or null), `report` a collected report envelope for |
| 626 | /// it (or null), `selfId` the viewing device. |
| 627 | /// dispatched -- posted, no lease seen yet: "Sent to your other devices." |
| 628 | /// no-peer-awake -- deadline passed, no lease ever taken: "No awake device…" |
| 629 | /// claimed -- lease mode 'claimed': "<machine> is picking this up." |
| 630 | /// running -- lease mode 'running': "<machine> is doing this. [Take back]" |
| 631 | /// done -- report status 'done': the answer, the badge clears |
| 632 | /// parked -- report status 'parked': "needs your permission — it will |
| 633 | /// re-run when you're back" (a survivable park, below the bound) |
| 634 | /// awaiting-consent -- an open consent-ask for this turn: the runner is blocked |
| 635 | /// on a live question the user must answer -- "<peer> needs your |
| 636 | /// permission to {act}", replacing "Sent to your other devices." |
| 637 | /// failed -- report a failure (a terminal park included), OR the lease |
| 638 | /// expired mid-run with no report (the peer stopped): the `why` |
| 639 | /// sentence + [Run here] |
| 640 | /// `ask` is a collected consent-ask envelope for this turn, or null. |
| 641 | function uiState(turn, lease, report, selfId, now, ask) { |
| 642 | if (!turn || turn.why !== REASON_DISPATCHED) return 'not-dispatched'; |
| 643 | var n = now == null ? Date.now() : now; |
| 644 | // A report settles it either way, and outlives the lease. |
| 645 | if (report && report.t === 'report') { |
| 646 | if (report.status === 'done') return 'done'; |
| 647 | if (report.status === 'parked') return 'parked'; // survivable: re-runs when a human is back |
| 648 | return 'failed'; // aborted / error / refused-spend: terminal |
| 649 | } |
| 650 | // A live question the runner is blocked on takes precedence over the lease |
| 651 | // state (which reads 'running' throughout the wait): the dispatch UI must stop |
| 652 | // saying "sent" and say what is actually holding it up. |
| 653 | if (ask && ask.t === T_ASK) return 'awaiting-consent'; |
| 654 | // No report yet: read the lease. |
| 655 | if (liveLease(lease, n)) { |
| 656 | return lease.mode === 'running' ? 'running' : 'claimed'; |
| 657 | } |
| 658 | // A lease that was taken then expired without a report is a peer that |
| 659 | // stopped mid-run (§6): failed, offer a local re-run. |
| 660 | if (lease && lease.mode !== 'released') return 'failed'; |
| 661 | // No live lease and none reported: waiting, or nobody picked it up in time. |
| 662 | var deadline = leaseMs(turn.deadline); |
| 663 | if (deadline && n > deadline) return 'no-peer-awake'; |
| 664 | return 'dispatched'; |
| 665 | } |
| 666 | |
| 667 | /// Should the DISPATCHING device RECOVER this turn locally now (on its return to |
| 668 | /// the foreground)? Pure. A backgrounded phone cannot run the deadline fallback, |
| 669 | /// so on return it must run any turn that was dispatched but that NO peer ran, or |
| 670 | /// the user comes back to nothing -- the reported "complete and utter failure". |
| 671 | /// |
| 672 | /// Recover when the turn is `dispatched`, is NOT finished (no done report, no |
| 673 | /// merged answer -- the app supplies this as `finished`), and is NOT held by a |
| 674 | /// LIVE FOREIGN lease. A live foreign lease means a peer IS on it: leave it, and |
| 675 | /// let the answer sync back (or the user take it back by hand). An expired lease, |
| 676 | /// no lease, or our own lease is reclaimable. This only decides whether to TRY; |
| 677 | /// the take-if-vacant lease is the money-safe arbiter at run time, so a peer that |
| 678 | /// claims between this decision and the local take still wins (and vice-versa). |
| 679 | function recoverDecision(turn, lease, finished, selfId, now) { |
| 680 | if (!turn || turn.why !== REASON_DISPATCHED) return false; |
| 681 | if (finished) return false; |
| 682 | var n = now == null ? Date.now() : now; |
| 683 | if (liveLease(lease, n) && lease.holder !== String(selfId)) return false; // a peer is on it |
| 684 | return true; |
| 685 | } |
| 686 | |
| 687 | // ── Presence (dev/PEER_DESIGN.md §4.2, §7 step 7) ────────── |
| 688 | // |
| 689 | // Each AWAKE, visible Daimond writes a heartbeat into the parcel so the phone |
| 690 | // knows which of its devices could take a turn -- and can NAME the machine |
| 691 | // ("waiting for argonaut"). Unlike the lease, presence is genuinely |
| 692 | // last-writer-wins: the freshest `lastSeen` per device is the truth, so it uses |
| 693 | // the FRESHEST-SCALAR merge (like the pause tree), NOT take-if-vacant. A stale |
| 694 | // beat is safe: the deadline and the lease catch a peer that actually slept, so |
| 695 | // the worst a stale beat costs is one dispatch that finds no runner and falls to |
| 696 | // `no-peer-awake` -- never a double run. |
| 697 | |
| 698 | var PRESENCE_BEAT_MS = 45000; // write a beat about this often while visible |
| 699 | var PRESENCE_FRESH_MS = 120000; // a beat older than this is not "awake" (≈ 2 min) |
| 700 | // The window the auto-dispatch DECISION uses -- deliberately TIGHTER than the |
| 701 | // display window above, and well under the gateway's 5-min presence TTL. A beat |
| 702 | // is written every 45 s, so two beats plus slack is a peer that is genuinely |
| 703 | // still beating; a peer last seen longer ago than this is treated as gone and |
| 704 | // the turn runs locally, rather than dispatched to a device that may have died |
| 705 | // since its last beat. The display can afford to name a peer as "awake" for |
| 706 | // longer; a DISPATCH cannot, because a dispatch into a dead peer is an orphan |
| 707 | // (the recovery-on-return catches it, but the tighter window avoids most). |
| 708 | var DISPATCH_FRESH_MS = 90000; // a peer beat older than this is not dispatched to (≈ 1.5 min) |
| 709 | |
| 710 | var _presence = {}; // deviceId -> { name, lastSeen } |
| 711 | |
| 712 | /// The freshest peer in a presence map that is NOT this device and whose beat is |
| 713 | /// within the window, or null. The one pure helper both the UI and the |
| 714 | /// auto-dispatch decision read, so "which peer is awake" has ONE answer. |
| 715 | function freshestPeer(presence, selfId, now, windowMs) { |
| 716 | var p = presence || {}, self = String(selfId || ''), n = now == null ? Date.now() : now; |
| 717 | var w = windowMs || PRESENCE_FRESH_MS, best = null; |
| 718 | for (var id in p) { |
| 719 | if (!Object.prototype.hasOwnProperty.call(p, id)) continue; |
| 720 | if (id === self) continue; |
| 721 | var rec = p[id]; |
| 722 | if (!rec || (n - leaseMs(rec.lastSeen)) > w) continue; // stale |
| 723 | if (!best || leaseMs(rec.lastSeen) > leaseMs(best.lastSeen)) { |
| 724 | best = { deviceId: id, name: (rec.name || ''), lastSeen: leaseMs(rec.lastSeen) }; |
| 725 | } |
| 726 | } |
| 727 | return best; |
| 728 | } |
| 729 | |
| 730 | /// Record this device's heartbeat. Answers whether the map changed (it always |
| 731 | /// does -- lastSeen moved -- which is what makes the beat a push). `attended` is |
| 732 | /// the attention signal (foreground + recent interaction) a live consent routes on: |
| 733 | /// `attendedAt` stamps when the device was last attended, so freshness is judged on |
| 734 | /// attention rather than on the bare beat. |
| 735 | function presenceBeat(deviceId, name, now, attended) { |
| 736 | var id = String(deviceId || ''); |
| 737 | if (!id) return false; |
| 738 | var n = now == null ? Date.now() : now; |
| 739 | var prev = _presence[id]; |
| 740 | var at = attended ? n : (prev ? leaseMs(prev.attendedAt) : 0); |
| 741 | _presence[id] = { name: String(name || ''), lastSeen: n, attended: !!attended, attendedAt: at }; |
| 742 | return true; |
| 743 | } |
| 744 | |
| 745 | /// Merge an arriving parcel's presence FRESHEST-SCALAR: the larger `lastSeen` |
| 746 | /// per device wins. Called by sync.js's reconcile. Answers whether anything |
| 747 | /// moved, so a pull that learned nothing schedules no push. |
| 748 | function presenceAdopt(incoming) { |
| 749 | if (!incoming) return false; |
| 750 | var moved = false; |
| 751 | for (var id in incoming) { |
| 752 | if (!Object.prototype.hasOwnProperty.call(incoming, id)) continue; |
| 753 | var inc = incoming[id]; |
| 754 | if (!inc) continue; |
| 755 | var cur = _presence[id]; |
| 756 | if (!cur || leaseMs(inc.lastSeen) > leaseMs(cur.lastSeen)) { |
| 757 | _presence[id] = { |
| 758 | name: String(inc.name || ''), lastSeen: leaseMs(inc.lastSeen), |
| 759 | attended: !!inc.attended, attendedAt: leaseMs(inc.attendedAt), |
| 760 | }; |
| 761 | moved = true; |
| 762 | } |
| 763 | } |
| 764 | return moved; |
| 765 | } |
| 766 | |
| 767 | /// Ingest an AUTHORITATIVE presence map from the gateway and REPLACE the local |
| 768 | /// view with it. The gateway is now the source of truth -- presence travels on |
| 769 | /// its own lightweight, non-waking path, not on the content parcel -- so this is |
| 770 | /// a replace, not the freshest-scalar merge `presenceAdopt` does. |
| 771 | /// |
| 772 | /// `serverMap` is deviceId -> { name, last_seen } with `last_seen` and |
| 773 | /// `serverNow` both stamped in the SERVER clock. Each `last_seen` is converted |
| 774 | /// into THIS client's frame -- `last_seen - (serverNow - now_at_receipt)` -- so |
| 775 | /// every existing freshness check that reads `Date.now()` (`awake`, |
| 776 | /// `freshestPeer`) keeps working unchanged and is immune to cross-device clock |
| 777 | /// skew. Answers whether the view moved, so a caller can skip a redraw that |
| 778 | /// learned nothing. |
| 779 | function presenceIngest(serverMap, serverNow) { |
| 780 | var recv = Date.now(); |
| 781 | var skew = leaseMs(serverNow) - recv; // how far the server clock leads ours |
| 782 | var next = {}, map = serverMap || {}; |
| 783 | for (var id in map) { |
| 784 | if (!Object.prototype.hasOwnProperty.call(map, id)) continue; |
| 785 | var rec = map[id]; |
| 786 | if (!rec) continue; |
| 787 | // The wire field is `last_seen`; tolerate `lastSeen` in case a caller |
| 788 | // hands an already-client-framed record straight in. |
| 789 | var seen = leaseMs(rec.last_seen != null ? rec.last_seen : rec.lastSeen); |
| 790 | // The attention signal a live consent routes on. The gateway relays it as |
| 791 | // `attended` / `attended_at` (skew-adjusted like last_seen); when a build's |
| 792 | // gateway does not yet carry it, it reads absent -> not attended, so a runner |
| 793 | // PARKS rather than routing a question to a device that may be unwatched (the |
| 794 | // fail-safe the design requires). |
| 795 | var atRaw = rec.attended_at != null ? rec.attended_at : rec.attendedAt; |
| 796 | next[String(id)] = { |
| 797 | name: String(rec.name || ''), lastSeen: seen - skew, |
| 798 | attended: !!rec.attended, attendedAt: atRaw != null ? (leaseMs(atRaw) - skew) : 0, |
| 799 | }; |
| 800 | } |
| 801 | var before = JSON.stringify(_presence); |
| 802 | _presence = next; |
| 803 | return JSON.stringify(_presence) !== before; |
| 804 | } |
| 805 | |
| 806 | /// The section as it rides the parcel, or null when empty. |
| 807 | function presenceSnapshot() { |
| 808 | return Object.keys(_presence).length ? _presence : null; |
| 809 | } |
| 810 | |
| 811 | /// The awake peers (not this device), freshest first, for the UI. |
| 812 | function presenceAwake(selfId, now, windowMs) { |
| 813 | var p = _presence, self = String(selfId || ''), n = now == null ? Date.now() : now; |
| 814 | var w = windowMs || PRESENCE_FRESH_MS, out = []; |
| 815 | for (var id in p) { |
| 816 | if (!Object.prototype.hasOwnProperty.call(p, id)) continue; |
| 817 | if (id === self) continue; |
| 818 | if ((n - leaseMs(p[id].lastSeen)) > w) continue; |
| 819 | out.push({ deviceId: id, name: p[id].name || '', lastSeen: leaseMs(p[id].lastSeen) }); |
| 820 | } |
| 821 | out.sort(function (a, b) { return b.lastSeen - a.lastSeen; }); |
| 822 | return out; |
| 823 | } |
| 824 | |
| 825 | /// This device's name for a peer's deviceId (for "waiting for argonaut"), or ''. |
| 826 | function presenceName(deviceId) { |
| 827 | var r = _presence[String(deviceId || '')]; |
| 828 | return (r && r.name) || ''; |
| 829 | } |
| 830 | |
| 831 | function presenceForget() { _presence = {}; } |
| 832 | |
| 833 | window.DaimondPresence = { |
| 834 | BEAT_MS: PRESENCE_BEAT_MS, |
| 835 | FRESH_MS: PRESENCE_FRESH_MS, |
| 836 | beat: presenceBeat, |
| 837 | /// sync.js's section contract -- freshest-scalar, NOT take-if-vacant. |
| 838 | snapshot: presenceSnapshot, |
| 839 | adopt: presenceAdopt, |
| 840 | /// Replace the local view from the gateway's authoritative map, converting |
| 841 | /// each last_seen into this client's clock frame (skew-immune). This is the |
| 842 | /// sync path now -- presence rides its own non-waking gateway route, not the |
| 843 | /// content parcel -- so it supersedes the freshest-scalar `adopt` above. |
| 844 | ingest: presenceIngest, |
| 845 | /// The awake peers, and one device's name, for the UI and auto-dispatch. |
| 846 | awake: presenceAwake, |
| 847 | name: presenceName, |
| 848 | forget: presenceForget, |
| 849 | }; |
| 850 | |
| 851 | // ── Remote consent — attention, routing, the park bound ──── |
| 852 | // |
| 853 | // A turn dispatched from the phone runs on a runner where nobody is. When it hits |
| 854 | // a genuinely per-turn consent the synced account policy does not cover, the |
| 855 | // runner cannot raise a dialog into an empty room: it routes the question to a |
| 856 | // device the user is ON, and awaits the answer while holding its turn in memory. |
| 857 | // Only the question and the answer travel; the turn never leaves the runner. |
| 858 | // |
| 859 | // The helpers here are the PURE decisions -- who is attended, whether to ask or |
| 860 | // park, and how far the park loop may run. daimond.js seals/posts/awaits over the |
| 861 | // real channel; the money-safety of the bound is decided here so a test drives it. |
| 862 | |
| 863 | // The default life of one live consent question. ~1 minute, tunable: the runner |
| 864 | // holds its lease straight to the errand deadline (~15 min) throughout, so the |
| 865 | // short consent wait sits INSIDE the long claim and no lease renewal is needed. |
| 866 | var CONSENT_DEADLINE_MS = 60 * 1000; |
| 867 | |
| 868 | // The hard cap on how many times ONE turn may park-and-re-dispatch. It is BOTH the |
| 869 | // liveness cap and the SPEND cap: each re-dispatch replays the turn's pre-consent |
| 870 | // account spend (the earlier LLM calls, a completed web_search), so N re-dispatches |
| 871 | // cost at most N× that spend. Small on purpose. The count is GLOBAL (synced on the |
| 872 | // errand / parked report / placeholder), so it is not multiplied by device count. |
| 873 | var MAX_PARKS = 2; |
| 874 | |
| 875 | /// Is a presence record ATTENDED -- a person is at this device now, not merely a |
| 876 | /// beating heartbeat? Attention is a foreground + recent-interaction signal the |
| 877 | /// beat carries (`attended`), distinct from `lastSeen`: routing a live question to |
| 878 | /// an awake-but-unwatched device would just relocate the invisible stall. Absent or |
| 879 | /// false is NOT attended, so an attention-indeterminate device is never asked. |
| 880 | function recAttended(rec, now, windowMs) { |
| 881 | if (!rec || !rec.attended) return false; |
| 882 | var w = windowMs || DISPATCH_FRESH_MS; |
| 883 | var seen = leaseMs(rec.attendedAt != null ? rec.attendedAt : rec.lastSeen); |
| 884 | return (now - seen) <= w; |
| 885 | } |
| 886 | |
| 887 | /// The freshest ATTENDED peer (not this device) in a presence map, or null. The one |
| 888 | /// pure answer to "is there a device the user is on that a live question can go to". |
| 889 | /// Attention fails SAFE: an indeterminate or stale-attention device is skipped, so |
| 890 | /// the caller parks rather than routing a question nobody will see. |
| 891 | function attendedPeer(presence, selfId, now, windowMs) { |
| 892 | var p = presence || {}, self = String(selfId || ''); |
| 893 | var n = now == null ? Date.now() : now, best = null; |
| 894 | for (var id in p) { |
| 895 | if (!Object.prototype.hasOwnProperty.call(p, id)) continue; |
| 896 | if (id === self) continue; |
| 897 | var rec = p[id]; |
| 898 | if (!recAttended(rec, n, windowMs)) continue; |
| 899 | if (!best || leaseMs(rec.lastSeen) > leaseMs(best.lastSeen)) { |
| 900 | best = { deviceId: id, name: (rec.name || ''), lastSeen: leaseMs(rec.lastSeen) }; |
| 901 | } |
| 902 | } |
| 903 | return best; |
| 904 | } |
| 905 | |
| 906 | /// Decide what a runner does with a per-turn consent it cannot answer locally: |
| 907 | /// ASK a specific attended peer, or PARK. Pure. `covered` is the policy-sync |
| 908 | /// short-circuit -- a standing account grant the runner already holds resolves the |
| 909 | /// act with no question at all, so the ask fires ONLY when the synced policy does |
| 910 | /// not already cover it (this composes with the shipped consent-sync, never |
| 911 | /// double-asking). Answers `{ action, peer?, verdict?, why? }`. |
| 912 | function consentRouteDecision(presence, selfId, now, opts) { |
| 913 | var o = opts || {}; |
| 914 | if (o.covered) return { action: 'allow', verdict: 'allow', why: 'policy' }; |
| 915 | var peer = attendedPeer(presence, selfId, now, o.windowMs); |
| 916 | if (peer) return { action: 'ask', peer: peer }; |
| 917 | return { action: 'park', why: 'no-attended-device' }; |
| 918 | } |
| 919 | |
| 920 | /// The outcome of a park, given the parkCount the errand carried. Pure and the |
| 921 | /// single arbiter of the spend bound. `next` is the new GLOBAL park total (this |
| 922 | /// park included); `terminal` is true once it reaches the bound, at which point the |
| 923 | /// turn fails clean rather than re-dispatching into another respend. |
| 924 | function parkOutcome(errandParkCount, maxParks) { |
| 925 | var mx = (maxParks == null) ? MAX_PARKS : (maxParks | 0); |
| 926 | var next = (errandParkCount | 0) + 1; |
| 927 | return { next: next, terminal: next >= mx }; |
| 928 | } |
| 929 | |
| 930 | // ── Smart auto-dispatch (dev/PEER_DESIGN.md §4.1, §4.2) ──── |
| 931 | // |
| 932 | // The pure decision: given the chat, the presence map and the moment, should |
| 933 | // this turn be handed to a peer, and which one? daimond.js only ACTS on it, at |
| 934 | // send-time while online, through the already-proven ordered dispatcher. The |
| 935 | // rules, in order: |
| 936 | // - NO fresh peer -> run locally, never dispatch into the void; |
| 937 | // - a per-chat OPT-OUT -> keep this chat on THIS device (toggle === false); |
| 938 | // - the toggle is on -> hand it off (blanket-when-awake for this chat); |
| 939 | // - MOBILE with a peer awake -> hand EVERY turn off (the phone is not where a |
| 940 | // turn should run when a persistent peer exists; |
| 941 | // sync brings the answer back). Desktop falls |
| 942 | // through -- it IS the persistent instance; |
| 943 | // - backgrounding in flight -> hand it off (the phone is about to sleep); |
| 944 | // - a long/agentic turn -> hand it off (tools/worker/expected-long); |
| 945 | // - otherwise (quick turn) -> run locally, for instant streaming. |
| 946 | // |
| 947 | // The order is money-safe by construction: the NO-fresh-peer guard is first, so |
| 948 | // nothing ever dispatches into the void, and the opt-out precedes every "hand it |
| 949 | // off" rule, so a chat pinned local cannot be routed away by the mobile default. |
| 950 | // Broadening WHICH turns route changes nothing about HOW routing works -- the |
| 951 | // single-runner guarantee is the lease's (below), never this decision's. |
| 952 | |
| 953 | /// Decide whether to dispatch, and to whom. Pure. Answers |
| 954 | /// `{ dispatch, peer, reason }`. |
| 955 | function autoDispatchDecision(chat, presence, opts, now) { |
| 956 | var o = opts || {}, c = chat || {}; |
| 957 | var peer = freshestPeer(presence, o.selfId, now, o.freshWindowMs); |
| 958 | if (!peer) return { dispatch: false, reason: 'no-fresh-peer' }; |
| 959 | // The per-chat choice: true = always hand off, false = keep on THIS device |
| 960 | // (the opt-out), null/undefined = decide by the policy below. |
| 961 | var toggle = (o.toggle != null) ? !!o.toggle : null; |
| 962 | // An explicit opt-out pins the chat here, even on a phone with a peer awake -- |
| 963 | // so it must be tested BEFORE the mobile default and before the global one. |
| 964 | if (toggle === false) return { dispatch: false, reason: 'chat-local' }; |
| 965 | // The toggle on, or the global default when the chat has not chosen. |
| 966 | if (toggle === true || o.globalDefault) return { dispatch: true, peer: peer, reason: 'toggle-on' }; |
| 967 | // MOBILE: a persistent peer is awake, so hand EVERY turn off -- quick or |
| 968 | // agentic -- rather than run it on the phone; the answer syncs back. Desktop |
| 969 | // (no isPhone) falls through: you are already on the persistent instance. |
| 970 | if (o.isPhone) return { dispatch: true, peer: peer, reason: 'mobile-peer' }; |
| 971 | // The phone is backgrounding with a turn still in flight -- move the running |
| 972 | // off it before it suspends. |
| 973 | if (o.backgrounding && o.turnInFlight) return { dispatch: true, peer: peer, reason: 'backgrounding-in-flight' }; |
| 974 | // A long or agentic turn is worth the round trip; a quick one is not. The |
| 975 | // worker signal must be a GENUINE one: daimond.js seeds `workerModel`/ |
| 976 | // `workerProvider` to the chat's OWN model for every active chat (newChat, |
| 977 | // startChat), so `c.workerModel` is truthy on a plain chat and a bare |
| 978 | // truthiness test dispatched EVERY turn (D2). A worker/Diamond chat is one |
| 979 | // whose worker pair DIFFERS from the chat's own -- the user deliberately chose |
| 980 | // a different model to fan work out to; a pair merely mirroring the default is |
| 981 | // not agentic and stays local for instant streaming. |
| 982 | var worker = (c.workerModel && String(c.workerModel) !== String(c.model || '')) |
| 983 | || (c.workerProvider && String(c.workerProvider) !== String(c.provider || '')); |
| 984 | var agentic = !!o.toolsEnabled || !!o.expectedLong || !!worker; |
| 985 | if (agentic) return { dispatch: true, peer: peer, reason: 'long-turn' }; |
| 986 | return { dispatch: false, reason: 'quick-local' }; |
| 987 | } |
| 988 | |
| 989 | // ── The nominated always-on runner (the claim guard) ─────── |
| 990 | // |
| 991 | // An account may name ONE device as the runner that should pick a dispatched |
| 992 | // turn up, so a laptop someone closes mid-turn does not race in and grab it. |
| 993 | // This decides only WHO attempts the lease claim; the take-if-vacant CAS below |
| 994 | // is still the single-runner arbiter, so a stale or racing decision here can at |
| 995 | // worst cost one extra claim attempt, never a double run. |
| 996 | |
| 997 | /// Should this device STAND DOWN from claiming a dispatched errand, deferring to |
| 998 | /// the account's NOMINATED always-on runner? Pure. True only when a nominee is |
| 999 | /// set, this device is NOT it, AND the nominee is presently, FRESHLY awake in the |
| 1000 | /// presence map -- judged by the SAME tight window a dispatch uses to call a peer |
| 1001 | /// awake (DISPATCH_FRESH_MS). A nominee whose last beat has aged out of that |
| 1002 | /// window is treated as OFFLINE, so this device claims (the owner's fallback: any |
| 1003 | /// awake device, over nobody-runs-it). |
| 1004 | /// |
| 1005 | /// Two stalls this guards against, both by construction: |
| 1006 | /// - STALE PRESENCE: a lagging map still showing a slept nominee as awake would |
| 1007 | /// have every fallback stand down for a device that is gone. The freshness |
| 1008 | /// bound is the defence -- an aged beat reads offline and the fallback claims. |
| 1009 | /// - A PERMANENT stand-down: the decision reads LIVE presence and is re-taken on |
| 1010 | /// every re-collect (the errand is HELD on the relay, not acked, while standing |
| 1011 | /// down -- post.js), so a nominee that slept just after its last beat stops |
| 1012 | /// being "fresh" within one window and a fallback then claims. The worst-case |
| 1013 | /// stall is one presence-sync lag plus DISPATCH_FRESH_MS. |
| 1014 | function nominationStandDown(nominatedId, selfId, presence, now, windowMs) { |
| 1015 | var nom = String(nominatedId || ''); |
| 1016 | if (!nom) return false; // no nomination -> first-come, unchanged |
| 1017 | if (nom === String(selfId || '')) return false; // this device IS the nominee -> claim |
| 1018 | var rec = (presence || {})[nom]; |
| 1019 | if (!rec) return false; // nominee absent from presence -> offline |
| 1020 | var n = now == null ? Date.now() : now; |
| 1021 | var w = windowMs || DISPATCH_FRESH_MS; |
| 1022 | return (n - leaseMs(rec.lastSeen)) <= w; // stand down only while the nominee is FRESH |
| 1023 | } |
| 1024 | |
| 1025 | // ════════════════════════════════════════════════════════════ |
| 1026 | // THE LEASE — the cross-device claim (dev/PEER_DESIGN.md §2). |
| 1027 | // ------------------------------------------------------------ |
| 1028 | // The MONEY-CRITICAL step. `withTurnLock` (daimond.js:17209) is |
| 1029 | // per-origin per-browser: it stops two TABS of one browser running a |
| 1030 | // turn twice, and says nothing about two DEVICES. Two awake desktops |
| 1031 | // would each collect the errand and each bill it. The lease is the |
| 1032 | // cross-device layer above the browser lock, and a wrong version of it |
| 1033 | // is a double bill. |
| 1034 | // |
| 1035 | // It lives in the sync parcel under a `leases` section keyed by turnId, |
| 1036 | // arbitrated by the parcel's compare-and-set. The section is the ONE |
| 1037 | // part of the parcel that is NOT append-only union / freshest-scalar: |
| 1038 | // a lease is a mutable claim, and freshest-scalar is exactly the |
| 1039 | // double-claim bug -- two devices claiming from one base version would |
| 1040 | // each write a lease, and last-write-by-clock could hand it to whichever |
| 1041 | // clock read a microsecond later. So it gets its OWN merge, take-if- |
| 1042 | // vacant, invoked by sync.js through `adopt` the same way every other |
| 1043 | // section's merge is (pause.js, trash.js). With the CAS this is |
| 1044 | // first-writer-by-version-wins. |
| 1045 | // ════════════════════════════════════════════════════════════ |
| 1046 | |
| 1047 | var LEASE_TTL_MS = 90000; // a lease past this is vacant (§2.4) |
| 1048 | var RENEW_EVERY_MS = 30000; // three renews per TTL, so one dropped renew is survivable |
| 1049 | var MAX_TAKE_TRIES = 10; // bound the CAS retry loop (was 6): more headroom under two-device churn |
| 1050 | var TAKE_BACKOFF_MS = 250; // jittered wait between take retries so a claim gets a clean window |
| 1051 | // The hard ceiling on how long ONE errand's liveness ticker may run before it gives |
| 1052 | // up and aborts. The ticker is READ-ONLY (it never renews the parcel -- the lease is |
| 1053 | // claimed straight to its deadline, so no renew is needed), so it is not itself a |
| 1054 | // churn source; but a runTurn whose promise never settles would keep the ticker (and |
| 1055 | // the errand) alive indefinitely, so this caps it -- after which the run is aborted |
| 1056 | // best-effort and the lease is left to expire at its deadline. Generous, because it |
| 1057 | // is a backstop for a hung turn, not a turn budget. |
| 1058 | var MAX_LEASE_LIFE_MS = 30 * 60 * 1000; // 30 min: a turn still 'running' past this is hung, not live |
| 1059 | |
| 1060 | /// The local materialised view of the parcel's `leases` section: turnId -> |
| 1061 | /// record. Snapshotted into the parcel and merged back through `adopt`. |
| 1062 | var _leases = {}; |
| 1063 | |
| 1064 | /// Epoch-ms as a NUMBER, never `| 0`. A lease timestamp is a real |
| 1065 | /// wall-clock ms -- ~1.7e12 in 2026 -- which overflows a 32-bit `| 0` to |
| 1066 | /// garbage, so every comparison below uses this. (This is the width bug that |
| 1067 | /// `| 0` on a timestamp always is; see the i64/BigInt note in the project log.) |
| 1068 | function leaseMs(x) { return typeof x === 'number' ? x : (parseFloat(x) || 0); } |
| 1069 | |
| 1070 | /// Is a lease record a LIVE claim at `now`? A released lease, or one past its |
| 1071 | /// expiry, is vacant -- it grants nothing and may be overwritten. |
| 1072 | function liveLease(r, now) { |
| 1073 | return !!r && r.mode !== 'released' && leaseMs(r.expiry) > now; |
| 1074 | } |
| 1075 | |
| 1076 | /// The ceiling an adopted expiry may reach on the ADOPTING device's clock: one TTL |
| 1077 | /// from now, OR the errand's own `deadline` when the record carries one. A running |
| 1078 | /// turn's lease is claimed with `expiry = deadline` (see leaseTakeFrom), because a |
| 1079 | /// busy device cannot propagate a 30s renew (sync.js:1077 suppresses the push over a |
| 1080 | /// live turn), so a TTL-capped lease would read EXPIRED on other devices after 90s |
| 1081 | /// while the turn is still running -- and the phone's recovery would then re-run and |
| 1082 | /// re-bill it (the >TTL double-run). Bounding to the deadline lets the claim stay |
| 1083 | /// live for the whole turn with no renew at all. The deadline is authored by the |
| 1084 | /// DISPATCHER (buildDispatch), not the holder, so it is not a fast-clock lever. |
| 1085 | function expiryCap(r, now) { |
| 1086 | var cap = now + LEASE_TTL_MS; |
| 1087 | var dl = leaseMs(r && r.deadline); |
| 1088 | return dl > cap ? dl : cap; |
| 1089 | } |
| 1090 | |
| 1091 | /// Clamp a record's expiry to `expiryCap`. A holder with a fast clock could |
| 1092 | /// otherwise write a far-future expiry and, if it then died, park the turn for up |
| 1093 | /// to that skew (QA defect a). Every merge clamps what it keeps to the ADOPTING |
| 1094 | /// device's clock, so no foreign expiry outlives the cap here. A deadline-bounded |
| 1095 | /// lease is already <= its deadline <= cap, so this is a no-op for it; an expiry |
| 1096 | /// ABOVE the cap (a fast clock, or a lease reaching past its own deadline) is |
| 1097 | /// clamped -- the fast-clock defence is preserved, now measured against the deadline |
| 1098 | /// rather than a bare TTL. A released record (expiry 0) is untouched. Returns a copy |
| 1099 | /// only when it must change the value, so an unchanged merge stays byte-identical. |
| 1100 | function clampExpiry(r, now) { |
| 1101 | if (!r) return r; |
| 1102 | var cap = expiryCap(r, now); |
| 1103 | if (leaseMs(r.expiry) <= cap) return r; |
| 1104 | var c = {}; |
| 1105 | for (var k in r) if (Object.prototype.hasOwnProperty.call(r, k)) c[k] = r[k]; |
| 1106 | c.expiry = cap; |
| 1107 | return c; |
| 1108 | } |
| 1109 | |
| 1110 | /// Merge two records for ONE turnId under take-if-vacant. `incoming` is the |
| 1111 | /// arriving/authoritative side (a pulled parcel, or the server leases a claim |
| 1112 | /// is folded against); `local` is this device's side. |
| 1113 | /// |
| 1114 | /// The whole money-safety property is these lines, so they are spelled out |
| 1115 | /// rather than compressed: |
| 1116 | /// - SAME holder -> the fresher `renewedAt` wins, live or not. This is decided |
| 1117 | /// FIRST, and it is the ONLY place `renewedAt` is consulted, so a device's |
| 1118 | /// own renew AND its own release supersede its earlier record -- a release |
| 1119 | /// that lost to its own still-live running lease would never land; |
| 1120 | /// - different holders, INCOMING live -> incoming wins, the local fresh claim |
| 1121 | /// drops. Under the CAS only one device commits a claim at a given version, |
| 1122 | /// so the loser -- pulling the winner's blob -- meets exactly this branch and |
| 1123 | /// stands down, and it never diverges because the loser's claim was refused |
| 1124 | /// by the CAS and so never reaches the winner as an incoming; |
| 1125 | /// - different holders, only LOCAL live -> local (the incoming is dead/vacant, |
| 1126 | /// e.g. reclaiming an expired lease); |
| 1127 | /// - both vacant -> the fresher record, for history only (a dead lease grants |
| 1128 | /// nothing, so this never decides a claim). |
| 1129 | function mergeOneLease(local, incoming, now) { |
| 1130 | return clampExpiry(pickLease(local, incoming, now), now); |
| 1131 | } |
| 1132 | |
| 1133 | /// The winner of two records for one turnId, BEFORE the expiry clamp. |
| 1134 | function pickLease(local, incoming, now) { |
| 1135 | if (local && incoming && local.holder === incoming.holder) { |
| 1136 | var ri = leaseMs(incoming.renewedAt), rl = leaseMs(local.renewedAt); |
| 1137 | if (ri !== rl) return ri > rl ? incoming : local; |
| 1138 | // EQUAL renewedAt between same-holder records: a 'released' wins, so a |
| 1139 | // stale 'running' can never resurrect a lease the holder let go (QA |
| 1140 | // defect b -- unreachable through the gateway today, latent otherwise). |
| 1141 | if (incoming.mode === 'released' && local.mode !== 'released') return incoming; |
| 1142 | if (local.mode === 'released' && incoming.mode !== 'released') return local; |
| 1143 | return incoming; // truly identical: either |
| 1144 | } |
| 1145 | var lLive = liveLease(local, now), iLive = liveLease(incoming, now); |
| 1146 | if (iLive) return incoming; // different holder, incoming live: incoming wins |
| 1147 | if (lLive) return local; // only local live: local holds |
| 1148 | if (!local) return incoming; // both vacant |
| 1149 | if (!incoming) return local; |
| 1150 | return leaseMs(incoming.renewedAt) >= leaseMs(local.renewedAt) ? incoming : local; |
| 1151 | } |
| 1152 | |
| 1153 | /// The named take-if-vacant merge for the whole `leases` section: the union of |
| 1154 | /// turnIds, each resolved by `mergeOneLease`. This is the rule sync.js routes |
| 1155 | /// the section through, distinct from the append-only union / freshest-scalar |
| 1156 | /// the rest of the parcel uses. NOT freshest-scalar -- that is the double claim. |
| 1157 | function mergeLeases(local, incoming, now) { |
| 1158 | var out = {}, a = local || {}, b = incoming || {}, k; |
| 1159 | for (k in a) if (Object.prototype.hasOwnProperty.call(a, k)) out[k] = a[k]; |
| 1160 | for (k in b) { |
| 1161 | if (!Object.prototype.hasOwnProperty.call(b, k)) continue; |
| 1162 | out[k] = mergeOneLease(a[k], b[k], now); |
| 1163 | } |
| 1164 | return out; |
| 1165 | } |
| 1166 | |
| 1167 | function leaseNow(nowFn) { return (typeof nowFn === 'function') ? nowFn() : Date.now(); } |
| 1168 | |
| 1169 | /// The section as it rides the parcel, or null when empty (a null section is |
| 1170 | /// one the other device leaves untouched, the same contract pause.js keeps). |
| 1171 | function leaseSnapshot() { |
| 1172 | return Object.keys(_leases).length ? _leases : null; |
| 1173 | } |
| 1174 | |
| 1175 | // A change listener the dispatching UI registers (D4). A lease learned through a |
| 1176 | // sync pull -- the phone seeing the peer claim, then run, its turn -- moves the |
| 1177 | // local view here but touches no message record, so nothing would otherwise |
| 1178 | // redraw the dispatched footer: it would sit on "Sent to your other devices" while |
| 1179 | // the peer held and ran the lease, and never show "[Take back]". `leaseAdopt` |
| 1180 | // fires this whenever the merge actually moved, so a sync update advances the |
| 1181 | // footer claimed -> running the same way a report does. |
| 1182 | var _onLeaseChange = null; |
| 1183 | function leaseOnChange(fn) { _onLeaseChange = fn; } |
| 1184 | |
| 1185 | /// Merge an arriving parcel's `leases` into the local view under take-if-vacant. |
| 1186 | /// Called by sync.js's reconcile. Answers whether anything MOVED, so a pull that |
| 1187 | /// agreed with us schedules no push -- the same quiet-on-no-change contract the |
| 1188 | /// pause tree keeps (sync.js:786) -- and, when it moved, notifies the UI so the |
| 1189 | /// dispatched footer re-renders against the fresh lease (D4). |
| 1190 | function leaseAdopt(incoming, nowFn) { |
| 1191 | if (!incoming) return false; |
| 1192 | var before = JSON.stringify(_leases); |
| 1193 | _leases = mergeLeases(_leases, incoming, leaseNow(nowFn)); |
| 1194 | var moved = JSON.stringify(_leases) !== before; |
| 1195 | if (moved && _onLeaseChange) { |
| 1196 | try { _onLeaseChange(); } catch (err) { /* a redraw must never break a sync */ } |
| 1197 | } |
| 1198 | return moved; |
| 1199 | } |
| 1200 | |
| 1201 | /// The live holder of a turn's lease at `now`, or null when it is vacant. |
| 1202 | function leaseHolder(turnId, now) { |
| 1203 | var r = _leases[String(turnId)]; |
| 1204 | return liveLease(r, now == null ? Date.now() : now) ? r.holder : null; |
| 1205 | } |
| 1206 | |
| 1207 | /// The full lease record for a turn, or null. What the guards and the UI state |
| 1208 | /// machine read (they need mode/expiry/holder, not just the live holder). |
| 1209 | function leaseRecord(turnId) { |
| 1210 | return _leases[String(turnId)] || null; |
| 1211 | } |
| 1212 | |
| 1213 | // ── The lifecycle, over a compare-and-set ────────────────── |
| 1214 | // |
| 1215 | // A CAS is `{ read, write }`: `read()` answers `{ version, leases }`; `write(base, |
| 1216 | // leases)` answers `{ ok:true, version }` when `base` was current (and bumps it), |
| 1217 | // or `{ ok:false, version, leases }` with the current blob when it was not -- the |
| 1218 | // 409. In production this is the parcel's own push/pull (sync.js:19-25); the |
| 1219 | // arbitration is identical, and modelling it as a CAS is what lets the race be |
| 1220 | // driven deterministically in a test. The ARBITRATION IS ENTIRELY IN |
| 1221 | // `mergeLeases`: a take folds its claim through the merge and stands down the |
| 1222 | // instant the merge does not keep it, so there is no second code path where a |
| 1223 | // wrong merge could still be caught -- swap the merge for freshest-scalar and the |
| 1224 | // take double-claims. That is on purpose. |
| 1225 | |
| 1226 | /// TAKE the lease for a turn, based on a parcel snapshot already read. Answers |
| 1227 | /// `{ won, holder, why }`. Stands down -- never runs -- when a live foreign |
| 1228 | /// lease exists (the merge drops the claim), when the deadline has passed, or |
| 1229 | /// when the CAS could not be won in bounds. |
| 1230 | /// |
| 1231 | /// The fold is `mergeLeases(MY claim /*local*/, server leases /*incoming*/)`: |
| 1232 | /// the server's existing foreign lease is the INCOMING that beats my fresh |
| 1233 | /// claim, so the merge -- and nothing else -- decides the race. This argument |
| 1234 | /// order is load-bearing; reversed, a loser would keep its own claim and double |
| 1235 | /// bill, which is exactly what the freshest-scalar mutation test proves. |
| 1236 | async function leaseTakeFrom(snap, turnId, opts, cas, nowFn) { |
| 1237 | var o = opts || {}; |
| 1238 | var holder = String(o.holder || ''); |
| 1239 | var tid = String(turnId); |
| 1240 | for (var attempt = 0; attempt < MAX_TAKE_TRIES; attempt++) { |
| 1241 | var now = leaseNow(nowFn); |
| 1242 | var deadline = leaseMs(o.deadline); |
| 1243 | if (deadline && now > deadline) { |
| 1244 | return { won: false, why: 'deadline' }; |
| 1245 | } |
| 1246 | // The claim expiry is the errand's DEADLINE, not now + TTL, so the lease |
| 1247 | // stays live for the whole turn WITHOUT a renew -- a busy turn cannot push a |
| 1248 | // renew (sync.js:1077), so a TTL-capped claim would read expired elsewhere |
| 1249 | // after 90s and be re-run (the >TTL double-run). A recovery errand carries |
| 1250 | // no deadline (deadline 0), so it falls back to a single TTL, which is right: |
| 1251 | // recovery is the owner running its own orphan, not a peer holding for long. |
| 1252 | // The record carries `deadline` so every merge/clamp honours the same bound. |
| 1253 | var claim = { |
| 1254 | turnId: tid, eid: String(o.eid || ''), holder: holder, |
| 1255 | mode: 'claimed', deadline: deadline || 0, |
| 1256 | expiry: (deadline && deadline > now) ? deadline : (now + LEASE_TTL_MS), |
| 1257 | renewedAt: now, |
| 1258 | }; |
| 1259 | var proposed = mergeLeases({ [tid]: claim }, snap.leases, now); |
| 1260 | if (!proposed[tid] || proposed[tid].holder !== holder) { |
| 1261 | _leases = mergeLeases(_leases, snap.leases, now); // adopt what we learned |
| 1262 | return { won: false, holder: proposed[tid] ? proposed[tid].holder : null }; |
| 1263 | } |
| 1264 | var res = await cas.write(snap.version, proposed); |
| 1265 | if (res.ok) { |
| 1266 | // A version bump is NOT proof our claim landed. The real sync resolves a |
| 1267 | // 409 mid-push by PULLING the concurrent winner's lease in, merging it |
| 1268 | // (take-if-vacant DROPS our claim), and pushing THAT -- yet the version |
| 1269 | // still advances, so a bare `ok` would let a loser believe it won and |
| 1270 | // double-run/double-charge (confirmed: two racers both `won` through the |
| 1271 | // pull-merge-retry commit). Trust the MERGE, never the version: re-read the |
| 1272 | // authoritative section and stand down unless it still names us as a LIVE |
| 1273 | // holder. A concurrent winner cannot be displaced by a later pull either -- |
| 1274 | // its lease is live and foreign, which `mergeLeases` keeps -- so a re-read |
| 1275 | // that names us is a true win. |
| 1276 | var conf; |
| 1277 | try { conf = await cas.read(); } |
| 1278 | catch (e) { conf = { version: res.version, leases: res.leases || {} }; } |
| 1279 | _leases = conf.leases || {}; |
| 1280 | var landed = _leases[tid]; |
| 1281 | if (landed && landed.holder === holder && liveLease(landed, leaseNow(nowFn))) { |
| 1282 | return { won: true, holder: holder }; |
| 1283 | } |
| 1284 | return { won: false, holder: landed ? landed.holder : null }; |
| 1285 | } |
| 1286 | // 409: the version CHURNED under us. Under active two-device sync the parcel |
| 1287 | // version keeps moving, so a stale `base` is refused by the commit BEFORE it |
| 1288 | // even pushes -- back-to-back tries then all fail and the claim never lands |
| 1289 | // (the live why:'exhausted' hand-off failure). Take a FRESH read so the next |
| 1290 | // base is current, and back off a jittered moment so the two devices do not |
| 1291 | // collide in lockstep -- giving the claim a real window. The fold above still |
| 1292 | // stands us down if a live foreign winner has appeared, so this stays |
| 1293 | // single-run safe: only the persistence changes, never the arbitration. |
| 1294 | try { snap = await cas.read(); } |
| 1295 | catch (e) { snap = { version: res.version, leases: res.leases || {} }; } |
| 1296 | if (attempt + 1 < MAX_TAKE_TRIES) { |
| 1297 | await new Promise(function (r) { |
| 1298 | setTimeout(r, Math.round(TAKE_BACKOFF_MS * (0.5 + Math.random()))); |
| 1299 | }); |
| 1300 | } |
| 1301 | } |
| 1302 | return { won: false, why: 'exhausted' }; |
| 1303 | } |
| 1304 | |
| 1305 | /// TAKE, reading the current parcel first. The ordinary entry point; the test |
| 1306 | /// uses `leaseTakeFrom` directly to race two takes from ONE base version. |
| 1307 | async function leaseTake(turnId, opts, cas, nowFn) { |
| 1308 | return leaseTakeFrom(await cas.read(), turnId, opts, cas, nowFn); |
| 1309 | } |
| 1310 | |
| 1311 | /// RENEW a lease this device holds, bumping its expiry. A healthy peer renews on |
| 1312 | /// journal progress; a dead one stops, and the lease expires. Answers |
| 1313 | /// `{ ok, why }`. Aborts (ok:false, why:'revoked') if the lease is no longer |
| 1314 | /// ours -- which is how a take-back (§3.3) reaches the running peer. |
| 1315 | async function leaseRenew(turnId, holder, cas, nowFn) { |
| 1316 | var tid = String(turnId), h = String(holder); |
| 1317 | for (var attempt = 0; attempt < MAX_TAKE_TRIES; attempt++) { |
| 1318 | var snap = await cas.read(); |
| 1319 | var now = leaseNow(nowFn); |
| 1320 | var cur = snap.leases[tid]; |
| 1321 | if (!cur || cur.holder !== h || cur.mode === 'released') { |
| 1322 | _leases = mergeLeases(_leases, snap.leases, now); |
| 1323 | return { ok: false, why: 'revoked' }; |
| 1324 | } |
| 1325 | // A renew never SHRINKS a deadline-bounded expiry: it holds to the later of |
| 1326 | // one TTL from now and the errand's deadline. Since a running turn is claimed |
| 1327 | // straight to its deadline and no longer renews on a ticker (runErrand only |
| 1328 | // transitions claimed -> running once), this is a no-op for a live turn; it |
| 1329 | // stays correct for a direct DaimondLease.renew of a TTL-only (no-deadline) |
| 1330 | // lease, where it is the old `now + TTL`. |
| 1331 | var bumped = { |
| 1332 | turnId: tid, eid: cur.eid, holder: h, |
| 1333 | mode: cur.mode === 'claimed' ? 'running' : cur.mode, |
| 1334 | deadline: leaseMs(cur.deadline), |
| 1335 | expiry: Math.max(now + LEASE_TTL_MS, leaseMs(cur.deadline)), renewedAt: now, |
| 1336 | }; |
| 1337 | var proposed = mergeLeases({ [tid]: bumped }, snap.leases, now); |
| 1338 | var res = await cas.write(snap.version, proposed); |
| 1339 | if (res.ok) { _leases = proposed; return { ok: true }; } |
| 1340 | } |
| 1341 | return { ok: false, why: 'exhausted' }; |
| 1342 | } |
| 1343 | |
| 1344 | /// COMPLETE (mode 'done') or RELEASE (mode 'released', which is vacant) a lease |
| 1345 | /// this device holds. `release` is also how the phone takes a turn back from a |
| 1346 | /// live peer (§3.3): the peer's read-only liveness check sees it released and aborts. |
| 1347 | async function leaseSet(turnId, holder, mode, cas, nowFn) { |
| 1348 | var tid = String(turnId), h = String(holder); |
| 1349 | for (var attempt = 0; attempt < MAX_TAKE_TRIES; attempt++) { |
| 1350 | var snap = await cas.read(); |
| 1351 | var now = leaseNow(nowFn); |
| 1352 | var cur = snap.leases[tid]; |
| 1353 | if (!cur || cur.holder !== h) { |
| 1354 | _leases = mergeLeases(_leases, snap.leases, now); |
| 1355 | return { ok: false, why: 'not_ours' }; |
| 1356 | } |
| 1357 | var next = { |
| 1358 | turnId: tid, eid: cur.eid, holder: h, mode: mode, |
| 1359 | deadline: leaseMs(cur.deadline), |
| 1360 | expiry: mode === 'released' ? 0 : cur.expiry, renewedAt: now, |
| 1361 | }; |
| 1362 | var proposed = mergeLeases({ [tid]: next }, snap.leases, now); |
| 1363 | var res = await cas.write(snap.version, proposed); |
| 1364 | if (res.ok) { _leases = proposed; return { ok: true }; } |
| 1365 | } |
| 1366 | return { ok: false, why: 'exhausted' }; |
| 1367 | } |
| 1368 | |
| 1369 | /// REVOKE a turn's lease whoever holds it -- the phone's take-back (§3.3). Unlike |
| 1370 | /// `release`, which is the holder letting go, this vacates a lease held by a |
| 1371 | /// DIFFERENT device: the running peer's read-only liveness check reads |
| 1372 | /// `mode:'released'` and hard-aborts its turn. CAS-written, so it races the peer |
| 1373 | /// cleanly. `renewedAt` is stamped now so the same-holder merge keeps the released |
| 1374 | /// record over the peer's live one. |
| 1375 | async function leaseRevoke(turnId, cas, nowFn) { |
| 1376 | var tid = String(turnId); |
| 1377 | for (var attempt = 0; attempt < MAX_TAKE_TRIES; attempt++) { |
| 1378 | var snap = await cas.read(); |
| 1379 | var now = leaseNow(nowFn); |
| 1380 | var cur = snap.leases[tid]; |
| 1381 | if (!cur || cur.mode === 'released') return { ok: true }; // already vacant |
| 1382 | // `renewedAt` at least the current record's, so the same-holder merge's |
| 1383 | // released-wins tie-break (or a strictly-greater renew) always keeps this |
| 1384 | // over the peer's live running record -- a fast-clock peer cannot outbid it. |
| 1385 | var revoked = { |
| 1386 | turnId: tid, eid: cur.eid, holder: cur.holder, |
| 1387 | mode: 'released', expiry: 0, deadline: leaseMs(cur.deadline), |
| 1388 | renewedAt: Math.max(now, leaseMs(cur.renewedAt)), |
| 1389 | }; |
| 1390 | var proposed = mergeLeases({ [tid]: revoked }, snap.leases, now); |
| 1391 | var res = await cas.write(snap.version, proposed); |
| 1392 | if (res.ok) { _leases = proposed; return { ok: true }; } |
| 1393 | } |
| 1394 | return { ok: false, why: 'exhausted' }; |
| 1395 | } |
| 1396 | |
| 1397 | /// Stage a leases section as the local view, for the sync shim ONLY: the CAS |
| 1398 | /// `commit` installs the proposed section here so the next `DaimondSync.push` |
| 1399 | /// sends it. Everything else reaches `_leases` through `adopt`/`take`/`renew`. |
| 1400 | function leaseInstall(leases) { _leases = leases || {}; } |
| 1401 | |
| 1402 | /// Drop the local view, for a test or an account switch. |
| 1403 | function leaseForget() { _leases = {}; } |
| 1404 | |
| 1405 | // The lease section provider, attached like pause.js so sync.js finds it by the |
| 1406 | // same `snapshot`/`adopt` contract every other section keeps. |
| 1407 | window.DaimondLease = { |
| 1408 | LEASE_TTL_MS: LEASE_TTL_MS, |
| 1409 | RENEW_EVERY_MS: RENEW_EVERY_MS, |
| 1410 | MAX_LEASE_LIFE_MS: MAX_LEASE_LIFE_MS, |
| 1411 | /// The named take-if-vacant merge for one turnId and for the whole section. |
| 1412 | /// Published so sync.js and a verifier drive the ONE implementation. |
| 1413 | mergeOne: mergeOneLease, |
| 1414 | merge: mergeLeases, |
| 1415 | live: liveLease, |
| 1416 | /// The two halves of sync.js's section contract. |
| 1417 | snapshot: leaseSnapshot, |
| 1418 | adopt: leaseAdopt, |
| 1419 | /// Register a redraw the UI wants run when a sync pull moves the lease view |
| 1420 | /// (D4): the dispatched footer advances claimed -> running -> "[Take back]". |
| 1421 | onChange: leaseOnChange, |
| 1422 | /// The live holder of a turn, or null; and the full record, for the guards |
| 1423 | /// and the UI state machine. |
| 1424 | holder: leaseHolder, |
| 1425 | record: leaseRecord, |
| 1426 | /// The lifecycle over a compare-and-set. |
| 1427 | take: leaseTake, |
| 1428 | /// TAKE from a snapshot already read -- lets a test race two takes from ONE |
| 1429 | /// base version, which is the concurrency the lease exists to arbitrate. |
| 1430 | takeFrom: leaseTakeFrom, |
| 1431 | renew: leaseRenew, |
| 1432 | complete: function (turnId, holder, cas, nowFn) { return leaseSet(turnId, holder, 'done', cas, nowFn); }, |
| 1433 | release: function (turnId, holder, cas, nowFn) { return leaseSet(turnId, holder, 'released', cas, nowFn); }, |
| 1434 | /// The phone's take-back: revoke whoever holds the lease (§3.3). |
| 1435 | revoke: leaseRevoke, |
| 1436 | /// Stage a section for the sync shim's CAS commit. Not for general use. |
| 1437 | install: leaseInstall, |
| 1438 | forget: leaseForget, |
| 1439 | }; |
| 1440 | |
| 1441 | // ── The runner (dev/PEER_DESIGN.md §4.3, step 5) ─────────── |
| 1442 | // |
| 1443 | // On a Channel::Post wake the errand routes through takeRow -> peek -> absorb |
| 1444 | // (step 2) to the runner registered here. `runErrand` is PURE over injected |
| 1445 | // deps, so the whole flow -- take, run, push, report, release, AND the |
| 1446 | // revoke->abort path -- is tested without daimond.js, which supplies the real |
| 1447 | // deps (the sync-bound lease CAS, reconstruct via ensureApp/scopeChatTo/chunks, |
| 1448 | // runTurn, the transcript push, the report post, the ack, chat.app.abort). |
| 1449 | |
| 1450 | /// Bind DaimondLease's abstract compare-and-set to a sync-like object. `sync` |
| 1451 | /// exposes `version()`, `leases()` and `commit(base, leases) -> { ok, version, |
| 1452 | /// leases }`; production wires `commit` onto DaimondSync -- install the leases |
| 1453 | /// section, push under CAS, report whether the version moved -- and this shim is |
| 1454 | /// what the lease lifecycle drives. The arbitration (push 409 -> adopt -> retry) |
| 1455 | /// is the lease's own; this only translates the interface, and both are proven |
| 1456 | /// against a fake sync in the tests. |
| 1457 | function syncCas(sync) { |
| 1458 | return { |
| 1459 | // A sync that offers an async `read` (the real lease door does; a test's |
| 1460 | // fake sync does not) reads through it; otherwise the synchronous |
| 1461 | // version()/leases() getters, which is what the tests drive. |
| 1462 | read: function () { |
| 1463 | return sync.read |
| 1464 | ? sync.read() |
| 1465 | : Promise.resolve({ version: sync.version(), leases: sync.leases() }); |
| 1466 | }, |
| 1467 | write: function (base, leases) { return Promise.resolve(sync.commit(base, leases)); }, |
| 1468 | }; |
| 1469 | } |
| 1470 | |
| 1471 | /// Run one errand end to end. Stands down -- never runs -- if a peer already |
| 1472 | /// holds it; HARD-ABORTS the instant the lease is revoked; and ACKS ONLY AFTER |
| 1473 | /// the result is pushed, so a crash before the push leaves the errand on the |
| 1474 | /// relay and the lease to expire (the phone reclaims, §2.5 -- nothing dropped). |
| 1475 | /// |
| 1476 | /// Pure over `deps`: |
| 1477 | /// selfId this device's id (the lease holder); |
| 1478 | /// cas the lease CAS (`syncCas` over the real sync); |
| 1479 | /// reconstruct async (errand) -> ctx: pull to >= parcelVersion, find the chat, |
| 1480 | /// `scopeChatTo`, apply `pause`, fetch chunks; |
| 1481 | /// runTurn async (ctx, prompt, { onProgress }): the ordinary turn engine, |
| 1482 | /// calling `onProgress` on journal events so the lease renews; |
| 1483 | /// abort (): hard-stop the in-flight turn (`chat.app.abort`); |
| 1484 | /// pushResult async () -> version: `captureSession` + parcel push (append merge); |
| 1485 | /// post async (reportEnvelope): post the report; |
| 1486 | /// ack async (): `DaimondPost.ack`, AFTER the push committed; |
| 1487 | /// now optional clock, for tests. |
| 1488 | /// |
| 1489 | /// Answers `{ ran, done?, aborted?, error?, why?, holder?, trace }`. `trace` is |
| 1490 | /// the ordered side effects, so a test asserts the sequence rather than guessing. |
| 1491 | /// |
| 1492 | /// PARK — abandon and re-run, never resume (there is no mid-turn checkpoint). When |
| 1493 | /// a per-turn consent could not be answered (no attended device, or the wait timed |
| 1494 | /// out), egressAllowed on the runner records the intent and aborts the turn; this |
| 1495 | /// runner then parks. Park REPORTS THEN RELEASES -- mirroring the reconstruct-fail |
| 1496 | /// order below -- so the lease is never stranded. Below MAX_PARKS the report is |
| 1497 | /// `parked` (the turn re-dispatches, fresh, when a human next surfaces); AT the |
| 1498 | /// bound it is a terminal failure the user is told about, and it does not re-run. |
| 1499 | /// The parkCount it carries is the GLOBAL total, so two devices cannot each drive |
| 1500 | /// the loop independently -- the count and the single-runner lease together cap the |
| 1501 | /// respend at MAX_PARKS. |
| 1502 | async function parkAndRelease(e, d, trace, pk) { |
| 1503 | var turnId = String(e.turnId); |
| 1504 | var out = parkOutcome(e.parkCount, d.maxParks); |
| 1505 | var terminal = out.terminal; |
| 1506 | var why = String((pk && pk.why) || (terminal |
| 1507 | ? 'This turn needed your permission and no device was available to grant it -- it did not run.' |
| 1508 | : 'This turn needs your permission and no device was available -- it will re-run when you are back.')); |
| 1509 | // REPORT then RELEASE (the reconstruct-fail order), so the lease is freed only |
| 1510 | // after the account of the stop is on its way. A terminal park is an `aborted` |
| 1511 | // report (no re-dispatch); a survivable one is `parked`, carrying the bumped |
| 1512 | // GLOBAL count so the re-dispatcher increments from the true total. |
| 1513 | try { |
| 1514 | if (d.post) await d.post(makeReport({ |
| 1515 | eid: e.eid, turnId: turnId, chatId: e.chatId, |
| 1516 | status: terminal ? 'aborted' : 'parked', why: why, parkCount: out.next })); |
| 1517 | trace.push('report'); |
| 1518 | } catch (err) { /* the release below still frees the turn */ } |
| 1519 | try { await leaseSet(turnId, d.selfId, 'released', d.cas, d.now); trace.push('release'); } |
| 1520 | catch (err) { /* an unreleased lease still expires at its deadline */ } |
| 1521 | return { ran: true, parked: !terminal, terminal: terminal, why: why, parkCount: out.next, trace: trace }; |
| 1522 | } |
| 1523 | |
| 1524 | async function runErrand(errand, deps) { |
| 1525 | var d = deps || {}, e = errand || {}; |
| 1526 | var turnId = String(e.turnId); |
| 1527 | var trace = []; |
| 1528 | |
| 1529 | // D1(a) — NEVER run an errand THIS device dispatched, EXCEPT on a deliberate |
| 1530 | // local recovery (`allowSelf`). The phone returns from the background and the |
| 1531 | // ordinary collect loop re-collects its OWN self-posted errand |
| 1532 | // (peerCollectOnReturn); routed here and run, it would re-take a released lease |
| 1533 | // and re-run a turn a peer already ran -- a second completion and a second |
| 1534 | // charge. So the AUTOMATIC path stands down on its own dispatch. Recovery is |
| 1535 | // different: it is the dispatching device DELIBERATELY running its own orphaned |
| 1536 | // turn because no peer did, and it has already confirmed the turn is not |
| 1537 | // finished and not held by a live foreign lease. It is STILL money-safe, because |
| 1538 | // it goes through the SAME `finished` (D1(b)) check and the SAME take-if-vacant |
| 1539 | // lease below -- a peer that took the lease first wins the merge and recovery |
| 1540 | // stands down; a peer that collects AFTER recovery's ack finds no errand and, |
| 1541 | // if it somehow does, `finished`/the released-with-answer lease stand it down. |
| 1542 | // `allowSelf` only lifts THIS blanket refusal; every other guard is untouched. |
| 1543 | // `dispatchedBy` names the dispatching device (the per-device id, not the |
| 1544 | // account key), so a match to this device is our own dispatch. |
| 1545 | if (!d.allowSelf && e.dispatchedBy && String(e.dispatchedBy) === String(d.selfId)) { |
| 1546 | trace.push('self-dispatched'); |
| 1547 | return { ran: false, why: 'self-dispatched', trace: trace }; |
| 1548 | } |
| 1549 | |
| 1550 | // D1(b) — a COMPLETED turn is not vacant-for-rerun. A released lease reads |
| 1551 | // vacant (`liveLease` false), and `done` is transient before `released`, so a |
| 1552 | // turn the peer already finished would be re-taken and re-billed by the next |
| 1553 | // device to collect the errand. A turn that already carries a done report or a |
| 1554 | // merged answer is FINISHED: stand down before the take. `finished` is supplied |
| 1555 | // by the app (it checks the report box and the transcript); absent -- the |
| 1556 | // runner-acceptance path -- this is a no-op. |
| 1557 | if (d.finished) { |
| 1558 | var already = false; |
| 1559 | try { already = await d.finished(e); } catch (err) { already = false; } |
| 1560 | if (already) { trace.push('already-done'); return { ran: false, why: 'already-done', trace: trace }; } |
| 1561 | } |
| 1562 | |
| 1563 | // D1(c) — DEFER TO THE NOMINATED RUNNER. When the account has named an always-on |
| 1564 | // runner and it is FRESHLY awake, a non-nominee stands down and leaves the claim |
| 1565 | // to it, so a laptop that may be closed mid-turn does not grab a turn the desktop |
| 1566 | // should run. Gated on the nominee's LIVE freshness (DISPATCH_FRESH_MS): a nominee |
| 1567 | // that has actually slept reads offline and this device claims instead -- fall back |
| 1568 | // to any awake device, the owner's explicit choice over nobody-runs-it. NOT applied |
| 1569 | // on a deliberate local recovery (`allowSelf`): recovery is the guaranteed net that |
| 1570 | // a turn NO peer ran is still run, and must never itself stall for the nominee. |
| 1571 | // Standing down does NOT ack -- takeRow HOLDs the errand on the relay (post.js) -- |
| 1572 | // so it is re-collected and re-decided against live presence until the nominee runs |
| 1573 | // it or its beat ages out. Only WHO attempts the claim changes; the take-if-vacant |
| 1574 | // lease below is still the single-runner arbiter. |
| 1575 | if (!d.allowSelf && nominationStandDown(d.nominatedId, d.selfId, d.presence, leaseNow(d.now), d.freshWindowMs)) { |
| 1576 | trace.push('stood-down-for-nominee'); |
| 1577 | return { ran: false, why: 'nominee', trace: trace }; |
| 1578 | } |
| 1579 | |
| 1580 | // A missing lease CAS cannot arbitrate a claim, so there is no safe way to run: |
| 1581 | // stand down cleanly rather than let `leaseTake` dereference a null `cas` and |
| 1582 | // throw the opaque "Cannot read properties of null (reading 'read')". |
| 1583 | if (!d.cas || typeof d.cas.read !== 'function') { |
| 1584 | trace.push('no-cas'); |
| 1585 | return { ran: false, why: 'no-cas', trace: trace }; |
| 1586 | } |
| 1587 | |
| 1588 | // 1. TAKE. Stand down -- never run -- if a peer already holds it. |
| 1589 | var took = await leaseTake(turnId, |
| 1590 | { holder: d.selfId, eid: e.eid, deadline: e.deadline }, d.cas, d.now); |
| 1591 | trace.push('take'); |
| 1592 | if (!took.won) return { ran: false, why: took.why || 'stood-down', holder: took.holder, trace: trace }; |
| 1593 | |
| 1594 | // THE LEASE DOES NOT RENEW. It is claimed straight to the errand's DEADLINE |
| 1595 | // (leaseTakeFrom), so it stays live for the whole turn with no periodic write -- |
| 1596 | // which is what keeps the parcel a fixed point during a running turn AND closes |
| 1597 | // the >LEASE_TTL_MS double-run: a busy turn cannot push a 30s renew (sync.js:1077 |
| 1598 | // suppresses the push over a live turn), so a TTL-capped lease read EXPIRED on |
| 1599 | // other devices after 90s while the turn ran on, and the phone's recovery re-ran |
| 1600 | // and re-billed it. What runs on a ticker now is a READ-ONLY liveness check: it |
| 1601 | // detects a take-back (the phone REVOKED the lease) and HARD-ABORTS, and it caps a |
| 1602 | // hung turn's lifetime -- it never writes the parcel, so a running turn causes no |
| 1603 | // churn. The check is owned HERE (not in the injected runTurn) and stopped on |
| 1604 | // EVERY exit (the finally), so it can neither outlive the errand nor leak a timer. |
| 1605 | var revoked = false, checkStopped = false, checkTimer = null; |
| 1606 | var checkStart = leaseNow(d.now); |
| 1607 | var maxLife = (d.maxLeaseLifeMs != null) ? d.maxLeaseLifeMs : MAX_LEASE_LIFE_MS; |
| 1608 | var setT = d.setTimer || (typeof setInterval === 'function' ? setInterval : null); |
| 1609 | var clrT = d.clearTimer || (typeof clearInterval === 'function' ? clearInterval : null); |
| 1610 | function stopCheck() { |
| 1611 | checkStopped = true; |
| 1612 | if (checkTimer != null && clrT) { try { clrT(checkTimer); } catch (err) {} checkTimer = null; } |
| 1613 | } |
| 1614 | // READ-ONLY: never writes the parcel (no renew, no churn). Aborts on a revoke |
| 1615 | // -- the lease is no longer ours, or was released, which a sync pull adopts into |
| 1616 | // the view this reads -- and on the lifetime cap, the backstop for a runTurn |
| 1617 | // whose promise never settles, after which the lease is simply left to expire. |
| 1618 | async function liveness() { |
| 1619 | if (checkStopped || revoked) return; |
| 1620 | if (leaseNow(d.now) - checkStart > maxLife) { |
| 1621 | trace.push('renew-capped'); |
| 1622 | stopCheck(); |
| 1623 | revoked = true; |
| 1624 | try { if (d.abort) d.abort(); } catch (err) { /* idempotent */ } |
| 1625 | return; |
| 1626 | } |
| 1627 | var snap; |
| 1628 | try { snap = await d.cas.read(); } catch (err) { return; } // a failed read is not a revoke |
| 1629 | var cur = (snap && snap.leases) ? snap.leases[turnId] : null; |
| 1630 | if (!cur || cur.holder !== String(d.selfId) || cur.mode === 'released') { |
| 1631 | revoked = true; |
| 1632 | trace.push('abort'); |
| 1633 | try { if (d.abort) d.abort(); } catch (err) { /* idempotent */ } |
| 1634 | } |
| 1635 | } |
| 1636 | try { |
| 1637 | // 2. RECONSTRUCT the chat and workspace at the errand's version. |
| 1638 | var ctx; |
| 1639 | try { ctx = await d.reconstruct(e); trace.push('reconstruct'); } |
| 1640 | catch (err) { |
| 1641 | // The lease was TAKEN above. A reconstruct that throws must NOT leave it |
| 1642 | // pinned at 'claimed' to the errand's deadline: on an iOS phone that cannot |
| 1643 | // fire the deadline fallback, that reads as a turn stuck on "picking this |
| 1644 | // up" for ever, with the engine's real sentence swallowed and nothing to act |
| 1645 | // on. So SURFACE it and HAND IT BACK -- post an error report carrying the |
| 1646 | // reason (the phone shows it and offers [Run here]) and release the lease so |
| 1647 | // the turn is reclaimable at once rather than after the deadline. Nothing ran, |
| 1648 | // so there is no charge and the release is money-safe. |
| 1649 | var rwhy = String((err && err.message) || err); |
| 1650 | trace.push('reconstruct-failed'); |
| 1651 | try { if (typeof console !== 'undefined') console.error('peer: reconstruct failed for turn ' + turnId + ' -- ' + rwhy); } catch (e2) {} |
| 1652 | try { if (d.post) await d.post(makeReport({ eid: e.eid, turnId: turnId, chatId: e.chatId, status: 'error', why: rwhy })); } |
| 1653 | catch (e2) { /* the release below still frees the turn */ } |
| 1654 | try { await leaseSet(turnId, d.selfId, 'released', d.cas, d.now); trace.push('release'); } |
| 1655 | catch (e2) { /* an unreleased lease still expires at its deadline */ } |
| 1656 | return { ran: false, error: true, why: rwhy, trace: trace }; |
| 1657 | } |
| 1658 | |
| 1659 | // 3. Transition claimed -> running ONCE -- a semantic state change for the UI |
| 1660 | // footer ("running" vs "picking this up"), one write, before the turn goes |
| 1661 | // busy. This keeps the deadline expiry (leaseRenew never shrinks it); it does |
| 1662 | // NOT start a heartbeat. A lease already revoked between take and here aborts. |
| 1663 | var mk = await leaseRenew(turnId, d.selfId, d.cas, d.now); |
| 1664 | if (!mk.ok && mk.why === 'revoked') { |
| 1665 | revoked = true; |
| 1666 | trace.push('abort'); |
| 1667 | try { if (d.abort) d.abort(); } catch (err) { /* idempotent */ } |
| 1668 | return { ran: true, aborted: true, why: 'revoked', trace: trace }; |
| 1669 | } |
| 1670 | // 4. RUN. A revoked lease HARD-ABORTS at once, via the read-only ticker and |
| 1671 | // the injected onProgress (kept so a real journal-event piggyback can check |
| 1672 | // liveness between ticks); chat.app.abort is the hard stop. |
| 1673 | if (setT) checkTimer = setT(function () { liveness(); if (d.heartbeat) d.heartbeat(); }, RENEW_EVERY_MS); |
| 1674 | try { |
| 1675 | // D3 — the prompt is ALREADY in the synced transcript (the dispatcher |
| 1676 | // persist-first pushed it before posting the errand, §4.1). Tell runTurn so, |
| 1677 | // so it runs the turn against the existing user message instead of appending |
| 1678 | // a second copy -- otherwise the prompt sits twice in `messages` AND is fed |
| 1679 | // to the model twice (seeded history + the re-sent turn). `turnId` names the |
| 1680 | // existing user message (mid === turnId) the runner anchors to. |
| 1681 | await d.runTurn(ctx, e.prompt, { onProgress: liveness, promptInTranscript: true, turnId: turnId }); |
| 1682 | trace.push('run'); |
| 1683 | } catch (err) { |
| 1684 | // Revoked -> the lease is already whoever took it back's; touch nothing, |
| 1685 | // do NOT ack -- the errand is theirs now. |
| 1686 | if (revoked) return { ran: true, aborted: true, why: 'revoked', trace: trace }; |
| 1687 | // PARK -> egressAllowed could not get consent and aborted the turn on |
| 1688 | // purpose (no attended device, or the wait timed out). This is not a |
| 1689 | // crash: report the park (or terminal failure at the bound) and RELEASE |
| 1690 | // the lease, so the turn re-dispatches or ends clean rather than stranding. |
| 1691 | var pkErr = d.parkRequested ? d.parkRequested(e) : null; |
| 1692 | if (pkErr) { stopCheck(); return await parkAndRelease(e, d, trace, pkErr); } |
| 1693 | // A genuine crash: do NOT ack and do NOT complete, so the relay keeps the |
| 1694 | // errand and the lease EXPIRES (at its deadline). The phone reclaims (§2.5). |
| 1695 | return { ran: true, error: true, why: String(err && err.message || err), trace: trace }; |
| 1696 | } |
| 1697 | if (revoked) return { ran: true, aborted: true, why: 'revoked', trace: trace }; |
| 1698 | // PARK could also be signalled without the turn throwing (egressAllowed |
| 1699 | // returned a refusal and the turn wound down on its own). Park takes |
| 1700 | // precedence over completing a half-answer: nothing past the consent point |
| 1701 | // ran, so there is nothing to push. |
| 1702 | var pkOk = d.parkRequested ? d.parkRequested(e) : null; |
| 1703 | if (pkOk) { stopCheck(); return await parkAndRelease(e, d, trace, pkOk); } |
| 1704 | // The turn produced a result: stop the liveness ticker BEFORE completing, so |
| 1705 | // nothing races the done/release writes below. |
| 1706 | stopCheck(); |
| 1707 | |
| 1708 | // 4. COMPLETE in order: push the transcript (append merges), post the report, |
| 1709 | // mark the lease done, ACK the errand (only now the push has committed), then |
| 1710 | // release the lease. |
| 1711 | var parcelVersion = 0; |
| 1712 | try { parcelVersion = (await d.pushResult()) | 0; trace.push('push'); } |
| 1713 | catch (err) { return { ran: true, error: true, why: 'push-failed', trace: trace }; } |
| 1714 | try { |
| 1715 | await d.post(makeReport({ eid: e.eid, turnId: turnId, chatId: e.chatId, |
| 1716 | status: 'done', parcelVersion: parcelVersion })); |
| 1717 | trace.push('report'); |
| 1718 | } catch (err) { /* the report is only the nudge; the answer is already pushed */ } |
| 1719 | await leaseSet(turnId, d.selfId, 'done', d.cas, d.now); trace.push('complete'); |
| 1720 | try { if (d.ack) { await d.ack(); trace.push('ack'); } } |
| 1721 | catch (err) { /* a missed ack costs one idempotent re-collect, never a drop */ } |
| 1722 | await leaseSet(turnId, d.selfId, 'released', d.cas, d.now); trace.push('release'); |
| 1723 | return { ran: true, done: true, parcelVersion: parcelVersion, trace: trace }; |
| 1724 | } finally { |
| 1725 | stopCheck(); // EVERY exit stops the liveness ticker -- no timer leaks. |
| 1726 | } |
| 1727 | } |
| 1728 | |
| 1729 | // ── Public surface ───────────────────────────────────────── |
| 1730 | /// Does this errand name THIS device as its dispatcher? The sender must NOT ack |
| 1731 | /// its own un-run errand off the shared relay -- only the peer that actually runs |
| 1732 | /// it may (post.js `collect`/`takeRow` hold it otherwise). `dispatchedBy` carries |
| 1733 | /// the dispatcher's peer id, which is `DaimondIdentity.deviceId()` -- the same |
| 1734 | /// `selfDeviceId` the runner's D1(a) self-dispatch guard compares against. |
| 1735 | function isOwnDispatch(env) { |
| 1736 | if (!env || env.t !== T_ERRAND || !env.dispatchedBy) return false; |
| 1737 | var self = ''; |
| 1738 | try { |
| 1739 | self = (window.DaimondIdentity && DaimondIdentity.deviceId) |
| 1740 | ? String(DaimondIdentity.deviceId() || '') : ''; |
| 1741 | } catch (e) { self = ''; } |
| 1742 | return !!self && String(env.dispatchedBy) === self; |
| 1743 | } |
| 1744 | |
| 1745 | window.DaimondPeer = { |
| 1746 | ENVELOPE_V: ENVELOPE_V, |
| 1747 | T_ERRAND: T_ERRAND, |
| 1748 | T_REPORT: T_REPORT, |
| 1749 | T_ASK: T_ASK, |
| 1750 | T_GRANT: T_GRANT, |
| 1751 | makeErrand: makeErrand, |
| 1752 | makeReport: makeReport, |
| 1753 | /// The two remote-consent envelopes: a runner's live question (`makeAsk`, fresh |
| 1754 | /// `cid` per ask) and an attended device's answer (`makeGrant`, which signs |
| 1755 | /// `cid`/`turnId`/`verdict` but NOT `tool`/`host`/`detail`, so the runner replays |
| 1756 | /// the EXACT act it held). |
| 1757 | makeAsk: makeAsk, |
| 1758 | makeGrant: makeGrant, |
| 1759 | /// Seal an envelope to this account and answer the `{ to, addr, envelope }` |
| 1760 | /// post body. Reuses `DaimondPost.seal`; no server is involved, which is |
| 1761 | /// also how it is tested. |
| 1762 | sealForSelf: sealForSelf, |
| 1763 | /// Open a sealed envelope, or throw if it is not this account's OR was not |
| 1764 | /// signed by this account. The same-account-only property lives here. |
| 1765 | openEnvelope: openEnvelope, |
| 1766 | /// Sign an envelope with this account's key / verify one was so signed. |
| 1767 | signEnvelope: signEnvelope, |
| 1768 | verifyEnvelope: verifyEnvelope, |
| 1769 | /// The collector's two doors: classify a row's sealed body without acting |
| 1770 | /// (`peek`), then verify-and-run it (`absorb`). `takeRow` calls these. |
| 1771 | peek: peek, |
| 1772 | absorb: absorb, |
| 1773 | /// Whether an errand is THIS device's own dispatch -- so the sender's collect |
| 1774 | /// leaves it on the relay for the peer rather than acking it away. |
| 1775 | isOwnDispatch: isOwnDispatch, |
| 1776 | /// Register the runners the collector hands a verified envelope to. Set by |
| 1777 | /// daimond.js; absent, `absorb` verifies and drops. |
| 1778 | onErrand: onErrand, |
| 1779 | onReport: onReport, |
| 1780 | /// Register the remote-consent handlers: `onAsk` raises a runner's question on |
| 1781 | /// an attended device; `onGrant` delivers the answer to the awaiting runner. |
| 1782 | onAsk: onAsk, |
| 1783 | onGrant: onGrant, |
| 1784 | /// Route a batch of collected rows by their sealed type tag (direct-drive). |
| 1785 | routeRows: routeRows, |
| 1786 | /// Fold a peer's answer into a transcript as an append. |
| 1787 | foldAssistant: foldAssistant, |
| 1788 | /// The dispatcher's pure core: assemble the ordered plan + full errand |
| 1789 | /// (`buildDispatch`), and classify a `why:'dispatched'` turn against the |
| 1790 | /// lease (`dispatchState`). daimond.js runs the order; these hold the logic. |
| 1791 | buildDispatch: buildDispatch, |
| 1792 | dispatchState: dispatchState, |
| 1793 | /// The §5 display state of a dispatched turn (dispatched/no-peer-awake/ |
| 1794 | /// claimed/running/done/failed). Pure; daimond.js only renders it. |
| 1795 | uiState: uiState, |
| 1796 | /// Should this turn be auto-handed to a peer, and which one? Pure; daimond.js |
| 1797 | /// acts on it at send-time. And the shared "which peer is awake" answer. |
| 1798 | autoDispatchDecision: autoDispatchDecision, |
| 1799 | freshestPeer: freshestPeer, |
| 1800 | /// The remote-consent decisions, pure so a test drives the money-safe bound. |
| 1801 | /// `attendedPeer` -- the freshest device the user is ON (foreground + recent |
| 1802 | /// interaction), or null; `consentRouteDecision` -- ask that peer, resolve from |
| 1803 | /// policy, or park; `parkOutcome` -- the new GLOBAL park total and whether it |
| 1804 | /// reaches the terminal spend bound. |
| 1805 | attendedPeer: attendedPeer, |
| 1806 | consentRouteDecision: consentRouteDecision, |
| 1807 | parkOutcome: parkOutcome, |
| 1808 | CONSENT_DEADLINE_MS: CONSENT_DEADLINE_MS, |
| 1809 | MAX_PARKS: MAX_PARKS, |
| 1810 | /// Whether the dispatching device should RECOVER an orphaned dispatched turn |
| 1811 | /// locally on its return -- dispatched, not finished, not held by a live peer. |
| 1812 | /// daimond.js acts on it through the same lease, so it is money-safe. |
| 1813 | recoverDecision: recoverDecision, |
| 1814 | /// Should this device stand down from claiming a dispatched turn, deferring to |
| 1815 | /// the account's nominated always-on runner? Pure; runErrand consults it before |
| 1816 | /// the lease take, and daimond.js supplies the nominee id + live presence. |
| 1817 | nominationStandDown: nominationStandDown, |
| 1818 | REASON_DISPATCHED: REASON_DISPATCHED, |
| 1819 | DISPATCH_DEADLINE_MS: DISPATCH_DEADLINE_MS, |
| 1820 | /// The TIGHTER window the auto-dispatch decision uses (under the display |
| 1821 | /// window and well under the gateway TTL), so a dispatch never goes to a peer |
| 1822 | /// last seen too long ago to still be beating. |
| 1823 | DISPATCH_FRESH_MS: DISPATCH_FRESH_MS, |
| 1824 | /// The runner: bind the lease CAS to the real sync (`syncCas`), then run an |
| 1825 | /// errand end to end (`runErrand`) -- take, run, push, report, release, with |
| 1826 | /// a hard-abort on revoke and ack only after commit. Pure over injected deps. |
| 1827 | syncCas: syncCas, |
| 1828 | runErrand: runErrand, |
| 1829 | /// The content address of some sealed bytes, exposed for a caller that |
| 1830 | /// seals by hand. |
| 1831 | addressOf: addressOf, |
| 1832 | }; |
| 1833 | })(); |