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 '], { 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 '], { 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 '], { 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 '], { 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 '], { 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 '], { 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((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((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) => Promise, ): Promise { 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, ): Promise> { 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 } 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() })