import { mkdir, readFile, rename, writeFile } from 'node:fs/promises' import { dirname, join } from 'node:path' import { newId } from '../shared/ids.ts' import { hasPendingHeartbeat, orderWakes, type Wake, type WakeKind } from '../shared/wake.ts' /** * Durable wake queue on the agent volume. * JSONL append with explicit done markers; compaction rewrites live wakes only. * A wake is acknowledged (fsynced to disk) before any side effect runs. */ interface QueueLine { op: 'add' | 'done' wake?: Wake id?: string } export class WakeQueue { readonly #file: string #pending: Wake[] = [] #nextSeq: number private constructor(file: string, pending: Wake[], nextSeq: number) { this.#file = file this.#pending = pending this.#nextSeq = nextSeq } static async open(file: string): Promise { let nextSeq = 1 const pending = new Map() let raw = '' try { raw = await readFile(file, 'utf8') } catch { // fresh queue } for (const line of raw.split('\n')) { if (line === '') continue const entry = JSON.parse(line) as QueueLine if (entry.op === 'add' && entry.wake) { pending.set(entry.wake.id, entry.wake) nextSeq = Math.max(nextSeq, entry.wake.seq + 1) } else if (entry.op === 'done' && entry.id) { pending.delete(entry.id) } } return new WakeQueue(file, [...pending.values()], nextSeq) } /** Durable enqueue; heartbeats coalesce to at most one pending (spec AC16). */ async enqueue( kind: WakeKind, payload: unknown, opts: { coalesceHeartbeat?: boolean } = {}, ): Promise { if ( kind === 'heartbeat' && opts.coalesceHeartbeat !== false && hasPendingHeartbeat(this.#pending) ) { return this.#pending.find((w) => w.kind === 'heartbeat') as Wake } const wake: Wake = { id: newId(), seq: this.#nextSeq++, kind, createdAt: new Date().toISOString(), payload, attempts: 0, } this.#pending.push(wake) await this.#append({ op: 'add', wake }) return wake } /** Re-enqueue an existing wake (retry backoff or defer). */ async defer(wake: Wake, deferredUntil: string, attempts: number): Promise { const idx = this.#pending.findIndex((w) => w.id === wake.id) if (idx === -1) this.#pending.push(wake) else this.#pending[idx] = wake const target = this.#pending.find((w) => w.id === wake.id) if (target) { target.deferredUntil = deferredUntil target.attempts = attempts } await this.#append({ op: 'add', wake: this.#pending.find((w) => w.id === wake.id) as Wake }) } next(nowMs = Date.now()): Wake | undefined { return orderWakes(this.#pending, nowMs)[0] } peek(): readonly Wake[] { return this.#pending } async complete(wake: Wake): Promise { this.#pending = this.#pending.filter((w) => w.id !== wake.id) await this.#append({ op: 'done', id: wake.id }) if (this.#appendsSinceCompact > 200) await this.compact() } #appendsSinceCompact = 0 async #append(line: QueueLine): Promise { await mkdir(dirname(this.#file), { recursive: true }) await writeFile(this.#file, `${JSON.stringify(line)}\n`, { flag: 'a' }) this.#appendsSinceCompact++ } /** Rewrite the file with only live wakes; crash-safe via rename. */ async compact(): Promise { const tmp = join(dirname(this.#file), `.queue-${process.pid}.tmp`) const lines: string[] = this.#pending.map((wake) => { const line: QueueLine = { op: 'add', wake } return JSON.stringify(line) }) await writeFile(tmp, lines.length > 0 ? `${lines.join('\n')}\n` : '', 'utf8') await rename(tmp, this.#file) this.#appendsSinceCompact = 0 } }