Luigit
repositories / will

will

owned by admin

src/agent/runtime.int.test.ts

Raw
import assert from 'node:assert/strict'
import { execFile } from 'node:child_process'
import { mkdir, mkdtemp, readFile, writeFile } from 'node:fs/promises'
import { createServer } from 'node:http'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { test } from 'node:test'
import { promisify } from 'node:util'
import { startService } from '../service/main.ts'
import { createLuciClient, pollCi } from './ciwatch.ts'
import { createDeployClient, FAILURE_EVENT_KINDS, pollDeployEvents } from './deploy-client.ts'
import { checkStorage } from './loops.ts'
import { Outbox } from './outbox.ts'
import { WakeQueue } from './queue-store.ts'
import { startReadinessServer } from './readiness.ts'
import { createServiceClient } from './service-client.ts'
import { AgentState } from './state.ts'
import { createWillTools } from './tools.ts'
import type { Tool } from './tools-shape.ts'
import { readWorktreeState, syncWorktree } from './worktree.ts'

const exec = promisify(execFile)

/** committer identity and signing must not depend on the machine's global git config (CI has none) */
const GIT_IDENTITY = ['-c', 'user.name=t', '-c', 'user.email=t@t', '-c', 'commit.gpgsign=false']

test('worktree: clone, ff-only when clean, preserved when dirty or ahead', async () => {
  const dir = await mkdtemp(join(tmpdir(), 'wt-'))
  const remote = join(dir, 'origin.git')
  const work = join(dir, 'work')
  await exec('git', ['init', '--bare', '-b', 'trunk', remote])
  const seed = join(dir, 'seed')
  await exec('git', ['clone', remote, seed])
  await writeFile(join(seed, 'a.txt'), 'one\n')
  await exec('git', ['add', '.'], { cwd: seed })
  await exec('git', [...GIT_IDENTITY, 'commit', '-m', 'one', '--author', 't <t@t>'], { cwd: seed })
  await exec('git', ['push', 'origin', 'trunk'], { cwd: seed })

  const s1 = await syncWorktree(work, remote, 'trunk')
  assert.equal(s1.cloned, true)
  assert.equal(s1.dirty, false)

  // advance remote; clean worktree fast-forwards
  await writeFile(join(seed, 'a.txt'), 'two\n')
  await exec('git', [...GIT_IDENTITY, 'commit', '-am', 'two', '--author', 't <t@t>'], { cwd: seed })
  await exec('git', ['push', 'origin', 'trunk'], { cwd: seed })
  const s2 = await syncWorktree(work, remote, 'trunk')
  assert.equal(s2.fastForwarded, true)
  assert.equal((await readFile(join(work, 'a.txt'), 'utf8')).trim(), 'two')

  // dirty worktree: preserved, no ff, reported
  await writeFile(join(work, 'b.txt'), 'local\n')
  await writeFile(join(seed, 'a.txt'), 'three\n')
  await exec('git', [...GIT_IDENTITY, 'commit', '-am', 'three', '--author', 't <t@t>'], {
    cwd: seed,
  })
  await exec('git', ['push', 'origin', 'trunk'], { cwd: seed })
  const s3 = await syncWorktree(work, remote, 'trunk')
  assert.equal(s3.dirty, true)
  assert.equal(s3.fastForwarded, false)
  assert.equal((await readFile(join(work, 'a.txt'), 'utf8')).trim(), 'two')

  // ahead commits: diverged, preserved
  await exec('git', ['add', '.'], { cwd: work })
  await exec('git', [...GIT_IDENTITY, 'commit', '-m', 'local', '--author', 't <t@t>'], {
    cwd: work,
  })
  const s4 = await syncWorktree(work, remote, 'trunk')
  assert.equal(s4.diverged, true)
})

test('readWorktreeState: fresh head and dirty flag without fetch or merge', async () => {
  const dir = await mkdtemp(join(tmpdir(), 'wt-read-'))
  const remote = join(dir, 'origin.git')
  const work = join(dir, 'work')
  await exec('git', ['init', '--bare', '-b', 'trunk', remote])
  const seed = join(dir, 'seed')
  await exec('git', ['clone', remote, seed])
  await writeFile(join(seed, 'a.txt'), 'one\n')
  await exec('git', ['add', '.'], { cwd: seed })
  await exec('git', [...GIT_IDENTITY, 'commit', '-m', 'one', '--author', 't <t@t>'], { cwd: seed })
  await exec('git', ['push', 'origin', 'trunk'], { cwd: seed })
  await exec('git', ['clone', '--quiet', remote, work])

  const head = () => exec('git', ['rev-parse', 'HEAD'], { cwd: work })
  const s1 = await readWorktreeState(work)
  assert.equal(s1.head, (await head()).stdout.trim())
  assert.equal(s1.dirty, false)

  // a local commit moves HEAD; the read reports it without merging origin
  await writeFile(join(work, 'b.txt'), 'local\n')
  await exec('git', ['add', '.'], { cwd: work })
  await exec('git', [...GIT_IDENTITY, 'commit', '-m', 'local', '--author', 't <t@t>'], {
    cwd: work,
  })
  const s2 = await readWorktreeState(work)
  assert.equal(s2.head, (await head()).stdout.trim())
  assert.notEqual(s2.head, s1.head)

  // uncommitted work is reported dirty
  await writeFile(join(work, 'c.txt'), 'uncommitted\n')
  const s3 = await readWorktreeState(work)
  assert.equal(s3.dirty, true)

  // missing worktree: error reported, never throws
  const s4 = await readWorktreeState(join(dir, 'missing'))
  assert.ok(s4.error)
  assert.equal(s4.head, undefined)
})

test('ci watch enqueues one failure wake per failed run of a tracked sha', async () => {
  const runs: unknown[] = []
  const server = createServer((_req, res) => {
    res.writeHead(200, { 'content-type': 'application/json' })
    res.end(JSON.stringify({ runs }))
  })
  await new Promise<void>((r) => server.listen(0, '127.0.0.1', () => r()))
  const url = `http://127.0.0.1:${(server.address() as { port: number }).port}`
  const dir = await mkdtemp(join(tmpdir(), 'ci-'))
  const state = await AgentState.load(join(dir, 'state.json'))
  await state.update((d) => {
    d.pushedShas.push('sha-1')
  })
  const queue = await WakeQueue.open(join(dir, 'wakeq.jsonl'))
  const client = createLuciClient(url, 'will')

  runs.push({ id: 'r1', revision: 'sha-1', status: 'success' })
  assert.equal((await pollCi(client, state, queue)).wakes, 0)

  runs.push({
    id: 'r2',
    revision: 'sha-1',
    status: 'failed',
    jobs: [{ name: 'test', status: 'failed' }],
  })
  assert.equal((await pollCi(client, state, queue)).wakes, 1)
  assert.equal((await pollCi(client, state, queue)).wakes, 0, 'handled once')

  runs.push({ id: 'r3', revision: 'sha-untracked', status: 'failed' })
  assert.equal((await pollCi(client, state, queue)).wakes, 0, 'untracked sha ignored')
  const payload = queue.peek()[0]?.payload as { sha: string; job: string }
  assert.deepEqual(payload, { sha: 'sha-1', job: 'test' })
  server.close()
})

test('deploy events poll mirrors into service and wakes on failures only', async () => {
  const events: unknown[] = []
  let _requests = 0
  const server = createServer((req, res) => {
    _requests++
    const url = new URL(req.url ?? '/', 'http://x')
    if (url.pathname === '/v1/rollback') {
      res.writeHead(200, { 'content-type': 'application/json' })
      res.end(JSON.stringify({ revision: 'rev-prev' }))
      return
    }
    const after = url.searchParams.get('after') ?? undefined
    const all = events as { cursor: string; id: string }[]
    const filtered = after === undefined ? all : all.filter((e) => e.cursor > after)
    res.writeHead(200, { 'content-type': 'application/json' })
    res.end(JSON.stringify({ events: filtered, nextCursor: filtered.at(-1)?.cursor }))
  })
  await new Promise<void>((r) => server.listen(0, '127.0.0.1', () => r()))
  const url = `http://127.0.0.1:${(server.address() as { port: number }).port}`
  const client = createDeployClient(url)

  const dataDir = await mkdtemp(join(tmpdir(), 'de-'))
  const webDist = join(dataDir, 'web-dist')
  await mkdir(webDist, { recursive: true })
  const service = await startService({
    dataDir,
    publicAddr: '127.0.0.1:0',
    internalAddr: '127.0.0.1:0',
    webDist,
    docsDist: webDist,
    mediaMaxBytes: 1000,
    storageWarnPct: 20,
    storageStopPct: 10,
    workerUrl: join(import.meta.dirname, '../service/worker.ts'),
  })
  try {
    const svc = createServiceClient(
      `http://127.0.0.1:${(service.internalServer.address() as { port: number }).port}`,
    )
    const queue = await WakeQueue.open(join(dataDir, 'agent', 'wakeq.jsonl'))
    events.push(
      {
        cursor: 'c1',
        id: 'e1',
        ts: new Date().toISOString(),
        kind: 'promote',
        revision: 'r5',
        detail: {},
      },
      {
        cursor: 'c2',
        id: 'e2',
        ts: new Date().toISOString(),
        kind: 'probation-failed',
        revision: 'r6',
        detail: { probe: 'healthz' },
      },
    )
    const seen: string[] = []
    const onEvent = async (e: { id: string; kind: string }) => {
      seen.push(e.id)
      if (FAILURE_EVENT_KINDS.has(e.kind)) await queue.enqueue('deploy_failure', { revision: 'r6' })
    }
    const page1 = await pollDeployEvents(client, { service: svc, onEvent })
    assert.equal(page1.imported, 2)
    assert.equal(queue.peek().length, 1)
    assert.equal(queue.peek()[0]?.kind, 'deploy_failure')
    assert.deepEqual(seen, ['e1', 'e2'])
    const mirror = await svc.recall({ kind: 'deploy', limit: 10 })
    assert.equal(mirror.rows.length, 2)

    const page2 = await pollDeployEvents(client, {
      cursor: page1.cursor,
      service: svc,
      onEvent: () => {},
    })
    assert.equal(page2.imported, 0, 'cursor prevents re-import')

    const rollback = await client.rollbackPrevious()
    assert.deepEqual(rollback, { ok: true, revision: 'rev-prev' })
    assert.ok(FAILURE_EVENT_KINDS.has('rollback'))
  } finally {
    await service.close()
    server.close()
  }
})

async function withTools(
  fn: (tools: Tool[], outbox: Outbox, svc: ReturnType<typeof createServiceClient>) => Promise<void>,
): Promise<void> {
  const dir = await mkdtemp(join(tmpdir(), 'tools-'))
  const webDist = join(dir, 'web-dist')
  await mkdir(webDist, { recursive: true })
  const service = await startService({
    dataDir: dir,
    publicAddr: '127.0.0.1:0',
    internalAddr: '127.0.0.1:0',
    webDist,
    docsDist: webDist,
    mediaMaxBytes: 1000 << 10,
    storageWarnPct: 20,
    storageStopPct: 10,
    workerUrl: join(import.meta.dirname, '../service/worker.ts'),
  })
  try {
    const svc = createServiceClient(
      `http://127.0.0.1:${(service.internalServer.address() as { port: number }).port}`,
    )
    const outbox = await Outbox.open(join(dir, 'agent', 'outbox.jsonl'))
    const tools = createWillTools({ outbox, service: svc, chatId: -100 }) as unknown as Tool[]
    await fn(tools, outbox, svc)
  } finally {
    await service.close()
  }
}

async function runTool(
  tool: Tool,
  args: Record<string, unknown>,
): Promise<Record<string, unknown>> {
  const res = await tool.execute('t', args, new AbortController().signal, () => {}, {
    cwd: process.cwd(),
  })
  const text =
    (res.content as { type: string; text?: string }[]).find((c) => c.type === 'text')?.text ?? '{}'
  return JSON.parse(text) as Record<string, unknown>
}

test('chat_send queues outbound; recall queries service; maintenance ops gated', async () => {
  await withTools(async (tools, outbox, svc) => {
    const chatSend = tools.find((t) => t.name === 'chat_send')
    const recall = tools.find((t) => t.name === 'recall')
    const maint = tools.find((t) => t.name === 'maintenance_run')
    const rollback = tools.find((t) => t.name === 'deployment_rollback')
    assert.ok(chatSend && recall && maint && rollback)

    const sent = await runTool(chatSend, { text: 'hello group' })
    assert.equal(sent.queued, true)
    assert.equal(outbox.pending().length, 1)

    await svc.ingestChat([
      {
        id: 'tg-1',
        chatId: -100,
        messageId: 1,
        direction: 'in',
        ts: new Date().toISOString(),
        text: 'remember this',
      },
    ])
    const rec = await runTool(recall, { kind: 'chat', q: 'remember' })
    assert.equal((rec.rows as unknown[]).length, 1)

    const quick = await runTool(maint, { op: 'quick_check' })
    assert.deepEqual(quick.result, { result: 'ok' })

    const rb = await runTool(rollback, { reason: 'test' })
    assert.equal(rb.ok, false)
    assert.match(String(rb.error), /WILL_DEPLOY_API_URL/)
  })
})

test('storage watch: transitions emit deterministic outbox notice', async () => {
  const dir = await mkdtemp(join(tmpdir(), 'sw-'))
  const webDist = join(dir, 'web-dist')
  await mkdir(webDist, { recursive: true })
  const service = await startService({
    dataDir: dir,
    publicAddr: '127.0.0.1:0',
    internalAddr: '127.0.0.1:0',
    webDist,
    docsDist: webDist,
    mediaMaxBytes: 1000,
    // 0/0 pins the initial state to 'ok'; the test forces 'warn' by patching
    // meta thresholds below, deterministic on any host filesystem.
    storageWarnPct: 0,
    storageStopPct: 0,
    workerUrl: join(import.meta.dirname, '../service/worker.ts'),
  })
  try {
    const svc = createServiceClient(
      `http://127.0.0.1:${(service.internalServer.address() as { port: number }).port}`,
    )
    const state = await AgentState.load(join(dir, 'state.json'))
    const outbox = await Outbox.open(join(dir, 'outbox.jsonl'))
    // ok -> ok: no transition
    const r1 = await checkStorage(svc, state, outbox, -100)
    assert.equal(r1.transitioned, false)
    // simulate warn by editing the state that produced storageState
    await state.update((d) => {
      d.storageNoticeState = 'ok'
    })
    // force warn: patch service meta thresholds so freePct < warn on any host
    const dbFile = join(dir, 'will.db')
    const { DatabaseSync } = await import('node:sqlite')
    const db = new DatabaseSync(dbFile)
    db.prepare("INSERT OR REPLACE INTO meta (key, value) VALUES ('storage_warn_pct', '101')").run()
    db.close()
    const r2 = await checkStorage(svc, state, outbox, -100)
    assert.equal(r2.state, 'warn')
    assert.equal(r2.transitioned, true)
    assert.equal(outbox.pending().length, 1)
    assert.match(outbox.pending()[0]?.text ?? '', /warning threshold/)
  } finally {
    await service.close()
  }
})

test('readiness server reports 503 until local init completes', async () => {
  const server = await startReadinessServer('127.0.0.1:0')
  const port = (server.server.address() as { port: number }).port
  const before = await fetch(`http://127.0.0.1:${port}/readyz`)
  assert.equal(before.status, 503)
  server.setReady({
    sessionReady: true,
    serviceReachable: true,
    telegramLoopReady: false,
    storageOk: false,
    configValid: true,
  })
  const after = await fetch(`http://127.0.0.1:${port}/readyz`)
  assert.equal(after.status, 200)
  const body = (await after.json()) as { checks: { telegramLoopReady: boolean } }
  assert.equal(body.checks.telegramLoopReady, false, 'degraded but ready')
  await server.close()
})