// Slice 1 S5 (#1522): the /api/bus/* routes and the Console's bus transport. // The routes run against an in-process broker through its reader session; // the transport runs a fake human CLI. Nothing here starts the real human // CLI or reaches a real broker. import { test } from 'node:test'; import assert from 'node:assert/strict'; import { mkdtempSync, mkdirSync, readFileSync, rmSync, writeFileSync } from 'node:fs'; import { spawnSync } from 'node:child_process'; import { fileURLToPath } from 'node:url'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { Store } from '../../bus/src/store.mjs'; import { Broker } from '../../bus/src/broker.mjs'; import { views } from '../../bus/src/views.mjs'; import { hostFile, startTimeOf } from '../../cli/src/host.mjs'; import { startServer } from '../src/serve.mjs'; import { BusReadError, busReader, humanCall } from '../src/bus.mjs'; import { close } from './fixture.mjs'; const HOSTILE = ' ‮evil'; const at = '2026-10-08T12:00:00.000Z'; 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: [] } }, coder: { authority: { withinRole: ['message.send'], crossRole: ['task.scope.change'] } }, }, }, }; const tmp = (t, prefix) => { const root = mkdtempSync(join(tmpdir(), prefix)); t.after(() => rmSync(root, { recursive: true, force: true })); return root; }; // A broker with one open decision, one claimed role and one externally changed task. function seeded(t) { const store = new Store(tmp(t, 'webui-bus-')); t.after(() => store.close()); const b = new Broker({ store, businesses }); const coder = b.bindLaunch({ business: 'demo', role: 'coder', run: 'coder-run', harness: 'pi', address: 'coder-run' }); b.request(coder, { verb: 'role.claim' }); const decision = b.request(coder, { verb: 'decision.raise', args: { action: 'git.push.protected', target: 'refactor', task_ref: 'vikunja:32/7', question: `Push the release? ${HOSTILE}`, options: [{ key: 'yes', text: 'Allow' }, { key: 'no', text: 'Decline' }], recommendation: 'no', blocking: true } }); b.recordTask({ business: 'demo', snapshots: [{ task_ref: 'vikunja:32/7', updated: at, digest: 'b'.repeat(64), fields: { title: HOSTILE, bucket: 11, done: false }, via: 'board', read_at: at }], events: [{ kind: 'task.changed.external', subject: 'vikunja:32/7', body: { changed: ['bucket'] } }] }); const reader = b.bindReader({ business: 'demo' }); const calls = []; const bus = views({ call: async (verb, args = {}) => { calls.push(verb); return structuredClone(b.request(reader, { verb, args })); } }); return { decision, bus, calls }; } test('bus routes serve the four reads from the reader view, unescaped JSON for the page to escape', async t => { const { decision, bus, calls } = seeded(t); const web = await startServer({ port: 0, board: 'http://127.0.0.1:9', bus }); t.after(() => close(web)); const base = `http://127.0.0.1:${web.address().port}`; const read = async path => { const r = await fetch(base + path); return { status: r.status, cache: r.headers.get('cache-control'), body: await r.json() }; }; const inbox = await read('/api/bus/inbox'); assert.equal(inbox.status, 200); assert.equal(inbox.cache, 'no-store'); assert.match(inbox.body.at, /^\d{4}-\d\d-\d\dT/); assert.equal(inbox.body.rows.length, 1); assert.equal(inbox.body.rows[0].id, decision.id); assert.deepEqual(inbox.body.rows[0].authorization, { action: 'git.push.protected', target: 'refactor', approvalChoice: 'yes' }); assert.ok(inbox.body.rows[0].question.endsWith(HOSTILE)); assert.equal(inbox.body.rows[0].blocking, true); assert.deepEqual(inbox.body.rows[0].options.map(o => o.key), ['yes', 'no']); const tasks = await read('/api/bus/tasks'); assert.equal(tasks.body.rows.length, 1); assert.equal(tasks.body.rows[0].source, 'poll'); assert.equal(tasks.body.rows[0].fields.title, HOSTILE); const agents = await read('/api/bus/agents'); assert.deepEqual(agents.body.rows.map(r => [r.role, r.op]), [['coder', 'claim']]); const id = inbox.body.rows[0].id; const trail = await read('/api/bus/trail?subject=' + encodeURIComponent(id)); assert.equal(trail.status, 200); assert.ok(trail.body.rows.some(r => r.table === 'decisions' && r.id === id)); const taskTrail = await read('/api/bus/trail?subject=vikunja:32/7'); assert.ok(taskTrail.body.rows.some(r => r.table === 'task_snapshots')); assert.ok(taskTrail.body.rows.some(r => r.table === 'events' && r.kind === 'task.changed.external')); // Bad subjects never reach the bus. const before = calls.length; for (const q of ['', '?subject=', '?subject=' + 'a'.repeat(161), '?subject=-x', '?subject=a%20b', '?subject=a%00b']) assert.deepEqual([q, (await read('/api/bus/trail' + q)).body.error], [q, 'invalid-request']); assert.equal(calls.length, before); assert.deepEqual([...new Set(calls)].sort(), ['agents', 'inbox', 'tasks', 'trail']); }); test('bus routes are GET-only, have no write verb, and keep the same-origin checks', async t => { const calls = []; const bus = Object.fromEntries(['inbox', 'tasks', 'agents', 'trail'].map(v => [v, async () => { calls.push(v); return []; }])); const web = await startServer({ port: 0, board: 'http://127.0.0.1:9', bus }); t.after(() => close(web)); const base = `http://127.0.0.1:${web.address().port}`; for (const method of ['POST', 'PUT', 'DELETE', 'HEAD', 'OPTIONS']) assert.equal((await fetch(base + '/api/bus/inbox', { method })).status, 405, method); for (const path of ['/api/bus/decide', '/api/bus/seen', '/api/bus/', '/api/bus/inbox/x', '/api/bus/request', '/api/bus/raise']) { assert.equal((await fetch(base + path)).status, 404, path); assert.equal((await fetch(base + path, { method: 'POST', headers: { 'content-type': 'application/json' }, body: '{}' })).status, 404, path); } assert.equal((await fetch(base + '/api/bus/inbox', { headers: { origin: 'https://evil.example' } })).status, 403); assert.deepEqual(calls, []); for (const path of ['/bus.js', '/bus.css']) assert.equal((await fetch(base + path)).status, 200, path); }); test('bus refusals map to statuses and never carry a path or a stack', async t => { let fail; const bus = { inbox: async () => { throw fail; }, tasks: async () => ({ not: 'rows' }), agents: async () => [], trail: async () => [] }; const web = await startServer({ port: 0, board: 'http://127.0.0.1:9', bus }); const none = await startServer({ port: 0, board: 'http://127.0.0.1:9' }); t.after(() => Promise.all([close(web), close(none)])); const base = `http://127.0.0.1:${web.address().port}`; const inbox = async () => { const r = await fetch(base + '/api/bus/inbox'); return [r.status, (await r.json())]; }; for (const [code, status] of [['human-required', 403], ['unauthenticated', 403], ['unknown-business', 403], ['read-only', 403], ['outcome-unknown', 503], ['no-bus-host', 503], ['invalid-response', 502], ['response-too-large', 502], ['something-new', 502]]) { fail = new BusReadError(code); const [s, body] = await inbox(); assert.deepEqual([code, s, body.error], [code, status, code]); } fail = Object.assign(new Error('ENOENT: /home/someone/.mosaic-dev/bus/broker.sock'), { stack: 'at secret (/x.mjs:1)' }); const [s, body] = await inbox(); assert.equal(s, 502); assert.equal(body.error, 'invalid-response'); assert.doesNotMatch(JSON.stringify(body), /home|mosaic-dev|secret|\.mjs/); const shape = await fetch(base + '/api/bus/tasks'); assert.deepEqual([shape.status, (await shape.json()).error], [502, 'invalid-response']); const unset = await fetch(`http://127.0.0.1:${none.address().port}/api/bus/agents`); assert.deepEqual([unset.status, (await unset.json()).error], [503, 'not-configured']); }); // A fake human CLI: it reads the request, then acts on args.mode. const FAKE = ` import { readFileSync, writeFileSync } from 'node:fs'; const request = JSON.parse(readFileSync(0, 'utf8')); const mode = request.args.mode; if (mode === 'refuse') { process.stderr.write('Some warning: Text\\nhuman-required\\n'); process.exit(2); } if (mode === 'garbage') { process.stdout.write('not json'); process.exit(0); } if (mode === 'silent-fail') process.exit(2); if (mode === 'big') { process.stdout.write('"' + 'x'.repeat(4096) + '"'); process.exit(0); } if (mode === 'slow') { setTimeout(() => process.stdout.write('[]'), 60000); } else if (mode === 'pid') { writeFileSync(request.args.pidFile, String(process.pid)); setInterval(() => {}, 60000); } else process.stdout.write(JSON.stringify({ request, socket: process.argv[2] })); `; function fake(t) { const root = tmp(t, 'webui-bus-cli-'); const cli = join(root, 'fake-human-cli.mjs'); writeFileSync(cli, FAKE); return { root, cli }; } test('humanCall speaks the human transport protocol without blocking the server', async t => { const { cli } = fake(t); const call = humanCall({ socket: '/run/fake.sock', business: 'demo', cli, timeoutMs: 1500, maxBytes: 1024 }); assert.deepEqual(await call('trail', { subject: 'd-1' }), { request: { business: 'demo', verb: 'trail', args: { subject: 'd-1' } }, socket: '/run/fake.sock' }); const code = async args => { try { await call('inbox', args); return 'resolved'; } catch (e) { assert.ok(e instanceof BusReadError); return e.code; } }; assert.equal(await code({ mode: 'refuse' }), 'human-required'); assert.equal(await code({ mode: 'garbage' }), 'invalid-response'); assert.equal(await code({ mode: 'silent-fail' }), 'invalid-response'); assert.equal(await code({ mode: 'big' }), 'response-too-large'); assert.equal(await humanCall({ socket: 's', business: 'demo', cli: join(tmpdir(), 'no-such-dir-x', 'cli.mjs') })('inbox').catch(e => e.code), 'invalid-response'); // A slow read stays slow on its own: the server keeps answering meanwhile. const web = await startServer({ port: 0, board: 'http://127.0.0.1:9', bus: views({ call: (verb) => call(verb, { mode: 'slow' }) }) }); t.after(() => close(web)); const base = `http://127.0.0.1:${web.address().port}`; const started = Date.now(); const slow = fetch(base + '/api/bus/inbox'); const health = await fetch(base + '/healthz'); assert.equal(health.status, 200); assert.ok(Date.now() - started < 1000, 'healthz waited for the bus read'); const r = await slow; assert.deepEqual([r.status, (await r.json()).error], [503, 'outcome-unknown']); assert.ok(Date.now() - started >= 1400); }); test('busReader takes --business, else the live bus host, else refuses', async t => { const { root, cli } = fake(t); const dataRoot = join(root, 'data'); const pinned = busReader({ dataRoot, socket: 'sock', business: 'pinned', cli }); assert.equal((await pinned.inbox()).request.business, 'pinned'); const reader = busReader({ dataRoot, socket: 'sock', cli }); const code = () => reader.agents().then(() => 'resolved', e => e.code); assert.equal(await code(), 'no-bus-host'); mkdirSync(join(dataRoot, 'bus-host'), { recursive: true }); writeFileSync(hostFile(dataRoot), '{not json'); assert.equal(await code(), 'no-bus-host'); writeFileSync(hostFile(dataRoot), JSON.stringify({ pid: process.pid, startTime: 'not-my-start', business: 'stale' })); assert.equal(await code(), 'no-bus-host'); writeFileSync(hostFile(dataRoot), JSON.stringify({ pid: process.pid, startTime: startTimeOf(process.pid), business: 'hosted' })); const reply = await reader.agents(); assert.deepEqual([reply.request.business, reply.request.verb, reply.socket], ['hosted', 'agents', 'sock']); // The host is read on every request. writeFileSync(hostFile(dataRoot), JSON.stringify({ pid: process.pid, startTime: startTimeOf(process.pid), business: 'moved' })); assert.equal((await reader.tasks()).request.business, 'moved'); }); test('humanCall kills a transport that runs past its timeout', async t => { const { root, cli } = fake(t); const pidFile = join(root, 'pid'); const call = humanCall({ socket: 's', business: 'demo', cli, timeoutMs: 500 }); assert.equal(await call('inbox', { mode: 'pid', pidFile }).catch(e => e.code), 'outcome-unknown'); const pid = Number(readFileSync(pidFile, 'utf8')); const gone = () => { try { process.kill(pid, 0); return false; } catch { return true; } }; // If it survived, stop it so the failure is this assertion, not a hang. t.after(() => { if (!gone()) process.kill(pid, 'SIGKILL'); }); // SIGKILL lands asynchronously; give the kernel a moment to reap it. for (let i = 0; i < 40 && !gone(); i++) await new Promise(r => setTimeout(r, 25)); assert.ok(gone(), `transport ${pid} still runs after the timeout`); }); test('serve refuses a --business value that is not a business id', () => { const cli = fileURLToPath(new URL('../src/cli.mjs', import.meta.url)); for (const business of ['../acme', '-x', 'acme corp']) { const result = spawnSync(process.execPath, [cli, 'serve', '--business', business, '--port', '0'], { encoding: 'utf8', timeout: 10000 }); assert.equal(result.status, 2, business); assert.match(result.stderr, /^refused: business must be a business id/, business); } });