repositories / will
will
owned by admin
src/agent/runtime.int.test.ts
Rawimport 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()
})