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:
2026-10-05 17:15:53 -05:00
co-authored by Claude Opus 5.5
parent 942dca9ea8
commit 38828a2cb3
24 changed files with 3646 additions and 0 deletions
+228
View File
@@ -0,0 +1,228 @@
# Bus and broker core — slice 1 S2
The broker owns `<dataRoot>/bus/` and is the only cooperative writer of
`bus.sqlite`. Consumers use the local socket. This package does not load or
write the system or business configuration, launch a harness, call a service,
or implement task verbs. S1 supplies validated/resolved definitions; S3 adds
task handlers; S6 supplies trusted launch records. No service credential is
returned to a client or written to the database.
Slice 1 shares an OS user. Files, capabilities, process ancestry checks and
hooks are **not a wall against a hostile same-UID process**. In particular,
Linux `/proc` checks here are not SO_PEERCRED authentication. The trusted host
and well-behaved clients must follow this protocol. Containers remain later work.
## Start and stop
Requires Linux and Node 24 or 26, local disk, no dependencies. The data root must
already exist and belong to the user. A broker creates `bus/` mode 0700 and the
DB/lock/socket mode 0600. It refuses symlinks or unsafe existing permissions.
The trusted host can import `startBroker` from `src/index.mjs`, or fork
`src/process.mjs` with an IPC channel. With `fork`, send exactly one boot message:
```js
{op: 'boot', config: {
dataRoot, businesses, launches, readers: ['business-id'], repoRoots: [repoRoot]
}}
```
`businesses` maps business IDs to normalized business definitions. Call
`busBusiness(business, resolvedByInstance)` to adapt S1's `loadBusiness` and
`resolveInstance` results. It copies resolved authority limits and credential
references into each role instance. It never expands an authority map or loads
files. Until S1 freezes its interface, callers may supply that normalized shape
explicitly: `{id, human, arbiters:{technical,delivery}, roles:{instance:{authority:
{withinRole:[],crossRole:[]}, credentials:{gitea?,vikunja?}}}, launch?}`. Empty
maps grant no actions; missing maps refuse. The runtime always excludes its own repository root and every declared project
root when loading credentials; `repoRoots` adds any other source checkouts. The host owns config/definition validation.
Launch records are `{business, role, run, harness, address?, pid, startTime}`.
`pid` and kernel `/proc/<pid>/stat` start time identify the managed process.
Harness is `pi` or `claude-code`. Bind only a launch the trusted host actually
performed. Binding appends `session.launched`; rebinding after broker restart
requires identical recorded identity. A changed role, harness, address or process
identity refuses. Never register an arbitrary client-supplied launch record.
The launcher must set `MOSAIC_RUN_ID` in every managed process's environment and
preserve it for descendants, including after the original process exits.
The boot reply is `{ok:true,path,launches:[{business,run,cap}],readers:
[{business,cap}]}`. Distribute each random capability only to its intended
run/reader, through an inherited private channel. It is a broker capability,
not a Gitea/Vikunja token. A capability binds business, role and run; socket
requests never declare those identities. Reconnect after restart requires a
new capability from the trusted host, not a socket re-registration.
For a later real launch the trusted host sends `{op:'bindLaunch',record}` over
its original IPC channel, receiving `{ok:true,launch:{business,run,cap}}`.
A refused binding returns `{ok:false,error}` and keeps the broker running.
The embedded runtime offers the same `bindLaunch(record)` wrapper, which also
updates the human ancestry registry. S6 must use that wrapper. Do not use the
underlying Broker's test-level binding API to bypass runtime process checks.
`{op:'close'}`, SIGTERM and SIGINT close the socket and DB and remove the owned
writer lock. Losing the host IPC channel closes with exit 2. Startup or host
protocol failures return a static error code and exit 2, never a raw exception.
Capabilities vanish with the process. A SIGKILL/crash leaves `writer.lock` and
possibly the socket: **there is no automatic stale-lock reclaim**. An operator
must establish that the old broker is gone, preserve the DB/WAL/SHM as evidence,
and explicitly remove only the stale lock/socket before restarting. Never delete
the database to get around an open-time refusal.
## Client protocol
Import `Client` from `@mosaic/bus/client` (or `src/client.mjs`).
`new Client({path,cap}).call(verb,args)` opens one Unix socket per request and
returns plain JSON or throws `BusError` with a static `.code`. One JSON line,
64 KiB request maximum, five-second timeout. There is no SQL verb and no
automatic retry. Responses are capped at 4 MiB by the client. Oversized reads
refuse; pagination is not part of this first core.
A disconnect or timeout reports `outcome-unknown`. A mutation may already have
committed. Callers must inspect its evidence before deciding what to do; they
must not resend automatically. `message.receive` records delivery when the
broker hands the batch to the connection, not when the harness acts on it.
A lost response leaves an uncertain delivery, never automatic replay to a new
holder. `message.read` is an explicit acknowledgement. `trail(messageId)` lets
an operator inspect content and delivery evidence after an uncertain result.
| Verb | Arguments and result |
|---|---|
| `role.claim` | No args; claims capability's role/run. One current holder per business/instance; a revoked run cannot reclaim that role, including after restart. |
| `role.release` | No args; only its current holder releases. |
| `role.revoke` | `{role,decision}`; needs human-approved gated decision for this exact target and holder incarnation. |
| `decision.raise` | `{action,question,options:[{key,text}],recommendation,blocking,domain?,target?,project?,task_ref?,requirement_ref?,supersedes?,choice?,approvalChoice?}`. Returns decision with decoded options and authorization context. |
| `decision.resolve` | `{id,choice,note?}`; option must exist, resolver must be routed role or outside-agent human CLI. |
| `decision.seen` | `{id,note?}`; acknowledgement by resolver, not a resolution. |
| `decision.withdraw` | `{id,note?}`; raiser or human. Never rewrites the decision. |
| `message.send` | `{to,body,class?,in_reply_to?,corrects?,decision?}`; returns `{id,request}`. Human instruction's `request` is its `human.input` event ID. |
| `message.receive` | No args; up to 100 undelivered messages for current holder, including originating `request` for human instructions. |
| `message.read` | `{id}`; recipient holder's acknowledgement; no new holder can acknowledge old holder's delivery. |
| `event.emit` | `{kind:'action.allowed',body:{action:'routine',target?},subject?}`; routine observations only. |
| `launch.revoke`, `launch.restore` | No args; outside-agent human CLI only. |
| `inbox`, `tasks`, `agents` | No args; plain JSON projections scoped to capability's business. |
| `trail` | `{subject}`: task ref, decision UUID or message UUID. Linked records in broker write order. |
Unknown keys, roles, references, actions and verbs refuse. If writing refusal
evidence fails, the caller receives the static `storage-unavailable` code. `human` is the reserved
recipient. Message classes are REQUEST, ASSIGNMENT, REVIEW-REQUEST, REVIEW-RESULT,
RESULT, INFO, DECISION and REACTION. References cannot cross businesses.
Routing is immutable at raise: within-role closes atomically with `choice`,
cross-role routes to the configured delivery arbiter (or technical arbiter with
`domain:'technical'`). If that arbiter is the raiser, resolution routes to human
while the decision keeps its cross-role class. Gated routes to human. Unlisted actions are gated and the
nine gated-only actions cannot appear in authority maps. A blocking decision
requires a strict task ref. `role.launch` is gated unless the business launch
block's `by` names the acting role and refused when launch state is revoked. S6 still owns launch allowlists,
capacity and actual spawning.
Authorization requires the resolved option identified by `approvalChoice`
(default `yes`, which must exist), the same raiser role/run, action and target.
Inbox projections expose `{action,target,approvalChoice}` as `authorization`:
**S4 must display this context when offering resolution**, not infer approval
from recommendation or option position. Any other option is a refusal, even if
the decision is resolved. Approval is context-bound, not an exactly-once service
operation receipt. S3/S6 must implement their own external-operation identity and
uncertain-outcome handling. Raising/superseding and resolving are transactional.
## Human and reader paths
`src/human-cli.mjs <socket>` is S4's transport shim, not another product CLI.
It reads `{business,verb,args}` from stdin and prints JSON. It re-execs once with
a fresh nonce in its launch environment; the broker checks its exact script
path, UID, PID/start time, nonce and ancestry on every request. Missing process
proof, ancestry loops/disappearance, a known launch ancestor, or inherited
Mosaic/Pi/Claude/Codex run markers refuse. No claimed human name is accepted from
the socket. The configured business human is bound only after these checks.
An ancestor whose environment returns EACCES skips only the environment-marker
check: command and launch-registry checks still apply. All other read failures
refuse, and the CLI itself must have a readable nonce environment.
Reading `/proc` is Linux-specific. Reparenting and malicious same-UID PID/nonce
impersonation remain within the explicit cooperative trust limit.
Human resolutions and launch toggles append `human.input`. Agent capabilities
and read-only WebUI capabilities cannot reach these operations. Reader
capabilities expose the human inbox plus tasks/agents/trail for one business;
they do not confer mutation or resolution authority.
## Shared reads and adapter boundary
```js
import {Client} from './src/client.mjs';
import {views} from './src/views.mjs';
const read = views(new Client({path,cap}));
await read.inbox(); await read.tasks(); await read.agents();
await read.trail('vikunja:3/41');
```
S4 and S5 must import this client-only module. Neither opens `bus.sqlite`.
The tasks projection reports the schema’s `task_current` view, including tombstones; the caller
must not pretend a missing/inaccessible task is open or completed. Trail joins
human request, messages/deliveries, decisions/resolutions, task snapshots and
session/claim evidence. Broker timestamps advance monotonically across all
its writes, including after reopen; per-table `seq` stays the primary table
order. Timestamps can run slightly ahead of wall time during bursts. S3's
service `updated` and `read_at` remain observations, not broker write order.
Trusted in-process adapters can call `broker.authorize(cap,action,{decision,
target})` and `broker.recordEvent(cap,{kind,body,subject})`. These APIs are never
socket verbs. `authorize` verifies the current holder and appends action evidence;
it is not itself an external service executor. An adapter must not treat a prior
check as authority after an await/holder change. Task/credential/poll registration
and serialized external operations are S3's integration work, not a generic
user-defined-handler endpoint in S2. `recordEvent` stamps identity, uses a
transaction, and cannot emit human input or human launch control events.
Task refs in action/review bodies require matching `subject`; all task events
are also guarded by the pinned schema. `task.created.body.request` must name an
existing same-business human input, with a requirement ID.
Credentials are keyed by `business/instance` in the runtime; the distinct read-only
tracker sync identity uses `business/@sync` and has no agent capability. Trusted service
adapters use `runtime.credentials.use('business/instance','gitea', callback)`;
only that in-process callback receives the token, never the socket client.
Callbacks must not spawn children with it, log it or persist it. Errors are
replaced with `service-failed`. Known-token content in returned values, keys,
request payloads and broker evidence refuses. This is not general-purpose secret
redaction or protection against a malicious adapter encoding/exfiltrating values.
File refs are absolute, regular non-symlink files, owned, exactly 0600, outside
repo roots and dataRoot, bounded to 16 KiB. Open uses O_NOFOLLOW and verifies
inode/device before reading. Environment refs are supported; neither kind is
forwarded to agents. Tokens use the opaque ASCII alphabet `[A-Za-z0-9._~-]`,
with a minimum of 16 characters and at most one terminal LF. Dates are strict YYYY-MM-DD: Gitea `rotateBy`,
Vikunja `expires`. Vikunja use refuses at midnight UTC on expiry. Gitea rotation
due is a warning state, not an expiry; `credentials.status()` reports metadata
only (valid, expiring, expired or rotation-due) for the S3 notifier. Tokens are read once
at startup and retained only in memory. Restart after a rotation. S3 owns live
scope probes and rotation/expiry notifications; S2 never probes real services.
## Schema and verification
`schema.sql` is unchanged v3b, SHA-256
`179ffe356d4ff19a49b5ebad39b6c6bfd7771deb8e1b7c55b5e746f039d69e65`.
It has eight tables, **36** triggers and five views. The prototype's displayed
35 is after deliberately dropping a guard. Source digest and a trusted full
sqlite_master digest are checked at open, alongside schema version `3b`, WAL,
quick_check and foreign-key integrity. Missing metadata refuses; it never
reinitializes an existing file. Full synchronous WAL transactions either commit
or roll back. SQL UPDATE, DELETE and INSERT OR REPLACE refuse on every table.
No migration or pruning command ships here.
Run `node --test packages/bus/tests/*.test.mjs`. Tests use private disposable
directories, fixture token files, fake process tables, real local child brokers,
real SQLite and real Unix sockets. The prototype cases are assertions, not
print-only observations. Crash recovery deliberately tests refusal, not automatic
repair. Positive human process classification uses synthetic process ancestry;
real `/proc` reading and real-process refusal are separate tests. There is no
claim that an agent's test run proves a live outside-agent human terminal.
Lead decision 60 also fixes timestamp precision. Before every commit (including
standalone trusted `Store.run` inserts) and on reopen, the store refuses a row
whose `at` or non-null `read_at` differs from SQLite's canonical millisecond UTC
form, `YYYY-MM-DDTHH:mm:ss.sssZ`. It rolls back the entire transaction, not just
the malformed row. This deliberately scans the timestamp columns in this first
core; an indexed/insert-scoped optimization requires equivalent checks and tests.
The schema itself remains byte-identical to v3b. S3 must normalize service read
start times to this form before storing snapshots. External service `updated`
is not a broker timestamp; its interpretation remains S3's responsibility.
+16
View File
@@ -0,0 +1,16 @@
{
"name": "@mosaic/bus",
"private": true,
"type": "module",
"exports": {
".": "./src/index.mjs",
"./client": "./src/client.mjs",
"./views": "./src/views.mjs"
},
"engines": {
"node": ">=24"
},
"scripts": {
"test": "node --test tests/*.test.mjs"
}
}
+211
View File
@@ -0,0 +1,211 @@
PRAGMA journal_mode = WAL;
PRAGMA foreign_keys = ON;
CREATE TABLE meta (key TEXT PRIMARY KEY, value TEXT NOT NULL) STRICT;
CREATE TABLE events (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
id TEXT NOT NULL UNIQUE,
at TEXT NOT NULL,
business TEXT NOT NULL,
kind TEXT NOT NULL CHECK (kind IN (
'session.launched',
'session.ended',
'action.allowed',
'action.refused',
'task.created',
'task.assigned',
'task.state',
'task.closed',
'task.changed.external',
'task.conflict',
'task.missing',
'review.requested',
'review.verdict',
'human.input',
'config.refused',
'credential.expiring',
'credential.expired',
'credential.changed',
'launch.revoked',
'launch.restored',
'digest.sent')),
actor_role TEXT, actor_run TEXT,
subject TEXT,
corrects TEXT REFERENCES events(id),
body TEXT NOT NULL CHECK (json_valid(body))
) STRICT;
CREATE TABLE role_claims (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
at TEXT NOT NULL,
business TEXT NOT NULL, role TEXT NOT NULL,
op TEXT NOT NULL CHECK (op IN ('claim','release','revoke')),
holder_run TEXT NOT NULL,
harness TEXT NOT NULL, address TEXT,
by TEXT NOT NULL, reason TEXT,
decision TEXT
) STRICT;
CREATE TABLE decisions (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
id TEXT NOT NULL UNIQUE,
at TEXT NOT NULL,
business TEXT NOT NULL, project TEXT,
raised_by_role TEXT NOT NULL, raised_by_run TEXT NOT NULL,
class TEXT NOT NULL CHECK (class IN ('routine','within-role','cross-role','gated')),
action TEXT NOT NULL,
route_to TEXT NOT NULL,
question TEXT NOT NULL,
options TEXT NOT NULL CHECK (json_valid(options) AND json_array_length(options) BETWEEN 2 AND 9),
recommendation TEXT NOT NULL,
task_ref TEXT CHECK (task_ref IS NULL OR (task_ref GLOB 'vikunja:[1-9]*/[1-9]*' AND NOT substr(task_ref, 9) GLOB '*[^0-9/]*' AND NOT substr(task_ref, 9) GLOB '*/*/*')), requirement_ref TEXT,
blocking INTEGER NOT NULL CHECK (blocking IN (0,1)),
supersedes TEXT REFERENCES decisions(id)
) STRICT;
CREATE TABLE decision_events (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
decision TEXT NOT NULL REFERENCES decisions(id),
at TEXT NOT NULL,
op TEXT NOT NULL CHECK (op IN ('seen','resolved','withdrawn','expired')),
by TEXT NOT NULL,
choice TEXT, note TEXT, via TEXT
) STRICT;
CREATE TABLE messages (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
id TEXT NOT NULL UNIQUE,
at TEXT NOT NULL,
business TEXT NOT NULL,
from_role TEXT NOT NULL, from_run TEXT NOT NULL,
to_role TEXT NOT NULL,
class TEXT NOT NULL,
in_reply_to TEXT REFERENCES messages(id),
decision TEXT REFERENCES decisions(id),
corrects TEXT REFERENCES messages(id),
body TEXT NOT NULL
) STRICT;
CREATE TABLE deliveries (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
message TEXT NOT NULL REFERENCES messages(id),
at TEXT NOT NULL,
op TEXT NOT NULL CHECK (op IN ('routed','delivered','failed','read')),
holder_run TEXT, transport TEXT, address TEXT, detail TEXT
) STRICT;
CREATE TABLE task_snapshots (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
at TEXT NOT NULL,
business TEXT NOT NULL,
task_ref TEXT NOT NULL CHECK (task_ref GLOB 'vikunja:[1-9]*/[1-9]*' AND NOT substr(task_ref, 9) GLOB '*[^0-9/]*' AND NOT substr(task_ref, 9) GLOB '*/*/*'),
updated TEXT NOT NULL,
etag TEXT,
digest TEXT NOT NULL CHECK (length(digest) = 64 AND NOT digest GLOB '*[^0-9a-f]*'),
fields TEXT NOT NULL CHECK (json_valid(fields)),
source TEXT NOT NULL CHECK (source IN ('self','poll')),
via TEXT CHECK (via IN ('board','cursor','task','reconcile')),
read_at TEXT,
role TEXT, run TEXT,
CHECK ((source = 'self') = (role IS NOT NULL AND run IS NOT NULL)),
CHECK ((source = 'poll') = (via IS NOT NULL AND read_at IS NOT NULL)),
CHECK (json_type(fields, '$.bucket') IS 'integer'
OR (source IS 'poll' AND via IS 'task' AND json_type(fields, '$.gone') IS 'text'))
) STRICT;
CREATE INDEX task_snapshots_ref ON task_snapshots (business, task_ref, seq);
CREATE TRIGGER decisions_resolve_once BEFORE INSERT ON decision_events
WHEN NEW.op IN ('resolved','withdrawn','expired') AND EXISTS (
SELECT 1 FROM decision_events WHERE decision = NEW.decision AND op IN ('resolved','withdrawn','expired'))
BEGIN SELECT RAISE(ABORT, 'decision already closed'); END;
CREATE TRIGGER decisions_resolved_choice BEFORE INSERT ON decision_events
WHEN NEW.op = 'resolved' AND (NEW.choice IS NULL OR NOT EXISTS (
SELECT 1 FROM decisions d, json_each(d.options) o WHERE d.id = NEW.decision AND json_extract(o.value,'$.key') = NEW.choice))
BEGIN SELECT RAISE(ABORT, 'resolution must name one of the options'); END;
CREATE TRIGGER role_one_holder BEFORE INSERT ON role_claims
WHEN NEW.op = 'claim' AND (SELECT op FROM role_claims WHERE business = NEW.business AND role = NEW.role ORDER BY seq DESC LIMIT 1) = 'claim'
BEGIN SELECT RAISE(ABORT, 'role already held'); END;
CREATE TRIGGER role_release_by_holder BEFORE INSERT ON role_claims
WHEN NEW.op IN ('release','revoke') AND COALESCE((SELECT op FROM role_claims WHERE business = NEW.business AND role = NEW.role ORDER BY seq DESC LIMIT 1),'') <> 'claim'
BEGIN SELECT RAISE(ABORT, 'role is not held'); END;
CREATE TRIGGER role_release_same_run BEFORE INSERT ON role_claims
WHEN NEW.op = 'release' AND (SELECT holder_run FROM role_claims WHERE business = NEW.business AND role = NEW.role ORDER BY seq DESC LIMIT 1) <> NEW.holder_run
BEGIN SELECT RAISE(ABORT, 'only the holder releases; others revoke'); END;
CREATE TRIGGER role_revoke_needs_decision BEFORE INSERT ON role_claims
WHEN NEW.op = 'revoke' AND NEW.decision IS NULL
BEGIN SELECT RAISE(ABORT, 'revoke needs a resolved decision'); END;
CREATE TRIGGER meta_no_update BEFORE UPDATE ON meta BEGIN SELECT RAISE(ABORT, 'meta is append-only'); END;
CREATE TRIGGER meta_no_delete BEFORE DELETE ON meta BEGIN SELECT RAISE(ABORT, 'meta is append-only'); END;
CREATE TRIGGER events_no_update BEFORE UPDATE ON events BEGIN SELECT RAISE(ABORT, 'events is append-only'); END;
CREATE TRIGGER events_no_delete BEFORE DELETE ON events BEGIN SELECT RAISE(ABORT, 'events is append-only'); END;
CREATE TRIGGER role_claims_no_update BEFORE UPDATE ON role_claims BEGIN SELECT RAISE(ABORT, 'role_claims is append-only'); END;
CREATE TRIGGER role_claims_no_delete BEFORE DELETE ON role_claims BEGIN SELECT RAISE(ABORT, 'role_claims is append-only'); END;
CREATE TRIGGER decisions_no_update BEFORE UPDATE ON decisions BEGIN SELECT RAISE(ABORT, 'decisions is append-only'); END;
CREATE TRIGGER decisions_no_delete BEFORE DELETE ON decisions BEGIN SELECT RAISE(ABORT, 'decisions is append-only'); END;
CREATE TRIGGER decision_events_no_update BEFORE UPDATE ON decision_events BEGIN SELECT RAISE(ABORT, 'decision_events is append-only'); END;
CREATE TRIGGER decision_events_no_delete BEFORE DELETE ON decision_events BEGIN SELECT RAISE(ABORT, 'decision_events is append-only'); END;
CREATE TRIGGER messages_no_update BEFORE UPDATE ON messages BEGIN SELECT RAISE(ABORT, 'messages is append-only'); END;
CREATE TRIGGER messages_no_delete BEFORE DELETE ON messages BEGIN SELECT RAISE(ABORT, 'messages is append-only'); END;
CREATE TRIGGER deliveries_no_update BEFORE UPDATE ON deliveries BEGIN SELECT RAISE(ABORT, 'deliveries is append-only'); END;
CREATE TRIGGER deliveries_no_delete BEFORE DELETE ON deliveries BEGIN SELECT RAISE(ABORT, 'deliveries is append-only'); END;
CREATE TRIGGER meta_no_replace BEFORE INSERT ON meta WHEN EXISTS (SELECT 1 FROM meta WHERE key = NEW.key) BEGIN SELECT RAISE(ABORT, 'meta is append-only'); END;
CREATE TRIGGER events_no_replace BEFORE INSERT ON events WHEN EXISTS (SELECT 1 FROM events WHERE seq = NEW.seq OR id = NEW.id) BEGIN SELECT RAISE(ABORT, 'events is append-only'); END;
CREATE TRIGGER role_claims_no_replace BEFORE INSERT ON role_claims WHEN EXISTS (SELECT 1 FROM role_claims WHERE seq = NEW.seq) BEGIN SELECT RAISE(ABORT, 'role_claims is append-only'); END;
CREATE TRIGGER decisions_no_replace BEFORE INSERT ON decisions WHEN EXISTS (SELECT 1 FROM decisions WHERE seq = NEW.seq OR id = NEW.id) BEGIN SELECT RAISE(ABORT, 'decisions is append-only'); END;
CREATE TRIGGER decision_events_no_replace BEFORE INSERT ON decision_events WHEN EXISTS (SELECT 1 FROM decision_events WHERE seq = NEW.seq) BEGIN SELECT RAISE(ABORT, 'decision_events is append-only'); END;
CREATE TRIGGER messages_no_replace BEFORE INSERT ON messages WHEN EXISTS (SELECT 1 FROM messages WHERE seq = NEW.seq OR id = NEW.id) BEGIN SELECT RAISE(ABORT, 'messages is append-only'); END;
CREATE TRIGGER deliveries_no_replace BEFORE INSERT ON deliveries WHEN EXISTS (SELECT 1 FROM deliveries WHERE seq = NEW.seq) BEGIN SELECT RAISE(ABORT, 'deliveries is append-only'); END;
CREATE TRIGGER task_snapshots_no_update BEFORE UPDATE ON task_snapshots BEGIN SELECT RAISE(ABORT, 'task_snapshots is append-only'); END;
CREATE TRIGGER task_snapshots_no_delete BEFORE DELETE ON task_snapshots BEGIN SELECT RAISE(ABORT, 'task_snapshots is append-only'); END;
CREATE TRIGGER task_snapshots_no_replace BEFORE INSERT ON task_snapshots WHEN EXISTS (SELECT 1 FROM task_snapshots WHERE seq = NEW.seq) BEGIN SELECT RAISE(ABORT, 'task_snapshots is append-only'); END;
CREATE TRIGGER decisions_blocking_needs_task BEFORE INSERT ON decisions
WHEN NEW.blocking = 1 AND NEW.task_ref IS NULL
BEGIN SELECT RAISE(ABORT, 'a blocking decision cites the task it blocks'); END;
CREATE TRIGGER events_launch_by_human BEFORE INSERT ON events
WHEN NEW.kind IN ('launch.revoked','launch.restored') AND (NEW.actor_role IS NOT NULL OR NEW.actor_run IS NOT NULL)
BEGIN SELECT RAISE(ABORT, 'only the human revokes or restores launching'); END;
CREATE TRIGGER events_credential_body BEFORE INSERT ON events
WHEN NEW.kind GLOB 'credential.*' AND (
json_extract(NEW.body, '$.service') IS NULL OR json_extract(NEW.body, '$.service') NOT IN ('gitea','vikunja')
OR json_extract(NEW.body, '$.instance') IS NULL)
BEGIN SELECT RAISE(ABORT, 'credential events name a service and a role instance'); END;
CREATE TRIGGER events_task_missing_body BEFORE INSERT ON events
WHEN NEW.kind = 'task.missing' AND (
json_extract(NEW.body, '$.reason') IS NULL OR json_extract(NEW.body, '$.reason') NOT IN ('moved','not-found','no-access')
OR (json_extract(NEW.body, '$.reason') = 'moved' AND json_type(NEW.body, '$.project') IS NOT 'integer'))
BEGIN SELECT RAISE(ABORT, 'task.missing carries a reason, and the new project when moved'); END;
CREATE TRIGGER events_task_subject BEFORE INSERT ON events
WHEN NEW.kind GLOB 'task.*' AND (NEW.subject IS NULL OR NOT (NEW.subject GLOB 'vikunja:[1-9]*/[1-9]*' AND NOT substr(NEW.subject, 9) GLOB '*[^0-9/]*' AND NOT substr(NEW.subject, 9) GLOB '*/*/*'))
BEGIN SELECT RAISE(ABORT, 'a task event names its task in subject'); END;
CREATE TRIGGER events_task_created_body BEFORE INSERT ON events
WHEN NEW.kind = 'task.created' AND (
json_type(NEW.body, '$.request') IS NOT 'text'
OR NOT EXISTS (SELECT 1 FROM events WHERE id = json_extract(NEW.body, '$.request')
AND kind = 'human.input' AND business = NEW.business)
OR json_type(NEW.body, '$.requirement') IS NOT 'text'
OR NOT json_extract(NEW.body, '$.requirement') GLOB 'REQ-[A-Z]*-[1-9]*'
OR json_extract(NEW.body, '$.requirement') GLOB 'REQ-*[^A-Z0-9-]*')
BEGIN SELECT RAISE(ABORT, 'task.created cites the human.input that asked for it and a requirement id'); END;
CREATE VIEW launch_state AS
SELECT business, CASE kind WHEN 'launch.revoked' THEN 'revoked' ELSE 'allowed' END AS state, at
FROM events e WHERE kind IN ('launch.revoked','launch.restored')
AND seq = (SELECT max(seq) FROM events WHERE business = e.business AND kind IN ('launch.revoked','launch.restored'));
CREATE VIEW task_external_changes AS
SELECT p.business, p.task_ref, p.seq, p.via, p.updated, p.digest, ls.digest AS self_digest
FROM task_snapshots p
LEFT JOIN task_snapshots ls ON ls.seq = (SELECT max(seq) FROM task_snapshots
WHERE business = p.business AND task_ref = p.task_ref AND source = 'self')
WHERE p.source = 'poll'
AND p.seq = (SELECT max(seq) FROM task_snapshots WHERE business = p.business AND task_ref = p.task_ref)
AND (ls.seq IS NULL OR (p.updated >= ls.updated AND p.read_at > ls.at AND p.digest <> ls.digest));
CREATE VIEW task_current AS
SELECT s.seq, s.at, s.business, s.task_ref, s.updated, s.etag, s.digest, s.fields, s.source, s.via, s.read_at, s.role, s.run
FROM task_snapshots s
WHERE s.seq = (SELECT max(c.seq) FROM task_snapshots c
WHERE c.business = s.business AND c.task_ref = s.task_ref
AND (c.source = 'self'
OR c.read_at > coalesce((SELECT ls.at FROM task_snapshots ls
WHERE ls.business = c.business AND ls.task_ref = c.task_ref AND ls.source = 'self'
ORDER BY ls.seq DESC LIMIT 1), '')));
CREATE VIEW tasks_open AS
SELECT s.business, s.task_ref, json_extract(s.fields, '$.bucket') AS bucket, s.seq
FROM task_current s
WHERE json_type(s.fields, '$.gone') IS NULL
AND json_extract(s.fields, '$.done') = 0;
CREATE VIEW urgent_inbox AS
SELECT d.id, d.business, d.task_ref, d.question, d.at
FROM decisions d
WHERE d.class = 'gated' AND d.blocking = 1
AND NOT EXISTS (SELECT 1 FROM decision_events x WHERE x.decision = d.id AND x.op IN ('resolved','withdrawn','expired'));
+806
View File
@@ -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,
);
}
}
+21
View File
@@ -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;
}
+44
View File
@@ -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'));
}
});
});
}
}
+143
View File
@@ -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();
}
}
+39
View File
@@ -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;
}
+101
View File
@@ -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();
}
+5
View File
@@ -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';
+57
View File
@@ -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));
}
+76
View File
@@ -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;
}
}
+90
View File
@@ -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;
}
},
};
}
+175
View File
@@ -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);
}
}
+9
View File
@@ -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 }),
});
}
+356
View File
@@ -0,0 +1,356 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Store } from '../src/store.mjs';
const { Broker } = await import('../src/broker.mjs').catch((e) => {
if (e.code === 'ERR_MODULE_NOT_FOUND') return {};
throw e;
});
const options = [
{ key: 'yes', text: 'Allow' },
{ key: 'no', text: 'Decline' },
];
export const businesses = {
demo: {
id: 'demo',
human: 'jason',
arbiters: { technical: 'cto', delivery: 'pm' },
roles: {
pm: { authority: { withinRole: ['message.send', 'task.create'], crossRole: [] } },
cto: { authority: { withinRole: ['message.send', 'decision.resolve.technical'], crossRole: [] } },
coder: { authority: { withinRole: ['message.send'], crossRole: ['task.scope.change'] } },
reviewer: { authority: { withinRole: ['message.send'], crossRole: [] } },
},
},
};
function setup(t, config = businesses) {
assert.equal(typeof Broker, 'function');
const root = mkdtempSync(join(tmpdir(), 'bus-broker-'));
const store = new Store(root);
const b = new Broker({ store, businesses: config });
t.after(() => {
store.close();
rmSync(root, { recursive: true, force: true });
});
const agent = (role, run = role + '-run') =>
b.bindLaunch({ business: 'demo', role, run, harness: 'pi', address: run });
const human = b.bindHuman({ business: 'demo', human: 'jason', via: 'cli', outsideAgent: true });
return { b, store, agent, human };
}
const call = (b, cap, verb, args = {}) => b.request(cap, { verb, args });
const raise = (b, cap, action, extra = {}) =>
call(b, cap, 'decision.raise', {
action,
question: 'Allow?',
options,
recommendation: 'no',
blocking: false,
...extra,
});
test('launch identity is stamped, payload identity is refused and stale holder cannot send', (t) => {
const { b, store, agent } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
assert.throws(() => call(b, c, 'message.send', { to: 'pm', body: 'hello', role: 'pm' }), /invalid-request/);
const m = call(b, c, 'message.send', { to: 'pm', body: 'hello' });
assert.equal(store.get('SELECT from_role FROM messages WHERE id=?', m.id).from_role, 'coder');
call(b, c, 'role.release');
assert.throws(() => call(b, c, 'message.send', { to: 'pm', body: 'late' }), /not-holder/);
assert.throws(() => call(b, 'bogus', 'role.claim'), /unauthenticated/);
});
test('decision classes route from policy; gated resolution is human-only, choice and target must match', (t) => {
const { b, agent, human } = setup(t),
c = agent('coder'),
cto = agent('cto');
call(b, c, 'role.claim');
call(b, cto, 'role.claim');
const cross = raise(b, c, 'task.scope.change', { domain: 'technical', target: 'task-1' });
assert.equal(cross.route_to, 'cto');
assert.throws(() => call(b, c, 'decision.resolve', { id: cross.id, choice: 'yes' }), /not-resolver/);
call(b, cto, 'decision.resolve', { id: cross.id, choice: 'yes' });
assert.equal(
b.authorize(c, 'task.scope.change', { decision: cross.id, target: 'task-1' }).class,
'cross-role',
);
assert.throws(
() => b.authorize(c, 'task.scope.change', { decision: cross.id, target: 'task-2' }),
/decision-mismatch/,
);
const d = raise(b, c, 'deploy', { target: 'release-1' });
assert.equal(d.route_to, 'human');
assert.throws(() => call(b, cto, 'decision.resolve', { id: d.id, choice: 'yes' }), /human-required/);
assert.throws(() => call(b, human, 'decision.resolve', { id: d.id, choice: 'other' }), /invalid-choice/);
call(b, human, 'decision.resolve', { id: d.id, choice: 'no' });
assert.throws(
() => b.authorize(c, 'deploy', { decision: d.id, target: 'release-1' }),
/decision-not-approved/,
);
assert.throws(() => call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' }), /decision-closed/);
});
test('claim exclusion, holder release, gated revoke and rerouting to a new holder are atomic', (t) => {
const { b, store, agent, human } = setup(t),
c = agent('coder'),
p = agent('pm'),
p2 = agent('pm', 'pm-next');
call(b, c, 'role.claim');
const m = call(b, c, 'message.send', { to: 'pm', body: 'waiting' });
assert.equal(store.get('SELECT holder_run FROM deliveries WHERE message=?', m.id).holder_run, null);
call(b, p, 'role.claim');
assert.throws(() => call(b, p2, 'role.claim'), /role-already-held/);
const delivered = call(b, p, 'message.receive');
assert.equal(delivered[0].id, m.id);
assert.equal(call(b, p, 'message.receive').length, 0, 'receiving twice does not replay');
assert.throws(() => call(b, p2, 'role.release'), /not-holder/);
const d = raise(b, c, 'role.revoke', { target: 'pm' });
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
call(b, c, 'role.revoke', { role: 'pm', decision: d.id });
call(b, p2, 'role.claim');
assert.throws(() => call(b, p, 'message.receive'), /not-holder/);
assert.equal(call(b, p2, 'message.receive').length, 0, 'old delivered message is not replayed');
});
test('launch events require a human CLI capability; generic emit cannot forge authority events', (t) => {
const { b, agent, human, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
assert.throws(() => call(b, c, 'launch.revoke'), /human-required/);
assert.throws(() => call(b, c, 'event.emit', { kind: 'launch.revoked', body: {} }), /reserved-event/);
call(b, human, 'launch.revoke');
assert.equal(store.get('SELECT state FROM launch_state').state, 'revoked');
call(b, human, 'launch.restore');
assert.equal(store.get('SELECT state FROM launch_state').state, 'allowed');
assert.throws(
() => b.bindHuman({ business: 'demo', human: 'jason', via: 'webui', outsideAgent: true }),
/human-required/,
);
assert.throws(
() => b.bindHuman({ business: 'demo', human: 'jason', via: 'cli', outsideAgent: false }),
/human-required/,
);
});
test('within-role decisions close atomically and invalid options or blocking omissions refuse', (t) => {
const { b, agent, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
const d = raise(b, c, 'message.send', { choice: 'yes' });
assert.equal(d.class, 'within-role');
assert.equal(store.get('SELECT choice FROM decision_events WHERE decision=?', d.id).choice, 'yes');
assert.throws(() => raise(b, c, 'deploy', { blocking: true }), /task-required/);
assert.throws(
() =>
raise(b, c, 'deploy', {
options: [
{ key: 'a', text: 'a' },
{ key: 'a', text: 'b' },
],
recommendation: 'a',
}),
/invalid-options/,
);
assert.throws(() => raise(b, c, 'made.up'), /unknown-action/);
assert.equal(call(b, c, 'inbox').length, 0);
});
test('observer capabilities read human inbox but cannot mutate or forge launch identity', (t) => {
const { b, agent } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
raise(b, c, 'deploy');
const view = b.bindReader({ business: 'demo' });
assert.equal(call(b, view, 'inbox').length, 1);
for (const [verb, args] of [
['role.claim', {}],
['launch.revoke', {}],
['message.send', { to: 'pm', body: 'x' }],
])
assert.throws(() => call(b, view, verb, args), /read-only/);
});
test('task action subjects and linked decision trail are complete and ordered', (t) => {
const { b, agent, human, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
const request = call(b, human, 'message.send', { to: 'coder', body: 'make this' });
const input = store.get("SELECT id FROM events WHERE kind='human.input' AND subject=?", request.id).id;
b.recordEvent(c, {
kind: 'task.created',
subject: 'vikunja:3/1',
body: { request: input, requirement: 'REQ-TASK-1' },
});
const d = raise(b, c, 'deploy', { target: 'vikunja:3/1', task_ref: 'vikunja:3/1' });
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
b.authorize(c, 'deploy', { target: 'vikunja:3/1', decision: d.id });
const rows = call(b, human, 'trail', { subject: 'vikunja:3/1' });
assert.ok(rows.some((r) => r.id === input));
assert.ok(rows.some((r) => r.id === request.id));
assert.ok(rows.some((r) => r.decision === d.id && r.op === 'resolved'));
assert.ok(
store
.all(
"SELECT subject FROM events WHERE kind='action.allowed' AND json_extract(body,'$.action')='deploy'",
)
.every((e) => e.subject === 'vikunja:3/1'),
);
for (let i = 1; i < rows.length; i++)
assert.ok(rows[i].at > rows[i - 1].at, 'broker write order has unique timestamps');
assert.throws(
() => b.recordEvent(c, { kind: 'review.verdict', body: { task_ref: 'vikunja:3/1' } }),
/task-subject-required/,
);
});
test('launch binding is durable and reconnecting requires the identical trusted record', (t) => {
const { b, agent, store } = setup(t);
agent('coder', 'durable');
assert.equal(store.get("SELECT count(*) n FROM events WHERE kind='session.launched'").n, 1);
const again = new Broker({ store, businesses });
again.bindLaunch({ business: 'demo', role: 'coder', run: 'durable', harness: 'pi', address: 'durable' });
assert.equal(store.get("SELECT count(*) n FROM events WHERE kind='session.launched'").n, 1);
const wrong = new Broker({ store, businesses });
assert.throws(
() => wrong.bindLaunch({ business: 'demo', role: 'cto', run: 'durable', harness: 'pi' }),
/launch-record-mismatch/,
);
});
test('business isolation includes inherited object names and cross-business message references', (t) => {
const { store } = setup(t);
const config = structuredClone(businesses);
config.other = { ...structuredClone(config.demo), id: 'other' };
const b = new Broker({ store, businesses: config });
assert.throws(
() => b.bindLaunch({ business: 'demo', role: 'toString', run: 'fake', harness: 'pi' }),
/unknown-role/,
);
assert.throws(() => b.bindReader({ business: 'constructor' }), /unknown-business/);
const caps = ['demo', 'other'].map((business) =>
b.bindLaunch({ business, role: 'coder', run: 'r', harness: 'pi' }),
);
for (const c of caps) call(b, c, 'role.claim');
const m = call(b, caps[0], 'message.send', { to: 'pm', body: 'private' });
assert.throws(
() => call(b, caps[1], 'message.send', { to: 'pm', body: 'reply', in_reply_to: m.id }),
/message-not-found/,
);
const d = raise(b, caps[0], 'deploy');
assert.throws(
() => call(b, caps[1], 'decision.resolve', { id: d.id, choice: 'yes' }),
/decision-not-found/,
);
assert.deepEqual(call(b, caps[1], 'trail', { subject: d.id }), []);
});
test('authority never transfers between action, run, target, unresolved or replaced role holder', (t) => {
const { b, agent, human } = setup(t),
c = agent('coder'),
p = agent('pm');
call(b, c, 'role.claim');
call(b, p, 'role.claim');
const d = raise(b, c, 'deploy', { target: 'v1' });
assert.throws(() => b.authorize(c, 'deploy', { target: 'v1', decision: d.id }), /decision-not-approved/);
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
assert.throws(() => b.authorize(c, 'spend', { target: 'v1', decision: d.id }), /decision-mismatch/);
call(b, c, 'role.release');
const c2 = agent('coder', 'r2');
call(b, c2, 'role.claim');
assert.throws(() => b.authorize(c2, 'deploy', { target: 'v1', decision: d.id }), /decision-mismatch/);
const revoke = raise(b, c2, 'role.revoke', { target: 'pm' });
call(b, human, 'decision.resolve', { id: revoke.id, choice: 'yes' });
call(b, p, 'role.release');
const p2 = agent('pm', 'p2');
call(b, p2, 'role.claim');
assert.throws(() => call(b, c2, 'role.revoke', { role: 'pm', decision: revoke.id }), /decision-mismatch/);
});
test('task projection uses schema current view, skipping earlier and equal-start polls', (t) => {
const { b, store, human } = setup(t);
const put = (source, at, read, bucket) =>
store.run(
'INSERT INTO task_snapshots(at,business,task_ref,updated,digest,fields,source,via,read_at,role,run) VALUES(?,?,?,?,?,?,?,?,?,?,?)',
at,
'demo',
'vikunja:3/46',
'2026-10-04T00:00:00.000Z',
'a'.repeat(64),
JSON.stringify({ bucket, done: 0 }),
source,
source === 'poll' ? 'board' : null,
read,
source === 'self' ? 'coder' : null,
source === 'self' ? 'r1' : null,
);
put('self', '2026-10-04T00:00:10.000Z', null, 13);
put('poll', '2026-10-04T00:00:11.000Z', '2026-10-04T00:00:09.000Z', 11);
assert.equal(call(b, human, 'tasks')[0].fields.bucket, 13);
put('poll', '2026-10-04T00:00:12.000Z', '2026-10-04T00:00:10.000Z', 11);
assert.equal(call(b, human, 'tasks')[0].fields.bucket, 13);
put('poll', '2026-10-04T00:00:13.000Z', '2026-10-04T00:00:12.000Z', 11);
assert.equal(call(b, human, 'tasks')[0].fields.bucket, 11);
assert.equal(store.get('SELECT bucket FROM tasks_open').bucket, 11);
});
test('revocation permanently bars the old run from reclaiming first, including after broker restart', (t) => {
const { b, agent, human, store } = setup(t),
coder = agent('coder'),
old = agent('pm', 'revoked-pm');
call(b, coder, 'role.claim');
call(b, old, 'role.claim');
const d = raise(b, coder, 'role.revoke', { target: 'pm' });
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
call(b, coder, 'role.revoke', { role: 'pm', decision: d.id });
assert.throws(() => call(b, old, 'role.claim'), /run-revoked/);
const reopened = new Broker({ store, businesses });
const oldAgain = reopened.bindLaunch({
business: 'demo',
role: 'pm',
run: 'revoked-pm',
harness: 'pi',
address: 'revoked-pm',
});
assert.throws(() => call(reopened, oldAgain, 'role.claim'), /run-revoked/);
const replacement = reopened.bindLaunch({
business: 'demo',
role: 'pm',
run: 'replacement',
harness: 'pi',
});
call(reopened, replacement, 'role.claim');
});
test('empty message references refuse before storage; refusal-evidence failure stays a typed error', (t) => {
const { b, agent, store } = setup(t),
c = agent('coder');
call(b, c, 'role.claim');
for (const key of ['in_reply_to', 'corrects'])
assert.throws(() => call(b, c, 'message.send', { to: 'pm', body: 'x', [key]: '' }), /invalid-request/);
const original = store.transaction;
store.transaction = () => {
throw Error('raw private database detail');
};
try {
assert.throws(
() => call(b, c, 'unknown.verb'),
(e) => e.code === 'storage-unavailable' && e.message === 'storage-unavailable',
);
} finally {
store.transaction = original;
}
});
test('both arbiters require human resolution when their cross-role route is themselves', (t) => {
const config = structuredClone(businesses);
config.demo.roles.cto.authority.crossRole = ['task.scope.change'];
config.demo.roles.pm.authority.crossRole = ['task.priority.change'];
const { b, agent, human } = setup(t, config);
for (const [role, action, domain] of [
['cto', 'task.scope.change', 'technical'],
['pm', 'task.priority.change', 'delivery'],
]) {
const cap = agent(role);
call(b, cap, 'role.claim');
const d = raise(b, cap, action, { domain, target: 'vikunja:3/41' });
assert.equal(d.class, 'cross-role');
assert.equal(d.route_to, 'human');
assert.throws(() => call(b, cap, 'decision.resolve', { id: d.id, choice: 'yes' }), /human-required/);
assert.throws(
() => b.authorize(cap, action, { decision: d.id, target: 'vikunja:3/41' }),
/decision-not-approved/,
);
call(b, human, 'decision.resolve', { id: d.id, choice: 'yes' });
assert.equal(b.authorize(cap, action, { decision: d.id, target: 'vikunja:3/41' }).class, 'cross-role');
}
});
+28
View File
@@ -0,0 +1,28 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { busBusiness } from '../src/business.mjs';
test('S1 adapter takes resolved limits and refs, rejects mismatched instance, never mutates input', () => {
const business = {
id: 'demo',
human: 'jason',
roles: {
coder: { definition: 'coder', credentials: { gitea: { file: '/unused', rotateBy: '2099-01-01' } } },
},
};
const resolved = {
coder: {
business: 'demo',
instance: 'coder',
definition: 'coder',
limits: { authority: { withinRole: ['message.send'], crossRole: [] } },
credentials: {},
},
};
const before = structuredClone(business),
b = busBusiness(business, resolved);
assert.deepEqual(b.roles.coder.authority, resolved.coder.limits.authority);
assert.deepEqual(b.roles.coder.credentials, {});
assert.deepEqual(business, before);
resolved.coder.instance = 'pm';
assert.throws(() => busBusiness(business, resolved), /invalid-resolution/);
});
+112
View File
@@ -0,0 +1,112 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, writeFileSync, chmodSync, symlinkSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
const { Credentials } = await import('../src/credentials.mjs').catch((e) => {
if (e.code === 'ERR_MODULE_NOT_FOUND') return {};
throw e;
});
function fixture(t) {
assert.equal(typeof Credentials, 'function');
const root = mkdtempSync(join(tmpdir(), 'bus-token-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const token = 'fixture_only_0123456789abcdef';
const path = join(root, 'token');
writeFileSync(path, token + '\n', { mode: 0o600 });
return { root, path, token, ref: { file: path, rotateBy: '2099-01-01' } };
}
test('only validated broker references load; returned data and exceptions cannot expose a known token', async (t) => {
const { ref, token } = fixture(t);
const c = new Credentials({
references: { coder: { gitea: ref } },
repoRoots: [],
dataRoot: '/nonexistent-data',
});
assert.equal(
await c.use('coder', 'gitea', async (value) => ({ matches: value === token })).then((x) => x.matches),
true,
);
await assert.rejects(
c.use('coder', 'gitea', async (value) => ({ echo: value })),
/credential-leak/,
);
await assert.rejects(
c.use('coder', 'gitea', async (value) => {
throw Error(value);
}),
/service-failed/,
);
assert.throws(() => c.assertClean({ [token]: 'x' }), /credential-leak/);
assert.throws(() => c.assertClean({ nested: ['prefix ' + token] }), /credential-leak/);
c.close();
await assert.rejects(
c.use('coder', 'gitea', () => true),
/credential-unavailable/,
);
});
test('bad file modes, symlinks, repository/data paths, malformed tokens and missing dates refuse', (t) => {
const { path, root, ref } = fixture(t);
const create = (r = ref, extra = {}) =>
new Credentials({ references: { coder: { gitea: r } }, repoRoots: [], dataRoot: '/no-data', ...extra });
chmodSync(path, 0o644);
assert.throws(() => create(), /credential-file/);
chmodSync(path, 0o600);
const link = join(root, 'link');
symlinkSync(path, link);
assert.throws(() => create({ ...ref, file: link }), /credential-file/);
assert.throws(() => create(ref, { repoRoots: [root] }), /credential-location/);
assert.throws(() => create(ref, { dataRoot: root }), /credential-location/);
assert.throws(() => create({ file: path }), /credential-date/);
for (const value of ['', 'value\nsecond', 'value\r\n', 'space value', '"value"']) {
writeFileSync(path, value);
assert.throws(() => create(), /credential-format/);
}
});
test('expiry refuses use and env references never become client data', async (t) => {
fixture(t);
const c = new Credentials({
references: { coder: { vikunja: { env: 'FIXTURE_BUS_TOKEN', expires: '2000-01-01' } } },
env: { FIXTURE_BUS_TOKEN: 'fixture_expired_0123456789' },
repoRoots: [],
dataRoot: '/no-data',
});
await assert.rejects(
c.use('coder', 'vikunja', () => true),
/credential-expired/,
);
await assert.rejects(
c.use('pm', 'vikunja', () => true),
/credential-unavailable/,
);
});
test('S1 parsed service refs work, service mismatch refuses, Gitea rotation due is a warning state', async (t) => {
const { ref } = fixture(t);
const c = new Credentials({
references: { coder: { gitea: { ...ref, service: 'gitea', rotateBy: '2000-01-01' } } },
});
assert.equal(await c.use('coder', 'gitea', () => true), true);
assert.deepEqual(c.status(), [
{ instance: 'coder', service: 'gitea', date: '2000-01-01', state: 'rotation-due' },
]);
assert.throws(
() => new Credentials({ references: { coder: { gitea: { ...ref, service: 'vikunja' } } } }),
/credential-reference/,
);
});
test('opaque tokens shorter than 16 characters refuse before use', () => {
for (const value of ['a', 'a'.repeat(15), 'a'.repeat(15) + '\n'])
assert.throws(
() =>
new Credentials({
references: { coder: { gitea: { env: 'FIXTURE_TOKEN', rotateBy: '2099-01-01' } } },
env: { FIXTURE_TOKEN: value },
}),
/credential-format/,
);
const c = new Credentials({
references: { coder: { gitea: { env: 'FIXTURE_TOKEN', rotateBy: '2099-01-01' } } },
env: { FIXTURE_TOKEN: 'a'.repeat(16) },
});
c.close();
});
+11
View File
@@ -0,0 +1,11 @@
import { writeSync } from 'node:fs';
import { Store } from '../../src/store.mjs';
const store = new Store(process.argv[2]);
store.transaction(() => {
store.run(
"INSERT INTO events(id,at,business,kind,body) VALUES('uncommitted','2026-10-04T00:00:00.000Z','demo','action.allowed','{}')",
);
writeSync(1, 'inserted\n');
Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 30000);
});
store.close();
+125
View File
@@ -0,0 +1,125 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { verifyHuman, readProcess } from '../src/human.mjs';
const cli = '/opt/mosaic/human-cli.mjs';
function processTable() {
return new Map([
[
50,
{
pid: 50,
ppid: 40,
startTime: '100',
uid: process.getuid(),
argv: ['node', cli],
env: { MOSAIC_BUS_CLI_NONCE: 'a'.repeat(64) },
},
],
[40, { pid: 40, ppid: 1, startTime: '90', uid: process.getuid(), argv: ['bash'], env: {} }],
]);
}
const proof = { pid: 50, startTime: '100', nonce: 'a'.repeat(64), business: 'demo' };
function check(table, extra = {}) {
return verifyHuman(proof, {
cliPath: cli,
readProcess: (pid) => {
if (!table.has(pid)) throw Error('gone');
return table.get(pid);
},
launches: [],
...extra,
});
}
test('human proof binds CLI entry, process start and nonce; agents and incomplete ancestry refuse', () => {
const table = processTable();
assert.equal(check(table).business, 'demo');
for (const change of [
{ startTime: '101' },
{ uid: process.getuid() + 1 },
{ argv: ['node', 'other.mjs'] },
{ env: { MOSAIC_BUS_CLI_NONCE: 'x' } },
]) {
const t = processTable();
Object.assign(t.get(50), change);
assert.throws(() => check(t), /human-required/);
}
for (const env of [{ MOSAIC_RUN_ID: 'run-x' }, { CLAUDECODE: '1' }, { CODEX_THREAD_ID: 'thread' }]) {
const t = processTable();
t.get(40).env = env;
assert.throws(() => check(t), /human-required/);
}
assert.throws(() => check(table, { launches: [{ pid: 40, startTime: '90' }] }), /human-required/);
const broken = processTable();
broken.delete(40);
assert.throws(() => check(broken), /human-required/);
const loop = processTable();
loop.get(40).ppid = 50;
assert.throws(() => check(loop), /human-required/);
});
test('process reader gets own kernel identity without exposing environment values', () => {
const p = readProcess(process.pid);
assert.equal(p.pid, process.pid);
assert.equal(p.uid, process.getuid());
assert.match(p.startTime, /^[0-9]+$/);
assert.equal(p.ppid, process.ppid);
});
test('EACCES ancestor environments skip only markers; commands and registered launches still refuse', () => {
const table = processTable();
table.get(40).ppid = 30;
table.set(30, {
pid: 30,
ppid: 20,
startTime: '80',
uid: process.getuid(),
argv: ['systemd', '--user'],
envError: 'EACCES',
});
table.set(20, {
pid: 20,
ppid: 1,
startTime: '70',
uid: 0,
argv: ['plasmalogin-helper'],
envError: 'EACCES',
});
const io = {
stat: (path) => ({ uid: table.get(Number(path.split('/')[2])).uid }),
readFile: (path) => {
const r = table.get(Number(path.split('/')[2]));
if (!r) throw Object.assign(Error(), { code: 'ENOENT' });
if (path.endsWith('/environ')) {
if (r.envError) throw Object.assign(Error(), { code: r.envError });
return Object.entries(r.env ?? {})
.map(([k, v]) => k + '=' + v)
.join('\0');
}
if (path.endsWith('/cmdline')) return r.argv.join('\0');
const fields = Array(20).fill('0');
fields[0] = 'S';
fields[1] = String(r.ppid);
fields[19] = r.startTime;
return `${r.pid} (fixture) ${fields.join(' ')}`;
},
};
const read = (pid) => readProcess(pid, io);
assert.equal(check(table, { readProcess: read }).business, 'demo');
assert.throws(
() => check(table, { readProcess: read, launches: [{ pid: 30, startTime: '80' }] }),
/human-required/,
);
table.get(30).argv = ['pi'];
assert.throws(() => check(table, { readProcess: read }), /human-required/);
table.get(30).argv = ['systemd'];
table.get(30).envError = 'ENOENT';
assert.throws(() => check(table, { readProcess: read }), /human-required/);
table.get(30).envError = 'EACCES';
table.get(50).envError = 'EACCES';
assert.throws(() => check(table, { readProcess: read }), /human-required/);
});
test('real pid 1 remains inspectable when its environment is protected', () => {
const p = readProcess(1);
assert.equal(p.pid, 1);
assert.match(p.startTime, /^[0-9]+$/);
assert.ok(Array.isArray(p.argv));
if (p.uid !== process.getuid()) assert.equal(p.env, null);
});
+193
View File
@@ -0,0 +1,193 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { fork } from 'node:child_process';
import { mkdtempSync, rmSync, existsSync, writeFileSync, mkdirSync, readFileSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { once } from 'node:events';
import { Client } from '../src/client.mjs';
const businesses = {
demo: {
id: 'demo',
human: 'jason',
arbiters: { technical: 'cto', delivery: 'cto' },
roles: { cto: { authority: { withinRole: ['message.send'], crossRole: [] } } },
},
};
function launch(t) {
const child = fork(new URL('../src/process.mjs', import.meta.url), [], {
stdio: ['ignore', 'pipe', 'pipe', 'ipc'],
});
let logs = '';
child.stdout.on('data', (b) => (logs += b));
child.stderr.on('data', (b) => (logs += b));
t.after(() => {
if (child.exitCode === null) child.kill('SIGKILL');
});
return { child, logs: () => logs };
}
async function boot(child, config) {
const reply = Promise.race([
once(child, 'message'),
once(child, 'exit').then(() => {
throw Error('exited-before-reply');
}),
]);
child.send({ op: 'boot', config });
return (await reply)[0];
}
test('broker process binds trusted launches, offers reader capabilities, refuses human mutation, closes cleanly', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-process-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const { child } = launch(t);
const ready = await boot(child, {
dataRoot: root,
businesses,
launches: [{ business: 'demo', role: 'cto', run: 'r1', harness: 'pi', pid: process.pid, startTime: '1' }],
readers: ['demo'],
});
assert.equal(ready.ok, true);
const agent = new Client({ path: ready.path, cap: ready.launches[0].cap });
await agent.call('role.claim');
const reader = new Client({ path: ready.path, cap: ready.readers[0].cap });
assert.equal((await reader.call('agents'))[0].holder_run, 'r1');
await assert.rejects(reader.call('launch.revoke'), /read-only/);
await assert.rejects(
new Client({
path: ready.path,
human: { business: 'demo', pid: process.pid, startTime: '1', nonce: 'f'.repeat(64) },
}).call('launch.revoke'),
/human-required/,
);
const exit = once(child, 'exit');
child.send({ op: 'close' });
assert.equal((await exit)[0], 0);
assert.equal(existsSync(join(root, 'bus/writer.lock')), false);
});
test('startup token refusal returns safe code without value or partial listening broker', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-boot-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const data = join(root, 'data');
mkdirSync(data);
const secret = 'fixture-token-never-in-db-948723';
const file = join(root, 'secret');
writeFileSync(file, secret, { mode: 0o644 });
const config = structuredClone(businesses);
config.demo.roles.cto.credentials = { gitea: { file, rotateBy: '2099-01-01' } };
const { child, logs } = launch(t);
const exit = once(child, 'exit');
const ready = await boot(child, { dataRoot: data, businesses: config, launches: [], readers: [] });
assert.equal(ready.ok, false);
assert.equal(ready.error, 'credential-file');
assert.equal((await exit)[0], 2);
assert.ok(!logs().includes(secret));
assert.equal(existsSync(join(data, 'bus/broker.sock')), false);
});
test('loaded fixture token is absent from socket replies and SQLite, including refusal evidence', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-secret-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const data = join(root, 'data');
mkdirSync(data);
const token = 'fixture-opaque-token-e9c39140';
const file = join(root, 'token');
writeFileSync(file, token, { mode: 0o600 });
const config = structuredClone(businesses);
config.demo.roles.cto.credentials = { gitea: { file, rotateBy: '2099-01-01' } };
const { child, logs } = launch(t);
const ready = await boot(child, {
dataRoot: data,
businesses: config,
launches: [{ business: 'demo', role: 'cto', run: 'r2', harness: 'pi', pid: process.pid, startTime: '1' }],
});
assert.equal(ready.ok, true);
const c = new Client({ path: ready.path, cap: ready.launches[0].cap });
await c.call('role.claim');
await assert.rejects(c.call('message.send', { to: 'cto', body: token }), /credential-leak/);
const exit = once(child, 'exit');
child.send({ op: 'close' });
await exit;
assert.ok(!readFileSync(join(data, 'bus/bus.sqlite')).includes(Buffer.from(token)));
assert.ok(!logs().includes(token));
});
test('killed broker leaves an explicit stale lock; another process cannot silently reclaim it', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-crash-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const first = launch(t);
const config = { dataRoot: root, businesses, launches: [] };
assert.equal((await boot(first.child, config)).ok, true);
const killed = once(first.child, 'exit');
first.child.kill('SIGKILL');
await killed;
assert.equal(existsSync(join(root, 'bus/writer.lock')), true);
const second = launch(t),
ended = once(second.child, 'exit');
const refusal = await boot(second.child, config);
assert.equal(refusal.ok, false);
assert.equal((await ended)[0], 2);
assert.equal(existsSync(join(root, 'bus/writer.lock')), true);
});
test('trusted host registers later launches; socket clients never have a registration verb', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-add-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const { child } = launch(t);
const ready = await boot(child, { dataRoot: root, businesses, launches: [] });
assert.equal(ready.ok, true);
const message = once(child, 'message');
child.send({
op: 'bindLaunch',
record: { business: 'demo', role: 'cto', run: 'later', harness: 'pi', pid: process.pid, startTime: '1' },
});
const bound = (await message)[0];
assert.equal(bound.ok, true);
const c = new Client({ path: ready.path, cap: bound.launch.cap });
await c.call('role.claim');
await assert.rejects(c.call('bindLaunch', { role: 'cto' }), /unknown-verb/);
const exit = once(child, 'exit');
child.send({ op: 'close' });
assert.equal((await exit)[0], 0);
});
test('runtime excludes declared project roots even when host supplies no repoRoots', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-project-token-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const data = join(root, 'data');
mkdirSync(data);
const project = join(root, 'project');
mkdirSync(project);
const file = join(project, 'token');
writeFileSync(file, 'fixture-do-not-load-from-project', { mode: 0o600 });
const config = structuredClone(businesses);
config.demo.projects = { stack: { root: project } };
config.demo.roles.cto.credentials = { gitea: { file, rotateBy: '2099-01-01' } };
const { child } = launch(t),
exit = once(child, 'exit');
const result = await boot(child, { dataRoot: data, businesses: config });
assert.equal(result.ok, false);
assert.equal(result.error, 'credential-location');
assert.equal((await exit)[0], 2);
});
test('a refused launch binding leaves the broker and existing capabilities alive; bad protocol stops it', async (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-bind-refusal-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
const { child } = launch(t);
const record = {
business: 'demo',
role: 'cto',
run: 'one',
harness: 'pi',
pid: process.pid,
startTime: '1',
};
const ready = await boot(child, { dataRoot: root, businesses, launches: [record] });
assert.equal(ready.ok, true);
const client = new Client({ path: ready.path, cap: ready.launches[0].cap });
await client.call('role.claim');
const reply = once(child, 'message');
child.send({ op: 'bindLaunch', record });
assert.deepEqual((await reply)[0], { ok: false, error: 'duplicate-run' });
assert.equal((await client.call('agents'))[0].holder_run, 'one');
const bad = once(child, 'message'),
exit = once(child, 'exit');
child.send({ op: 'not-a-protocol-verb' });
assert.equal((await bad)[0].ok, false);
assert.equal((await exit)[0], 2);
});
+522
View File
@@ -0,0 +1,522 @@
import test from 'node:test';
import assert from 'node:assert/strict';
// Slice 1 prototype, v3a schema (v3 plus the task event rules of lead decision 56, Q3). Same pattern as proto-v3.mjs.
import { DatabaseSync } from 'node:sqlite';
import { readFileSync, mkdtempSync, rmSync } from 'node:fs';
import { join } from 'node:path';
import { tmpdir } from 'node:os';
import { createHash } from 'node:crypto';
const root = mkdtempSync(join(tmpdir(), 'bus-schema-'));
const f = join(root, 'bus.sqlite');
test('v3b prototype refusals, views and append-only mutations', (t) => {
t.after(() => {
db.close();
rmSync(root, { recursive: true, force: true });
});
let db = new DatabaseSync(f, { timeout: 5000 });
db.exec(readFileSync(new URL('../schema.sql', import.meta.url), 'utf8'));
let tick = 0;
const now = () => new Date(Date.UTC(2026, 9, 4, 12, 0, tick++)).toISOString();
const refused = new Set([
'raise without blocking',
'raise with blocking 2',
'raise blocking without task_ref',
'unknown kind task.deleted',
'unknown kind task.deleted with a subject',
'credential.changed without instance',
'credential.expiring service github',
'task.missing without reason',
'task.missing reason deleted',
'task.missing without subject',
'task.missing moved without project',
'task.state without subject',
'task.state subject PROJ-41',
'task.state subject vikunja:3/41x',
'task.state subject vikunja:03/41',
'task.state subject vikunja:3/4/1',
'task.created without request',
'task.created request names an unknown event',
'task.created request names a credential event',
'task.created request from another business',
'task.created request as a number',
'task.created without requirement',
'task.created requirement REQ-task-1',
'task.created requirement REQ-TASK-0',
'task.created without subject',
'decision with task_ref vikunja:3/41x',
'launch.revoked by pm run',
'self without role and run',
'poll with a role',
'poll without via',
'poll without read_at',
'self with via',
'unknown via webhook',
'fields without bucket',
'bucket as text',
'gone on a board read',
'gone on self',
'bad task_ref',
'task_ref vikunja:3/41x',
'bad digest',
]);
const tested = new Set();
const tryit = (label, fn) => {
tested.add(label);
if (refused.has(label) || / (UPDATE|DELETE|INSERT OR REPLACE)$/.test(label))
assert.throws(fn, undefined, label);
else assert.doesNotThrow(fn, label);
};
const expected = {
urgent_inbox: ['d-3'],
'urgent_inbox after resolve': [],
'external after cursor X': [],
'external after cursor Y': ['vikunja:3/41'],
'external after self Z, stale board read': [],
'external after stale cursor read': [],
'external after board agrees': [],
"external after person's move": ['vikunja:3/41'],
tasks_open: ['vikunja:3/41', 'vikunja:3/42'],
'tasks_open after tombstone': ['vikunja:3/41'],
};
const show = (label, sql) => {
const rows = db.prepare(sql).all();
if (Object.hasOwn(expected, label))
assert.deepEqual(
rows.map((r) => r.id ?? r.task_ref),
expected[label],
label,
);
};
const hex = (s) => createHash('sha256').update(s).digest('hex');
const schemaDigest = (d) =>
hex(
d
.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'),
);
db.prepare("INSERT INTO meta (key, value) VALUES ('schema_digest', ?)").run(schemaDigest(db));
const check = () =>
db.prepare("SELECT value FROM meta WHERE key = 'schema_digest'").get().value === schemaDigest(db)
? 'match'
: 'MISMATCH';
assert.equal(check(), 'match');
const dec = db.prepare(
'INSERT INTO decisions (id,at,business,raised_by_role,raised_by_run,class,action,route_to,question,options,recommendation,task_ref,blocking) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)',
);
const opts = JSON.stringify([
{ key: 'A', text: 'rotate' },
{ key: 'B', text: 'wait' },
]);
tryit('raise without blocking', () =>
db.exec(
`INSERT INTO decisions (id,at,business,raised_by_role,raised_by_run,class,action,route_to,question,options,recommendation) VALUES ('d-0','${now()}','mosaic-stack','coder','run-C','gated','credential.mint','human','?','${opts}','A')`,
),
);
tryit('raise with blocking 2', () =>
dec.run(
'd-1',
now(),
'mosaic-stack',
'coder',
'run-C',
'gated',
'credential.mint',
'human',
'Rotate?',
opts,
'A',
'vikunja:3/41',
2,
),
);
tryit('raise blocking without task_ref', () =>
dec.run(
'd-2',
now(),
'mosaic-stack',
'coder',
'run-C',
'gated',
'credential.mint',
'human',
'Rotate?',
opts,
'A',
null,
1,
),
);
tryit('raise blocking gated with task_ref', () =>
dec.run(
'd-3',
now(),
'mosaic-stack',
'coder',
'run-C',
'gated',
'credential.mint',
'human',
'Rotate coder vikunja token?',
opts,
'A',
'vikunja:3/41',
1,
),
);
tryit('raise non-blocking gated', () =>
dec.run(
'd-4',
now(),
'mosaic-stack',
'pm',
'run-P',
'gated',
'deploy',
'human',
'Deploy?',
opts,
'B',
null,
0,
),
);
tryit('raise blocking cross-role', () =>
dec.run(
'd-5',
now(),
'mosaic-stack',
'coder',
'run-C',
'cross-role',
'task.scope.change',
'pm',
'Widen scope?',
opts,
'B',
'vikunja:3/41',
1,
),
);
show('urgent_inbox', 'SELECT id, task_ref FROM urgent_inbox');
tryit('resolve d-3 with A', () =>
db
.prepare('INSERT INTO decision_events (decision,at,op,by,choice,via) VALUES (?,?,?,?,?,?)')
.run('d-3', now(), 'resolved', 'jason', 'A', 'cli'),
);
show('urgent_inbox after resolve', 'SELECT id FROM urgent_inbox');
const ev = db.prepare(
'INSERT INTO events (id,at,business,kind,actor_role,actor_run,subject,body) VALUES (?,?,?,?,?,?,?,?)',
);
let n = 0;
const e = (kind, role, run, body, subject = null) =>
ev.run(`e-${++n}`, now(), 'mosaic-stack', kind, role, run, subject, JSON.stringify(body));
tryit('unknown kind task.deleted', () => e('task.deleted', 'pm', 'run-P', {}));
tryit('unknown kind task.deleted with a subject', () =>
e('task.deleted', 'pm', 'run-P', {}, 'vikunja:3/41'),
);
tryit('credential.expiring vikunja coder', () =>
e('credential.expiring', null, null, { service: 'vikunja', instance: 'coder', expires: '2026-10-11' }),
);
tryit('credential.expired vikunja coder', () =>
e('credential.expired', null, null, { service: 'vikunja', instance: 'coder', decision: 'd-3' }),
);
tryit('credential.changed gitea pm', () =>
e('credential.changed', null, null, { service: 'gitea', instance: 'pm', stat: { inode: 1, size: 41 } }),
);
tryit('credential.changed without instance', () =>
e('credential.changed', null, null, { service: 'gitea' }),
);
tryit('credential.expiring service github', () =>
e('credential.expiring', null, null, { service: 'github', instance: 'pm' }),
);
tryit('task.missing without reason', () =>
e('task.missing', null, null, { reconcile: 'r-1' }, 'vikunja:3/40'),
);
tryit('task.missing reason deleted', () =>
e('task.missing', null, null, { reason: 'deleted' }, 'vikunja:3/40'),
);
tryit('task.missing without subject', () => e('task.missing', null, null, { reason: 'not-found' }));
tryit('task.missing moved without project', () =>
e('task.missing', null, null, { reason: 'moved' }, 'vikunja:3/40'),
);
tryit('task.missing not-found', () =>
e('task.missing', null, null, { reason: 'not-found' }, 'vikunja:3/40'),
);
tryit('task.missing moved to project 9', () =>
e('task.missing', null, null, { reason: 'moved', project: 9 }, 'vikunja:3/43'),
);
tryit('task.state without subject', () => e('task.state', 'coder', 'run-C', { bucket: 12 }));
tryit('task.state subject PROJ-41', () => e('task.state', 'coder', 'run-C', { bucket: 12 }, 'PROJ-41'));
tryit('task.state subject vikunja:3/41x', () =>
e('task.state', 'coder', 'run-C', { bucket: 12 }, 'vikunja:3/41x'),
);
tryit('task.state subject vikunja:03/41', () =>
e('task.state', 'coder', 'run-C', { bucket: 12 }, 'vikunja:03/41'),
);
tryit('task.state subject vikunja:3/4/1', () =>
e('task.state', 'coder', 'run-C', { bucket: 12 }, 'vikunja:3/4/1'),
);
tryit('task.state subject vikunja:3/41', () =>
e('task.state', 'coder', 'run-C', { bucket: 12 }, 'vikunja:3/41'),
);
tryit('human.input from the cli', () =>
ev.run(
'h-1',
now(),
'mosaic-stack',
'human.input',
null,
null,
null,
JSON.stringify({ via: 'cli', text: 'Add the broker push.' }),
),
);
tryit('human.input in another business', () =>
ev.run('h-2', now(), 'other', 'human.input', null, null, null, JSON.stringify({ via: 'cli', text: 'x' })),
);
const tc = (body, subject = 'vikunja:3/50') => e('task.created', 'pm', 'run-P', body, subject);
tryit('task.created without request', () => tc({ requirement: 'REQ-TASK-1' }));
tryit('task.created request names an unknown event', () =>
tc({ request: 'h-9', requirement: 'REQ-TASK-1' }),
);
tryit('task.created request names a credential event', () =>
tc({ request: 'e-3', requirement: 'REQ-TASK-1' }),
);
tryit('task.created request from another business', () =>
tc({ request: 'h-2', requirement: 'REQ-TASK-1' }),
);
tryit('task.created request as a number', () => tc({ request: 1, requirement: 'REQ-TASK-1' }));
tryit('task.created without requirement', () => tc({ request: 'h-1' }));
tryit('task.created requirement REQ-task-1', () => tc({ request: 'h-1', requirement: 'REQ-task-1' }));
tryit('task.created requirement REQ-TASK-0', () => tc({ request: 'h-1', requirement: 'REQ-TASK-0' }));
tryit('task.created without subject', () => tc({ request: 'h-1', requirement: 'REQ-TASK-1' }, null));
tryit('task.created cites h-1 and REQ-TASK-1', () => tc({ request: 'h-1', requirement: 'REQ-TASK-1' }));
tryit('decision with task_ref vikunja:3/41x', () =>
dec.run(
'd-6',
now(),
'mosaic-stack',
'coder',
'run-C',
'gated',
'deploy',
'human',
'?',
opts,
'A',
'vikunja:3/41x',
0,
),
);
show('e-3 is', "SELECT kind FROM events WHERE id = 'e-3'");
show(
'trail from h-1',
"SELECT e.kind, e.subject FROM events e WHERE json_extract(e.body, '$.request') = 'h-1'",
);
tryit('digest.sent', () => e('digest.sent', null, null, { decisions: ['d-4'], transport: 'discord-dm' }));
tryit('launch.revoked by pm run', () => e('launch.revoked', 'pm', 'run-P', {}));
tryit('launch.revoked by human', () => e('launch.revoked', null, null, { via: 'cli' }));
show('launch_state', 'SELECT business, state FROM launch_state');
tryit('launch.restored by human', () => e('launch.restored', null, null, { via: 'cli' }));
show('launch_state', 'SELECT business, state FROM launch_state');
const snap = db.prepare(
'INSERT INTO task_snapshots (at,business,task_ref,updated,etag,digest,fields,source,via,read_at,role,run) VALUES (?,?,?,?,?,?,?,?,?,?,?,?)',
);
const at = (sec) => new Date(Date.UTC(2026, 9, 4, 13, 0, sec)).toISOString();
const self = (ref, updated, fields, sec, role, run) =>
snap.run(
at(sec),
'mosaic-stack',
ref,
updated,
`"${hex(JSON.stringify(fields)).slice(0, 8)}"`,
hex(JSON.stringify(fields)),
JSON.stringify(fields),
'self',
null,
null,
role,
run,
);
const poll = (ref, updated, fields, via, readSec) =>
snap.run(
at(readSec + 1),
'mosaic-stack',
ref,
updated,
null,
hex(JSON.stringify(fields)),
JSON.stringify(fields),
'poll',
via,
at(readSec),
null,
null,
);
const raw = (source, via, readAt, role, run, fields) =>
snap.run(
at(59),
'mosaic-stack',
'vikunja:3/49',
'2026-10-04T12:00:00Z',
null,
hex(JSON.stringify(fields)),
JSON.stringify(fields),
source,
via,
readAt,
role,
run,
);
// Buckets: 11 todo, 12 in-progress, 13 in-review, 14 blocked, 15 done.
const X = { title: 'Add broker push', bucket: 12, done: 0 },
Y = { ...X, title: 'Add broker push (Jason edit)' },
Z = { ...Y, bucket: 13 },
B = { ...Z, bucket: 14 },
W = { title: 'Add broker push', bucket: 11, done: 0 };
const U = '2026-10-04T12:00:00Z';
tryit('self without role and run', () => raw('self', null, null, null, null, X));
tryit('poll with a role', () => raw('poll', 'board', at(58), 'pm', 'run-P', X));
tryit('poll without via', () => raw('poll', null, at(58), null, null, X));
tryit('poll without read_at', () => raw('poll', 'board', null, null, null, X));
tryit('self with via', () => raw('self', 'board', at(58), 'coder', 'run-C', X));
tryit('unknown via webhook', () => raw('poll', 'webhook', at(58), null, null, X));
tryit('fields without bucket', () => raw('poll', 'cursor', at(58), null, null, { title: 'x', done: 0 }));
tryit('bucket as text', () =>
raw('poll', 'cursor', at(58), null, null, { title: 'x', bucket: 'in-progress', done: 0 }),
);
tryit('gone on a board read', () => raw('poll', 'board', at(58), null, null, { gone: 'not-found' }));
tryit('gone on self', () => raw('self', null, null, 'pm', 'run-P', { gone: 'not-found' }));
tryit('bad task_ref', () =>
snap.run(
at(59),
'mosaic-stack',
'PROJ-41',
U,
null,
hex('x'),
JSON.stringify(X),
'poll',
'board',
at(58),
null,
null,
),
);
tryit('task_ref vikunja:3/41x', () =>
snap.run(
at(59),
'mosaic-stack',
'vikunja:3/41x',
U,
null,
hex('x'),
JSON.stringify(X),
'poll',
'board',
at(58),
null,
null,
),
);
tryit('bad digest', () =>
snap.run(
at(59),
'mosaic-stack',
'vikunja:3/41',
U,
null,
'abc',
JSON.stringify(X),
'poll',
'board',
at(58),
null,
null,
),
);
tryit('self X by coder, response :02', () => self('vikunja:3/41', U, X, 2, 'coder', 'run-C'));
tryit('cursor X sent :04 (unchanged)', () => poll('vikunja:3/41', U, X, 'cursor', 4));
show('external after cursor X', 'SELECT task_ref FROM task_external_changes');
tryit('cursor Y sent :06, same second (person edit)', () => poll('vikunja:3/41', U, Y, 'cursor', 6));
show('external after cursor Y', 'SELECT task_ref, via FROM task_external_changes');
tryit('self Z by coder (move to in-review), response :10', () =>
self('vikunja:3/41', U, Z, 10, 'coder', 'run-C'),
);
tryit('stale board read sent :09 shows Y', () => poll('vikunja:3/41', U, Y, 'board', 9));
show('external after self Z, stale board read', 'SELECT task_ref FROM task_external_changes');
tryit('stale cursor read, updated 11:59:00', () =>
poll('vikunja:3/41', '2026-10-04T11:59:00Z', W, 'cursor', 20),
);
show('external after stale cursor read', 'SELECT task_ref FROM task_external_changes');
tryit('board read sent :30 agrees with Z', () => poll('vikunja:3/41', U, Z, 'board', 30));
show('external after board agrees', 'SELECT task_ref FROM task_external_changes');
tryit('person moves to blocked, updated unchanged, board sent :40', () =>
poll('vikunja:3/41', U, B, 'board', 40),
);
show("external after person's move", 'SELECT task_ref, via FROM task_external_changes');
tryit('board read on vikunja:3/42 with no self row', () =>
poll('vikunja:3/42', '2026-10-04T12:01:00Z', W, 'board', 41),
);
tryit('cursor read on vikunja:3/44, done', () =>
poll('vikunja:3/44', '2026-10-04T12:01:00Z', { ...W, bucket: 15, done: 1 }, 'cursor', 41),
);
show('tasks_open', 'SELECT task_ref, bucket FROM tasks_open ORDER BY task_ref');
tryit('tombstone for vikunja:3/42 (GET 404)', () =>
poll('vikunja:3/42', '2026-10-04T12:01:00Z', { gone: 'not-found' }, 'task', 45),
);
show('tasks_open after tombstone', 'SELECT task_ref, bucket FROM tasks_open ORDER BY task_ref');
show('external, all', 'SELECT task_ref, via FROM task_external_changes ORDER BY task_ref');
const keys = {
meta: "key = 'schema_digest'",
events: "id = 'h-1'",
role_claims: '1',
decisions: "id = 'd-3'",
decision_events: '1',
messages: '1',
deliveries: '1',
task_snapshots: 'seq = 1',
};
db.exec(
"INSERT INTO role_claims (at,business,role,op,holder_run,harness,by) VALUES ('x','mosaic-stack','pm','claim','run-P','pi','run-P')",
);
db.exec(
"INSERT INTO messages (id,at,business,from_role,from_run,to_role,class,decision,body) VALUES ('m-1','x','mosaic-stack','pm','run-P','human','RESULT','d-3','rotate')",
);
db.exec("INSERT INTO deliveries (message,at,op,transport) VALUES ('m-1','x','delivered','discord-dm')");
for (const [tbl, where] of Object.entries(keys)) {
const row = db.prepare(`SELECT * FROM ${tbl} WHERE ${where} LIMIT 1`).get();
const cols = Object.keys(row);
const col = cols.find((c) => !['seq', 'id', 'key'].includes(c));
tryit(`${tbl} UPDATE`, () => db.exec(`UPDATE ${tbl} SET ${col} = ${col} WHERE ${where}`));
tryit(`${tbl} DELETE`, () => db.exec(`DELETE FROM ${tbl} WHERE ${where}`));
tryit(`${tbl} INSERT OR REPLACE`, () =>
db
.prepare(
`INSERT OR REPLACE INTO ${tbl} (${cols.join(',')}) VALUES (${cols.map(() => '?').join(',')})`,
)
.run(...cols.map((c) => row[c])),
);
}
db.close();
db = new DatabaseSync(f, { timeout: 5000 });
assert.equal(check(), 'match');
db.exec('DROP TRIGGER task_snapshots_no_update');
db.close();
db = new DatabaseSync(f, { timeout: 5000 });
assert.equal(check(), 'MISMATCH');
const count = (type) => db.prepare('SELECT count(*) n FROM sqlite_master WHERE type = ?').get(type).n;
for (const label of refused) assert.ok(tested.has(label), label);
assert.equal(count('table') - 1, 8);
assert.equal(count('trigger'), 35);
assert.equal(count('view'), 5);
});
+159
View File
@@ -0,0 +1,159 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, rmSync, statSync, symlinkSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { DatabaseSync } from 'node:sqlite';
const { Store } = await import('../src/store.mjs').catch((e) => {
if (e.code === 'ERR_MODULE_NOT_FOUND') return {};
throw e;
});
const scratch = (t) => {
const root = mkdtempSync(join(tmpdir(), 'bus-store-'));
t.after(() => rmSync(root, { recursive: true, force: true }));
return root;
};
test('creates private WAL store and excludes a second writer until explicit close', (t) => {
assert.equal(typeof Store, 'function', 'Store implementation exists');
const root = scratch(t);
const s = new Store(root);
t.after(() => s.close());
assert.equal(statSync(join(root, 'bus')).mode & 0o777, 0o700);
assert.equal(statSync(join(root, 'bus/bus.sqlite')).mode & 0o777, 0o600);
assert.equal(s.get('PRAGMA journal_mode').journal_mode, 'wal');
assert.throws(() => new Store(root), /writer-locked/);
s.close();
const next = new Store(root);
next.close();
});
test('rollback is atomic and schema metadata is checked against trusted DDL, not just itself', (t) => {
assert.equal(typeof Store, 'function');
const root = scratch(t);
const s = new Store(root);
assert.throws(
() =>
s.transaction(() => {
s.run("INSERT INTO events(id,at,business,kind,body) VALUES('e','x','b','action.allowed','{}')");
throw Error('rollback');
}),
/rollback/,
);
assert.equal(s.get('SELECT count(*) n FROM events').n, 0);
s.close();
const db = new DatabaseSync(join(root, 'bus/bus.sqlite'));
db.exec('DROP TRIGGER events_no_update');
db.close();
assert.throws(() => new Store(root), /schema-mismatch/);
});
test('existing empty database and symlink runtime directory refuse, never initialize over damage', (t) => {
assert.equal(typeof Store, 'function');
const root = scratch(t);
const s = new Store(root);
s.close();
const db = new DatabaseSync(join(root, 'bus/bus.sqlite'));
db.exec('DROP TRIGGER meta_no_delete; DELETE FROM meta');
db.close();
assert.throws(() => new Store(root), /schema-mismatch/);
const other = scratch(t);
symlinkSync(join(root, 'bus'), join(other, 'bus'));
assert.throws(() => new Store(other), /unsafe-path/);
});
test('crash during a transaction recovers no partial event after explicit fixture-only lock removal', async (t) => {
const { spawn } = await import('node:child_process');
const { once } = await import('node:events');
const { unlinkSync } = await import('node:fs');
const root = scratch(t);
const child = spawn(
process.execPath,
[new URL('./fixtures/crash-writer.mjs', import.meta.url).pathname, root],
{ stdio: ['ignore', 'pipe', 'pipe'] },
);
t.after(() => {
if (child.exitCode === null) child.kill('SIGKILL');
});
const ready = await Promise.race([
once(child.stdout, 'data'),
once(child, 'exit').then(() => {
throw Error('crash fixture exited early');
}),
]);
assert.match(String(ready[0]), /inserted/);
const exit = once(child, 'exit');
child.kill('SIGKILL');
await exit;
assert.throws(() => new Store(root), /writer-locked/);
unlinkSync(join(root, 'bus/writer.lock')); // Known dead test child, never production recovery.
const recovered = new Store(root);
try {
assert.equal(recovered.get("SELECT count(*) n FROM events WHERE id='uncommitted'").n, 0);
} finally {
recovered.close();
}
});
test('writer refuses mixed at/read_at forms atomically, even through trusted SQL helpers', (t) => {
const root = scratch(t),
store = new Store(root);
t.after(() => store.close());
const insert = (at) =>
store.run(
'INSERT INTO events(id,at,business,kind,body) VALUES(?,?,?,?,?)',
'time-' + at,
at,
'demo',
'action.allowed',
'{}',
);
for (const at of [
'2026-10-04T00:00:00Z',
'2026-10-04T01:00:00.000+01:00',
'2026-10-04T00:00:00.000000Z',
'2026-02-30T00:00:00.000Z',
])
assert.throws(() => insert(at), /invalid-timestamp/);
assert.equal(store.get('SELECT count(*) n FROM events').n, 0);
insert('2026-10-04T00:00:00.000Z');
assert.throws(
() =>
store.run(
"INSERT INTO events(seq,id,at,business,kind,body) VALUES(-1,'low','2026-10-04T00:00:00Z','demo','action.allowed','{}')",
),
/invalid-timestamp/,
);
assert.throws(
() =>
store.transaction(() => {
insert('2026-10-04T00:00:01.000Z');
store.run(
'INSERT INTO task_snapshots(at,business,task_ref,updated,digest,fields,source,via,read_at) VALUES(?,?,?,?,?,?,?,?,?)',
'2026-10-04T00:00:02.000Z',
'demo',
'vikunja:1/1',
'external',
'a'.repeat(64),
'{"bucket":1}',
'poll',
'board',
'2026-10-04T00:00:01Z',
);
}),
/invalid-timestamp/,
);
assert.equal(store.get('SELECT count(*) n FROM events').n, 1, 'whole transaction rolls back');
assert.equal(store.get('SELECT count(*) n FROM task_snapshots').n, 0);
});
test('async transactions refuse before invoking their function', async (t) => {
const root = scratch(t),
store = new Store(root);
t.after(() => store.close());
let ran = false;
assert.throws(
() =>
store.transaction(async () => {
ran = true;
await Promise.resolve();
}),
/async-transaction-refused/,
);
await Promise.resolve();
assert.equal(ran, false);
});
+119
View File
@@ -0,0 +1,119 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, rmSync, lstatSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { connect } from 'node:net';
import { Store } from '../src/store.mjs';
import { Broker } from '../src/broker.mjs';
import { serve } from '../src/server.mjs';
import { Client } from '../src/client.mjs';
import { views } from '../src/views.mjs';
const businesses = {
demo: {
id: 'demo',
human: 'jason',
arbiters: { technical: 'cto', delivery: 'cto' },
roles: { cto: { authority: { withinRole: ['message.send'], crossRole: [] } } },
},
};
async function setup(t) {
const root = mkdtempSync(join(tmpdir(), 'bus-wire-')),
store = new Store(root),
broker = new Broker({ store, businesses });
const path = join(store.directory, 'broker.sock');
const server = await serve({ broker, path });
t.after(async () => {
await server.close();
store.close();
rmSync(root, { recursive: true, force: true });
});
return { store, broker, path };
}
test('socket capability stamps launch identity; shared views use wire, no SQL client', async (t) => {
const { broker, path } = await setup(t),
cap = broker.bindLaunch({ business: 'demo', role: 'cto', run: 'one', harness: 'pi' });
assert.equal(lstatSync(path).mode & 0o777, 0o600);
const c = new Client({ path, cap });
await c.call('role.claim');
assert.equal((await views(c).agents())[0].holder_run, 'one');
await assert.rejects(new Client({ path, cap: 'forged' }).call('agents'), /unauthenticated/);
await assert.rejects(c.call('decision.resolve', { id: 'none', role: 'human' }), /invalid-request/);
await assert.rejects(c.call('sql', { query: 'SELECT * FROM meta' }), /unknown-verb/);
});
test('two wire claims serialize; a lost reply never automatically retries', async (t) => {
const { broker, path, store } = await setup(t),
caps = ['a', 'b'].map((run) => broker.bindLaunch({ business: 'demo', role: 'cto', run, harness: 'pi' }));
const outcomes = await Promise.allSettled(caps.map((cap) => new Client({ path, cap }).call('role.claim')));
assert.equal(outcomes.filter((x) => x.status === 'fulfilled').length, 1);
assert.equal(store.get("SELECT count(*) n FROM role_claims WHERE op='claim'").n, 1);
const cap = caps[outcomes.findIndex((x) => x.status === 'fulfilled')];
await new Client({ path, cap }).call('message.send', { to: 'cto', body: 'once' });
assert.equal((await new Client({ path, cap }).call('message.receive')).length, 1);
assert.deepEqual(await new Client({ path, cap }).call('message.receive'), []);
});
test('malformed, oversized and identity-forging envelopes refuse without echoing input', async (t) => {
const { path } = await setup(t);
async function raw(data) {
return new Promise((resolve, reject) => {
const s = connect(path);
let out = '';
s.on('connect', () => s.end(data));
s.on('data', (b) => (out += b));
s.on('end', () => resolve(out));
s.on('error', reject);
});
}
for (const input of [
'not-json\n',
JSON.stringify({ cap: 'x', verb: 'agents', role: 'human' }) + '\n',
'x'.repeat(70000) + '\n',
]) {
const out = await raw(input);
assert.equal(JSON.parse(out).ok, false);
assert.ok(!out.includes('not-json'));
assert.ok(out.length < 150);
}
});
test('client preserves UTF-8 when a response divides a multibyte character', async (t) => {
const { createServer } = await import('node:net');
const root = mkdtempSync(join(tmpdir(), 'bus-utf8-')),
path = join(root, 's');
const server = createServer((s) =>
s.once('data', () => {
const bytes = Buffer.from('{"ok":true,"result":"🙂"}\n');
const i = bytes.indexOf(Buffer.from('🙂')) + 1;
s.write(bytes.subarray(0, i));
setTimeout(() => s.end(bytes.subarray(i)), 10);
}),
);
await new Promise((resolve) => server.listen(path, resolve));
t.after(async () => {
await new Promise((resolve) => server.close(resolve));
rmSync(root, { recursive: true, force: true });
});
assert.equal(await new Client({ path, cap: 'fixture' }).call('agents'), '🙂');
});
test('committed mutation followed by dropped reply reports unknown and is never retried', async (t) => {
const { broker, store, path } = await setup(t);
const cap = broker.bindLaunch({ business: 'demo', role: 'cto', run: 'drop', harness: 'pi' });
broker.request(cap, { verb: 'role.claim' });
const { createServer } = await import('node:net');
let calls = 0;
const proxy = createServer((s) =>
s.once('data', (data) => {
calls++;
const r = JSON.parse(data);
broker.request(r.cap, { verb: r.verb, args: r.args });
s.destroy();
}),
);
await new Promise((resolve) => proxy.listen(path + '.drop', resolve));
t.after(() => new Promise((resolve) => proxy.close(resolve)));
await assert.rejects(
new Client({ path: path + '.drop', cap }).call('message.send', { to: 'cto', body: 'only once' }),
/outcome-unknown/,
);
assert.equal(calls, 1);
assert.equal(store.get('SELECT count(*) n FROM messages').n, 1);
});