diff --git a/lib/host/store.js b/lib/host/store.js index 3317c31..ac18aaa 100644 --- a/lib/host/store.js +++ b/lib/host/store.js @@ -1,6 +1,6 @@ import { asBoardSettings, emptyLedger, isPlausibleTaskRecord, pruneExecutions } from "../shared/protocol.js"; import { dirname, join } from "node:path"; -import { mkdir, open, readFile, rename } from "node:fs/promises"; +import { mkdir, open, readFile, rename, stat, unlink } from "node:fs/promises"; //#region src/host/store.ts /** * Host-side task ledger: one JSON file under the active data directory, mutated through a @@ -15,7 +15,9 @@ import { mkdir, open, readFile, rename } from "node:fs/promises"; * validates the resulting document, bumps the global revision, persists * atomically (temp file + rename), and only then notifies subscribers. */ -var TaskStore = class { +var TaskStore = class TaskStore { + /** Stores from old and new plugin generations share one write lane per file. */ + static fileQueues = /* @__PURE__ */ new Map(); file; storageQueue; ledger = emptyLedger(); @@ -37,10 +39,25 @@ var TaskStore = class { await this.load(); await persistAtomic(file, JSON.stringify(this.ledger)); } + /** Re-read through the normalizer while holding the cross-instance write lock. */ + async reload() { + this.loaded = false; + this.loadPromise = void 0; + await this.load(); + } /** Switch future writes after a prepared migration commits. */ setLocation(file) { this.file = file; } + /** Serialize one operation with every TaskStore in this host for this file. */ + inProcessWriteLane(file, run) { + const next = (TaskStore.fileQueues.get(file) ?? Promise.resolve()).then(run, run); + TaskStore.fileQueues.set(file, next); + next.finally(() => { + if (TaskStore.fileQueues.get(file) === next) TaskStore.fileQueues.delete(file); + }).catch(() => {}); + return next; + } /** Load (once) from disk; a missing file starts empty; a corrupt file is quarantined, not thrown. */ load() { if (this.loaded) return Promise.resolve(); @@ -131,29 +148,32 @@ var TaskStore = class { async mutate(kind, mutator) { const run = async () => { await this.load(); - const draft = structuredClone(this.ledger); - const changed = mutator(draft); - if (changed === void 0) return { - ledger: deepFreeze(structuredClone(this.ledger)), - changed: [] - }; - for (const task of changed) pruneExecutions(task); - draft.revision += 1; - const json = JSON.stringify(draft); - await persistAtomic(this.file, json); - this.ledger = draft; - const change = { - revision: draft.revision, - tasks: changed, - kind - }; - for (const fn of this.subscribers) try { - fn(change); - } catch {} - return { - ledger: deepFreeze(structuredClone(draft)), - changed: changed.map((t) => deepFreeze(structuredClone(t))) - }; + const file = this.file; + return this.inProcessWriteLane(file, async () => withLedgerWriteLock(file, async () => { + await this.reload(); + const draft = structuredClone(this.ledger); + const changed = mutator(draft); + if (changed === void 0) return { + ledger: deepFreeze(structuredClone(this.ledger)), + changed: [] + }; + for (const task of changed) pruneExecutions(task); + draft.revision += 1; + await persistAtomic(file, JSON.stringify(draft)); + this.ledger = draft; + const change = { + revision: draft.revision, + tasks: changed, + kind + }; + for (const fn of this.subscribers) try { + fn(change); + } catch {} + return { + ledger: deepFreeze(structuredClone(draft)), + changed: changed.map((t) => deepFreeze(structuredClone(t))) + }; + })); }; if (this.storageQueue !== void 0) return this.storageQueue.run(run); return this.queue = this.queue.then(run, run); @@ -173,6 +193,42 @@ var TaskStore = class { return this.queue = this.queue.then(run, run); } }; +const LOCK_RETRY_MS = 25; +const LOCK_TIMEOUT_MS = 3e4; +const STALE_LOCK_MS = 5 * 6e4; +/** +* `rename()` makes one replacement atomic, but not the preceding +* read-modify-write sequence. This exclusive sidecar lock protects that full +* sequence across hosts; a crashed owner is recoverable after five minutes. +*/ +async function withLedgerWriteLock(file, run) { + const lock = `${file}.lock`; + const token = `${process.pid}:${Date.now()}:${Math.random().toString(36).slice(2)}`; + const deadline = Date.now() + LOCK_TIMEOUT_MS; + await mkdir(dirname(file), { recursive: true }); + for (;;) try { + const handle = await open(lock, "wx"); + try { + await handle.writeFile(token, "utf8"); + } finally { + await handle.close(); + } + try { + return await run(); + } finally { + if (await readFile(lock, "utf8").catch(() => void 0) === token) await unlink(lock).catch(() => {}); + } + } catch (error) { + if (error.code !== "EEXIST") throw error; + const age = await stat(lock).then((info) => Date.now() - info.mtimeMs, () => void 0); + if (age !== void 0 && age > STALE_LOCK_MS) { + await unlink(lock).catch(() => {}); + continue; + } + if (Date.now() >= deadline) throw new Error(`timed out waiting for task ledger lock: ${file}`); + await new Promise((resolve) => setTimeout(resolve, LOCK_RETRY_MS)); + } +} /** Recursively freeze a plain-data value (defense in depth for handed-out snapshots). */ function deepFreeze(value) { if (value !== null && typeof value === "object") { diff --git a/lib/host/store.js.map b/lib/host/store.js.map index 585ef20..0b96fff 100644 --- a/lib/host/store.js.map +++ b/lib/host/store.js.map @@ -1 +1 @@ -{"version":3,"file":"store.js","names":[],"sources":["../../src/host/store.ts"],"sourcesContent":["/**\n * Host-side task ledger: one JSON file under the active data directory, mutated through a\n * serial write queue, published as immutable snapshots with a global\n * monotonic revision. Change subscribers (P2: SSE route) observe every\n * committed mutation.\n *\n * @module dsh-taskboard/host/store\n */\nimport { mkdir, open, readFile, rename } from 'node:fs/promises'\nimport { dirname, join } from 'node:path'\nimport {\n LEDGER_SCHEMA_VERSION,\n asBoardSettings,\n emptyLedger,\n isPlausibleTaskRecord,\n pruneExecutions,\n type TaskLedger,\n type TaskRecord,\n} from '../shared/protocol.ts'\nimport type { StorageQueue } from './storage-queue.ts'\n\n/** One committed ledger mutation, handed to change subscribers. */\nexport interface LedgerChange {\n /** Revision after the mutation. */\n revision: number\n /** The mutated tasks, if any (a comment purge may touch none). */\n tasks: readonly TaskRecord[]\n /** What kind of mutation this was (for SSE event naming later). */\n kind: 'task-created' | 'task-updated' | 'task-moved' | 'task-deleted' | 'comment-added' | 'execution-recorded' | 'settings-updated' | 'ledger-replaced'\n}\n\n/** Options for {@link TaskStore}. */\nexport interface TaskStoreOptions {\n /** Absolute ledger file path. */\n file: string\n /** Optional queue shared with templates/assets and storage migration. */\n queue?: StorageQueue\n}\n\n/**\n * The durable ledger. All mutations run through {@link mutate}, which:\n * validates the resulting document, bumps the global revision, persists\n * atomically (temp file + rename), and only then notifies subscribers.\n */\nexport class TaskStore {\n private file: string\n private readonly storageQueue?: StorageQueue\n private ledger: TaskLedger = emptyLedger()\n private readonly subscribers = new Set<(change: LedgerChange) => void>()\n private queue: Promise = Promise.resolve()\n private loaded = false\n private loadPromise: Promise | undefined\n\n /** @param options - file location. */\n constructor(options: TaskStoreOptions) {\n this.file = options.file\n this.storageQueue = options.queue\n }\n\n /** Current absolute ledger path. */\n location(): string { return this.file }\n\n /** Persist the live in-memory ledger to another file without switching. */\n async writeCopy(file: string): Promise {\n await this.load()\n await persistAtomic(file, JSON.stringify(this.ledger))\n }\n\n /** Switch future writes after a prepared migration commits. */\n setLocation(file: string): void { this.file = file }\n\n /** Load (once) from disk; a missing file starts empty; a corrupt file is quarantined, not thrown. */\n load(): Promise {\n if (this.loaded) return Promise.resolve()\n if (this.loadPromise !== undefined) return this.loadPromise\n this.loadPromise = this.loadOnce()\n return this.loadPromise\n }\n\n /** Perform the single physical ledger read shared by all startup callers. */\n private async loadOnce(): Promise {\n try {\n const raw = await readFile(this.file, 'utf8')\n const parsed = JSON.parse(raw) as TaskLedger\n if (typeof parsed.revision === 'number' && Array.isArray(parsed.tasks)) {\n // S11: trust no record wholesale — drop structurally broken entries\n // (including R4's traversal-shaped ids from a hand-edited file) with\n // a notice instead of letting them reach the path-building layers.\n const plausible: TaskRecord[] = []\n for (const entry of parsed.tasks as unknown[]) {\n if (!isPlausibleTaskRecord(entry)) {\n const rawId = (entry as { id?: unknown })?.id\n const id = typeof rawId === 'string' ? rawId.slice(0, 60) : String(rawId)\n console.warn('[dsh-taskboard] dropping implausible ledger entry on load:', id)\n continue\n }\n plausible.push(entry as TaskRecord)\n }\n const tasks = plausible\n // Migration from pre-claim-field ledgers: an agent-held in_progress\n // task carried its holder in updatedBy — backfill the explicit claim\n // fields so the hold survives user edits (updatedBy is audit-only).\n for (const task of tasks) {\n if (task.status === 'in_progress' && task.claimedBy === undefined\n && task.updatedBy?.kind === 'agent' && typeof task.updatedBy.sessionId === 'string') {\n task.claimedBy = task.updatedBy.sessionId\n task.claimedAt = task.updatedAt\n }\n }\n let settings = undefined\n if (parsed.settings !== undefined) {\n try {\n settings = asBoardSettings(parsed.settings)\n } catch {\n console.warn('[dsh-taskboard] dropping invalid board settings on load')\n }\n }\n this.ledger = {\n schemaVersion: LEDGER_SCHEMA_VERSION,\n revision: parsed.revision,\n tasks,\n ...(settings !== undefined ? { settings } : {}),\n }\n }\n } catch (error) {\n const code = (error as NodeJS.ErrnoException).code\n if (code !== 'ENOENT') {\n // Quarantine a corrupt ledger: rename it aside, start fresh. Never\n // take the host down over ledger damage.\n try {\n await rename(this.file, `${this.file}.corrupt-${Date.now()}`)\n } catch { /* best effort */ }\n }\n }\n this.loaded = true\n }\n\n /**\n * The current snapshot — a deep-frozen clone. Mutating the returned value\n * throws (strict mode) instead of silently bypassing the revision/persist\n * path; internal state is never handed out.\n */\n snapshot(): TaskLedger {\n return deepFreeze(structuredClone(this.ledger))\n }\n\n /** Find a task by id (frozen clone; internal state is never handed out). */\n get(id: string): TaskRecord | undefined {\n const task = this.ledger.tasks.find(t => t.id === id)\n return task === undefined ? undefined : deepFreeze(structuredClone(task))\n }\n\n /** Subscribe to committed changes; returns the unsubscribe. */\n subscribe(fn: (change: LedgerChange) => void): () => void {\n this.subscribers.add(fn)\n return () => this.subscribers.delete(fn)\n }\n\n /**\n * Write a timestamped backup copy of the current ledger next to the live\n * file (import-replace safety, 0.4.0). Never throws the caller's flow —\n * a backup failure fails the import itself.\n * @returns the backup file path.\n */\n async backup(): Promise {\n const run = async (): Promise => {\n await this.load()\n const target = `${this.file}.backup-${Date.now()}`\n await persistAtomic(target, JSON.stringify(this.ledger, null, 2))\n return target\n }\n return this.storageQueue === undefined ? run() : this.storageQueue.run(run)\n }\n\n /**\n * Run one mutation inside the serial queue. The mutator works on a\n * structured clone; returning `undefined` aborts with no write.\n * @param kind - change kind for subscribers.\n * @param mutator - receives the cloned ledger; mutate tasks in place; return the touched tasks.\n */\n async mutate(\n kind: LedgerChange['kind'],\n mutator: (ledger: TaskLedger) => TaskRecord[] | undefined,\n ): Promise<{ ledger: TaskLedger; changed: readonly TaskRecord[] }> {\n const run = async (): Promise<{ ledger: TaskLedger; changed: readonly TaskRecord[] }> => {\n await this.load()\n const draft: TaskLedger = structuredClone(this.ledger)\n const changed = mutator(draft)\n if (changed === undefined) {\n // S9 parity: even a no-op mutation hands out a frozen clone — never\n // the live internal ledger.\n return { ledger: deepFreeze(structuredClone(this.ledger)), changed: [] }\n }\n // Retention cap: every committed mutation re-checks the touched tasks,\n // so execution history can never grow unbounded (SSE state payload).\n for (const task of changed) pruneExecutions(task)\n draft.revision += 1\n const json = JSON.stringify(draft)\n await persistAtomic(this.file, json)\n this.ledger = draft\n const change: LedgerChange = { revision: draft.revision, tasks: changed, kind }\n for (const fn of this.subscribers) {\n try {\n fn(change)\n } catch { /* subscriber errors never abort the write */ }\n }\n // S9: hand out frozen clones — the return value used to BE the new\n // internal ledger; callers must never mutate internal state in place.\n return {\n ledger: deepFreeze(structuredClone(draft)),\n changed: changed.map(t => deepFreeze(structuredClone(t))),\n }\n }\n if (this.storageQueue !== undefined) return this.storageQueue.run(run)\n return (this.queue = this.queue.then(run, run)) as ReturnType\n }\n\n /**\n * Run a read INSIDE the serial queue (R3): observes exactly the ledger\n * state after all previously enqueued mutations — immune to the\n * write-then-publish window around `mutate`'s persistence. Read-only: the\n * callback receives a frozen deep clone and nothing is written.\n */\n async read(fn: (ledger: TaskLedger) => T): Promise {\n const run = async (): Promise => {\n await this.load()\n return fn(deepFreeze(structuredClone(this.ledger)))\n }\n if (this.storageQueue !== undefined) return this.storageQueue.run(run)\n return (this.queue = this.queue.then(run, run)) as Promise\n }\n}\n\n/** Recursively freeze a plain-data value (defense in depth for handed-out snapshots). */\nfunction deepFreeze(value: T): T {\n if (value !== null && typeof value === 'object') {\n if (!Object.isFrozen(value)) Object.freeze(value)\n for (const key of Object.keys(value as Record)) {\n deepFreeze((value as Record)[key])\n }\n }\n return value\n}\n\n/**\n * Atomic file persist: write temp, fsync, then rename over the target (S10:\n * without the sync, a power loss after rename can leave a zero-length file —\n * the next load would quarantine the ledger and start empty).\n */\nasync function persistAtomic(file: string, contents: string): Promise {\n await mkdir(dirname(file), { recursive: true })\n const temp = join(dirname(file), `.${Math.random().toString(36).slice(2)}.tmp`)\n const fh = await open(temp, 'w')\n try {\n await fh.writeFile(contents, 'utf8')\n await fh.sync()\n } finally {\n await fh.close()\n }\n await rename(temp, file)\n}\n"],"mappings":";;;;;;;;;;;;;;;;;AA4CA,IAAa,YAAb,MAAuB;CACrB;CACA;CACA,SAA6B,YAAY;CACzC,8BAA+B,IAAI,IAAoC;CACvE,QAAkC,QAAQ,QAAQ;CAClD,SAAiB;CACjB;;CAGA,YAAY,SAA2B;EACrC,KAAK,OAAO,QAAQ;EACpB,KAAK,eAAe,QAAQ;CAC9B;;CAGA,WAAmB;EAAE,OAAO,KAAK;CAAK;;CAGtC,MAAM,UAAU,MAA6B;EAC3C,MAAM,KAAK,KAAK;EAChB,MAAM,cAAc,MAAM,KAAK,UAAU,KAAK,MAAM,CAAC;CACvD;;CAGA,YAAY,MAAoB;EAAE,KAAK,OAAO;CAAK;;CAGnD,OAAsB;EACpB,IAAI,KAAK,QAAQ,OAAO,QAAQ,QAAQ;EACxC,IAAI,KAAK,gBAAgB,KAAA,GAAW,OAAO,KAAK;EAChD,KAAK,cAAc,KAAK,SAAS;EACjC,OAAO,KAAK;CACd;;CAGA,MAAc,WAA0B;EACtC,IAAI;GACF,MAAM,MAAM,MAAM,SAAS,KAAK,MAAM,MAAM;GAC5C,MAAM,SAAS,KAAK,MAAM,GAAG;GAC7B,IAAI,OAAO,OAAO,aAAa,YAAY,MAAM,QAAQ,OAAO,KAAK,GAAG;IAItE,MAAM,YAA0B,CAAC;IACjC,KAAK,MAAM,SAAS,OAAO,OAAoB;KAC7C,IAAI,CAAC,sBAAsB,KAAK,GAAG;MACjC,MAAM,QAAS,OAA4B;MAC3C,MAAM,KAAK,OAAO,UAAU,WAAW,MAAM,MAAM,GAAG,EAAE,IAAI,OAAO,KAAK;MACxE,QAAQ,KAAK,8DAA8D,EAAE;MAC7E;KACF;KACA,UAAU,KAAK,KAAmB;IACpC;IACA,MAAM,QAAQ;IAId,KAAK,MAAM,QAAQ,OACjB,IAAI,KAAK,WAAW,iBAAiB,KAAK,cAAc,KAAA,KACnD,KAAK,WAAW,SAAS,WAAW,OAAO,KAAK,UAAU,cAAc,UAAU;KACrF,KAAK,YAAY,KAAK,UAAU;KAChC,KAAK,YAAY,KAAK;IACxB;IAEF,IAAI,WAAW,KAAA;IACf,IAAI,OAAO,aAAa,KAAA,GACtB,IAAI;KACF,WAAW,gBAAgB,OAAO,QAAQ;IAC5C,QAAQ;KACN,QAAQ,KAAK,yDAAyD;IACxE;IAEF,KAAK,SAAS;KACZ,eAAA;KACA,UAAU,OAAO;KACjB;KACA,GAAI,aAAa,KAAA,IAAY,EAAE,SAAS,IAAI,CAAC;IAC/C;GACF;EACF,SAAS,OAAO;GAEd,IADc,MAAgC,SACjC,UAGX,IAAI;IACF,MAAM,OAAO,KAAK,MAAM,GAAG,KAAK,KAAK,WAAW,KAAK,IAAI,GAAG;GAC9D,QAAQ,CAAoB;EAEhC;EACA,KAAK,SAAS;CAChB;;;;;;CAOA,WAAuB;EACrB,OAAO,WAAW,gBAAgB,KAAK,MAAM,CAAC;CAChD;;CAGA,IAAI,IAAoC;EACtC,MAAM,OAAO,KAAK,OAAO,MAAM,MAAK,MAAK,EAAE,OAAO,EAAE;EACpD,OAAO,SAAS,KAAA,IAAY,KAAA,IAAY,WAAW,gBAAgB,IAAI,CAAC;CAC1E;;CAGA,UAAU,IAAgD;EACxD,KAAK,YAAY,IAAI,EAAE;EACvB,aAAa,KAAK,YAAY,OAAO,EAAE;CACzC;;;;;;;CAQA,MAAM,SAA0B;EAC9B,MAAM,MAAM,YAA6B;GACvC,MAAM,KAAK,KAAK;GAChB,MAAM,SAAS,GAAG,KAAK,KAAK,UAAU,KAAK,IAAI;GAC/C,MAAM,cAAc,QAAQ,KAAK,UAAU,KAAK,QAAQ,MAAM,CAAC,CAAC;GAChE,OAAO;EACT;EACA,OAAO,KAAK,iBAAiB,KAAA,IAAY,IAAI,IAAI,KAAK,aAAa,IAAI,GAAG;CAC5E;;;;;;;CAQA,MAAM,OACJ,MACA,SACiE;EACjE,MAAM,MAAM,YAA6E;GACvF,MAAM,KAAK,KAAK;GAChB,MAAM,QAAoB,gBAAgB,KAAK,MAAM;GACrD,MAAM,UAAU,QAAQ,KAAK;GAC7B,IAAI,YAAY,KAAA,GAGd,OAAO;IAAE,QAAQ,WAAW,gBAAgB,KAAK,MAAM,CAAC;IAAG,SAAS,CAAC;GAAE;GAIzE,KAAK,MAAM,QAAQ,SAAS,gBAAgB,IAAI;GAChD,MAAM,YAAY;GAClB,MAAM,OAAO,KAAK,UAAU,KAAK;GACjC,MAAM,cAAc,KAAK,MAAM,IAAI;GACnC,KAAK,SAAS;GACd,MAAM,SAAuB;IAAE,UAAU,MAAM;IAAU,OAAO;IAAS;GAAK;GAC9E,KAAK,MAAM,MAAM,KAAK,aACpB,IAAI;IACF,GAAG,MAAM;GACX,QAAQ,CAAgD;GAI1D,OAAO;IACL,QAAQ,WAAW,gBAAgB,KAAK,CAAC;IACzC,SAAS,QAAQ,KAAI,MAAK,WAAW,gBAAgB,CAAC,CAAC,CAAC;GAC1D;EACF;EACA,IAAI,KAAK,iBAAiB,KAAA,GAAW,OAAO,KAAK,aAAa,IAAI,GAAG;EACrE,OAAQ,KAAK,QAAQ,KAAK,MAAM,KAAK,KAAK,GAAG;CAC/C;;;;;;;CAQA,MAAM,KAAQ,IAA2C;EACvD,MAAM,MAAM,YAAwB;GAClC,MAAM,KAAK,KAAK;GAChB,OAAO,GAAG,WAAW,gBAAgB,KAAK,MAAM,CAAC,CAAC;EACpD;EACA,IAAI,KAAK,iBAAiB,KAAA,GAAW,OAAO,KAAK,aAAa,IAAI,GAAG;EACrE,OAAQ,KAAK,QAAQ,KAAK,MAAM,KAAK,KAAK,GAAG;CAC/C;AACF;;AAGA,SAAS,WAAc,OAAa;CAClC,IAAI,UAAU,QAAQ,OAAO,UAAU,UAAU;EAC/C,IAAI,CAAC,OAAO,SAAS,KAAK,GAAG,OAAO,OAAO,KAAK;EAChD,KAAK,MAAM,OAAO,OAAO,KAAK,KAAgC,GAC5D,WAAY,MAAkC,IAAI;CAEtD;CACA,OAAO;AACT;;;;;;AAOA,eAAe,cAAc,MAAc,UAAiC;CAC1E,MAAM,MAAM,QAAQ,IAAI,GAAG,EAAE,WAAW,KAAK,CAAC;CAC9C,MAAM,OAAO,KAAK,QAAQ,IAAI,GAAG,IAAI,KAAK,OAAO,CAAC,CAAC,SAAS,EAAE,CAAC,CAAC,MAAM,CAAC,EAAE,KAAK;CAC9E,MAAM,KAAK,MAAM,KAAK,MAAM,GAAG;CAC/B,IAAI;EACF,MAAM,GAAG,UAAU,UAAU,MAAM;EACnC,MAAM,GAAG,KAAK;CAChB,UAAU;EACR,MAAM,GAAG,MAAM;CACjB;CACA,MAAM,OAAO,MAAM,IAAI;AACzB"} \ No newline at end of file +{"version":3,"file":"store.js","names":[],"sources":["../../src/host/store.ts"],"sourcesContent":["/**\n * Host-side task ledger: one JSON file under the active data directory, mutated through a\n * serial write queue, published as immutable snapshots with a global\n * monotonic revision. Change subscribers (P2: SSE route) observe every\n * committed mutation.\n *\n * @module dsh-taskboard/host/store\n */\nimport { mkdir, open, readFile, rename, stat, unlink } from 'node:fs/promises'\nimport { dirname, join } from 'node:path'\nimport {\n LEDGER_SCHEMA_VERSION,\n asBoardSettings,\n emptyLedger,\n isPlausibleTaskRecord,\n pruneExecutions,\n type TaskLedger,\n type TaskRecord,\n} from '../shared/protocol.ts'\nimport type { StorageQueue } from './storage-queue.ts'\n\n/** One committed ledger mutation, handed to change subscribers. */\nexport interface LedgerChange {\n /** Revision after the mutation. */\n revision: number\n /** The mutated tasks, if any (a comment purge may touch none). */\n tasks: readonly TaskRecord[]\n /** What kind of mutation this was (for SSE event naming later). */\n kind: 'task-created' | 'task-updated' | 'task-moved' | 'task-deleted' | 'comment-added' | 'execution-recorded' | 'settings-updated' | 'ledger-replaced'\n}\n\n/** Options for {@link TaskStore}. */\nexport interface TaskStoreOptions {\n /** Absolute ledger file path. */\n file: string\n /** Optional queue shared with templates/assets and storage migration. */\n queue?: StorageQueue\n}\n\n/**\n * The durable ledger. All mutations run through {@link mutate}, which:\n * validates the resulting document, bumps the global revision, persists\n * atomically (temp file + rename), and only then notifies subscribers.\n */\nexport class TaskStore {\n /** Stores from old and new plugin generations share one write lane per file. */\n private static readonly fileQueues = new Map>()\n\n private file: string\n private readonly storageQueue?: StorageQueue\n private ledger: TaskLedger = emptyLedger()\n private readonly subscribers = new Set<(change: LedgerChange) => void>()\n private queue: Promise = Promise.resolve()\n private loaded = false\n private loadPromise: Promise | undefined\n\n /** @param options - file location. */\n constructor(options: TaskStoreOptions) {\n this.file = options.file\n this.storageQueue = options.queue\n }\n\n /** Current absolute ledger path. */\n location(): string { return this.file }\n\n /** Persist the live in-memory ledger to another file without switching. */\n async writeCopy(file: string): Promise {\n await this.load()\n await persistAtomic(file, JSON.stringify(this.ledger))\n }\n\n /** Re-read through the normalizer while holding the cross-instance write lock. */\n private async reload(): Promise {\n this.loaded = false\n this.loadPromise = undefined\n await this.load()\n }\n\n /** Switch future writes after a prepared migration commits. */\n setLocation(file: string): void { this.file = file }\n\n /** Serialize one operation with every TaskStore in this host for this file. */\n private inProcessWriteLane(file: string, run: () => Promise): Promise {\n const previous = TaskStore.fileQueues.get(file) ?? Promise.resolve()\n const next = previous.then(run, run)\n TaskStore.fileQueues.set(file, next)\n void next.finally(() => {\n if (TaskStore.fileQueues.get(file) === next) TaskStore.fileQueues.delete(file)\n }).catch(() => { /* the caller receives the original rejection */ })\n return next\n }\n\n /** Load (once) from disk; a missing file starts empty; a corrupt file is quarantined, not thrown. */\n load(): Promise {\n if (this.loaded) return Promise.resolve()\n if (this.loadPromise !== undefined) return this.loadPromise\n this.loadPromise = this.loadOnce()\n return this.loadPromise\n }\n\n /** Perform the single physical ledger read shared by all startup callers. */\n private async loadOnce(): Promise {\n try {\n const raw = await readFile(this.file, 'utf8')\n const parsed = JSON.parse(raw) as TaskLedger\n if (typeof parsed.revision === 'number' && Array.isArray(parsed.tasks)) {\n // S11: trust no record wholesale — drop structurally broken entries\n // (including R4's traversal-shaped ids from a hand-edited file) with\n // a notice instead of letting them reach the path-building layers.\n const plausible: TaskRecord[] = []\n for (const entry of parsed.tasks as unknown[]) {\n if (!isPlausibleTaskRecord(entry)) {\n const rawId = (entry as { id?: unknown })?.id\n const id = typeof rawId === 'string' ? rawId.slice(0, 60) : String(rawId)\n console.warn('[dsh-taskboard] dropping implausible ledger entry on load:', id)\n continue\n }\n plausible.push(entry as TaskRecord)\n }\n const tasks = plausible\n // Migration from pre-claim-field ledgers: an agent-held in_progress\n // task carried its holder in updatedBy — backfill the explicit claim\n // fields so the hold survives user edits (updatedBy is audit-only).\n for (const task of tasks) {\n if (task.status === 'in_progress' && task.claimedBy === undefined\n && task.updatedBy?.kind === 'agent' && typeof task.updatedBy.sessionId === 'string') {\n task.claimedBy = task.updatedBy.sessionId\n task.claimedAt = task.updatedAt\n }\n }\n let settings = undefined\n if (parsed.settings !== undefined) {\n try {\n settings = asBoardSettings(parsed.settings)\n } catch {\n console.warn('[dsh-taskboard] dropping invalid board settings on load')\n }\n }\n this.ledger = {\n schemaVersion: LEDGER_SCHEMA_VERSION,\n revision: parsed.revision,\n tasks,\n ...(settings !== undefined ? { settings } : {}),\n }\n }\n } catch (error) {\n const code = (error as NodeJS.ErrnoException).code\n if (code !== 'ENOENT') {\n // Quarantine a corrupt ledger: rename it aside, start fresh. Never\n // take the host down over ledger damage.\n try {\n await rename(this.file, `${this.file}.corrupt-${Date.now()}`)\n } catch { /* best effort */ }\n }\n }\n this.loaded = true\n }\n\n /**\n * The current snapshot — a deep-frozen clone. Mutating the returned value\n * throws (strict mode) instead of silently bypassing the revision/persist\n * path; internal state is never handed out.\n */\n snapshot(): TaskLedger {\n return deepFreeze(structuredClone(this.ledger))\n }\n\n /** Find a task by id (frozen clone; internal state is never handed out). */\n get(id: string): TaskRecord | undefined {\n const task = this.ledger.tasks.find(t => t.id === id)\n return task === undefined ? undefined : deepFreeze(structuredClone(task))\n }\n\n /** Subscribe to committed changes; returns the unsubscribe. */\n subscribe(fn: (change: LedgerChange) => void): () => void {\n this.subscribers.add(fn)\n return () => this.subscribers.delete(fn)\n }\n\n /**\n * Write a timestamped backup copy of the current ledger next to the live\n * file (import-replace safety, 0.4.0). Never throws the caller's flow —\n * a backup failure fails the import itself.\n * @returns the backup file path.\n */\n async backup(): Promise {\n const run = async (): Promise => {\n await this.load()\n const target = `${this.file}.backup-${Date.now()}`\n await persistAtomic(target, JSON.stringify(this.ledger, null, 2))\n return target\n }\n return this.storageQueue === undefined ? run() : this.storageQueue.run(run)\n }\n\n /**\n * Run one mutation inside the serial queue. The mutator works on a\n * structured clone; returning `undefined` aborts with no write.\n * @param kind - change kind for subscribers.\n * @param mutator - receives the cloned ledger; mutate tasks in place; return the touched tasks.\n */\n async mutate(\n kind: LedgerChange['kind'],\n mutator: (ledger: TaskLedger) => TaskRecord[] | undefined,\n ): Promise<{ ledger: TaskLedger; changed: readonly TaskRecord[] }> {\n const run = async (): Promise<{ ledger: TaskLedger; changed: readonly TaskRecord[] }> => {\n await this.load()\n const file = this.file\n return this.inProcessWriteLane(file, async () => withLedgerWriteLock(file, async () => {\n // A different TaskStore may have committed while this instance waited.\n // Reload under the lock: a cached-revision check before it leaves a\n // TOCTOU window in which two writers emit the same next revision.\n await this.reload()\n const draft: TaskLedger = structuredClone(this.ledger)\n const changed = mutator(draft)\n if (changed === undefined) return { ledger: deepFreeze(structuredClone(this.ledger)), changed: [] }\n for (const task of changed) pruneExecutions(task)\n draft.revision += 1\n await persistAtomic(file, JSON.stringify(draft))\n this.ledger = draft\n const change: LedgerChange = { revision: draft.revision, tasks: changed, kind }\n for (const fn of this.subscribers) {\n try {\n fn(change)\n } catch { /* subscriber errors never abort the write */ }\n }\n return {\n ledger: deepFreeze(structuredClone(draft)),\n changed: changed.map(t => deepFreeze(structuredClone(t))),\n }\n }))\n }\n if (this.storageQueue !== undefined) return this.storageQueue.run(run)\n return (this.queue = this.queue.then(run, run)) as ReturnType\n }\n\n /**\n * Run a read INSIDE the serial queue (R3): observes exactly the ledger\n * state after all previously enqueued mutations — immune to the\n * write-then-publish window around `mutate`'s persistence. Read-only: the\n * callback receives a frozen deep clone and nothing is written.\n */\n async read(fn: (ledger: TaskLedger) => T): Promise {\n const run = async (): Promise => {\n await this.load()\n return fn(deepFreeze(structuredClone(this.ledger)))\n }\n if (this.storageQueue !== undefined) return this.storageQueue.run(run)\n return (this.queue = this.queue.then(run, run)) as Promise\n }\n}\n\nconst LOCK_RETRY_MS = 25\nconst LOCK_TIMEOUT_MS = 30_000\nconst STALE_LOCK_MS = 5 * 60_000\n\n/**\n * `rename()` makes one replacement atomic, but not the preceding\n * read-modify-write sequence. This exclusive sidecar lock protects that full\n * sequence across hosts; a crashed owner is recoverable after five minutes.\n */\nasync function withLedgerWriteLock(file: string, run: () => Promise): Promise {\n const lock = `${file}.lock`\n const token = `${process.pid}:${Date.now()}:${Math.random().toString(36).slice(2)}`\n const deadline = Date.now() + LOCK_TIMEOUT_MS\n await mkdir(dirname(file), { recursive: true })\n for (;;) {\n try {\n const handle = await open(lock, 'wx')\n try {\n await handle.writeFile(token, 'utf8')\n } finally {\n await handle.close()\n }\n try {\n return await run()\n } finally {\n // Do not delete a new owner's lock if ours was reclaimed while a\n // slow filesystem operation was in flight.\n const held = await readFile(lock, 'utf8').catch(() => undefined)\n if (held === token) await unlink(lock).catch(() => { /* best effort */ })\n }\n } catch (error) {\n if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error\n const age = await stat(lock).then(info => Date.now() - info.mtimeMs, () => undefined)\n if (age !== undefined && age > STALE_LOCK_MS) {\n await unlink(lock).catch(() => { /* another writer may have recovered it */ })\n continue\n }\n if (Date.now() >= deadline) throw new Error(`timed out waiting for task ledger lock: ${file}`)\n await new Promise(resolve => setTimeout(resolve, LOCK_RETRY_MS))\n }\n }\n}\n\n/** Recursively freeze a plain-data value (defense in depth for handed-out snapshots). */\nfunction deepFreeze(value: T): T {\n if (value !== null && typeof value === 'object') {\n if (!Object.isFrozen(value)) Object.freeze(value)\n for (const key of Object.keys(value as Record)) {\n deepFreeze((value as Record)[key])\n }\n }\n return value\n}\n\n/**\n * Atomic file persist: write temp, fsync, then rename over the target (S10:\n * without the sync, a power loss after rename can leave a zero-length file —\n * the next load would quarantine the ledger and start empty).\n */\nasync function persistAtomic(file: string, contents: string): Promise {\n await mkdir(dirname(file), { recursive: true })\n const temp = join(dirname(file), `.${Math.random().toString(36).slice(2)}.tmp`)\n const fh = await open(temp, 'w')\n try {\n await fh.writeFile(contents, 'utf8')\n await fh.sync()\n } finally {\n await fh.close()\n }\n await rename(temp, file)\n}\n"],"mappings":";;;;;;;;;;;;;;;;;AA4CA,IAAa,YAAb,MAAa,UAAU;;CAErB,OAAwB,6BAAa,IAAI,IAA8B;CAEvE;CACA;CACA,SAA6B,YAAY;CACzC,8BAA+B,IAAI,IAAoC;CACvE,QAAkC,QAAQ,QAAQ;CAClD,SAAiB;CACjB;;CAGA,YAAY,SAA2B;EACrC,KAAK,OAAO,QAAQ;EACpB,KAAK,eAAe,QAAQ;CAC9B;;CAGA,WAAmB;EAAE,OAAO,KAAK;CAAK;;CAGtC,MAAM,UAAU,MAA6B;EAC3C,MAAM,KAAK,KAAK;EAChB,MAAM,cAAc,MAAM,KAAK,UAAU,KAAK,MAAM,CAAC;CACvD;;CAGA,MAAc,SAAwB;EACpC,KAAK,SAAS;EACd,KAAK,cAAc,KAAA;EACnB,MAAM,KAAK,KAAK;CAClB;;CAGA,YAAY,MAAoB;EAAE,KAAK,OAAO;CAAK;;CAGnD,mBAA8B,MAAc,KAAmC;EAE7E,MAAM,QADW,UAAU,WAAW,IAAI,IAAI,KAAK,QAAQ,QAAQ,EAAA,CAC7C,KAAK,KAAK,GAAG;EACnC,UAAU,WAAW,IAAI,MAAM,IAAI;EACnC,KAAU,cAAc;GACtB,IAAI,UAAU,WAAW,IAAI,IAAI,MAAM,MAAM,UAAU,WAAW,OAAO,IAAI;EAC/E,CAAC,CAAC,CAAC,YAAY,CAAmD,CAAC;EACnE,OAAO;CACT;;CAGA,OAAsB;EACpB,IAAI,KAAK,QAAQ,OAAO,QAAQ,QAAQ;EACxC,IAAI,KAAK,gBAAgB,KAAA,GAAW,OAAO,KAAK;EAChD,KAAK,cAAc,KAAK,SAAS;EACjC,OAAO,KAAK;CACd;;CAGA,MAAc,WAA0B;EACtC,IAAI;GACF,MAAM,MAAM,MAAM,SAAS,KAAK,MAAM,MAAM;GAC5C,MAAM,SAAS,KAAK,MAAM,GAAG;GAC7B,IAAI,OAAO,OAAO,aAAa,YAAY,MAAM,QAAQ,OAAO,KAAK,GAAG;IAItE,MAAM,YAA0B,CAAC;IACjC,KAAK,MAAM,SAAS,OAAO,OAAoB;KAC7C,IAAI,CAAC,sBAAsB,KAAK,GAAG;MACjC,MAAM,QAAS,OAA4B;MAC3C,MAAM,KAAK,OAAO,UAAU,WAAW,MAAM,MAAM,GAAG,EAAE,IAAI,OAAO,KAAK;MACxE,QAAQ,KAAK,8DAA8D,EAAE;MAC7E;KACF;KACA,UAAU,KAAK,KAAmB;IACpC;IACA,MAAM,QAAQ;IAId,KAAK,MAAM,QAAQ,OACjB,IAAI,KAAK,WAAW,iBAAiB,KAAK,cAAc,KAAA,KACnD,KAAK,WAAW,SAAS,WAAW,OAAO,KAAK,UAAU,cAAc,UAAU;KACrF,KAAK,YAAY,KAAK,UAAU;KAChC,KAAK,YAAY,KAAK;IACxB;IAEF,IAAI,WAAW,KAAA;IACf,IAAI,OAAO,aAAa,KAAA,GACtB,IAAI;KACF,WAAW,gBAAgB,OAAO,QAAQ;IAC5C,QAAQ;KACN,QAAQ,KAAK,yDAAyD;IACxE;IAEF,KAAK,SAAS;KACZ,eAAA;KACA,UAAU,OAAO;KACjB;KACA,GAAI,aAAa,KAAA,IAAY,EAAE,SAAS,IAAI,CAAC;IAC/C;GACF;EACF,SAAS,OAAO;GAEd,IADc,MAAgC,SACjC,UAGX,IAAI;IACF,MAAM,OAAO,KAAK,MAAM,GAAG,KAAK,KAAK,WAAW,KAAK,IAAI,GAAG;GAC9D,QAAQ,CAAoB;EAEhC;EACA,KAAK,SAAS;CAChB;;;;;;CAOA,WAAuB;EACrB,OAAO,WAAW,gBAAgB,KAAK,MAAM,CAAC;CAChD;;CAGA,IAAI,IAAoC;EACtC,MAAM,OAAO,KAAK,OAAO,MAAM,MAAK,MAAK,EAAE,OAAO,EAAE;EACpD,OAAO,SAAS,KAAA,IAAY,KAAA,IAAY,WAAW,gBAAgB,IAAI,CAAC;CAC1E;;CAGA,UAAU,IAAgD;EACxD,KAAK,YAAY,IAAI,EAAE;EACvB,aAAa,KAAK,YAAY,OAAO,EAAE;CACzC;;;;;;;CAQA,MAAM,SAA0B;EAC9B,MAAM,MAAM,YAA6B;GACvC,MAAM,KAAK,KAAK;GAChB,MAAM,SAAS,GAAG,KAAK,KAAK,UAAU,KAAK,IAAI;GAC/C,MAAM,cAAc,QAAQ,KAAK,UAAU,KAAK,QAAQ,MAAM,CAAC,CAAC;GAChE,OAAO;EACT;EACA,OAAO,KAAK,iBAAiB,KAAA,IAAY,IAAI,IAAI,KAAK,aAAa,IAAI,GAAG;CAC5E;;;;;;;CAQA,MAAM,OACJ,MACA,SACiE;EACjE,MAAM,MAAM,YAA6E;GACvF,MAAM,KAAK,KAAK;GAChB,MAAM,OAAO,KAAK;GAClB,OAAO,KAAK,mBAAmB,MAAM,YAAY,oBAAoB,MAAM,YAAY;IAIrF,MAAM,KAAK,OAAO;IAClB,MAAM,QAAoB,gBAAgB,KAAK,MAAM;IACrD,MAAM,UAAU,QAAQ,KAAK;IAC7B,IAAI,YAAY,KAAA,GAAW,OAAO;KAAE,QAAQ,WAAW,gBAAgB,KAAK,MAAM,CAAC;KAAG,SAAS,CAAC;IAAE;IAClG,KAAK,MAAM,QAAQ,SAAS,gBAAgB,IAAI;IAChD,MAAM,YAAY;IAClB,MAAM,cAAc,MAAM,KAAK,UAAU,KAAK,CAAC;IAC/C,KAAK,SAAS;IACd,MAAM,SAAuB;KAAE,UAAU,MAAM;KAAU,OAAO;KAAS;IAAK;IAC9E,KAAK,MAAM,MAAM,KAAK,aACpB,IAAI;KACF,GAAG,MAAM;IACX,QAAQ,CAAgD;IAE1D,OAAO;KACL,QAAQ,WAAW,gBAAgB,KAAK,CAAC;KACzC,SAAS,QAAQ,KAAI,MAAK,WAAW,gBAAgB,CAAC,CAAC,CAAC;IAC1D;GACF,CAAC,CAAC;EACJ;EACA,IAAI,KAAK,iBAAiB,KAAA,GAAW,OAAO,KAAK,aAAa,IAAI,GAAG;EACrE,OAAQ,KAAK,QAAQ,KAAK,MAAM,KAAK,KAAK,GAAG;CAC/C;;;;;;;CAQA,MAAM,KAAQ,IAA2C;EACvD,MAAM,MAAM,YAAwB;GAClC,MAAM,KAAK,KAAK;GAChB,OAAO,GAAG,WAAW,gBAAgB,KAAK,MAAM,CAAC,CAAC;EACpD;EACA,IAAI,KAAK,iBAAiB,KAAA,GAAW,OAAO,KAAK,aAAa,IAAI,GAAG;EACrE,OAAQ,KAAK,QAAQ,KAAK,MAAM,KAAK,KAAK,GAAG;CAC/C;AACF;AAEA,MAAM,gBAAgB;AACtB,MAAM,kBAAkB;AACxB,MAAM,gBAAgB,IAAI;;;;;;AAO1B,eAAe,oBAAuB,MAAc,KAAmC;CACrF,MAAM,OAAO,GAAG,KAAK;CACrB,MAAM,QAAQ,GAAG,QAAQ,IAAI,GAAG,KAAK,IAAI,EAAE,GAAG,KAAK,OAAO,CAAC,CAAC,SAAS,EAAE,CAAC,CAAC,MAAM,CAAC;CAChF,MAAM,WAAW,KAAK,IAAI,IAAI;CAC9B,MAAM,MAAM,QAAQ,IAAI,GAAG,EAAE,WAAW,KAAK,CAAC;CAC9C,SACE,IAAI;EACF,MAAM,SAAS,MAAM,KAAK,MAAM,IAAI;EACpC,IAAI;GACF,MAAM,OAAO,UAAU,OAAO,MAAM;EACtC,UAAU;GACR,MAAM,OAAO,MAAM;EACrB;EACA,IAAI;GACF,OAAO,MAAM,IAAI;EACnB,UAAU;GAIR,IAAI,MADe,SAAS,MAAM,MAAM,CAAC,CAAC,YAAY,KAAA,CAAS,MAClD,OAAO,MAAM,OAAO,IAAI,CAAC,CAAC,YAAY,CAAoB,CAAC;EAC1E;CACF,SAAS,OAAO;EACd,IAAK,MAAgC,SAAS,UAAU,MAAM;EAC9D,MAAM,MAAM,MAAM,KAAK,IAAI,CAAC,CAAC,MAAK,SAAQ,KAAK,IAAI,IAAI,KAAK,eAAe,KAAA,CAAS;EACpF,IAAI,QAAQ,KAAA,KAAa,MAAM,eAAe;GAC5C,MAAM,OAAO,IAAI,CAAC,CAAC,YAAY,CAA6C,CAAC;GAC7E;EACF;EACA,IAAI,KAAK,IAAI,KAAK,UAAU,MAAM,IAAI,MAAM,2CAA2C,MAAM;EAC7F,MAAM,IAAI,SAAc,YAAW,WAAW,SAAS,aAAa,CAAC;CACvE;AAEJ;;AAGA,SAAS,WAAc,OAAa;CAClC,IAAI,UAAU,QAAQ,OAAO,UAAU,UAAU;EAC/C,IAAI,CAAC,OAAO,SAAS,KAAK,GAAG,OAAO,OAAO,KAAK;EAChD,KAAK,MAAM,OAAO,OAAO,KAAK,KAAgC,GAC5D,WAAY,MAAkC,IAAI;CAEtD;CACA,OAAO;AACT;;;;;;AAOA,eAAe,cAAc,MAAc,UAAiC;CAC1E,MAAM,MAAM,QAAQ,IAAI,GAAG,EAAE,WAAW,KAAK,CAAC;CAC9C,MAAM,OAAO,KAAK,QAAQ,IAAI,GAAG,IAAI,KAAK,OAAO,CAAC,CAAC,SAAS,EAAE,CAAC,CAAC,MAAM,CAAC,EAAE,KAAK;CAC9E,MAAM,KAAK,MAAM,KAAK,MAAM,GAAG;CAC/B,IAAI;EACF,MAAM,GAAG,UAAU,UAAU,MAAM;EACnC,MAAM,GAAG,KAAK;CAChB,UAAU;EACR,MAAM,GAAG,MAAM;CACjB;CACA,MAAM,OAAO,MAAM,IAAI;AACzB"} \ No newline at end of file diff --git a/src/host/store.ts b/src/host/store.ts index 0b95f5f..099f65a 100644 --- a/src/host/store.ts +++ b/src/host/store.ts @@ -6,7 +6,7 @@ * * @module dsh-taskboard/host/store */ -import { mkdir, open, readFile, rename } from 'node:fs/promises' +import { mkdir, open, readFile, rename, stat, unlink } from 'node:fs/promises' import { dirname, join } from 'node:path' import { LEDGER_SCHEMA_VERSION, @@ -43,6 +43,9 @@ export interface TaskStoreOptions { * atomically (temp file + rename), and only then notifies subscribers. */ export class TaskStore { + /** Stores from old and new plugin generations share one write lane per file. */ + private static readonly fileQueues = new Map>() + private file: string private readonly storageQueue?: StorageQueue private ledger: TaskLedger = emptyLedger() @@ -66,9 +69,27 @@ export class TaskStore { await persistAtomic(file, JSON.stringify(this.ledger)) } + /** Re-read through the normalizer while holding the cross-instance write lock. */ + private async reload(): Promise { + this.loaded = false + this.loadPromise = undefined + await this.load() + } + /** Switch future writes after a prepared migration commits. */ setLocation(file: string): void { this.file = file } + /** Serialize one operation with every TaskStore in this host for this file. */ + private inProcessWriteLane(file: string, run: () => Promise): Promise { + const previous = TaskStore.fileQueues.get(file) ?? Promise.resolve() + const next = previous.then(run, run) + TaskStore.fileQueues.set(file, next) + void next.finally(() => { + if (TaskStore.fileQueues.get(file) === next) TaskStore.fileQueues.delete(file) + }).catch(() => { /* the caller receives the original rejection */ }) + return next + } + /** Load (once) from disk; a missing file starts empty; a corrupt file is quarantined, not thrown. */ load(): Promise { if (this.loaded) return Promise.resolve() @@ -184,32 +205,30 @@ export class TaskStore { ): Promise<{ ledger: TaskLedger; changed: readonly TaskRecord[] }> { const run = async (): Promise<{ ledger: TaskLedger; changed: readonly TaskRecord[] }> => { await this.load() - const draft: TaskLedger = structuredClone(this.ledger) - const changed = mutator(draft) - if (changed === undefined) { - // S9 parity: even a no-op mutation hands out a frozen clone — never - // the live internal ledger. - return { ledger: deepFreeze(structuredClone(this.ledger)), changed: [] } - } - // Retention cap: every committed mutation re-checks the touched tasks, - // so execution history can never grow unbounded (SSE state payload). - for (const task of changed) pruneExecutions(task) - draft.revision += 1 - const json = JSON.stringify(draft) - await persistAtomic(this.file, json) - this.ledger = draft - const change: LedgerChange = { revision: draft.revision, tasks: changed, kind } - for (const fn of this.subscribers) { - try { - fn(change) - } catch { /* subscriber errors never abort the write */ } - } - // S9: hand out frozen clones — the return value used to BE the new - // internal ledger; callers must never mutate internal state in place. - return { - ledger: deepFreeze(structuredClone(draft)), - changed: changed.map(t => deepFreeze(structuredClone(t))), - } + const file = this.file + return this.inProcessWriteLane(file, async () => withLedgerWriteLock(file, async () => { + // A different TaskStore may have committed while this instance waited. + // Reload under the lock: a cached-revision check before it leaves a + // TOCTOU window in which two writers emit the same next revision. + await this.reload() + const draft: TaskLedger = structuredClone(this.ledger) + const changed = mutator(draft) + if (changed === undefined) return { ledger: deepFreeze(structuredClone(this.ledger)), changed: [] } + for (const task of changed) pruneExecutions(task) + draft.revision += 1 + await persistAtomic(file, JSON.stringify(draft)) + this.ledger = draft + const change: LedgerChange = { revision: draft.revision, tasks: changed, kind } + for (const fn of this.subscribers) { + try { + fn(change) + } catch { /* subscriber errors never abort the write */ } + } + return { + ledger: deepFreeze(structuredClone(draft)), + changed: changed.map(t => deepFreeze(structuredClone(t))), + } + })) } if (this.storageQueue !== undefined) return this.storageQueue.run(run) return (this.queue = this.queue.then(run, run)) as ReturnType @@ -231,6 +250,49 @@ export class TaskStore { } } +const LOCK_RETRY_MS = 25 +const LOCK_TIMEOUT_MS = 30_000 +const STALE_LOCK_MS = 5 * 60_000 + +/** + * `rename()` makes one replacement atomic, but not the preceding + * read-modify-write sequence. This exclusive sidecar lock protects that full + * sequence across hosts; a crashed owner is recoverable after five minutes. + */ +async function withLedgerWriteLock(file: string, run: () => Promise): Promise { + const lock = `${file}.lock` + const token = `${process.pid}:${Date.now()}:${Math.random().toString(36).slice(2)}` + const deadline = Date.now() + LOCK_TIMEOUT_MS + await mkdir(dirname(file), { recursive: true }) + for (;;) { + try { + const handle = await open(lock, 'wx') + try { + await handle.writeFile(token, 'utf8') + } finally { + await handle.close() + } + try { + return await run() + } finally { + // Do not delete a new owner's lock if ours was reclaimed while a + // slow filesystem operation was in flight. + const held = await readFile(lock, 'utf8').catch(() => undefined) + if (held === token) await unlink(lock).catch(() => { /* best effort */ }) + } + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error + const age = await stat(lock).then(info => Date.now() - info.mtimeMs, () => undefined) + if (age !== undefined && age > STALE_LOCK_MS) { + await unlink(lock).catch(() => { /* another writer may have recovered it */ }) + continue + } + if (Date.now() >= deadline) throw new Error(`timed out waiting for task ledger lock: ${file}`) + await new Promise(resolve => setTimeout(resolve, LOCK_RETRY_MS)) + } + } +} + /** Recursively freeze a plain-data value (defense in depth for handed-out snapshots). */ function deepFreeze(value: T): T { if (value !== null && typeof value === 'object') { diff --git a/tests/store-staleness.spec.ts b/tests/store-staleness.spec.ts new file mode 100644 index 0000000..c816ebf --- /dev/null +++ b/tests/store-staleness.spec.ts @@ -0,0 +1,117 @@ +import { mkdtemp, readFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { describe, expect, it } from 'vitest' +import { TaskStore } from '../src/host/store.ts' +import type { TaskRecord } from '../src/shared/protocol.ts' + +/** + * The ledger file is shared state: two store instances can point at the same + * file (a plugin hot reload leaves the previous generation's long-lived + * callbacks running). Because a store caches its snapshot forever and + * `mutate` persists the whole document, a stale instance used to roll a newer + * commit back silently — same revision, no error, invisible to subscribers. + * These tests pin the guard that prevents exactly that. + */ + +/** A minimal task that satisfies the ledger's plausibility rules. */ +function task(id: string): TaskRecord { + return { + id, title: 'shared-ledger fixture', description: '', prompt: 'Do the thing', + workspaceId: 'project', urgency: 'normal', status: 'todo', blocked: false, + isolation: 'none', execution: { mode: 'claim' }, + version: 1, createdAt: 1, updatedAt: 1, createdBy: { kind: 'user' }, updatedBy: { kind: 'user' }, + comments: [], executions: [], + } +} + +/** Create a store backed by a real file that already contains one task. */ +async function fixture() { + const root = await mkdtemp(join(tmpdir(), 'taskboard-store-')) + const file = join(root, 'dsh-taskboard.json') + const seed = new TaskStore({ file }) + await seed.load() + const first = task('t-shared-1') + await seed.mutate('task-created', ledger => { ledger.tasks.push(first); return [first] }) + return { file } +} + +/** Read the ledger straight from disk (bypassing any instance's cache). */ +async function onDisk(file: string) { + return JSON.parse(await readFile(file, 'utf8')) as { + revision: number + tasks: Array<{ comments: Array<{ id: string }>; version: number }> + } +} + +/** Append a marker comment to the first task. */ +function appendMarker(store: TaskStore, id: string) { + return store.mutate('comment-added', ledger => { + const first = ledger.tasks[0] + if (first === undefined) return undefined + first.comments.push({ id } as (typeof first.comments)[number]) + return [first] + }) +} + +describe('TaskStore shared-file staleness', () => { + it('does not roll back a commit made by another instance since we loaded', async () => { + const { file } = await fixture() + const older = new TaskStore({ file }) + const newer = new TaskStore({ file }) + await older.load() + await newer.load() + + // The newer instance commits a marker… + await appendMarker(newer, 'c-newer') + expect((await onDisk(file)).tasks[0]!.comments.map(c => c.id)).toContain('c-newer') + + // …then the STALE instance writes. It must re-read first, not clobber. + await appendMarker(older, 'c-older') + + const ids = (await onDisk(file)).tasks[0]!.comments.map(c => c.id) + expect(ids).toContain('c-newer') + expect(ids).toContain('c-older') + }) + + it('keeps the revision monotonic across interleaved writers', async () => { + const { file } = await fixture() + const a = new TaskStore({ file }) + const b = new TaskStore({ file }) + await a.load() + await b.load() + + const revisions: number[] = [] + for (let i = 0; i < 4; i++) { + // Alternate writers; neither re-loads from scratch in between. + const store = i % 2 === 0 ? a : b + const result = await store.mutate('task-updated', ledger => { + const first = ledger.tasks[0] + if (first === undefined) return undefined + first.version += 1 + return [first] + }) + revisions.push(result.ledger.revision) + } + + // Strictly increasing ⇒ no writer reused a stale counter. + for (let i = 1; i < revisions.length; i++) { + expect(revisions[i]!).toBeGreaterThan(revisions[i - 1]!) + } + }) + + it('preserves both updates when stale stores mutate concurrently', async () => { + const { file } = await fixture() + const a = new TaskStore({ file }) + const b = new TaskStore({ file }) + await Promise.all([a.load(), b.load()]) + + // Both stores start from the same revision. A write-before-refresh guard + // loses one marker here because both reads can finish before either write. + await Promise.all([appendMarker(a, 'c-a'), appendMarker(b, 'c-b')]) + + const disk = await onDisk(file) + expect(disk.revision).toBe(3) + expect(disk.tasks[0]!.comments.map(c => c.id)).toEqual(expect.arrayContaining(['c-a', 'c-b'])) + }) +})