repositories / will
will
owned by admin
src/agent/loops.ts
Rawimport { 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)
}