import { 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 } const PI_EVENT_KINDS: Record = { 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 { 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 | 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) | 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 { 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()) }