import assert from 'node:assert/strict' import { mkdtemp } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { test } from 'node:test' import { startService } from '../service/main.ts' import type { AuditEvent } from '../shared/events.ts' import { WakeQueue } from './queue-store.ts' import { createServiceClient } from './service-client.ts' import { AuditSpool } from './spool.ts' import { AgentState } from './state.ts' test('wake queue persists across reopen; heartbeat coalesces; done compacts', async () => { const dir = await mkdtemp(join(tmpdir(), 'wq-')) const file = join(dir, 'wakeq.jsonl') const q = await WakeQueue.open(file) const chat = await q.enqueue('chat', { text: 'hi' }) await q.enqueue('chat', { text: 'second' }) const hb1 = await q.enqueue('heartbeat', {}) const hb2 = await q.enqueue('heartbeat', {}) assert.equal(hb1.id, hb2.id, 'heartbeat coalesces into the pending one') assert.equal(q.peek().length, 3) const next = q.next() assert.equal(next?.kind, 'chat') assert.equal(next?.seq, chat.seq) await q.complete(next) assert.equal(q.peek().length, 2) // reopen: durability const q2 = await WakeQueue.open(file) assert.equal(q2.peek().length, 2) assert.deepEqual( q2 .peek() .map((w) => w.kind) .sort(), ['chat', 'heartbeat'], ) }) test('wake queue defer reorders wake with backoff', async () => { const dir = await mkdtemp(join(tmpdir(), 'wq-')) const q = await WakeQueue.open(join(dir, 'wakeq.jsonl')) const w = await q.enqueue('chat', {}) const future = new Date(Date.now() + 60_000).toISOString() await q.defer(w, future, 1) assert.equal(q.next(), undefined, 'deferred wake not claimable yet') const q2 = await WakeQueue.open(join(dir, 'wakeq.jsonl')) assert.equal(q2.peek()[0]?.attempts, 1) }) test('spool drains idempotently into the running service', async () => { const dataDir = await mkdtemp(join(tmpdir(), 'spool-svc-')) const webDist = join(dataDir, 'web-dist') await import('node:fs/promises').then((fs) => fs.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 << 10, storageWarnPct: 20, storageStopPct: 10, workerUrl: join(import.meta.dirname, '../service/worker.ts'), }) try { const addr = service.internalServer.address() if (addr === null || typeof addr !== 'object') throw new Error('no addr') const client = createServiceClient(`http://127.0.0.1:${addr.port}`) const spool = await AuditSpool.open(join(dataDir, 'agent', 'spool')) const event: AuditEvent = { id: 'ev-1', ts: new Date().toISOString(), kind: 'tool_execution_start', payload: { tool: 'bash' }, } await spool.append(event, { revision: 'rev-1', worktreeDirty: false }) assert.equal((await spool.list()).length, 1) const submitted = await spool.drain((events) => client.ingestAudit(events)) assert.equal(submitted, 1) assert.equal((await spool.list()).length, 0, 'drained events are removed') const status = await client.status() assert.equal(status.counts.audit, 1) // service crash between submit and remove -> drain resubmits, service dedups await spool.append(event, { revision: 'rev-1', worktreeDirty: false }) await spool.append( { id: 'ev-2', ts: new Date().toISOString(), kind: 'message_end', payload: {} }, {}, ) const second = await spool.drain((events) => client.ingestAudit(events)) assert.equal(second, 2) const status2 = await client.status() assert.equal(status2.counts.audit, 2, 'duplicate id ignored') } finally { await service.close() } }) test('agent state persists mutations atomically', async () => { const dir = await mkdtemp(join(tmpdir(), 'st-')) const path = join(dir, 'state.json') const st = await AgentState.load(path) await st.update((d) => { d.telegramOffset = 41 d.pushedShas.push('abc') }) const st2 = await AgentState.load(path) assert.equal(st2.data.telegramOffset, 41) assert.deepEqual(st2.data.pushedShas, ['abc']) })