feat(tasks): the Vikunja v2 adapter, broker task verbs and sync (row 38, S3, darkwing)

packages/tasks adds the Vikunja v2 client, the eight task verbs, the
board-plus-cursor poll with its 60 s window and the digest. The broker
gains the task verbs and boots trackers from the boot config (lead
decisions 66 to 68). Due dates are truncated to the second and recorded
as truncated (B1). A write that lands but whose final read fails counts
as landed, in update and in create (B2).

Candidate agents/darkwing/work/slice1-s3, build-r2.patch 71ce87e6,
manifest e10e30e3 (28 files). Filbert approved round 2 on #1520
(comment 26853). Darkwing's post-reset rerun: test-release 14/14,
test-task 98/98 (comment 26857).

Integration gate in a worktree on c4baf779 with the patch applied:
bus 67, business 60, control-board 124, discord 173, ledger 78,
mosaic 69, queue 148, seat 19, tasks 51 and webui 14, all with no
failures. Conversation is 149/3. The three cohort kill cases (K1, K3,
K10) fail the same on the unpatched base, and the patch doesn't touch
the package. Every scripts/test-*.sh is green. test-release 14/14 and
test-task 98/98 ran on the existing gate2 compose network, because the
host's Docker address pools are exhausted. No network was created or
pruned.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
This commit is contained in:
2026-10-09 07:40:48 -05:00
co-authored by Claude Opus 5.5
parent c4baf77916
commit 7e73c2cd13
28 changed files with 6504 additions and 25 deletions
+121 -19
View File
@@ -38,6 +38,19 @@ const GATED = new Set([
'role.revoke',
]);
const CLOSED = "('resolved','withdrawn','expired')";
export const TASK_VERBS = Object.freeze(ACTIONS.filter((a) => a.startsWith('task.')));
const SELF_KINDS = new Set(['task.created', 'task.assigned', 'task.state', 'task.closed', 'task.conflict']);
const POLL_KINDS = new Set([
'task.changed.external',
'task.missing',
'credential.expiring',
'credential.expired',
'credential.changed',
]);
const canonical = (x) =>
typeof x === 'string' && /^\d{4}-\d\d-\d\dT\d\d:\d\d:\d\d\.\d{3}Z$/.test(x) &&
Number.isFinite(Date.parse(x)) &&
new Date(x).toISOString() === x;
export class BusError extends Error {
constructor(code) {
super(code);
@@ -321,6 +334,95 @@ export class Broker {
return this.#event(s, kind, body, subject);
});
}
// Trusted S3 adapter only. With a cap it records a role's own completed write: the external call
// already happened, so this records and does not re-authorize. Without one it records sync reads.
recordTask({ cap = null, business, snapshots = [], events = [] }) {
const s = cap === null ? { business, role: null, run: null } : this.#session(cap);
this.#business(business);
if (cap !== null && (s.human || s.reader || s.business !== business)) fail('agent-required');
if (!Array.isArray(snapshots) || !Array.isArray(events)) fail('invalid-request');
const kinds = cap === null ? POLL_KINDS : SELF_KINDS;
this.#secretCheck({ snapshots, events });
return this.#store.transaction(() => ({
snapshots: snapshots.map((x) => this.#snapshot(s, x)),
events: events.map((e) => {
keys(e, ['kind', 'body', 'subject'], ['kind', 'body']);
if (!kinds.has(e.kind)) fail('reserved-event');
if (!object(e.body)) fail('invalid-request');
return this.#event(s, e.kind, e.body, e.subject ?? null);
}),
}));
}
#snapshot(s, x) {
const poll = s.role === null;
keys(
x,
['task_ref', 'updated', 'etag', 'digest', 'fields', ...(poll ? ['via', 'read_at'] : [])],
['task_ref', 'updated', 'digest', 'fields', ...(poll ? ['via', 'read_at'] : [])],
);
if (!taskRef(x.task_ref) || !object(x.fields) || !/^[0-9a-f]{64}$/.test(x.digest ?? ''))
fail('invalid-snapshot');
string(x.updated, 64);
if (x.etag !== undefined && x.etag !== null) string(x.etag, 256);
if (poll && (!['board', 'cursor', 'task', 'reconcile'].includes(x.via) || !canonical(x.read_at)))
fail('invalid-snapshot');
return this.#store.run(
'INSERT INTO task_snapshots(at,business,task_ref,updated,etag,digest,fields,source,via,read_at,role,run) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)',
this.#now(),
s.business,
x.task_ref,
x.updated,
x.etag ?? null,
x.digest,
JSON.stringify(x.fields),
poll ? 'poll' : 'self',
poll ? x.via : null,
poll ? x.read_at : null,
s.role,
s.run,
).lastInsertRowid;
}
// Trusted S3 reads. Agents read the same rows through the 'tasks' view verb.
taskView(business, view, arg = null) {
this.#business(business);
const parse = (r) => ({ ...r, fields: JSON.parse(r.fields) });
if (view === 'current' && arg === null)
return this.#store.all('SELECT * FROM task_current WHERE business=? ORDER BY task_ref', business).map(parse);
if (view === 'current') {
if (!taskRef(arg)) fail('invalid-request');
const r = this.#store.get('SELECT * FROM task_current WHERE business=? AND task_ref=?', business, arg);
return r ? parse(r) : null;
}
if (view === 'open')
return this.#store.all('SELECT task_ref,bucket FROM tasks_open WHERE business=? ORDER BY task_ref', business);
if (view === 'input')
return Boolean(
typeof arg === 'string' &&
this.#store.get("SELECT 1 FROM events WHERE business=? AND kind='human.input' AND id=?", business, arg),
);
fail('invalid-request');
}
#refused(s, e, request, other = 'storage-refused') {
const error = e instanceof BusError ? e : new BusError(other);
// No request content or raw exception text in refusal evidence.
try {
this.#store.transaction(() =>
this.#event(
s,
'action.refused',
{ code: error.code },
taskRef(request?.args?.task_ref)
? request.args.task_ref
: taskRef(request?.args?.target)
? request.args.target
: null,
),
);
} catch {
return new BusError('storage-unavailable');
}
return error;
}
request(cap, request) {
const s = this.#session(cap);
try {
@@ -336,25 +438,25 @@ export class Broker {
return result;
});
} catch (e) {
const error = e instanceof BusError ? e : new BusError('storage-refused');
// No request content or raw exception text in refusal evidence.
try {
this.#store.transaction(() =>
this.#event(
s,
'action.refused',
{ code: error.code },
taskRef(request?.args?.task_ref)
? request.args.task_ref
: taskRef(request?.args?.target)
? request.args.target
: null,
),
);
} catch {
throw new BusError('storage-unavailable');
}
throw error;
throw this.#refused(s, e, request);
}
}
// The eight task verbs reach the trusted S3 handler here. It runs outside any transaction, gets the
// cap to authorize and record through authorize() and recordTask(); refusals are recorded as in request().
async requestTask(cap, request, handler) {
const s = this.#session(cap);
try {
keys(request, ['verb', 'args'], ['verb']);
if (!TASK_VERBS.includes(request.verb)) fail('unknown-verb');
const args = request.args ?? {};
if (!object(args)) fail('invalid-request');
this.#secretCheck(args);
this.#store.transaction(() => this.#agent(s));
const result = await handler(cap, request.verb, args);
this.#secretCheck(result);
return result;
} catch (e) {
throw this.#refused(s, e, request, 'adapter-failed');
}
}
#dispatch(s, verb, a) {
+1 -1
View File
@@ -1,4 +1,4 @@
export { Broker, BusError, ACTIONS } from './broker.mjs';
export { Broker, BusError, ACTIONS, TASK_VERBS } from './broker.mjs';
export { startBroker } from './runtime.mjs';
export { Client } from './client.mjs';
export { views } from './views.mjs';
+7 -1
View File
@@ -39,7 +39,13 @@ if (!process.send) {
if (message?.op !== 'boot' || booted) throw new BusError('invalid-host-request');
booted = true;
clearTimeout(timer);
runtime = await startBroker(message.config);
// `trackers` is S3's plain-data boot config; the adapter is loaded only when a business has one.
const { trackers, ...config } = message.config ?? {};
if (trackers) {
const { tasksAdapter } = await import('../../tasks/src/adapter.mjs');
config.tasks = tasksAdapter({ trackers });
}
runtime = await startBroker(config);
if (closing) {
await runtime.close();
return;
+18 -2
View File
@@ -7,8 +7,16 @@ import { serve } from './server.mjs';
import { verifyHuman } from './human.mjs';
// Trusted host API. S1 supplies resolved definitions, S6 supplies launch records.
// Neither a business-file writer nor a socket-accessible configuration endpoint.
export async function startBroker({ dataRoot, businesses, launches = [], readers = [], repoRoots = [] }) {
let store, credentials, server;
// `tasks` is S3's adapter factory: ({broker, credentials, businesses}) => {handle, timeout, close}.
export async function startBroker({
dataRoot,
businesses,
launches = [],
readers = [],
repoRoots = [],
tasks = null,
}) {
let store, credentials, server, adapter;
try {
const references = {};
for (const [business, b] of Object.entries(businesses)) {
@@ -40,8 +48,14 @@ export async function startBroker({ dataRoot, businesses, launches = [], readers
}
const bound = launches.map(bindLaunch);
const readCaps = readers.map((business) => ({ business, cap: broker.bindReader({ business }) }));
if (tasks) {
adapter = await tasks({ broker, credentials, businesses });
if (typeof adapter?.handle !== 'function' || !Number.isSafeInteger(adapter.timeout) || adapter.timeout < 1)
throw new BusError('invalid-adapter');
}
const path = join(store.directory, 'broker.sock');
server = await serve({
tasks: adapter ? { handle: adapter.handle, timeout: adapter.timeout } : null,
broker,
path,
authenticateHuman: (proof) => {
@@ -63,12 +77,14 @@ export async function startBroker({ dataRoot, businesses, launches = [], readers
readers: readCaps,
async close() {
await server.close();
await adapter?.close?.();
store.close();
credentials.close();
},
};
} catch (e) {
await server?.close();
await adapter?.close?.();
store?.close();
credentials?.close();
throw e;
+19 -2
View File
@@ -1,9 +1,10 @@
import { createServer } from 'node:net';
import { lstatSync, chmodSync, unlinkSync } from 'node:fs';
import { BusError } from './broker.mjs';
import { BusError, TASK_VERBS } from './broker.mjs';
const LIMIT = 65536;
// One request per connection. There is deliberately no reconnect/retry or SQL verb.
export async function serve({ broker, path, authenticateHuman = null, timeout = 5000 }) {
// `tasks` ({handle, timeout}) is the trusted S3 adapter; only the eight task verbs reach it.
export async function serve({ broker, path, authenticateHuman = null, timeout = 5000, tasks = null }) {
try {
lstatSync(path);
throw new BusError('socket-exists');
@@ -54,6 +55,22 @@ export async function serve({ broker, path, authenticateHuman = null, timeout =
transient = cap = authenticateHuman(r.human);
}
if (typeof cap !== 'string') throw new BusError('unauthenticated');
if (tasks && TASK_VERBS.includes(r.verb)) {
const own = transient;
transient = null;
// A Vikunja write can outlast the idle timeout; a late reply is still outcome-unknown to the client.
socket.setTimeout(tasks.timeout);
broker
.requestTask(cap, { verb: r.verb, args: r.args ?? {} }, tasks.handle)
.then(
(result) => finish({ ok: true, result }),
(e) => finish({ ok: false, error: e instanceof BusError ? e.code : 'adapter-failed' }),
)
.finally(() => {
if (own) broker.disconnect(own);
});
return;
}
const result = broker.request(cap, { verb: r.verb, args: r.args ?? {} });
finish({ ok: true, result });
} catch (e) {
+305
View File
@@ -0,0 +1,305 @@
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);
});