repositories / will
will
owned by admin
src/agent/queue-store.ts
Rawimport { 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<WakeQueue> {
let nextSeq = 1
const pending = new Map<string, Wake>()
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<Wake> {
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<void> {
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<void> {
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<void> {
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<void> {
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
}
}