repositories / will
will
owned by admin
src/agent/main.ts
Rawimport { mkdir } from 'node:fs/promises'
import { join } from 'node:path'
import { pathToFileURL } from 'node:url'
import { createLogger } from '../shared/log.ts'
import { nowIso } from '../shared/time.ts'
import { createLuciClient, pollCi } from './ciwatch.ts'
import { type AgentConfig, loadAgentConfig, resolveSessionMode } from './config.ts'
import { createDeployClient, FAILURE_EVENT_KINDS, pollDeployEvents } from './deploy-client.ts'
import { createEngine, type SessionPort } from './engine.ts'
import { startHeartbeatLoop, startStorageWatch } from './loops.ts'
import { Outbox } from './outbox.ts'
import { createFakeSession, PiSession } from './pi-session.ts'
import { buildSystemPrompt } from './prompt.ts'
import { WakeQueue } from './queue-store.ts'
import { type ReadinessChecks, startReadinessServer } from './readiness.ts'
import { createServiceClient } from './service-client.ts'
import { AuditSpool } from './spool.ts'
import { AgentState } from './state.ts'
import { TelegramClient } from './telegram-client.ts'
import { createTelegramRuntime } from './telegram-runtime.ts'
import { createWillTools } from './tools.ts'
import { sumAssistantUsage } from './usage.ts'
import { readWorktreeState, syncWorktree, type WorktreeStatus } from './worktree.ts'
import { ZaiRestClient } from './zai/client.ts'
import { createZaiTools } from './zai/tools.ts'
const log = createLogger({ component: 'agent' })
export interface RunningAgent {
config: AgentConfig
session: SessionPort
readiness: { setReady(checks: ReadinessChecks): void }
close(): Promise<void>
}
const PI_EVENT_KINDS: Record<string, string> = {
agent_start: 'agent_start',
agent_end: 'agent_end',
turn_start: 'turn_start',
turn_end: 'turn_end',
message_start: 'message_start',
message_end: 'message_end',
tool_execution_start: 'tool_execution_start',
tool_execution_end: 'tool_execution_end',
compaction_start: 'compaction_start',
compaction_end: 'compaction_end',
auto_retry_start: 'auto_retry_start',
auto_retry_end: 'auto_retry_end',
}
export async function startAgent(config: AgentConfig): Promise<RunningAgent> {
await mkdir(join(config.agentDir, 'spool'), { recursive: true })
const state = await AgentState.load(join(config.agentDir, 'state.json'))
const queue = await WakeQueue.open(join(config.agentDir, 'wakeq.jsonl'))
const spool = await AuditSpool.open(join(config.agentDir, 'spool'))
const service = createServiceClient(config.serviceUrl)
const outbox = await Outbox.open(join(config.agentDir, 'outbox.jsonl'))
const provenance = () => ({
revision: config.revision,
agentDigest: config.agentDigest,
})
const audit = async (kind: string, payload: unknown) => {
await spool.append({ id: crypto.randomUUID(), ts: nowIso(), kind, payload }, provenance())
kickDelivery()
}
const drainSpool = () =>
spool
.drain((events) => service.ingestAudit(events))
.then((n) => {
if (n > 0) log.debug('spool drained', { events: n })
})
.catch((err: unknown) => {
log.warn('spool drain failed', { error: err instanceof Error ? err.message : String(err) })
})
const telegramFlusher = () => {
if (telegramFlush === undefined) return Promise.resolve()
return telegramFlush().catch(() => undefined)
}
/**
* Delivery runs between wakes AND during long agent runs: every audited
* event and every chat_send kicks a debounced drain+flush, so the timeline
* stays live and replies go out while the turn is still streaming.
*/
let deliveryTimer: ReturnType<typeof setTimeout> | undefined
const kickDelivery = () => {
if (deliveryTimer !== undefined) return
deliveryTimer = setTimeout(() => {
deliveryTimer = undefined
void drainSpool().then(telegramFlusher)
}, 300)
deliveryTimer.unref?.()
}
// --- tools ---
const deploy =
config.deployApiUrl !== undefined ? createDeployClient(config.deployApiUrl) : undefined
const willTools = createWillTools({
outbox,
service,
...(deploy !== undefined ? { deploy } : {}),
...(config.telegram.chatId !== undefined ? { chatId: config.telegram.chatId } : {}),
onChatQueued: kickDelivery,
})
const zaiClient =
config.zai.apiKey !== undefined
? new ZaiRestClient(config.zai.apiKey, config.zai.baseUrl)
: undefined
const zaiTools =
zaiClient !== undefined
? createZaiTools({ client: zaiClient, service, analysisModel: config.zai.analysisModel })
: []
const customTools = [...willTools, ...zaiTools]
// --- session ---
const systemPrompt = await buildSystemPrompt(join(config.resourceDir, 'persona'))
const onPiEvent = (event: unknown) => {
const e = event as { type?: string }
const kind = e.type !== undefined ? PI_EVENT_KINDS[e.type] : undefined
if (kind === undefined) return
void audit(kind, event).then(() => undefined)
if (e.type === 'agent_end') {
// Derived cost-view rows: one per (wake, provider/model) (spec: tokens
// and projected cost surface). Raw events stay untouched.
for (const usage of sumAssistantUsage((event as { messages?: unknown }).messages)) {
void audit('turn_usage', usage).then(() => undefined)
}
}
}
let session: SessionPort
const sessionMode = resolveSessionMode()
if (sessionMode === 'fake') {
session = createFakeSession()
log.warn('running with a FAKE session (no model calls; set ZAI_API_KEY for real)', {})
} else {
const pi = await PiSession.create({ config, systemPrompt, customTools, onEvent: onPiEvent })
session = pi
}
// --- engine ---
// Mutable worktree snapshot: the boot sync below assigns it; wake facts
// refresh it with a cheap local read so they never report a stale HEAD.
let worktreeStatus: WorktreeStatus = {
dirty: false,
diverged: false,
fastForwarded: false,
cloned: false,
}
const engine = createEngine({
queue,
state,
session,
budgetConfig: config.budget,
buildFacts: async () => {
let storageState: string | undefined
try {
storageState = (await service.status()).storageState
} catch {
storageState = undefined
}
// Cheap local read (no fetch/ff); the boot sync owns fetch + fast-forward.
worktreeStatus = await readWorktreeState(config.workDir)
return {
timeUtc: new Date().toISOString(),
revision: config.revision,
agentDigest: config.agentDigest,
...(storageState !== undefined ? { serviceStorageState: storageState } : {}),
queueDepth: queue.peek().length,
budgetUsed: state.data.budget.used,
budgetCap: config.budget.dailyCapCalls,
worktreeHead: worktreeStatus.head,
...(worktreeStatus.dirty !== undefined ? { worktreeDirty: worktreeStatus.dirty } : {}),
}
},
audit,
turnDeadlineMs: config.maxTurnMinutes * 60_000,
})
// --- worktree ---
worktreeStatus = await syncWorktree(config.workDir, config.vcs.gitUrl, config.vcs.branch)
if (worktreeStatus.error) log.warn('worktree not ready', { error: worktreeStatus.error })
else
log.info('worktree synced', {
head: worktreeStatus.head,
dirty: worktreeStatus.dirty,
ff: worktreeStatus.fastForwarded,
})
// --- telegram ---
const checks: ReadinessChecks = {
sessionReady: true,
serviceReachable: false,
telegramLoopReady: false,
storageOk: true,
configValid: true,
}
let telegramStop: (() => void) | undefined
let telegramFlush: (() => Promise<void>) | undefined
if (
config.telegram.token !== undefined &&
config.telegram.chatId !== undefined &&
config.telegram.oliverId !== undefined
) {
const tg = new TelegramClient(
config.telegram.token,
config.telegram.apiBase !== undefined
? { apiBase: config.telegram.apiBase, fileBase: config.telegram.apiBase }
: {},
)
const rt = createTelegramRuntime({
tg,
service,
queue,
spool,
state,
chatId: config.telegram.chatId,
oliverId: config.telegram.oliverId,
})
telegramStop = () => rt.stop()
telegramFlush = async () => {
await outbox.flush({ tg, service })
}
void rt.run().then(() => undefined)
checks.telegramLoopReady = true
} else {
log.warn('telegram not configured; chat wakes disabled', {})
}
// --- ci watch + deploy events ---
let ciStop: (() => void) | undefined
if (config.luci.baseUrl !== undefined && config.luci.repo !== undefined) {
const luci = createLuciClient(config.luci.baseUrl, config.luci.repo)
const timer = setInterval(() => {
void pollCi(luci, state, queue).catch(() => undefined)
}, 120_000)
timer.unref()
ciStop = () => clearInterval(timer)
}
let deployStop: (() => void) | undefined
if (deploy !== undefined) {
const timer = setInterval(() => {
void pollDeployEvents(deploy, {
cursor: state.data.deployEventCursor,
service,
onEvent: (event) => {
if (!FAILURE_EVENT_KINDS.has(event.kind)) return
return queue
.enqueue('deploy_failure', {
revision: event.revision,
kind: event.kind,
detail: event.detail,
})
.then(() => undefined)
},
})
.then(async ({ cursor }) => {
if (cursor !== undefined && cursor !== state.data.deployEventCursor) {
await state.update((d) => {
d.deployEventCursor = cursor
})
}
})
.catch(() => undefined)
}, 60_000)
timer.unref()
deployStop = () => clearInterval(timer)
}
// --- readiness + loops ---
const readiness = await startReadinessServer(config.readyAddr)
const stopHeartbeat = startHeartbeatLoop(state, queue)
const stopStorageWatch = startStorageWatch(
service,
state,
outbox,
config.telegram.chatId,
(paused) => engine.setStoragePaused(paused),
)
// --- main loop ---
let stopped = false
const mainLoop = async () => {
while (!stopped) {
await drainSpool()
await telegramFlusher()
try {
await service.status()
checks.serviceReachable = true
} catch {
checks.serviceReachable = false
}
readiness.setReady(checks)
const outcome = await engine.tick()
if (
outcome.kind === 'idle' ||
outcome.kind === 'budget-blocked' ||
outcome.kind === 'storage-blocked'
) {
await sleep(outcome.kind === 'idle' ? 1000 : 5000)
}
}
}
void mainLoop()
return {
config,
session,
readiness: {
setReady: (c) => readiness.setReady(c),
},
async close() {
stopped = true
telegramStop?.()
ciStop?.()
deployStop?.()
stopHeartbeat()
stopStorageWatch()
await drainSpool()
await session.abort().catch(() => undefined)
await readiness.close()
log.info('agent stopped cleanly', {})
},
}
}
function sleep(ms: number): Promise<void> {
return new Promise((res) => setTimeout(res, ms))
}
const isEntryPoint =
process.argv[1] !== undefined && import.meta.url === pathToFileURL(process.argv[1]).href
if (isEntryPoint) {
const config = loadAgentConfig()
const agent = await startAgent(config)
const shutdown = async () => {
const hard = setTimeout(() => process.exit(1), 29_000)
hard.unref()
await agent.close()
clearTimeout(hard)
process.exit(0)
}
process.on('SIGTERM', () => void shutdown())
process.on('SIGINT', () => void shutdown())
}