Files
stack/packages/bus/src/broker.mjs
T
jason.woltjeandClaude Opus 5.5 d27042faf2 feat(bus): a within-role message cites a decision without consuming it (row 44, S2c, rocko)
A within-role message.send stores its decision as a citation: the broker
checks it exists in the business, and doesn't class-match it, consume it
or put it in action.allowed (lead decision 65). Sends that aren't
within-role keep the S2b authority rules. The README wording on
within-role sends is fixed.

Candidate agents/rocko/work/slice1-s2c-r1, build.patch 2c4f8d9f,
manifest aa249ad9. Darkwing approved round 1 on #1526 (comment 26778).
Integration gate in a worktree on ac4a6499: every package green on
Node 26; bus 58/58 on Node 24; every scripts/test-*.sh green. Node 24
failures in conversation, ledger, queue, seat and webui match the base.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-10-05 18:13:46 -05:00

825 lines
29 KiB
JavaScript

import { randomUUID, randomBytes } from 'node:crypto';
// Kept closed while S1 lands; its vocabulary is the integration authority.
export const ACTIONS = Object.freeze([
'task.create',
'task.assign',
'task.schedule',
'task.update.assigned',
'task.close',
'task.reassign',
'task.scope.change',
'task.priority.change',
'git.push.working',
'git.push.protected',
'git.merge.protected',
'review.request',
'review.verdict',
'message.send',
'message.external',
'role.launch',
'role.revoke',
'credential.mint',
'spend',
'deploy',
'policy.change',
'prd.approve',
'decision.resolve.technical',
]);
const GATED = new Set([
'credential.mint',
'git.merge.protected',
'git.push.protected',
'deploy',
'spend',
'message.external',
'policy.change',
'prd.approve',
'role.revoke',
]);
const CLOSED = "('resolved','withdrawn','expired')";
export class BusError extends Error {
constructor(code) {
super(code);
this.code = code;
}
}
const fail = (code) => {
throw new BusError(code);
};
const object = (x) => x !== null && typeof x === 'object' && !Array.isArray(x);
const text = (x, max = 4096) => typeof x === 'string' && x.length > 0 && x.length <= max && !x.includes('\0');
const id = (x) => text(x, 160) && /^[a-zA-Z0-9][a-zA-Z0-9_.:/-]*$/.test(x);
function keys(x, allowed, required = []) {
if (!object(x) || Object.keys(x).some((k) => !allowed.includes(k)) || required.some((k) => !(k in x)))
fail('invalid-request');
}
function string(x, max) {
if (!text(x, max)) fail('invalid-request');
return x;
}
function identifier(x) {
if (!id(x)) fail('invalid-request');
return x;
}
const taskRef = (x) => typeof x === 'string' && /^vikunja:[1-9][0-9]*\/[1-9][0-9]*$/.test(x);
export class Broker {
#store;
#businesses;
#sessions = new Map();
#runs = new Set();
#secretCheck;
#clock = 0;
constructor({ store, businesses, secretCheck = () => {} }) {
if (!object(businesses) || !Object.keys(businesses).length) fail('invalid-business');
this.#businesses = structuredClone(businesses);
this.#store = store;
this.#secretCheck = secretCheck;
for (const table of [
'events',
'role_claims',
'decisions',
'decision_events',
'messages',
'deliveries',
'task_snapshots',
]) {
const at = store.get(`SELECT max(at) at FROM ${table}`).at;
if (at) {
const ms = Date.parse(at);
if (!Number.isFinite(ms)) fail('invalid-timestamp');
this.#clock = Math.max(this.#clock, ms);
}
}
for (const [name, b] of Object.entries(this.#businesses)) {
if (!id(name) || b.id !== name || !id(b.human) || !object(b.roles) || !object(b.arbiters))
fail('invalid-business');
for (const role of Object.keys(b.roles)) {
if (!id(role) || role === 'human') fail('invalid-business');
const a = b.roles[role].authority;
if (!object(a) || !Array.isArray(a.withinRole) || !Array.isArray(a.crossRole))
fail('invalid-authority');
const all = [...a.withinRole, ...a.crossRole];
if (new Set(all).size !== all.length || all.some((v) => !ACTIONS.includes(v) || GATED.has(v)))
fail('invalid-authority');
}
if (!Object.hasOwn(b.roles, b.arbiters.technical) || !Object.hasOwn(b.roles, b.arbiters.delivery))
fail('invalid-arbiter');
}
}
#now() {
this.#clock = Math.max(Date.now(), this.#clock + 1);
return new Date(this.#clock).toISOString();
}
#business(name) {
if (!Object.hasOwn(this.#businesses, name)) fail('unknown-business');
return this.#businesses[name];
}
#session(cap) {
return this.#sessions.get(cap) ?? fail('unauthenticated');
}
#holder(business, role) {
return this.#store.get(
'SELECT * FROM role_claims WHERE business=? AND role=? ORDER BY seq DESC LIMIT 1',
business,
role,
);
}
#agent(s) {
if (s.reader) fail('read-only');
if (s.human) fail('agent-required');
const h = this.#holder(s.business, s.role);
if (h?.op !== 'claim' || h.holder_run !== s.run) fail('not-holder');
return h;
}
#human(s) {
if (s.reader) fail('read-only');
if (!s.human || s.via !== 'cli' || s.outsideAgent !== true) fail('human-required');
}
#cap(s) {
const cap = randomBytes(32).toString('hex');
this.#sessions.set(cap, Object.freeze(s));
return cap;
}
// Trusted launcher API, deliberately NOT a socket verb. S6 supplies immutable launch records.
bindLaunch(record) {
keys(
record,
['business', 'role', 'run', 'harness', 'address', 'pid', 'startTime'],
['business', 'role', 'run', 'harness'],
);
const b = this.#business(identifier(record.business));
if (!Object.hasOwn(b.roles, record.role)) fail('unknown-role');
identifier(record.run);
if (!['pi', 'claude-code'].includes(record.harness)) fail('invalid-harness');
if (record.address !== undefined) string(record.address, 1024);
const key = record.business + ':' + record.run;
if (this.#runs.has(key)) fail('duplicate-run');
this.#runs.add(key);
const session = { ...record, address: record.address ?? null, human: false };
this.#secretCheck(record);
const old = this.#store.get(
"SELECT actor_role,body FROM events WHERE business=? AND kind='session.launched' AND subject=? ORDER BY seq LIMIT 1",
record.business,
record.run,
);
const body = {
harness: record.harness,
address: record.address ?? null,
pid: record.pid ?? null,
startTime: record.startTime ?? null,
};
if (old) {
if (old.actor_role !== record.role || old.body !== JSON.stringify(body)) fail('launch-record-mismatch');
} else this.#store.transaction(() => this.#event(session, 'session.launched', body, record.run));
return this.#cap(session);
}
// Trusted human transport invokes only AFTER checking its CLI process and launch ancestry.
bindHuman(record) {
keys(record, ['business', 'human', 'via', 'outsideAgent'], ['business', 'human', 'via', 'outsideAgent']);
const b = this.#business(record.business);
if (record.human !== b.human || record.via !== 'cli' || record.outsideAgent !== true)
fail('human-required');
return this.#cap({ ...record, role: null, run: null, human: b.human });
}
bindReader({ business }) {
this.#business(business);
return this.#cap({ business, reader: true, role: null, run: null, human: false });
}
disconnect(cap) {
this.#sessions.delete(cap);
}
identity(cap) {
const s = this.#session(cap);
return { business: s.business, role: s.role, run: s.run, human: s.human || null };
}
#event(s, kind, body, subject = null) {
this.#secretCheck({ kind, body, subject });
if (
(kind.startsWith('action.') || kind.startsWith('review.')) &&
taskRef(body.task_ref ?? body.target) &&
subject !== (body.task_ref ?? body.target)
)
fail('task-subject-required');
const eid = randomUUID();
this.#store.run(
'INSERT INTO events(id,at,business,kind,actor_role,actor_run,subject,body) VALUES(?,?,?,?,?,?,?,?)',
eid,
this.#now(),
s.business,
kind,
s.role,
s.run,
subject,
JSON.stringify(body),
);
return eid;
}
#classification(s, action) {
if (!ACTIONS.includes(action)) fail('unknown-action');
if (GATED.has(action) || (action === 'role.launch' && this.#business(s.business).launch?.by !== s.role))
return 'gated';
const a = this.#business(s.business).roles[s.role]?.authority;
return a?.withinRole.includes(action)
? 'within-role'
: a?.crossRole.includes(action)
? 'cross-role'
: 'gated';
}
#decision(s, id) {
const d = this.#store.get('SELECT * FROM decisions WHERE id=? AND business=?', id, s.business);
if (!d) fail('decision-not-found');
return d;
}
#decisionView(d) {
const row = this.#store.get(
"SELECT body FROM events WHERE business=? AND kind='action.allowed' AND json_extract(body,'$.operation')='decision.raise' AND json_extract(body,'$.decision')=? ORDER BY seq LIMIT 1",
d.business,
d.id,
);
const context = row ? JSON.parse(row.body) : null;
return {
...d,
options: JSON.parse(d.options),
blocking: !!d.blocking,
authorization: context
? { action: d.action, target: context.target, approvalChoice: context.approvalChoice }
: null,
};
}
#closed(id) {
return this.#store.get(
`SELECT * FROM decision_events WHERE decision=? AND op IN ${CLOSED} ORDER BY seq DESC LIMIT 1`,
id,
);
}
#checkAuthority(s, action, { decision = null, target = null } = {}) {
this.#agent(s);
if (
action === 'role.launch' &&
this.#store.get('SELECT state FROM launch_state WHERE business=?', s.business)?.state === 'revoked'
)
fail('launch-revoked');
const cls = this.#classification(s, action);
if (cls === 'within-role' && !decision) return { class: cls };
if (!decision) fail('decision-required');
const d = this.#decision(s, decision),
r = this.#closed(d.id);
const evidence = this.#store.get(
"SELECT body FROM events WHERE business=? AND json_extract(body,'$.decision')=? AND json_extract(body,'$.operation')='decision.raise' AND kind='action.allowed' ORDER BY seq LIMIT 1",
s.business,
d.id,
);
const context = evidence ? JSON.parse(evidence.body) : null;
if (
d.class !== cls ||
d.action !== action ||
d.raised_by_role !== s.role ||
d.raised_by_run !== s.run ||
!context ||
context.target !== target
)
fail('decision-mismatch');
if (r?.op !== 'resolved' || r.choice !== context.approvalChoice) fail('decision-not-approved');
if (action === 'role.revoke' && context.holderRun !== this.#holder(s.business, target)?.holder_run)
fail('decision-mismatch');
return { class: cls, decision: d.id };
}
// Caller owns the transaction: authority, consumption and effect commit together.
#consumeAuthority(s, action, context = {}) {
const result = this.#checkAuthority(s, action, context);
// decision.raise records context, not permission consumed by an executor.
if (
result.decision &&
this.#store.get(
"SELECT 1 FROM events WHERE business=? AND kind='action.allowed' AND json_extract(body,'$.decision')=? AND json_extract(body,'$.operation') IS NOT 'decision.raise' LIMIT 1",
s.business,
result.decision,
)
)
fail('decision-consumed');
this.#event(
s,
'action.allowed',
{ action, ...result, target: context.target ?? null },
context.target ?? null,
);
return result;
}
// Trusted S3 handler calls inside its own operation; not a generic socket action executor.
authorize(cap, action, context = {}) {
const s = this.#session(cap);
return this.#store.transaction(() => this.#consumeAuthority(s, action, context));
}
// Trusted adapters only. No agent socket route reaches this method.
recordEvent(cap, { kind, body, subject = null }) {
const s = this.#session(cap);
return this.#store.transaction(() => {
this.#agent(s);
if (['human.input', 'launch.revoked', 'launch.restored'].includes(kind)) fail('reserved-event');
return this.#event(s, kind, body, subject);
});
}
request(cap, request) {
const s = this.#session(cap);
try {
keys(request, ['verb', 'args'], ['verb']);
string(request.verb, 64);
const args = request.args ?? {};
if (!object(args)) fail('invalid-request');
this.#secretCheck(args);
return this.#store.transaction(() => {
if (s.reader && !['inbox', 'agents', 'tasks', 'trail'].includes(request.verb)) fail('read-only');
const result = this.#dispatch(s, request.verb, args);
this.#secretCheck(result);
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;
}
}
#dispatch(s, verb, a) {
const b = this.#business(s.business),
by = s.human || s.run;
switch (verb) {
case 'role.claim': {
keys(a, []);
if (s.human) fail('agent-required');
if (
this.#store.get(
"SELECT 1 FROM role_claims WHERE business=? AND role=? AND holder_run=? AND op='revoke' LIMIT 1",
s.business,
s.role,
s.run,
)
)
fail('run-revoked');
if (this.#holder(s.business, s.role)?.op === 'claim') fail('role-already-held');
this.#store.run(
'INSERT INTO role_claims(at,business,role,op,holder_run,harness,address,by) VALUES(?,?,?,?,?,?,?,?)',
this.#now(),
s.business,
s.role,
'claim',
s.run,
s.harness,
s.address,
by,
);
this.#routeWaiting(s.business, s.role);
return { role: s.role, run: s.run };
}
case 'role.release': {
keys(a, []);
const h = this.#agent(s);
this.#claimEnd(s, h, 'release', null);
return { released: true };
}
case 'role.revoke': {
keys(a, ['role', 'decision'], ['role', 'decision']);
identifier(a.role);
identifier(a.decision);
this.#consumeAuthority(s, 'role.revoke', { decision: a.decision, target: a.role });
const h = this.#holder(s.business, a.role);
if (h?.op !== 'claim') fail('not-held');
this.#claimEnd(s, h, 'revoke', a.decision);
return { revoked: true };
}
case 'decision.raise': {
keys(
a,
[
'action',
'domain',
'target',
'project',
'question',
'options',
'recommendation',
'blocking',
'task_ref',
'requirement_ref',
'supersedes',
'choice',
'approvalChoice',
],
['action', 'question', 'options', 'recommendation', 'blocking'],
);
this.#agent(s);
const cls = this.#classification(s, a.action);
string(a.question);
if (typeof a.blocking !== 'boolean') fail('invalid-request');
if (a.blocking && !a.task_ref) fail('task-required');
if (a.task_ref !== undefined && !/^vikunja:[1-9][0-9]*\/[1-9][0-9]*$/.test(a.task_ref))
fail('invalid-task-ref');
for (const k of ['target', 'project', 'requirement_ref', 'supersedes'])
if (a[k] !== undefined) identifier(a[k]);
if (!Array.isArray(a.options) || a.options.length < 2 || a.options.length > 9)
fail('invalid-options');
for (const o of a.options) {
keys(o, ['key', 'text'], ['key', 'text']);
if (!id(o.key) || !text(o.text, 1024)) fail('invalid-options');
}
const choices = a.options.map((o) => o.key);
if (new Set(choices).size !== choices.length || !choices.includes(a.recommendation))
fail('invalid-options');
const approvalChoice = a.approvalChoice ?? 'yes';
if (!choices.includes(approvalChoice)) fail('invalid-approval-choice');
const domain = a.domain ?? 'delivery';
if (!['technical', 'delivery'].includes(domain)) fail('invalid-domain');
const proposedRoute = cls === 'gated' ? 'human' : cls === 'cross-role' ? b.arbiters[domain] : s.role;
const route = cls === 'cross-role' && proposedRoute === s.role ? 'human' : proposedRoute;
if (a.supersedes) {
const old = this.#decision(s, a.supersedes);
if (old.raised_by_run !== s.run || this.#closed(old.id)) fail('supersede-refused');
this.#store.run(
'INSERT INTO decision_events(decision,at,op,by) VALUES(?,?,?,?)',
old.id,
this.#now(),
'withdrawn',
by,
);
}
const did = randomUUID();
this.#store.run(
'INSERT INTO decisions(id,at,business,project,raised_by_role,raised_by_run,class,action,route_to,question,options,recommendation,task_ref,requirement_ref,blocking,supersedes) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)',
did,
this.#now(),
s.business,
a.project ?? null,
s.role,
s.run,
cls,
a.action,
route,
a.question,
JSON.stringify(a.options),
a.recommendation,
a.task_ref ?? null,
a.requirement_ref ?? null,
Number(a.blocking),
a.supersedes ?? null,
);
this.#event(
s,
'action.allowed',
{
decision: did,
action: a.action,
operation: 'decision.raise',
task_ref: a.task_ref ?? null,
target: a.target ?? null,
approvalChoice,
holderRun:
a.action === 'role.revoke' ? (this.#holder(s.business, a.target)?.holder_run ?? null) : null,
},
a.task_ref ?? (taskRef(a.target) ? a.target : did),
);
if (cls === 'within-role') {
if (!choices.includes(a.choice)) fail('invalid-choice');
this.#store.run(
'INSERT INTO decision_events(decision,at,op,by,choice,via) VALUES(?,?,?,?,?,?)',
did,
this.#now(),
'resolved',
by,
a.choice,
'broker',
);
}
return this.#decisionView(this.#decision(s, did));
}
case 'decision.resolve':
case 'decision.seen':
case 'decision.withdraw': {
keys(a, ['id', 'choice', 'note'], ['id']);
identifier(a.id);
if (a.note !== undefined) string(a.note);
const d = this.#decision(s, a.id);
if (this.#closed(d.id)) fail('decision-closed');
if (verb === 'decision.withdraw') {
if (s.human) {
this.#human(s);
} else {
this.#agent(s);
if (d.raised_by_run !== s.run) fail('not-resolver');
}
} else if (d.route_to === 'human') this.#human(s);
else {
this.#agent(s);
if (s.role !== d.route_to) fail('not-resolver');
}
const op = verb === 'decision.resolve' ? 'resolved' : verb === 'decision.seen' ? 'seen' : 'withdrawn';
if (op === 'resolved' && !JSON.parse(d.options).some((o) => o.key === a.choice))
fail('invalid-choice');
this.#store.run(
'INSERT INTO decision_events(decision,at,op,by,choice,note,via) VALUES(?,?,?,?,?,?,?)',
d.id,
this.#now(),
op,
by,
op === 'resolved' ? a.choice : null,
a.note ?? null,
s.human ? 'cli' : 'broker',
);
if (s.human)
this.#event(
s,
'human.input',
{ kind: 'answer', decision: d.id, op, choice: a.choice ?? null },
d.id,
);
return { id: d.id, op };
}
case 'message.send': {
keys(a, ['to', 'body', 'class', 'in_reply_to', 'corrects', 'decision'], ['to', 'body']);
if (!s.human) {
const withinRole = this.#classification(s, 'message.send') === 'within-role';
this.#consumeAuthority(s, 'message.send', {
decision: withinRole ? null : (a.decision ?? null),
target: a.to,
});
} else this.#human(s);
if (a.to !== 'human' && !Object.hasOwn(b.roles, a.to)) fail('unknown-role');
string(a.body, 32768);
if (
a.class !== undefined &&
![
'REQUEST',
'ASSIGNMENT',
'REVIEW-REQUEST',
'REVIEW-RESULT',
'RESULT',
'INFO',
'DECISION',
'REACTION',
].includes(a.class)
)
fail('invalid-class');
for (const k of ['in_reply_to', 'corrects']) {
if (a[k] !== undefined) {
identifier(a[k]);
if (!this.#store.get('SELECT id FROM messages WHERE id=? AND business=?', a[k], s.business))
fail('message-not-found');
}
}
if (a.decision) this.#decision(s, a.decision);
const mid = randomUUID();
this.#store.run(
'INSERT INTO messages(id,at,business,from_role,from_run,to_role,class,in_reply_to,decision,corrects,body) VALUES(?,?,?,?,?,?,?,?,?,?,?)',
mid,
this.#now(),
s.business,
s.human ? 'human' : s.role,
s.human ? 'human-cli' : s.run,
a.to,
a.class ?? 'INFO',
a.in_reply_to ?? null,
a.decision ?? null,
a.corrects ?? null,
a.body,
);
const h = this.#holder(s.business, a.to);
this.#delivery(mid, 'routed', h?.op === 'claim' ? h : null);
const request = s.human
? this.#event(s, 'human.input', { kind: 'instruction', message: mid }, mid)
: null;
return { id: mid, request };
}
case 'message.receive': {
keys(a, []);
if (s.human) this.#human(s);
else this.#agent(s);
const role = s.human ? 'human' : s.role;
const rows = this.#store.all(
"SELECT * FROM messages m WHERE business=? AND to_role=? AND NOT EXISTS(SELECT 1 FROM deliveries d WHERE d.message=m.id AND d.op IN ('delivered','read')) ORDER BY seq LIMIT 100",
s.business,
role,
);
for (const m of rows)
this.#delivery(
m.id,
'delivered',
s.human ? null : this.#holder(s.business, s.role),
s.human ? 'cli' : 'broker',
);
return rows.map((m) => ({
...m,
request:
this.#store.get(
"SELECT id FROM events WHERE business=? AND kind='human.input' AND subject=? ORDER BY seq LIMIT 1",
s.business,
m.id,
)?.id ?? null,
}));
}
case 'message.read': {
keys(a, ['id'], ['id']);
if (s.human) this.#human(s);
else this.#agent(s);
const m = this.#store.get('SELECT * FROM messages WHERE business=? AND id=?', s.business, a.id);
if (!m || m.to_role !== (s.human ? 'human' : s.role)) fail('message-not-found');
const delivery = this.#store.get(
"SELECT * FROM deliveries WHERE message=? AND op='delivered' ORDER BY seq DESC LIMIT 1",
m.id,
);
if (!delivery || delivery.holder_run !== (s.run ?? null)) fail('not-recipient');
if (!this.#store.get("SELECT seq FROM deliveries WHERE message=? AND op='read'", m.id))
this.#delivery(m.id, 'read', s.human ? null : this.#holder(s.business, s.role));
return { read: true };
}
case 'launch.revoke':
case 'launch.restore': {
keys(a, []);
this.#human(s);
const kind = verb === 'launch.revoke' ? 'launch.revoked' : 'launch.restored';
const eid = this.#event(s, kind, { via: 'cli' });
this.#event(s, 'human.input', { kind: 'admin', event: eid });
return { id: eid };
}
case 'event.emit': {
keys(a, ['kind', 'body', 'subject'], ['kind', 'body']);
this.#agent(s);
// Authoritative lifecycle/task/credential events are emitted by trusted handlers, not arbitrary clients.
if (a.kind !== 'action.allowed') fail('reserved-event');
keys(a.body, ['action', 'target'], ['action']);
if (a.body.action !== 'routine') fail('reserved-event');
return {
id: this.#event(
s,
'action.allowed',
{ action: 'routine', target: a.body.target ?? null },
a.subject ?? null,
),
};
}
case 'inbox':
case 'agents':
case 'tasks':
case 'trail': {
keys(a, verb === 'trail' ? ['subject'] : [], verb === 'trail' ? ['subject'] : []);
if (!s.human && !s.reader) this.#agent(s);
if (verb === 'inbox')
return this.#store
.all(
`SELECT d.* FROM decisions d WHERE business=? AND route_to=? AND NOT EXISTS(SELECT 1 FROM decision_events x WHERE x.decision=d.id AND x.op IN ${CLOSED}) ORDER BY CASE class WHEN 'gated' THEN 0 ELSE 1 END,seq`,
s.business,
s.human || s.reader ? 'human' : s.role,
)
.map((d) => this.#decisionView(d));
if (verb === 'agents')
return this.#store.all(
"SELECT * FROM role_claims c WHERE business=? AND op='claim' AND seq=(SELECT max(seq) FROM role_claims WHERE business=c.business AND role=c.role) ORDER BY role",
s.business,
);
if (verb === 'tasks')
return this.#store
.all('SELECT * FROM task_current WHERE business=? ORDER BY task_ref', s.business)
.map((r) => ({ ...r, fields: JSON.parse(r.fields) }));
return this.#trail(s, identifier(a.subject));
}
default:
fail('unknown-verb');
}
}
#claimEnd(s, h, op, decision) {
this.#store.run(
'INSERT INTO role_claims(at,business,role,op,holder_run,harness,address,by,decision) VALUES(?,?,?,?,?,?,?,?,?)',
this.#now(),
s.business,
h.role,
op,
h.holder_run,
h.harness,
h.address,
s.run ?? s.human,
decision,
);
}
#delivery(mid, op, h, transport = null) {
this.#store.run(
'INSERT INTO deliveries(message,at,op,holder_run,transport,address) VALUES(?,?,?,?,?,?)',
mid,
this.#now(),
op,
h?.holder_run ?? null,
transport ?? h?.harness ?? null,
h?.address ?? null,
);
}
#routeWaiting(business, role) {
const h = this.#holder(business, role);
const rows = this.#store.all(
"SELECT id FROM messages m WHERE business=? AND to_role=? AND NOT EXISTS(SELECT 1 FROM deliveries WHERE message=m.id AND op IN ('delivered','read'))",
business,
role,
);
for (const m of rows) this.#delivery(m.id, 'routed', h);
}
#trail(s, subject) {
const rows = new Map();
const add = (table, list) => {
for (const row of list) {
const r = { ...row, table };
for (const k of ['body', 'options', 'fields'])
if (k in r && table !== 'messages')
try {
r[k] = JSON.parse(r[k]);
} catch {}
rows.set(table + ':' + row.seq, r);
}
};
add(
'events',
this.#store.all('SELECT * FROM events WHERE business=? AND subject=?', s.business, subject),
);
add('messages', this.#store.all('SELECT * FROM messages WHERE business=? AND id=?', s.business, subject));
const decisions = this.#store.all(
'SELECT * FROM decisions WHERE business=? AND (id=? OR task_ref=?)',
s.business,
subject,
subject,
);
add('decisions', decisions);
for (const d of decisions) {
add('decision_events', this.#store.all('SELECT * FROM decision_events WHERE decision=?', d.id));
add(
'events',
this.#store.all(
"SELECT * FROM events WHERE business=? AND (subject=? OR json_extract(body,'$.decision')=?)",
s.business,
d.id,
d.id,
),
);
add(
'messages',
this.#store.all('SELECT * FROM messages WHERE business=? AND decision=?', s.business, d.id),
);
}
for (const e of [...rows.values()])
if (e.table === 'events' && e.kind === 'task.created') {
const input = this.#store.get(
"SELECT * FROM events WHERE business=? AND kind='human.input' AND id=?",
s.business,
e.body.request,
);
if (input) {
add('events', [input]);
const body = JSON.parse(input.body);
if (body.message)
add(
'messages',
this.#store.all('SELECT * FROM messages WHERE business=? AND id=?', s.business, body.message),
);
}
}
for (const m of [...rows.values()])
if (m.table === 'messages')
add('deliveries', this.#store.all('SELECT * FROM deliveries WHERE message=?', m.id));
add(
'task_snapshots',
this.#store.all('SELECT * FROM task_snapshots WHERE business=? AND task_ref=?', s.business, subject),
);
const runs = new Set(
[...rows.values()].map((r) => r.actor_run ?? r.raised_by_run ?? r.run).filter(Boolean),
);
for (const run of runs) {
add(
'events',
this.#store.all(
"SELECT * FROM events WHERE business=? AND actor_run=? AND kind IN ('session.launched','session.ended')",
s.business,
run,
),
);
add(
'role_claims',
this.#store.all('SELECT * FROM role_claims WHERE business=? AND holder_run=?', s.business, run),
);
}
return [...rows.values()].sort(
(a, b) => a.at.localeCompare(b.at) || a.table.localeCompare(b.table) || a.seq - b.seq,
);
}
}