import test from 'node:test'; import assert from 'node:assert/strict'; import { fork } from 'node:child_process'; import { mkdtempSync, rmSync, existsSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { once } from 'node:events'; import { Store } from '../src/store.mjs'; import { Broker, BusError, TASK_VERBS } from '../src/broker.mjs'; import { serve } from '../src/server.mjs'; import { startBroker } from '../src/runtime.mjs'; import { Client } from '../src/client.mjs'; // The S3 hooks: recordTask, taskView and requestTask on the broker, the server's task branch, the // runtime's adapter factory and the process's `trackers` boot field. const businesses = { demo: { id: 'demo', human: 'jason', arbiters: { technical: 'cto', delivery: 'pm' }, roles: { pm: { authority: { withinRole: ['message.send', 'task.create'], crossRole: [] } }, cto: { authority: { withinRole: ['message.send'], crossRole: [] } }, }, }, other: { id: 'other', human: 'jason', arbiters: { technical: 'pm', delivery: 'pm' }, roles: { pm: { authority: { withinRole: ['message.send'], crossRole: [] } } }, }, }; const SECRET = 'tk_' + 'z'.repeat(40); const digest = 'b'.repeat(64); const at = '2026-10-08T12:00:00.000Z'; function setup(t) { const root = mkdtempSync(join(tmpdir(), 'bus-tasks-')); const store = new Store(root); const b = new Broker({ store, businesses, secretCheck: (v) => { if (JSON.stringify(v ?? null).includes(SECRET)) throw new BusError('secret-detected'); }, }); t.after(() => { store.close(); rmSync(root, { recursive: true, force: true }); }); const agent = (role, business = 'demo') => { const cap = b.bindLaunch({ business, role, run: `${business}-${role}`, harness: 'pi' }); b.request(cap, { verb: 'role.claim' }); return cap; }; const human = b.bindHuman({ business: 'demo', human: 'jason', via: 'cli', outsideAgent: true }); const count = (table, where = '1') => store.get(`SELECT count(*) n FROM ${table} WHERE ${where}`).n; return { root, store, b, agent, human, count }; } const poll = (extra = {}) => ({ task_ref: 'vikunja:3/7', updated: at, digest, fields: { bucket: 11, done: false }, via: 'board', read_at: at, ...extra, }); const self = (extra = {}) => ({ task_ref: 'vikunja:3/7', updated: at, digest, fields: { bucket: 12, done: false }, ...extra }); test('recordTask keeps sync reads and a role write apart', (t) => { const { b, agent, human, store, count } = setup(t); const pm = agent('pm'); b.recordTask({ business: 'demo', snapshots: [poll()], events: [{ kind: 'task.changed.external', subject: 'vikunja:3/7', body: { changed: ['bucket'] } }], }); const row = store.get('SELECT * FROM task_snapshots'); assert.deepEqual([row.source, row.via, row.read_at, row.role, row.run], ['poll', 'board', at, null, null]); assert.equal(store.get("SELECT actor_role FROM events WHERE kind='task.changed.external'").actor_role, null); b.recordTask({ cap: pm, business: 'demo', snapshots: [self()], events: [{ kind: 'task.assigned', subject: 'vikunja:3/7', body: { role: 'coder' } }] }); const mine = store.get("SELECT * FROM task_snapshots WHERE source='self'"); assert.deepEqual([mine.via, mine.read_at, mine.role, mine.run], [null, null, 'pm', 'demo-pm']); // Each side records only its own event kinds and snapshot fields. const refuse = (record, code) => assert.throws(() => b.recordTask(record), (e) => e.code === code, JSON.stringify(record)); refuse({ business: 'demo', events: [{ kind: 'task.created', body: {} }] }, 'reserved-event'); refuse({ cap: pm, business: 'demo', events: [{ kind: 'task.changed.external', body: {} }] }, 'reserved-event'); refuse({ business: 'demo', events: [{ kind: 'human.input', body: {} }] }, 'reserved-event'); refuse({ cap: pm, business: 'demo', snapshots: [{ ...self(), via: 'board' }] }, 'invalid-request'); refuse({ business: 'demo', snapshots: [poll({ via: 'ui' })] }, 'invalid-snapshot'); refuse({ business: 'demo', snapshots: [poll({ digest: 'B'.repeat(64) })] }, 'invalid-snapshot'); refuse({ business: 'demo', snapshots: [poll({ task_ref: 'jira:3/7' })] }, 'invalid-snapshot'); refuse({ business: 'demo', snapshots: [poll({ fields: [] })] }, 'invalid-snapshot'); refuse({ business: 'demo', events: [{ kind: 'task.missing', body: 'gone' }] }, 'invalid-request'); refuse({ business: 'demo', events: {} }, 'invalid-request'); refuse({ business: 'nope' }, 'unknown-business'); // A cap must be an agent of the business it records for. refuse({ cap: human, business: 'demo' }, 'agent-required'); refuse({ cap: agent('pm', 'other'), business: 'demo' }, 'agent-required'); refuse({ business: 'demo', events: [{ kind: 'task.missing', body: { note: SECRET } }] }, 'secret-detected'); assert.equal(count('task_snapshots'), 2); assert.equal(count('events', "kind LIKE 'task.%'"), 2); }); test('read_at must be one canonical UTC format, so the projection compares strings safely', (t) => { const { b, count } = setup(t); // task_current orders read_at as text; a second format would sort wrong against the broker's `at`. for (const read_at of [ '2026-10-08T12:00:00Z', '2026-10-08T12:00:00.000+00:00', '2026-10-08 12:00:00.000Z', '2026-02-30T12:00:00.000Z', Date.parse(at), null, ]) assert.throws( () => b.recordTask({ business: 'demo', snapshots: [poll({ read_at })] }), (e) => e.code === 'invalid-snapshot', String(read_at), ); assert.equal(count('task_snapshots'), 0); }); test('a bad entry refuses the whole record', (t) => { const { b, count } = setup(t); assert.throws(() => b.recordTask({ business: 'demo', snapshots: [poll(), poll({ task_ref: 'vikunja:3/8' })], events: [ { kind: 'task.changed.external', body: {} }, { kind: 'task.state', body: {} }, ], }), ); assert.equal(count('task_snapshots'), 0); assert.equal(count('events', "kind LIKE 'task.%'"), 0); }); test('taskView reads the projection for one business', (t) => { const { b, human } = setup(t); b.recordTask({ business: 'demo', snapshots: [poll(), poll({ task_ref: 'vikunja:3/8', fields: { bucket: 14, done: true } })] }); assert.deepEqual( b.taskView('demo', 'current').map((r) => [r.task_ref, r.fields.bucket]), [ ['vikunja:3/7', 11], ['vikunja:3/8', 14], ], ); assert.equal(b.taskView('demo', 'current', 'vikunja:3/8').fields.done, true); assert.equal(b.taskView('demo', 'current', 'vikunja:3/9'), null); assert.deepEqual(b.taskView('other', 'current'), []); assert.deepEqual( b.taskView('demo', 'open').map((r) => ({ ...r })), [{ task_ref: 'vikunja:3/7', bucket: 11 }], ); const { request } = b.request(human, { verb: 'message.send', args: { to: 'pm', body: 'please' } }); assert.equal(b.taskView('demo', 'input', request), true); assert.equal(b.taskView('other', 'input', request), false); assert.equal(b.taskView('demo', 'input', 'evt-none'), false); assert.throws(() => b.taskView('demo', 'current', 'vikunja:x'), (e) => e.code === 'invalid-request'); assert.throws(() => b.taskView('demo', 'snapshots'), (e) => e.code === 'invalid-request'); assert.throws(() => b.taskView('nope', 'current'), (e) => e.code === 'unknown-business'); }); test('requestTask hands only a holder and a task verb to the handler, and records refusals', async (t) => { const { b, agent, human, store } = setup(t); const pm = agent('pm'); const calls = []; const handler = async (cap, verb, args) => { calls.push([cap === pm, verb, args]); if (args.fail === 'typed') throw new BusError('tracker-unavailable'); if (args.fail === 'raw') throw new Error('socket hang up at http://vikunja.test'); if (args.fail === 'leak') return { echo: SECRET }; return { task_ref: 'vikunja:3/7' }; }; assert.deepEqual(TASK_VERBS.length, 8); assert.deepEqual(await b.requestTask(pm, { verb: 'task.create', args: { title: 'x' } }, handler), { task_ref: 'vikunja:3/7' }); const refused = () => store.all("SELECT body, subject FROM events WHERE kind='action.refused' ORDER BY seq").map((e) => [JSON.parse(e.body).code, e.subject]); const refuse = async (cap, request, code) => assert.rejects(b.requestTask(cap, request, handler), (e) => e instanceof BusError && e.code === code, JSON.stringify(request)); await refuse(pm, { verb: 'message.send', args: {} }, 'unknown-verb'); await refuse(pm, { verb: 'task.create', args: [] }, 'invalid-request'); await refuse(pm, { verb: 'task.create', extra: 1 }, 'invalid-request'); await refuse(pm, { verb: 'task.create', args: { title: SECRET } }, 'secret-detected'); await refuse(human, { verb: 'task.close', args: { task_ref: 'vikunja:3/7' } }, 'agent-required'); const stale = b.bindLaunch({ business: 'demo', role: 'cto', run: 'unclaimed', harness: 'pi' }); await refuse(stale, { verb: 'task.close', args: { task_ref: 'vikunja:3/7' } }, 'not-holder'); assert.equal(calls.length, 1); await refuse(pm, { verb: 'task.close', args: { task_ref: 'vikunja:3/7', fail: 'typed' } }, 'tracker-unavailable'); await refuse(pm, { verb: 'task.close', args: { task_ref: 'vikunja:3/7', fail: 'raw' } }, 'adapter-failed'); await refuse(pm, { verb: 'task.close', args: { task_ref: 'vikunja:3/7', fail: 'leak' } }, 'secret-detected'); assert.equal(calls.length, 4); assert.deepEqual(refused(), [ ['unknown-verb', null], ['invalid-request', null], ['invalid-request', null], ['secret-detected', null], ['agent-required', 'vikunja:3/7'], ['not-holder', 'vikunja:3/7'], ['tracker-unavailable', 'vikunja:3/7'], ['adapter-failed', 'vikunja:3/7'], ['secret-detected', 'vikunja:3/7'], ]); // Neither the raw error text nor the secret reaches the database. const all = JSON.stringify(store.all('SELECT * FROM events')); assert.ok(!all.includes('hang up') && !all.includes(SECRET)); }); test('the server sends task verbs to the adapter with its own timeout; other verbs stay synchronous', async (t) => { const { b, agent, store } = setup(t); const pm = agent('pm'); const path = join(store.directory, 'broker.sock'); let seen = null; const tasks = { timeout: 2000, handle: async (cap, verb, args) => { seen = verb; // Longer than the server's idle timeout of 50 ms below. await new Promise((ok) => setTimeout(ok, 200)); return { verb, args }; }, }; const server = await serve({ broker: b, path, timeout: 50, tasks }); t.after(() => server.close()); const c = new Client({ path, cap: pm }); assert.deepEqual(await c.call('task.close', { task_ref: 'vikunja:3/7' }), { verb: 'task.close', args: { task_ref: 'vikunja:3/7' } }); assert.equal(seen, 'task.close'); seen = null; assert.equal((await c.call('message.send', { to: 'pm', body: 'hi' })).request, null); assert.equal(seen, null); await assert.rejects(c.call('task.sql', {}), /unknown-verb/); // A client that gives up first gets outcome-unknown, never a retry. await assert.rejects(new Client({ path, cap: pm, timeout: 50 }).call('task.create', { title: 'x' }), /outcome-unknown/); }); test('without an adapter the server refuses every task verb', async (t) => { const { b, agent, store } = setup(t); const path = join(store.directory, 'broker.sock'); const server = await serve({ broker: b, path }); t.after(() => server.close()); const c = new Client({ path, cap: agent('pm') }); for (const verb of TASK_VERBS) await assert.rejects(c.call(verb, { task_ref: 'vikunja:3/7' }), /unknown-verb/, verb); }); test('the runtime refuses an invalid adapter and closes a valid one', async (t) => { const root = mkdtempSync(join(tmpdir(), 'bus-tasks-rt-')); t.after(() => rmSync(root, { recursive: true, force: true })); let closed = 0; for (const made of [null, {}, { handle: () => {}, timeout: 0 }, { handle: () => {}, timeout: 1.5 }, { handle: 'x', timeout: 10 }]) await assert.rejects( startBroker({ dataRoot: root, businesses, tasks: async () => (made ? { ...made, close: () => closed++ } : made) }), (e) => e.code === 'invalid-adapter', JSON.stringify(made), ); // Each refusal closed what it had opened, the adapter included; the writer lock is free. assert.equal(closed, 4); assert.equal(existsSync(join(root, 'bus/writer.lock')), false); let given = null; const runtime = await startBroker({ dataRoot: root, businesses, tasks: async (deps) => { given = Object.keys(deps).sort(); return { handle: async () => ({ ok: 1 }), timeout: 1000, close: async () => closed++ }; }, }); assert.deepEqual(given, ['broker', 'businesses', 'credentials']); await runtime.close(); assert.equal(closed, 5); await assert.rejects( startBroker({ dataRoot: root, businesses, tasks: async () => { throw new BusError('tracker-config'); }, }), (e) => e.code === 'tracker-config', ); assert.equal(existsSync(join(root, 'bus/writer.lock')), false); }); async function boot(t, config) { const child = fork(new URL('../src/process.mjs', import.meta.url), [], { stdio: ['ignore', 'pipe', 'pipe', 'ipc'] }); t.after(() => { if (child.exitCode === null) child.kill('SIGKILL'); }); const reply = Promise.race([ once(child, 'message'), once(child, 'exit').then(() => { throw Error('exited-before-reply'); }), ]); child.send({ op: 'boot', config }); return { child, reply: (await reply)[0] }; } test('the process loads the S3 adapter from plain-data trackers', async (t) => { const root = mkdtempSync(join(tmpdir(), 'bus-tasks-proc-')); t.after(() => rmSync(root, { recursive: true, force: true })); const launches = [{ business: 'demo', role: 'pm', run: 'r1', harness: 'pi', pid: process.pid, startTime: '1' }]; // A tracker for a business the host doesn't define refuses the boot with the adapter's code. const bad = await boot(t, { dataRoot: root, businesses, launches, trackers: { nope: { baseUrl: 'http://127.0.0.1:9', project: 1 } } }); assert.deepEqual(bad.reply, { ok: false, error: 'tracker-config' }); await once(bad.child, 'exit'); // An empty tracker map loads the adapter; a business without an entry refuses task verbs. const { child, reply } = await boot(t, { dataRoot: root, businesses, launches, trackers: {} }); assert.equal(reply.ok, true); const pm = new Client({ path: reply.path, cap: reply.launches[0].cap }); await pm.call('role.claim'); await assert.rejects(pm.call('task.create', { title: 'x' }), /tracker-unconfigured/); const exit = once(child, 'exit'); child.send({ op: 'close' }); assert.equal((await exit)[0], 0); });