Luigit
repositories / will

will

owned by admin

src/agent/queue-store.ts

Raw
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<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
  }
}