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]>
This commit is contained in:
@@ -0,0 +1,806 @@
|
||||
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') 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.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 };
|
||||
}
|
||||
// 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(() => {
|
||||
const result = this.#checkAuthority(s, action, context);
|
||||
this.#event(
|
||||
s,
|
||||
'action.allowed',
|
||||
{ action, ...result, target: context.target ?? null },
|
||||
context.target ?? null,
|
||||
);
|
||||
return result;
|
||||
});
|
||||
}
|
||||
// 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.#checkAuthority(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) this.#checkAuthority(s, 'message.send', { decision: 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,
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
import { BusError } from './broker.mjs';
|
||||
// Adapt S1's validated loadBusiness + resolveInstance outputs, without loading or writing files.
|
||||
export function busBusiness(business, resolved) {
|
||||
if (!business || typeof business.id !== 'string' || !business.roles || !resolved)
|
||||
throw new BusError('invalid-business');
|
||||
const out = structuredClone(business);
|
||||
for (const role of Object.keys(out.roles)) {
|
||||
const r = resolved[role];
|
||||
if (
|
||||
!r ||
|
||||
r.business !== business.id ||
|
||||
r.instance !== role ||
|
||||
r.definition !== out.roles[role].definition ||
|
||||
!r.limits?.authority
|
||||
)
|
||||
throw new BusError('invalid-resolution');
|
||||
out.roles[role].authority = structuredClone(r.limits.authority);
|
||||
out.roles[role].credentials = structuredClone(r.credentials ?? {});
|
||||
}
|
||||
return out;
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
import { connect } from 'node:net';
|
||||
import { BusError } from './broker.mjs';
|
||||
export class Client {
|
||||
constructor({ path, cap, human, timeout = 5000 }) {
|
||||
this.path = path;
|
||||
this.cap = cap;
|
||||
this.human = human;
|
||||
this.timeout = timeout;
|
||||
}
|
||||
call(verb, args = {}) {
|
||||
const request =
|
||||
JSON.stringify({ ...(this.human ? { human: this.human } : { cap: this.cap }), verb, args }) + '\n';
|
||||
if (Buffer.byteLength(request) > 65536) return Promise.reject(new BusError('request-too-large'));
|
||||
return new Promise((resolve, reject) => {
|
||||
const socket = connect(this.path);
|
||||
let input = Buffer.alloc(0),
|
||||
settled = false;
|
||||
const end = (error, value) => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
socket.destroy();
|
||||
error ? reject(error) : resolve(value);
|
||||
};
|
||||
socket.setTimeout(this.timeout, () => end(new BusError('outcome-unknown')));
|
||||
socket.on('connect', () => socket.write(request));
|
||||
socket.on('error', () => end(new BusError('outcome-unknown')));
|
||||
socket.on('end', () => end(new BusError('outcome-unknown')));
|
||||
socket.on('data', (b) => {
|
||||
input = Buffer.concat([input, b]);
|
||||
if (input.length > 4 * 1024 * 1024) return end(new BusError('response-too-large'));
|
||||
if (!input.includes(10)) return;
|
||||
try {
|
||||
const r = JSON.parse(input.toString('utf8'));
|
||||
if (r.ok === true) end(null, r.result);
|
||||
else if (r.ok === false && typeof r.error === 'string' && /^[a-z-]{1,64}$/.test(r.error))
|
||||
end(new BusError(r.error));
|
||||
else end(new BusError('invalid-response'));
|
||||
} catch {
|
||||
end(new BusError('invalid-response'));
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,143 @@
|
||||
import { openSync, closeSync, fstatSync, readFileSync, lstatSync, realpathSync, constants } from 'node:fs';
|
||||
import { isAbsolute, relative, resolve } from 'node:path';
|
||||
import { BusError } from './broker.mjs';
|
||||
const deny = (code) => {
|
||||
throw new BusError(code);
|
||||
};
|
||||
const inside = (root, path) => {
|
||||
const r = relative(root, path);
|
||||
return r === '' || (!r.startsWith('..') && !isAbsolute(r));
|
||||
};
|
||||
const day = (value) =>
|
||||
typeof value === 'string' &&
|
||||
/^\d{4}-\d{2}-\d{2}$/.test(value) &&
|
||||
!Number.isNaN(Date.parse(value)) &&
|
||||
new Date(value).toISOString().slice(0, 10) === value;
|
||||
export class Credentials {
|
||||
#tokens = new Map();
|
||||
constructor({ references = {}, repoRoots = [], dataRoot, env = process.env }) {
|
||||
try {
|
||||
for (const [instance, services] of Object.entries(references))
|
||||
for (const [service, ref] of Object.entries(services)) {
|
||||
if (!['gitea', 'vikunja'].includes(service) || !ref || typeof ref !== 'object')
|
||||
deny('credential-reference');
|
||||
const date = service === 'vikunja' ? ref.expires : ref.rotateBy;
|
||||
if (!day(date)) deny('credential-date');
|
||||
if (
|
||||
(ref.service !== undefined && ref.service !== service) ||
|
||||
Object.keys(ref).some(
|
||||
(k) => !['service', 'file', 'env', service === 'vikunja' ? 'expires' : 'rotateBy'].includes(k),
|
||||
) ||
|
||||
Boolean(ref.file) === Boolean(ref.env)
|
||||
)
|
||||
deny('credential-reference');
|
||||
let token;
|
||||
if (ref.file) {
|
||||
if (!isAbsolute(ref.file)) deny('credential-location');
|
||||
let resolved, original;
|
||||
try {
|
||||
original = lstatSync(ref.file);
|
||||
resolved = realpathSync(ref.file);
|
||||
} catch {
|
||||
deny('credential-file');
|
||||
}
|
||||
if (original.isSymbolicLink() || !original.isFile()) deny('credential-file');
|
||||
if (
|
||||
[...repoRoots, dataRoot].filter(Boolean).some((root) =>
|
||||
inside(
|
||||
(() => {
|
||||
try {
|
||||
return realpathSync(root);
|
||||
} catch {
|
||||
return resolve(root);
|
||||
}
|
||||
})(),
|
||||
resolved,
|
||||
),
|
||||
)
|
||||
)
|
||||
deny('credential-location');
|
||||
let fd;
|
||||
try {
|
||||
fd = openSync(ref.file, constants.O_RDONLY | constants.O_NOFOLLOW);
|
||||
const s = fstatSync(fd);
|
||||
if (
|
||||
!s.isFile() ||
|
||||
s.uid !== process.getuid() ||
|
||||
(s.mode & 0o777) !== 0o600 ||
|
||||
s.size > 16384 ||
|
||||
s.ino !== original.ino ||
|
||||
s.dev !== original.dev
|
||||
)
|
||||
deny('credential-file');
|
||||
token = readFileSync(fd, 'utf8');
|
||||
} catch (e) {
|
||||
if (e instanceof BusError) throw e;
|
||||
deny('credential-file');
|
||||
} finally {
|
||||
if (fd !== undefined) closeSync(fd);
|
||||
}
|
||||
} else {
|
||||
if (!/^[A-Z_][A-Z0-9_]*$/.test(ref.env)) deny('credential-reference');
|
||||
token = env[ref.env];
|
||||
}
|
||||
if (typeof token !== 'string' || token.length > 16384 || !/^[-A-Za-z0-9._~]+\n?$/.test(token))
|
||||
deny('credential-format');
|
||||
token = token.replace(/\n$/, '');
|
||||
if (token.length < 16) deny('credential-format');
|
||||
this.#tokens.set(instance + ':' + service, { token, date });
|
||||
}
|
||||
} catch (e) {
|
||||
this.close();
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
assertClean(value) {
|
||||
let encoded;
|
||||
try {
|
||||
encoded = JSON.stringify(value);
|
||||
} catch {
|
||||
deny('invalid-request');
|
||||
}
|
||||
for (const { token } of this.#tokens.values()) if (encoded?.includes(token)) deny('credential-leak');
|
||||
}
|
||||
async use(instance, service, operation) {
|
||||
const entry = this.#tokens.get(instance + ':' + service);
|
||||
if (!entry) deny('credential-unavailable');
|
||||
if (service === 'vikunja' && Date.now() >= Date.parse(entry.date + 'T00:00:00Z'))
|
||||
deny('credential-expired');
|
||||
let result;
|
||||
try {
|
||||
result = await operation(entry.token);
|
||||
} catch {
|
||||
deny('service-failed');
|
||||
}
|
||||
this.assertClean(result);
|
||||
return result;
|
||||
}
|
||||
status() {
|
||||
return [...this.#tokens.entries()].map(([key, { date }]) => {
|
||||
const i = key.lastIndexOf(':'),
|
||||
service = key.slice(i + 1),
|
||||
left = Date.parse(date + 'T00:00:00Z') - Date.now();
|
||||
return {
|
||||
instance: key.slice(0, i),
|
||||
service,
|
||||
date,
|
||||
state:
|
||||
service === 'gitea'
|
||||
? left <= 0
|
||||
? 'rotation-due'
|
||||
: 'valid'
|
||||
: left <= 0
|
||||
? 'expired'
|
||||
: left < 7 * 86400000
|
||||
? 'expiring'
|
||||
: 'valid',
|
||||
};
|
||||
});
|
||||
}
|
||||
close() {
|
||||
this.#tokens.clear();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
// Transport shim used by S4, not a second product CLI. Input and proof never carry service tokens.
|
||||
import { randomBytes } from 'node:crypto';
|
||||
import { spawnSync } from 'node:child_process';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { readFileSync } from 'node:fs';
|
||||
import { Client } from './client.mjs';
|
||||
import { readProcess } from './human.mjs';
|
||||
const file = fileURLToPath(import.meta.url);
|
||||
if (!process.env.MOSAIC_BUS_CLI_NONCE) {
|
||||
const child = spawnSync(process.execPath, [file, ...process.argv.slice(2)], {
|
||||
stdio: 'inherit',
|
||||
env: { ...process.env, MOSAIC_BUS_CLI_NONCE: randomBytes(32).toString('hex') },
|
||||
});
|
||||
process.exit(child.status ?? 2);
|
||||
}
|
||||
try {
|
||||
const request = JSON.parse(readFileSync(0, 'utf8'));
|
||||
if (
|
||||
!request ||
|
||||
typeof request !== 'object' ||
|
||||
Object.keys(request).some((k) => !['business', 'verb', 'args'].includes(k))
|
||||
)
|
||||
throw Error();
|
||||
const me = readProcess(process.pid);
|
||||
const client = new Client({
|
||||
path: process.argv[2],
|
||||
human: {
|
||||
business: request.business,
|
||||
pid: me.pid,
|
||||
startTime: me.startTime,
|
||||
nonce: process.env.MOSAIC_BUS_CLI_NONCE,
|
||||
},
|
||||
});
|
||||
const result = await client.call(request.verb, request.args ?? {});
|
||||
process.stdout.write(JSON.stringify(result) + '\n');
|
||||
} catch (e) {
|
||||
process.stderr.write((/^[a-z-]{1,64}$/.test(e.code ?? '') ? e.code : 'invalid-request') + '\n');
|
||||
process.exitCode = 2;
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
import { readFileSync, lstatSync } from 'node:fs';
|
||||
import { basename, resolve } from 'node:path';
|
||||
import { BusError } from './broker.mjs';
|
||||
const MARKERS = [
|
||||
'MOSAIC_BUS_CAP',
|
||||
'MOSAIC_RUN_ID',
|
||||
'MOSAIC_AGENT_RUN',
|
||||
'CLAUDECODE',
|
||||
'CLAUDE_CODE_ENTRYPOINT',
|
||||
'CODEX_THREAD_ID',
|
||||
'PI_AGENT_DIR',
|
||||
];
|
||||
const refuse = () => {
|
||||
throw new BusError('human-required');
|
||||
};
|
||||
export function readProcess(pid, { readFile = readFileSync, stat = lstatSync } = {}) {
|
||||
if (!Number.isSafeInteger(pid) || pid < 1) refuse();
|
||||
try {
|
||||
const dir = `/proc/${pid}`,
|
||||
raw = readFile(dir + '/stat', 'utf8'),
|
||||
fields = raw.slice(raw.lastIndexOf(')') + 2).split(' ');
|
||||
let env = {};
|
||||
let environment;
|
||||
try {
|
||||
environment = readFile(dir + '/environ', 'utf8');
|
||||
} catch (e) {
|
||||
if (e.code !== 'EACCES') throw e;
|
||||
env = null;
|
||||
}
|
||||
if (environment !== undefined)
|
||||
for (const entry of environment.split('\0')) {
|
||||
const i = entry.indexOf('=');
|
||||
const key = entry.slice(0, i);
|
||||
if (key === 'MOSAIC_BUS_CLI_NONCE' || MARKERS.includes(key)) env[key] = entry.slice(i + 1);
|
||||
}
|
||||
return {
|
||||
pid,
|
||||
ppid: Number(fields[1]),
|
||||
startTime: fields[19],
|
||||
uid: stat(dir).uid,
|
||||
argv: readFile(dir + '/cmdline', 'utf8')
|
||||
.split('\0')
|
||||
.filter(Boolean),
|
||||
env,
|
||||
};
|
||||
} catch {
|
||||
refuse();
|
||||
}
|
||||
}
|
||||
// Cooperative same-UID process checks, not peer credentials or a hostile-process wall.
|
||||
// A nonce in the CLI's launch environment binds the submitted PID to a live CLI.
|
||||
export function verifyHuman(proof, { cliPath, launches = [], readProcess: read = readProcess }) {
|
||||
if (
|
||||
!proof ||
|
||||
typeof proof !== 'object' ||
|
||||
Object.keys(proof).some((k) => !['pid', 'startTime', 'nonce', 'business'].includes(k)) ||
|
||||
!Number.isSafeInteger(proof.pid) ||
|
||||
proof.pid < 2 ||
|
||||
typeof proof.business !== 'string' ||
|
||||
typeof proof.nonce !== 'string' ||
|
||||
!/^[a-f0-9]{64}$/.test(proof.nonce)
|
||||
)
|
||||
refuse();
|
||||
try {
|
||||
const p = read(proof.pid);
|
||||
if (
|
||||
p.uid !== process.getuid() ||
|
||||
p.startTime !== proof.startTime ||
|
||||
p.env?.MOSAIC_BUS_CLI_NONCE !== proof.nonce ||
|
||||
p.argv[1] !== resolve(cliPath)
|
||||
)
|
||||
refuse();
|
||||
const seen = new Set();
|
||||
let current = p;
|
||||
for (let depth = 0; depth < 128; depth++) {
|
||||
if (seen.has(current.pid)) refuse();
|
||||
seen.add(current.pid);
|
||||
if (
|
||||
(current.env !== null && MARKERS.some((k) => current.env[k])) ||
|
||||
launches.some((r) => r.pid === current.pid && r.startTime === current.startTime)
|
||||
)
|
||||
refuse();
|
||||
const command = basename(current.argv[0] ?? '');
|
||||
if (
|
||||
/^(pi|claude|claude-code|codex)(?:\.js)?$/.test(command) ||
|
||||
current.argv.some((v) => /\/(?:pi-coding-agent|codex)\/(?:dist|bin)\//.test(v))
|
||||
)
|
||||
refuse();
|
||||
if (current.ppid === 1) {
|
||||
const again = read(proof.pid);
|
||||
if (again.startTime !== proof.startTime || again.ppid !== p.ppid) refuse();
|
||||
return { business: proof.business };
|
||||
}
|
||||
if (current.ppid < 2) refuse();
|
||||
current = read(current.ppid);
|
||||
}
|
||||
} catch {
|
||||
refuse();
|
||||
}
|
||||
refuse();
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
export { Broker, BusError, ACTIONS } from './broker.mjs';
|
||||
export { startBroker } from './runtime.mjs';
|
||||
export { Client } from './client.mjs';
|
||||
export { views } from './views.mjs';
|
||||
export { busBusiness } from './business.mjs';
|
||||
@@ -0,0 +1,57 @@
|
||||
// Started with fork() by the trusted host; boot refs/capabilities travel over IPC,
|
||||
// never command arguments, stdout or service-token-bearing environment variables.
|
||||
import { startBroker } from './runtime.mjs';
|
||||
import { BusError } from './broker.mjs';
|
||||
let runtime,
|
||||
booted = false,
|
||||
closing = false;
|
||||
async function close(code) {
|
||||
if (closing) return;
|
||||
closing = true;
|
||||
try {
|
||||
await runtime?.close();
|
||||
} catch {
|
||||
code = 2;
|
||||
}
|
||||
process.exitCode = code;
|
||||
if (process.connected) process.disconnect();
|
||||
}
|
||||
if (!process.send) {
|
||||
process.stderr.write('trusted-host-required\n');
|
||||
process.exitCode = 2;
|
||||
} else {
|
||||
const timer = setTimeout(() => close(2), 10000);
|
||||
process.on('message', async (message) => {
|
||||
try {
|
||||
if (message?.op === 'bindLaunch' && runtime) {
|
||||
try {
|
||||
process.send({ ok: true, launch: runtime.bindLaunch(message.record) });
|
||||
} catch (e) {
|
||||
process.send({ ok: false, error: e instanceof BusError ? e.code : 'bind-refused' });
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (message?.op === 'close') {
|
||||
clearTimeout(timer);
|
||||
await close(0);
|
||||
return;
|
||||
}
|
||||
if (message?.op !== 'boot' || booted) throw new BusError('invalid-host-request');
|
||||
booted = true;
|
||||
clearTimeout(timer);
|
||||
runtime = await startBroker(message.config);
|
||||
if (closing) {
|
||||
await runtime.close();
|
||||
return;
|
||||
}
|
||||
process.send({ ok: true, path: runtime.path, launches: runtime.launches, readers: runtime.readers });
|
||||
} catch (e) {
|
||||
process.send?.({ ok: false, error: e instanceof BusError ? e.code : 'startup-refused' }, () =>
|
||||
close(2),
|
||||
);
|
||||
}
|
||||
});
|
||||
process.on('disconnect', () => close(2));
|
||||
process.on('SIGTERM', () => close(0));
|
||||
process.on('SIGINT', () => close(0));
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
import { join } from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { Store } from './store.mjs';
|
||||
import { Broker, BusError } from './broker.mjs';
|
||||
import { Credentials } from './credentials.mjs';
|
||||
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;
|
||||
try {
|
||||
const references = {};
|
||||
for (const [business, b] of Object.entries(businesses)) {
|
||||
for (const [role, r] of Object.entries(b.roles))
|
||||
if (r.credentials) references[business + '/' + role] = r.credentials;
|
||||
if (b.tracker?.sync?.credentials) references[business + '/@sync'] = b.tracker.sync.credentials;
|
||||
}
|
||||
const protectedRoots = [
|
||||
fileURLToPath(new URL('../../../', import.meta.url)),
|
||||
...repoRoots,
|
||||
...Object.values(businesses).flatMap((b) => Object.values(b.projects ?? {}).map((p) => p.root)),
|
||||
];
|
||||
credentials = new Credentials({ references, repoRoots: protectedRoots, dataRoot });
|
||||
store = new Store(dataRoot);
|
||||
const broker = new Broker({ store, businesses, secretCheck: (value) => credentials.assertClean(value) });
|
||||
// Process identities must be supplied for every managed run when using the human transport.
|
||||
const launchRecords = [];
|
||||
function bindLaunch(r) {
|
||||
if (
|
||||
!Number.isSafeInteger(r.pid) ||
|
||||
r.pid < 2 ||
|
||||
typeof r.startTime !== 'string' ||
|
||||
!/^\d+$/.test(r.startTime)
|
||||
)
|
||||
throw new BusError('invalid-launch-process');
|
||||
const cap = broker.bindLaunch(r);
|
||||
launchRecords.push({ ...r });
|
||||
return { business: r.business, run: r.run, cap };
|
||||
}
|
||||
const bound = launches.map(bindLaunch);
|
||||
const readCaps = readers.map((business) => ({ business, cap: broker.bindReader({ business }) }));
|
||||
const path = join(store.directory, 'broker.sock');
|
||||
server = await serve({
|
||||
broker,
|
||||
path,
|
||||
authenticateHuman: (proof) => {
|
||||
const { business } = verifyHuman(proof, {
|
||||
cliPath: fileURLToPath(new URL('./human-cli.mjs', import.meta.url)),
|
||||
launches: launchRecords,
|
||||
});
|
||||
const human = businesses[business]?.human;
|
||||
if (!human) throw new BusError('unknown-business');
|
||||
return broker.bindHuman({ business, human, via: 'cli', outsideAgent: true });
|
||||
},
|
||||
});
|
||||
return {
|
||||
broker,
|
||||
credentials,
|
||||
path,
|
||||
bindLaunch,
|
||||
launches: bound,
|
||||
readers: readCaps,
|
||||
async close() {
|
||||
await server.close();
|
||||
store.close();
|
||||
credentials.close();
|
||||
},
|
||||
};
|
||||
} catch (e) {
|
||||
await server?.close();
|
||||
store?.close();
|
||||
credentials?.close();
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
import { createServer } from 'node:net';
|
||||
import { lstatSync, chmodSync, unlinkSync } from 'node:fs';
|
||||
import { BusError } 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 }) {
|
||||
try {
|
||||
lstatSync(path);
|
||||
throw new BusError('socket-exists');
|
||||
} catch (e) {
|
||||
if (e.code !== 'ENOENT') throw e;
|
||||
}
|
||||
const sockets = new Set();
|
||||
const server = createServer({ allowHalfOpen: true }, (socket) => {
|
||||
sockets.add(socket);
|
||||
socket.on('close', () => sockets.delete(socket));
|
||||
socket.on('error', () => {});
|
||||
let input = Buffer.alloc(0),
|
||||
done = false;
|
||||
const finish = (result) => {
|
||||
if (done) return;
|
||||
done = true;
|
||||
socket.end(JSON.stringify(result) + '\n');
|
||||
};
|
||||
socket.setTimeout(timeout, () => {
|
||||
finish({ ok: false, error: 'request-timeout' });
|
||||
socket.destroySoon();
|
||||
});
|
||||
socket.on('data', (chunk) => {
|
||||
if (done) return;
|
||||
if (input.length + chunk.length > LIMIT) {
|
||||
finish({ ok: false, error: 'request-too-large' });
|
||||
return;
|
||||
}
|
||||
input = Buffer.concat([input, chunk]);
|
||||
const end = input.indexOf(10);
|
||||
if (end < 0) return;
|
||||
let transient;
|
||||
try {
|
||||
if (end !== input.length - 1) throw new BusError('invalid-envelope');
|
||||
const r = JSON.parse(input.subarray(0, end).toString('utf8'));
|
||||
if (
|
||||
!r ||
|
||||
Array.isArray(r) ||
|
||||
typeof r !== 'object' ||
|
||||
Object.keys(r).some((k) => !['cap', 'human', 'verb', 'args'].includes(k)) ||
|
||||
typeof r.verb !== 'string' ||
|
||||
Boolean(r.cap) === Boolean(r.human)
|
||||
)
|
||||
throw new BusError('invalid-envelope');
|
||||
let cap = r.cap;
|
||||
if (r.human) {
|
||||
if (!authenticateHuman) throw new BusError('human-required');
|
||||
transient = cap = authenticateHuman(r.human);
|
||||
}
|
||||
if (typeof cap !== 'string') throw new BusError('unauthenticated');
|
||||
const result = broker.request(cap, { verb: r.verb, args: r.args ?? {} });
|
||||
finish({ ok: true, result });
|
||||
} catch (e) {
|
||||
finish({ ok: false, error: e instanceof BusError ? e.code : 'invalid-envelope' });
|
||||
} finally {
|
||||
if (transient) broker.disconnect(transient);
|
||||
}
|
||||
});
|
||||
socket.on('end', () => {
|
||||
if (!done) finish({ ok: false, error: 'incomplete-request' });
|
||||
});
|
||||
});
|
||||
await new Promise((resolve, reject) => {
|
||||
server.once('error', reject);
|
||||
server.listen(path, () => {
|
||||
server.removeListener('error', reject);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
chmodSync(path, 0o600);
|
||||
const owned = lstatSync(path);
|
||||
return {
|
||||
async close() {
|
||||
for (const s of sockets) s.destroy();
|
||||
await new Promise((resolve) => server.close(resolve));
|
||||
try {
|
||||
const cur = lstatSync(path);
|
||||
if (cur.ino === owned.ino && cur.dev === owned.dev) unlinkSync(path);
|
||||
} catch (e) {
|
||||
if (e.code !== 'ENOENT') throw e;
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,175 @@
|
||||
import { DatabaseSync } from 'node:sqlite';
|
||||
import {
|
||||
readFileSync,
|
||||
mkdirSync,
|
||||
lstatSync,
|
||||
openSync,
|
||||
closeSync,
|
||||
writeFileSync,
|
||||
fsyncSync,
|
||||
unlinkSync,
|
||||
} from 'node:fs';
|
||||
import { join, isAbsolute } from 'node:path';
|
||||
import { createHash } from 'node:crypto';
|
||||
|
||||
const ddl = readFileSync(new URL('../schema.sql', import.meta.url), 'utf8');
|
||||
export const SCHEMA_SHA256 = '179ffe356d4ff19a49b5ebad39b6c6bfd7771deb8e1b7c55b5e746f039d69e65';
|
||||
const timedTables = [
|
||||
'events',
|
||||
'role_claims',
|
||||
'decisions',
|
||||
'decision_events',
|
||||
'messages',
|
||||
'deliveries',
|
||||
'task_snapshots',
|
||||
];
|
||||
const hex = (value) => createHash('sha256').update(value).digest('hex');
|
||||
if (hex(ddl) !== SCHEMA_SHA256) throw Error('schema-source-mismatch');
|
||||
export const schemaDigest = (db) =>
|
||||
hex(
|
||||
db
|
||||
.prepare('SELECT type,name,sql FROM sqlite_master WHERE sql IS NOT NULL ORDER BY type,name')
|
||||
.all()
|
||||
.map((r) => `${r.type}|${r.name}|${r.sql}`)
|
||||
.join('\n'),
|
||||
);
|
||||
const reference = new DatabaseSync(':memory:');
|
||||
reference.exec(ddl);
|
||||
const EXPECTED_DIGEST = schemaDigest(reference);
|
||||
reference.close();
|
||||
function exists(path) {
|
||||
try {
|
||||
return lstatSync(path);
|
||||
} catch (e) {
|
||||
if (e.code === 'ENOENT') return null;
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
function safe(path, directory = false) {
|
||||
const s = lstatSync(path);
|
||||
if (
|
||||
s.isSymbolicLink() ||
|
||||
!(directory ? s.isDirectory() : s.isFile()) ||
|
||||
s.uid !== process.getuid() ||
|
||||
s.mode & 0o077
|
||||
)
|
||||
throw Error('unsafe-path');
|
||||
return s;
|
||||
}
|
||||
|
||||
// Internal trusted SQL surface. Only the broker owns a Store; never expose SQL over the socket.
|
||||
export class Store {
|
||||
#db;
|
||||
#lock;
|
||||
#lockStat;
|
||||
#closed = false;
|
||||
#inTransaction = false;
|
||||
constructor(dataRoot) {
|
||||
if (typeof dataRoot !== 'string' || !isAbsolute(dataRoot)) throw Error('unsafe-path');
|
||||
const rootStat = lstatSync(dataRoot);
|
||||
if (!rootStat.isDirectory() || rootStat.isSymbolicLink() || rootStat.uid !== process.getuid())
|
||||
throw Error('unsafe-path');
|
||||
this.directory = join(dataRoot, 'bus');
|
||||
if (!exists(this.directory)) mkdirSync(this.directory, { mode: 0o700 });
|
||||
safe(this.directory, true);
|
||||
this.#lock = join(this.directory, 'writer.lock');
|
||||
let fd;
|
||||
try {
|
||||
fd = openSync(this.#lock, 'wx', 0o600);
|
||||
} catch (e) {
|
||||
if (e.code === 'EEXIST') throw Error('writer-locked');
|
||||
throw e;
|
||||
}
|
||||
this.#lockStat = lstatSync(this.#lock);
|
||||
try {
|
||||
writeFileSync(fd, JSON.stringify({ pid: process.pid, at: new Date().toISOString() }));
|
||||
fsyncSync(fd);
|
||||
closeSync(fd);
|
||||
fd = undefined;
|
||||
const dirfd = openSync(this.directory, 'r');
|
||||
try {
|
||||
fsyncSync(dirfd);
|
||||
} finally {
|
||||
closeSync(dirfd);
|
||||
}
|
||||
this.path = join(this.directory, 'bus.sqlite');
|
||||
const fresh = !exists(this.path);
|
||||
if (fresh) closeSync(openSync(this.path, 'wx', 0o600));
|
||||
else safe(this.path);
|
||||
for (const suffix of ['-wal', '-shm']) if (exists(this.path + suffix)) safe(this.path + suffix);
|
||||
this.#db = new DatabaseSync(this.path, { timeout: 5000, defensive: true });
|
||||
this.#db.exec('PRAGMA foreign_keys=ON; PRAGMA synchronous=FULL');
|
||||
if (fresh) {
|
||||
// journal_mode must be set outside a transaction.
|
||||
this.#db.exec('PRAGMA journal_mode=WAL');
|
||||
this.transaction(() => {
|
||||
this.#db.exec(ddl.replace('PRAGMA journal_mode = WAL;', ''));
|
||||
this.run('INSERT INTO meta(key,value) VALUES (?,?)', 'schema_version', '3b');
|
||||
this.run('INSERT INTO meta(key,value) VALUES (?,?)', 'schema_digest', EXPECTED_DIGEST);
|
||||
});
|
||||
}
|
||||
if (
|
||||
schemaDigest(this.#db) !== EXPECTED_DIGEST ||
|
||||
this.get("SELECT value FROM meta WHERE key='schema_version'")?.value !== '3b' ||
|
||||
this.get("SELECT value FROM meta WHERE key='schema_digest'")?.value !== EXPECTED_DIGEST ||
|
||||
this.get('PRAGMA journal_mode').journal_mode !== 'wal'
|
||||
)
|
||||
throw Error('schema-mismatch');
|
||||
if (this.get('PRAGMA quick_check').quick_check !== 'ok' || this.all('PRAGMA foreign_key_check').length)
|
||||
throw Error('integrity-failed');
|
||||
this.#verifyTimes();
|
||||
} catch (e) {
|
||||
if (fd !== undefined) closeSync(fd);
|
||||
this.close();
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
get(sql, ...args) {
|
||||
return this.#db.prepare(sql).get(...args);
|
||||
}
|
||||
all(sql, ...args) {
|
||||
return this.#db.prepare(sql).all(...args);
|
||||
}
|
||||
run(sql, ...args) {
|
||||
if (!this.#inTransaction) return this.transaction(() => this.run(sql, ...args));
|
||||
return this.#db.prepare(sql).run(...args);
|
||||
}
|
||||
#verifyTimes() {
|
||||
for (const table of timedTables) {
|
||||
const badAt =
|
||||
"length(at)<>24 OR at NOT GLOB '[0-9][0-9][0-9][0-9]-*' OR strftime('%Y-%m-%dT%H:%M:%fZ',at) IS NOT at";
|
||||
const badRead =
|
||||
table === 'task_snapshots'
|
||||
? " OR (read_at IS NOT NULL AND (length(read_at)<>24 OR read_at NOT GLOB '[0-9][0-9][0-9][0-9]-*' OR strftime('%Y-%m-%dT%H:%M:%fZ',read_at) IS NOT read_at))"
|
||||
: '';
|
||||
if (this.get(`SELECT 1 FROM ${table} WHERE ${badAt}${badRead} LIMIT 1`))
|
||||
throw Error('invalid-timestamp');
|
||||
}
|
||||
}
|
||||
transaction(fn) {
|
||||
if (typeof fn !== 'function' || fn.constructor.name === 'AsyncFunction')
|
||||
throw Error('async-transaction-refused');
|
||||
if (this.#inTransaction) throw Error('nested-transaction');
|
||||
this.#db.exec('BEGIN IMMEDIATE');
|
||||
this.#inTransaction = true;
|
||||
try {
|
||||
const result = fn();
|
||||
if (result?.then) throw Error('async-transaction-refused');
|
||||
this.#verifyTimes();
|
||||
this.#db.exec('COMMIT');
|
||||
return result;
|
||||
} catch (e) {
|
||||
this.#db.exec('ROLLBACK');
|
||||
throw e;
|
||||
} finally {
|
||||
this.#inTransaction = false;
|
||||
}
|
||||
}
|
||||
close() {
|
||||
if (this.#closed) return;
|
||||
this.#closed = true;
|
||||
this.#db?.close();
|
||||
const current = exists(this.#lock);
|
||||
if (current?.ino === this.#lockStat?.ino && current?.dev === this.#lockStat?.dev) unlinkSync(this.#lock);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
// Shared by the S4 CLI and S5 WebUI. Client-only: never opens bus.sqlite.
|
||||
export function views(client) {
|
||||
return Object.freeze({
|
||||
inbox: () => client.call('inbox'),
|
||||
tasks: () => client.call('tasks'),
|
||||
agents: () => client.call('agents'),
|
||||
trail: (subject) => client.call('trail', { subject }),
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user