Files
stack/packages/bus/tests/broker.test.mjs
T
jason.woltjeandClaude Opus 5.5 38828a2cb3 feat(bus): the bus and the broker core (row 37, S2, rocko)
Rocko's round 2 candidate, approved by Darkwing (#1519 comment 26757).
build.patch 40d7e838, manifest 61519059, 24 files under packages/bus,
schema v3b (179ffe35, lead decision 60). Integration gate in a git
worktree of 942dca9e (S1 in the tree) plus the patch: bus 43/43 and
business 60/60 on Node 24 and 26, every package test and every
scripts/test-*.sh green, test-task 98/98 with the live-provider cases.
Rulings from lead decisions 62 and 63: the human proof is cooperative in
slice 1, and a self-raised cross-role decision routes to the human.
Single-use gated approvals follow in row 43.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-10-05 17:15:53 -05:00

357 lines
15 KiB
JavaScript

import test from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Store } from '../src/store.mjs';
const { Broker } = await import('../src/broker.mjs').catch((e) => {
if (e.code === 'ERR_MODULE_NOT_FOUND') return {};
throw e;
});
const options = [
{ key: 'yes', text: 'Allow' },
{ key: 'no', text: 'Decline' },
];
export 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', 'decision.resolve.technical'], crossRole: [] } },
coder: { authority: { withinRole: ['message.send'], crossRole: ['task.scope.change'] } },
reviewer: { authority: { withinRole: ['message.send'], crossRole: [] } },
},
},
};
function setup(t, config = businesses) {
assert.equal(typeof Broker, 'function');
const root = mkdtempSync(join(tmpdir(), 'bus-broker-'));
const store = new Store(root);
const b = new Broker({ store, businesses: config });
t.after(() => {
store.close();
rmSync(root, { recursive: true, force: true });
});
const agent = (role, run = role + '-run') =>
b.bindLaunch({ business: 'demo', role, run, harness: 'pi', address: run });
const human = b.bindHuman({ business: 'demo', human: 'jason', via: 'cli', outsideAgent: true });
return { b, store, agent, human };
}
const call = (b, cap, verb, args = {}) => b.request(cap, { verb, args });
const raise = (b, cap, action, extra = {}) =>
call(b, cap, 'decision.raise', {
action,
question: 'Allow?',
options,
recommendation: 'no',
blocking: false,
...extra,
});
test('launch identity is stamped, payload identity is refused and stale holder cannot send', (t) => {
const { b, store, agent } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
assert.throws(() => call(b, c, 'message.send', { to: 'pm', body: 'hello', role: 'pm' }), /invalid-request/);
const m = call(b, c, 'message.send', { to: 'pm', body: 'hello' });
assert.equal(store.get('SELECT from_role FROM messages WHERE id=?', m.id).from_role, 'coder');
call(b, c, 'role.release');
assert.throws(() => call(b, c, 'message.send', { to: 'pm', body: 'late' }), /not-holder/);
assert.throws(() => call(b, 'bogus', 'role.claim'), /unauthenticated/);
});
test('decision classes route from policy; gated resolution is human-only, choice and target must match', (t) => {
const { b, agent, human } = setup(t),
c = agent('coder'),
cto = agent('cto');
call(b, c, 'role.claim');
call(b, cto, 'role.claim');
const cross = raise(b, c, 'task.scope.change', { domain: 'technical', target: 'task-1' });
assert.equal(cross.route_to, 'cto');
assert.throws(() => call(b, c, 'decision.resolve', { id: cross.id, choice: 'yes' }), /not-resolver/);
call(b, cto, 'decision.resolve', { id: cross.id, choice: 'yes' });
assert.equal(
b.authorize(c, 'task.scope.change', { decision: cross.id, target: 'task-1' }).class,
'cross-role',
);
assert.throws(
() => b.authorize(c, 'task.scope.change', { decision: cross.id, target: 'task-2' }),
/decision-mismatch/,
);
const d = raise(b, c, 'deploy', { target: 'release-1' });
assert.equal(d.route_to, 'human');
assert.throws(() => call(b, cto, 'decision.resolve', { id: d.id, choice: 'yes' }), /human-required/);
assert.throws(() => call(b, human, 'decision.resolve', { id: d.id, choice: 'other' }), /invalid-choice/);
call(b, human, 'decision.resolve', { id: d.id, choice: 'no' });
assert.throws(
() => b.authorize(c, 'deploy', { decision: d.id, target: 'release-1' }),
/decision-not-approved/,
);
assert.throws(() => call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' }), /decision-closed/);
});
test('claim exclusion, holder release, gated revoke and rerouting to a new holder are atomic', (t) => {
const { b, store, agent, human } = setup(t),
c = agent('coder'),
p = agent('pm'),
p2 = agent('pm', 'pm-next');
call(b, c, 'role.claim');
const m = call(b, c, 'message.send', { to: 'pm', body: 'waiting' });
assert.equal(store.get('SELECT holder_run FROM deliveries WHERE message=?', m.id).holder_run, null);
call(b, p, 'role.claim');
assert.throws(() => call(b, p2, 'role.claim'), /role-already-held/);
const delivered = call(b, p, 'message.receive');
assert.equal(delivered[0].id, m.id);
assert.equal(call(b, p, 'message.receive').length, 0, 'receiving twice does not replay');
assert.throws(() => call(b, p2, 'role.release'), /not-holder/);
const d = raise(b, c, 'role.revoke', { target: 'pm' });
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
call(b, c, 'role.revoke', { role: 'pm', decision: d.id });
call(b, p2, 'role.claim');
assert.throws(() => call(b, p, 'message.receive'), /not-holder/);
assert.equal(call(b, p2, 'message.receive').length, 0, 'old delivered message is not replayed');
});
test('launch events require a human CLI capability; generic emit cannot forge authority events', (t) => {
const { b, agent, human, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
assert.throws(() => call(b, c, 'launch.revoke'), /human-required/);
assert.throws(() => call(b, c, 'event.emit', { kind: 'launch.revoked', body: {} }), /reserved-event/);
call(b, human, 'launch.revoke');
assert.equal(store.get('SELECT state FROM launch_state').state, 'revoked');
call(b, human, 'launch.restore');
assert.equal(store.get('SELECT state FROM launch_state').state, 'allowed');
assert.throws(
() => b.bindHuman({ business: 'demo', human: 'jason', via: 'webui', outsideAgent: true }),
/human-required/,
);
assert.throws(
() => b.bindHuman({ business: 'demo', human: 'jason', via: 'cli', outsideAgent: false }),
/human-required/,
);
});
test('within-role decisions close atomically and invalid options or blocking omissions refuse', (t) => {
const { b, agent, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
const d = raise(b, c, 'message.send', { choice: 'yes' });
assert.equal(d.class, 'within-role');
assert.equal(store.get('SELECT choice FROM decision_events WHERE decision=?', d.id).choice, 'yes');
assert.throws(() => raise(b, c, 'deploy', { blocking: true }), /task-required/);
assert.throws(
() =>
raise(b, c, 'deploy', {
options: [
{ key: 'a', text: 'a' },
{ key: 'a', text: 'b' },
],
recommendation: 'a',
}),
/invalid-options/,
);
assert.throws(() => raise(b, c, 'made.up'), /unknown-action/);
assert.equal(call(b, c, 'inbox').length, 0);
});
test('observer capabilities read human inbox but cannot mutate or forge launch identity', (t) => {
const { b, agent } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
raise(b, c, 'deploy');
const view = b.bindReader({ business: 'demo' });
assert.equal(call(b, view, 'inbox').length, 1);
for (const [verb, args] of [
['role.claim', {}],
['launch.revoke', {}],
['message.send', { to: 'pm', body: 'x' }],
])
assert.throws(() => call(b, view, verb, args), /read-only/);
});
test('task action subjects and linked decision trail are complete and ordered', (t) => {
const { b, agent, human, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
const request = call(b, human, 'message.send', { to: 'coder', body: 'make this' });
const input = store.get("SELECT id FROM events WHERE kind='human.input' AND subject=?", request.id).id;
b.recordEvent(c, {
kind: 'task.created',
subject: 'vikunja:3/1',
body: { request: input, requirement: 'REQ-TASK-1' },
});
const d = raise(b, c, 'deploy', { target: 'vikunja:3/1', task_ref: 'vikunja:3/1' });
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
b.authorize(c, 'deploy', { target: 'vikunja:3/1', decision: d.id });
const rows = call(b, human, 'trail', { subject: 'vikunja:3/1' });
assert.ok(rows.some((r) => r.id === input));
assert.ok(rows.some((r) => r.id === request.id));
assert.ok(rows.some((r) => r.decision === d.id && r.op === 'resolved'));
assert.ok(
store
.all(
"SELECT subject FROM events WHERE kind='action.allowed' AND json_extract(body,'$.action')='deploy'",
)
.every((e) => e.subject === 'vikunja:3/1'),
);
for (let i = 1; i < rows.length; i++)
assert.ok(rows[i].at > rows[i - 1].at, 'broker write order has unique timestamps');
assert.throws(
() => b.recordEvent(c, { kind: 'review.verdict', body: { task_ref: 'vikunja:3/1' } }),
/task-subject-required/,
);
});
test('launch binding is durable and reconnecting requires the identical trusted record', (t) => {
const { b, agent, store } = setup(t);
agent('coder', 'durable');
assert.equal(store.get("SELECT count(*) n FROM events WHERE kind='session.launched'").n, 1);
const again = new Broker({ store, businesses });
again.bindLaunch({ business: 'demo', role: 'coder', run: 'durable', harness: 'pi', address: 'durable' });
assert.equal(store.get("SELECT count(*) n FROM events WHERE kind='session.launched'").n, 1);
const wrong = new Broker({ store, businesses });
assert.throws(
() => wrong.bindLaunch({ business: 'demo', role: 'cto', run: 'durable', harness: 'pi' }),
/launch-record-mismatch/,
);
});
test('business isolation includes inherited object names and cross-business message references', (t) => {
const { store } = setup(t);
const config = structuredClone(businesses);
config.other = { ...structuredClone(config.demo), id: 'other' };
const b = new Broker({ store, businesses: config });
assert.throws(
() => b.bindLaunch({ business: 'demo', role: 'toString', run: 'fake', harness: 'pi' }),
/unknown-role/,
);
assert.throws(() => b.bindReader({ business: 'constructor' }), /unknown-business/);
const caps = ['demo', 'other'].map((business) =>
b.bindLaunch({ business, role: 'coder', run: 'r', harness: 'pi' }),
);
for (const c of caps) call(b, c, 'role.claim');
const m = call(b, caps[0], 'message.send', { to: 'pm', body: 'private' });
assert.throws(
() => call(b, caps[1], 'message.send', { to: 'pm', body: 'reply', in_reply_to: m.id }),
/message-not-found/,
);
const d = raise(b, caps[0], 'deploy');
assert.throws(
() => call(b, caps[1], 'decision.resolve', { id: d.id, choice: 'yes' }),
/decision-not-found/,
);
assert.deepEqual(call(b, caps[1], 'trail', { subject: d.id }), []);
});
test('authority never transfers between action, run, target, unresolved or replaced role holder', (t) => {
const { b, agent, human } = setup(t),
c = agent('coder'),
p = agent('pm');
call(b, c, 'role.claim');
call(b, p, 'role.claim');
const d = raise(b, c, 'deploy', { target: 'v1' });
assert.throws(() => b.authorize(c, 'deploy', { target: 'v1', decision: d.id }), /decision-not-approved/);
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
assert.throws(() => b.authorize(c, 'spend', { target: 'v1', decision: d.id }), /decision-mismatch/);
call(b, c, 'role.release');
const c2 = agent('coder', 'r2');
call(b, c2, 'role.claim');
assert.throws(() => b.authorize(c2, 'deploy', { target: 'v1', decision: d.id }), /decision-mismatch/);
const revoke = raise(b, c2, 'role.revoke', { target: 'pm' });
call(b, human, 'decision.resolve', { id: revoke.id, choice: 'yes' });
call(b, p, 'role.release');
const p2 = agent('pm', 'p2');
call(b, p2, 'role.claim');
assert.throws(() => call(b, c2, 'role.revoke', { role: 'pm', decision: revoke.id }), /decision-mismatch/);
});
test('task projection uses schema current view, skipping earlier and equal-start polls', (t) => {
const { b, store, human } = setup(t);
const put = (source, at, read, bucket) =>
store.run(
'INSERT INTO task_snapshots(at,business,task_ref,updated,digest,fields,source,via,read_at,role,run) VALUES(?,?,?,?,?,?,?,?,?,?,?)',
at,
'demo',
'vikunja:3/46',
'2026-10-04T00:00:00.000Z',
'a'.repeat(64),
JSON.stringify({ bucket, done: 0 }),
source,
source === 'poll' ? 'board' : null,
read,
source === 'self' ? 'coder' : null,
source === 'self' ? 'r1' : null,
);
put('self', '2026-10-04T00:00:10.000Z', null, 13);
put('poll', '2026-10-04T00:00:11.000Z', '2026-10-04T00:00:09.000Z', 11);
assert.equal(call(b, human, 'tasks')[0].fields.bucket, 13);
put('poll', '2026-10-04T00:00:12.000Z', '2026-10-04T00:00:10.000Z', 11);
assert.equal(call(b, human, 'tasks')[0].fields.bucket, 13);
put('poll', '2026-10-04T00:00:13.000Z', '2026-10-04T00:00:12.000Z', 11);
assert.equal(call(b, human, 'tasks')[0].fields.bucket, 11);
assert.equal(store.get('SELECT bucket FROM tasks_open').bucket, 11);
});
test('revocation permanently bars the old run from reclaiming first, including after broker restart', (t) => {
const { b, agent, human, store } = setup(t),
coder = agent('coder'),
old = agent('pm', 'revoked-pm');
call(b, coder, 'role.claim');
call(b, old, 'role.claim');
const d = raise(b, coder, 'role.revoke', { target: 'pm' });
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
call(b, coder, 'role.revoke', { role: 'pm', decision: d.id });
assert.throws(() => call(b, old, 'role.claim'), /run-revoked/);
const reopened = new Broker({ store, businesses });
const oldAgain = reopened.bindLaunch({
business: 'demo',
role: 'pm',
run: 'revoked-pm',
harness: 'pi',
address: 'revoked-pm',
});
assert.throws(() => call(reopened, oldAgain, 'role.claim'), /run-revoked/);
const replacement = reopened.bindLaunch({
business: 'demo',
role: 'pm',
run: 'replacement',
harness: 'pi',
});
call(reopened, replacement, 'role.claim');
});
test('empty message references refuse before storage; refusal-evidence failure stays a typed error', (t) => {
const { b, agent, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
for (const key of ['in_reply_to', 'corrects'])
assert.throws(() => call(b, c, 'message.send', { to: 'pm', body: 'x', [key]: '' }), /invalid-request/);
const original = store.transaction;
store.transaction = () => {
throw Error('raw private database detail');
};
try {
assert.throws(
() => call(b, c, 'unknown.verb'),
(e) => e.code === 'storage-unavailable' && e.message === 'storage-unavailable',
);
} finally {
store.transaction = original;
}
});
test('both arbiters require human resolution when their cross-role route is themselves', (t) => {
const config = structuredClone(businesses);
config.demo.roles.cto.authority.crossRole = ['task.scope.change'];
config.demo.roles.pm.authority.crossRole = ['task.priority.change'];
const { b, agent, human } = setup(t, config);
for (const [role, action, domain] of [
['cto', 'task.scope.change', 'technical'],
['pm', 'task.priority.change', 'delivery'],
]) {
const cap = agent(role);
call(b, cap, 'role.claim');
const d = raise(b, cap, action, { domain, target: 'vikunja:3/41' });
assert.equal(d.class, 'cross-role');
assert.equal(d.route_to, 'human');
assert.throws(() => call(b, cap, 'decision.resolve', { id: d.id, choice: 'yes' }), /human-required/);
assert.throws(
() => b.authorize(cap, action, { decision: d.id, target: 'vikunja:3/41' }),
/decision-not-approved/,
);
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
assert.equal(b.authorize(cap, action, { decision: d.id, target: 'vikunja:3/41' }).class, 'cross-role');
}
});