Luigit
repositories / will

will

owned by admin

src/agent/main.ts

Raw
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<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())
}