Files
stack/packages/bus/tests/transport.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

120 lines
5.0 KiB
JavaScript

import test from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, rmSync, lstatSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { connect } from 'node:net';
import { Store } from '../src/store.mjs';
import { Broker } from '../src/broker.mjs';
import { serve } from '../src/server.mjs';
import { Client } from '../src/client.mjs';
import { views } from '../src/views.mjs';
const businesses = {
demo: {
id: 'demo',
human: 'jason',
arbiters: { technical: 'cto', delivery: 'cto' },
roles: { cto: { authority: { withinRole: ['message.send'], crossRole: [] } } },
},
};
async function setup(t) {
const root = mkdtempSync(join(tmpdir(), 'bus-wire-')),
store = new Store(root),
broker = new Broker({ store, businesses });
const path = join(store.directory, 'broker.sock');
const server = await serve({ broker, path });
t.after(async () => {
await server.close();
store.close();
rmSync(root, { recursive: true, force: true });
});
return { store, broker, path };
}
test('socket capability stamps launch identity; shared views use wire, no SQL client', async (t) => {
const { broker, path } = await setup(t),
cap = broker.bindLaunch({ business: 'demo', role: 'cto', run: 'one', harness: 'pi' });
assert.equal(lstatSync(path).mode & 0o777, 0o600);
const c = new Client({ path, cap });
await c.call('role.claim');
assert.equal((await views(c).agents())[0].holder_run, 'one');
await assert.rejects(new Client({ path, cap: 'forged' }).call('agents'), /unauthenticated/);
await assert.rejects(c.call('decision.resolve', { id: 'none', role: 'human' }), /invalid-request/);
await assert.rejects(c.call('sql', { query: 'SELECT * FROM meta' }), /unknown-verb/);
});
test('two wire claims serialize; a lost reply never automatically retries', async (t) => {
const { broker, path, store } = await setup(t),
caps = ['a', 'b'].map((run) => broker.bindLaunch({ business: 'demo', role: 'cto', run, harness: 'pi' }));
const outcomes = await Promise.allSettled(caps.map((cap) => new Client({ path, cap }).call('role.claim')));
assert.equal(outcomes.filter((x) => x.status === 'fulfilled').length, 1);
assert.equal(store.get("SELECT count(*) n FROM role_claims WHERE op='claim'").n, 1);
const cap = caps[outcomes.findIndex((x) => x.status === 'fulfilled')];
await new Client({ path, cap }).call('message.send', { to: 'cto', body: 'once' });
assert.equal((await new Client({ path, cap }).call('message.receive')).length, 1);
assert.deepEqual(await new Client({ path, cap }).call('message.receive'), []);
});
test('malformed, oversized and identity-forging envelopes refuse without echoing input', async (t) => {
const { path } = await setup(t);
async function raw(data) {
return new Promise((resolve, reject) => {
const s = connect(path);
let out = '';
s.on('connect', () => s.end(data));
s.on('data', (b) => (out += b));
s.on('end', () => resolve(out));
s.on('error', reject);
});
}
for (const input of [
'not-json\n',
JSON.stringify({ cap: 'x', verb: 'agents', role: 'human' }) + '\n',
'x'.repeat(70000) + '\n',
]) {
const out = await raw(input);
assert.equal(JSON.parse(out).ok, false);
assert.ok(!out.includes('not-json'));
assert.ok(out.length < 150);
}
});
test('client preserves UTF-8 when a response divides a multibyte character', async (t) => {
const { createServer } = await import('node:net');
const root = mkdtempSync(join(tmpdir(), 'bus-utf8-')),
path = join(root, 's');
const server = createServer((s) =>
s.once('data', () => {
const bytes = Buffer.from('{"ok":true,"result":"🙂"}\n');
const i = bytes.indexOf(Buffer.from('🙂')) + 1;
s.write(bytes.subarray(0, i));
setTimeout(() => s.end(bytes.subarray(i)), 10);
}),
);
await new Promise((resolve) => server.listen(path, resolve));
t.after(async () => {
await new Promise((resolve) => server.close(resolve));
rmSync(root, { recursive: true, force: true });
});
assert.equal(await new Client({ path, cap: 'fixture' }).call('agents'), '🙂');
});
test('committed mutation followed by dropped reply reports unknown and is never retried', async (t) => {
const { broker, store, path } = await setup(t);
const cap = broker.bindLaunch({ business: 'demo', role: 'cto', run: 'drop', harness: 'pi' });
broker.request(cap, { verb: 'role.claim' });
const { createServer } = await import('node:net');
let calls = 0;
const proxy = createServer((s) =>
s.once('data', (data) => {
calls++;
const r = JSON.parse(data);
broker.request(r.cap, { verb: r.verb, args: r.args });
s.destroy();
}),
);
await new Promise((resolve) => proxy.listen(path + '.drop', resolve));
t.after(() => new Promise((resolve) => proxy.close(resolve)));
await assert.rejects(
new Client({ path: path + '.drop', cap }).call('message.send', { to: 'cto', body: 'only once' }),
/outcome-unknown/,
);
assert.equal(calls, 1);
assert.equal(store.get('SELECT count(*) n FROM messages').n, 1);
});