// notes.json on disk: read, lock, apply one op, write. Shared by the app's // /api/notes and `umtool notes`, so the operator's page and an agent's CLI go // through ONE writer with one set of rules (./shape.mjs). // // Three writers can race on one file -- the page, an agent's `umtool notes // reply`, a second agent -- and they are different PROCESSES, so the in-process // queues lib/state.ts and lib/report/manifest.mjs use are not enough here: // // * a LOCKFILE (`notes.json.lock`, created O_EXCL) serialises // read-modify-write across processes; one left by a dead writer is stale // after 30 s and is taken over; // * the write is tmp + rename, so a reader never sees half a file; // * every write from the page carries the TOKEN it read (the file's mtime in // ns, or "absent"); a stale one is a 409, never a silent overwrite of a // reply an agent wrote in between; // * a notes.json that exists and does not parse is NEVER overwritten (the // takes.mjs rule) -- writing over it would erase every note in it; // * the last note deleted deletes the file: an empty notes.json says nothing. // // WHERE a write may land is the caller's to decide BEFORE it gets here // (lib/annotations/targets.mjs: isCorpusNotesFile for an article, a project // directory for a video) -- `writeOp` takes the file it is given. import { open, readFile, rename, stat, unlink, writeFile } from "node:fs/promises"; import path from "node:path"; import { NoteError, applyOp, emptyDoc, parseNotesDoc } from "./shape.mjs"; export { NoteError }; export const LOCK_STALE_MS = 30_000; const LOCK_WAIT_MS = 10_000; /** A write whose token no longer matches the file: someone wrote in between. */ export class StaleNotes extends Error { constructor(expected, got) { super(`the notes changed since you read them (${got} vs ${expected})`); this.name = "StaleNotes"; this.status = 409; } } /** A notes.json that exists and is not one: refused, never overwritten. */ export class NotesUnreadable extends Error { constructor(file, why) { super(`${path.basename(file)} does not parse (${why}); not overwriting it`); this.name = "NotesUnreadable"; this.status = 409; } } /** The file's identity for a write guard: mtime in ns plus size, or "absent". */ export async function notesToken(file) { const st = await stat(/* turbopackIgnore: true */ file, { bigint: true }).catch(() => null); return st ? `${st.mtimeNs}-${st.size}` : "absent"; } /** * `{ doc, token }` -- doc null when there is no file. `{ error }` too when the * file is there and is not a notes doc (the page shows it; writes refuse). * * @param {string} file * @returns {Promise<{ doc: import("./types").NotesDoc | null, token: string, error?: string }>} */ export async function readNotes(file) { const token = await notesToken(file); let text; try { text = await readFile(/* turbopackIgnore: true */ file, "utf8"); } catch (err) { if (/** @type {NodeJS.ErrnoException} */ (err).code === "ENOENT") return { doc: null, token: "absent" }; return { doc: null, token, error: String(/** @type {Error} */ (err).message ?? err) }; } let raw; try { raw = JSON.parse(text); } catch (err) { return { doc: null, token, error: `not JSON: ${/** @type {Error} */ (err).message}` }; } const r = parseNotesDoc(raw); if ("error" in r) return { doc: null, token, error: r.error }; return { doc: r.doc, token }; } const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); /** * Run `fn` holding `.lock`. A lock older than LOCK_STALE_MS is a dead * writer's and is removed; otherwise wait (up to 10 s) and retry. * * @template T * @param {string} file * @param {() => Promise} fn * @param {{ waitMs?: number, staleMs?: number }} [opts] * @returns {Promise} */ export async function withNotesLock(file, fn, { waitMs = LOCK_WAIT_MS, staleMs = LOCK_STALE_MS } = {}) { const lock = `${file}.lock`; const deadline = Date.now() + waitMs; let delay = 15; for (;;) { try { const h = await open(/* turbopackIgnore: true */ lock, "wx"); await h.writeFile(`${process.pid} ${new Date().toISOString()}\n`).catch(() => {}); await h.close(); break; } catch (err) { if (/** @type {NodeJS.ErrnoException} */ (err).code !== "EEXIST") throw err; const st = await stat(/* turbopackIgnore: true */ lock).catch(() => null); if (st && Date.now() - st.mtimeMs > staleMs) { await unlink(/* turbopackIgnore: true */ lock).catch(() => {}); continue; } if (Date.now() > deadline) throw new Error(`${path.basename(lock)} is held; try again`); await sleep(delay); delay = Math.min(delay * 2, 250); } } try { return await fn(); } finally { await unlink(/* turbopackIgnore: true */ lock).catch(() => {}); } } /** * Apply one op to the notes at `file` and write it back. `init` is the * subject (and source) a NEW file starts with; an existing file keeps its own * subject, and its `source` unless the op is `source`. `token` (when given) * must match the file as it is now, or StaleNotes. `by` is "operator" (the * app) or "agent" (the CLI). * * Returns the doc as written (null when the file was deleted), the new token, * and the note the op touched. * * @param {string} file * @param {{ subject: Record, source?: Record | null }} init * @param {Record} op * @param {{ by: "operator" | "agent", token?: string | null }} opts * @returns {Promise<{ doc: import("./types").NotesDoc | null, token: string, note: import("./types").Note | null }>} */ export async function writeOp(file, init, op, { by, token = null }) { return withNotesLock(file, async () => { const current = await readNotes(file); if (current.error) throw new NotesUnreadable(file, current.error); if (token !== null && token !== undefined && token !== current.token) throw new StaleNotes(token, current.token); const doc = /** @type {any} */ (current.doc ?? emptyDoc(init.subject, init.source ?? undefined)); const note = applyOp(doc, op, by); if (doc.notes.length === 0) { await unlink(/* turbopackIgnore: true */ file).catch((err) => { if (err.code !== "ENOENT") throw err; }); return { doc: null, token: "absent", note }; } const tmp = `${file}.tmp-${process.pid}-${Math.random().toString(36).slice(2, 8)}`; await writeFile(/* turbopackIgnore: true */ tmp, JSON.stringify(doc, null, 2) + "\n", "utf8"); await rename(/* turbopackIgnore: true */ tmp, file); return { doc, token: await notesToken(file), note }; }); }