Dewey's round 3 candidate, manifest
agents/dewey/work/queue-40/candidate-manifest-r3.sha256 (d0aa0ded,
27 files, checked OK in the canonical tree).
- WebUI inbox, tasks, agents and trail views, read-only over /api/bus.
The README says the bus proof ends at the Console process.
- CHAT-03 seal: the engine command is fixed, the engine environment is
explicit, SEAL_FLAGS has --no-approve, escalating is cleared on throw.
- Terminal input typed after Ctrl-T or Ctrl-O is held. Only the run whose
own parse set held drains it (T1), and #run catches errors per action.
- DEFERRED keeps N2 and moves F2 to done, citing T1.
Reviews: Filbert approve (comment 27011, rev 260), Darkwing approve
(27013, rev 264). Landing gate on 8cad7722 plus the candidate: webui 22,
conversation 161, control-board 124, every scripts/test-*.sh green,
test-task 98/0. Mutant Mr survives; its flows test is the first
follow-up row.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
218 lines
13 KiB
JavaScript
218 lines
13 KiB
JavaScript
// 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 = '<img src=x onerror="window.injected=true"> 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);
|
||
}
|
||
});
|