Luigit
repositories / will

will

owned by admin

src/agent/engine.ts

Raw
import { type BudgetConfig, canSpend, isRecoveryWake, rollover, spend } from '../shared/budget.ts'
import { createLogger } from '../shared/log.ts'
import { type Clock, systemClock } from '../shared/time.ts'
import type { Wake } from '../shared/wake.ts'
import type { WakeFacts } from './prompt.ts'
import { renderWakePrompt } from './prompt.ts'
import type { WakeQueue } from './queue-store.ts'
import type { AgentState } from './state.ts'

const _log = createLogger({ component: 'engine' })

/** Thin seam over the pi session so the wake engine is testable without a model. */
export interface SessionPort {
  prompt(text: string): Promise<void>
  abort(): Promise<void>
}

export interface EngineDeps {
  queue: WakeQueue
  state: AgentState
  session: SessionPort
  budgetConfig: BudgetConfig
  buildFacts: () => Promise<WakeFacts>
  audit: (kind: string, payload: unknown) => Promise<void>
  clock?: Clock
  /** wall-clock watchdog bound for one turn (spec AC17) */
  turnDeadlineMs?: number
  /** retries before a wake is deferred to the next heartbeat (spec AC17) */
  maxAttempts?: number
  backoffMs?: (attempt: number) => number
}

export interface TurnOutcome {
  kind: 'completed' | 'failed' | 'timeout' | 'budget-blocked' | 'storage-blocked' | 'idle'
  wake?: Wake
}

/**
 * The wake engine: claims the highest-priority due wake, drives one session
 * turn under the watchdog, and applies budget/retry/timeout policy.
 */
export function createEngine(deps: EngineDeps) {
  const clock = deps.clock ?? systemClock
  const turnDeadlineMs = deps.turnDeadlineMs ?? 30 * 60_000
  const maxAttempts = deps.maxAttempts ?? 3
  const backoffMs =
    deps.backoffMs ?? ((attempt: number) => Math.min(60_000 * 2 ** (attempt - 1), 15 * 60_000))
  let storagePaused = false
  let bootPending = true

  const setStoragePaused = (paused: boolean) => {
    storagePaused = paused
  }

  const runTurn = async (wake: Wake): Promise<TurnOutcome> => {
    const budget = rollover(deps.state.data.budget, clock)
    if (!canSpend(deps.budgetConfig, budget, wake.kind)) {
      if (budget !== deps.state.data.budget) await deps.state.update(() => {})
      await deps.audit('wake_budget_blocked', { wakeId: wake.id, kind: wake.kind })
      return { kind: 'budget-blocked', wake }
    }
    if (storagePaused && !isRecoveryWake(wake.kind)) {
      return { kind: 'storage-blocked', wake }
    }

    const facts = await deps.buildFacts()
    const promptText = renderWakePrompt(wake, facts, { boot: bootPending })
    bootPending = false
    await deps.audit('wake_start', { wakeId: wake.id, kind: wake.kind, attempts: wake.attempts })

    let timedOut = false
    let onDeadline: ((value: 'deadline') => void) | undefined
    const deadline = new Promise<'deadline'>((resolve) => {
      onDeadline = resolve
    })
    const timer = setTimeout(() => {
      timedOut = true
      onDeadline?.('deadline')
    }, turnDeadlineMs)
    timer.unref?.()

    try {
      const settled = await Promise.race([
        deps.session.prompt(promptText).then(
          () => 'done' as const,
          (err: unknown) => {
            throw err
          },
        ),
        deadline,
      ])
      if (settled === 'deadline') throw new Error('turn exceeded deadline')
      const spent = spend(deps.budgetConfig, rollover(deps.state.data.budget, clock), wake.kind)
      await deps.state.update((d) => {
        d.budget = spent
      })
      await deps.queue.complete(wake)
      await deps.audit('wake_end', { wakeId: wake.id, kind: wake.kind, outcome: 'completed' })
      return { kind: 'completed', wake }
    } catch (err) {
      if (timedOut) {
        // Timeout: cancel, then exactly one recovery wake; never blind replay (spec AC17).
        await Promise.race([
          deps.session.abort().catch(() => undefined),
          new Promise((res) => setTimeout(res, 5000)),
        ])
        await deps.queue.complete(wake)
        await deps.queue.enqueue('recovery', {
          interruptedWakeId: wake.id,
          note: 'turn watchdog fired',
        })
        await deps.audit('wake_timeout', { wakeId: wake.id, kind: wake.kind })
        return { kind: 'timeout', wake }
      }
      const attempts = wake.attempts + 1
      const message = err instanceof Error ? err.message : String(err)
      if (attempts < maxAttempts) {
        const until = new Date(clock.now().getTime() + backoffMs(attempts)).toISOString()
        await deps.queue.defer(wake, until, attempts)
        await deps.audit('wake_retry', { wakeId: wake.id, attempts, error: message })
        return { kind: 'failed', wake }
      }
      // Deferred: retry on a later heartbeat-triggered pass (spec: bounded attempts, then defer).
      const until = new Date(clock.now().getTime() + 60 * 60_000).toISOString()
      await deps.queue.defer(wake, until, 0)
      await deps.audit('wake_deferred', { wakeId: wake.id, attempts, error: message })
      return { kind: 'failed', wake }
    } finally {
      clearTimeout(timer)
    }
  }

  /** Try one wake; returns idle when nothing is due. */
  const tick = async (): Promise<TurnOutcome> => {
    const wake = deps.queue.next(clock.now().getTime())
    if (!wake) return { kind: 'idle' }
    return runTurn(wake)
  }

  return { tick, runTurn, setStoragePaused }
}