Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
106 changes: 81 additions & 25 deletions lib/host/store.js

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion lib/host/store.js.map

Large diffs are not rendered by default.

116 changes: 89 additions & 27 deletions src/host/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<string, Promise<unknown>>()

private file: string
private readonly storageQueue?: StorageQueue
private ledger: TaskLedger = emptyLedger()
Expand All @@ -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<void> {
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<T>(file: string, run: () => Promise<T>): Promise<T> {
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<void> {
if (this.loaded) return Promise.resolve()
Expand Down Expand Up @@ -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<typeof run>
Expand All @@ -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<T>(file: string, run: () => Promise<T>): Promise<T> {
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<void>(resolve => setTimeout(resolve, LOCK_RETRY_MS))
}
}
}

/** Recursively freeze a plain-data value (defense in depth for handed-out snapshots). */
function deepFreeze<T>(value: T): T {
if (value !== null && typeof value === 'object') {
Expand Down
Loading
Loading