Luigit
repositories / will

will

owned by admin

src/agent/loops.ts

Raw
import { createLogger } from '../shared/log.ts'
import { msUntilNextHour, utcHourKey } from '../shared/time.ts'
import { claimHeartbeatHour, hasPendingHeartbeat } from '../shared/wake.ts'
import type { Outbox } from './outbox.ts'
import type { WakeQueue } from './queue-store.ts'
import type { ServiceClient } from './service-client.ts'
import type { AgentState } from './state.ts'

const log = createLogger({ component: 'loop' })

/**
 * Heartbeat scheduler (spec AC16): at most once per UTC hour, no backfill.
 * Storage watch (spec AC13): deterministic notices on threshold transitions,
 * ordinary model work paused below the stop threshold.
 */

export function startHeartbeatLoop(
  state: AgentState,
  queue: WakeQueue,
  intervalMs = 60_000,
): () => void {
  const timer = setInterval(() => {
    void (async () => {
      const { claimed } = claimHeartbeatHour(state.data.heartbeat, utcHourKey(new Date()))
      if (!claimed) return
      if (hasPendingHeartbeat(queue.peek())) return
      await state.update((d) => {
        d.heartbeat = { lastClaimedHour: utcHourKey(new Date()) }
      })
      await queue.enqueue('heartbeat', {})
      log.info('heartbeat enqueued', { hour: utcHourKey(new Date()) })
    })()
  }, intervalMs)
  timer.unref()
  return () => clearInterval(timer)
}

export interface StorageWatchResult {
  transitioned: boolean
  state: 'ok' | 'warn' | 'stop'
}

export async function checkStorage(
  service: ServiceClient,
  state: AgentState,
  outbox: Outbox,
  chatId: number | undefined,
): Promise<StorageWatchResult> {
  const status = await service.status()
  const current = status.storageState
  const previous = state.data.storageNoticeState
  if (current !== previous) {
    await state.update((d) => {
      d.storageNoticeState = current
    })
    if (chatId !== undefined && (current === 'warn' || current === 'stop')) {
      // Deterministic notice: no model call (spec AC13).
      await outbox.enqueue({
        id: `out-storage-${Date.now()}`,
        chatId,
        text:
          current === 'stop'
            ? 'will: storage below the stop threshold; pausing new model work and media. Queued work is safe.'
            : 'will: storage below the warning threshold.',
      })
    }
    log.warn('storage threshold transition', { from: previous, to: current })
    return { transitioned: true, state: current }
  }
  return { transitioned: false, state: current }
}

export function startStorageWatch(
  service: ServiceClient,
  state: AgentState,
  outbox: Outbox,
  chatId: number | undefined,
  onPauseChange: (paused: boolean) => void,
  intervalMs = 60_000,
): () => void {
  const timer = setInterval(() => {
    void checkStorage(service, state, outbox, chatId)
      .then((r) => onPauseChange(r.state === 'stop'))
      .catch(() => undefined)
  }, intervalMs)
  timer.unref()
  return () => clearInterval(timer)
}

export function msToNextHourFrom(now: Date): number {
  return msUntilNextHour(now)
}