From 38828a2cb35d47bd207337b5857a4c3b62b6e533 Mon Sep 17 00:00:00 2001 From: Jason Woltje Date: Mon, 5 Oct 2026 17:15:53 -0500 Subject: [PATCH] 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 --- packages/bus/README.md | 228 ++++++ packages/bus/package.json | 16 + packages/bus/schema.sql | 211 +++++ packages/bus/src/broker.mjs | 806 +++++++++++++++++++ packages/bus/src/business.mjs | 21 + packages/bus/src/client.mjs | 44 + packages/bus/src/credentials.mjs | 143 ++++ packages/bus/src/human-cli.mjs | 39 + packages/bus/src/human.mjs | 101 +++ packages/bus/src/index.mjs | 5 + packages/bus/src/process.mjs | 57 ++ packages/bus/src/runtime.mjs | 76 ++ packages/bus/src/server.mjs | 90 +++ packages/bus/src/store.mjs | 175 ++++ packages/bus/src/views.mjs | 9 + packages/bus/tests/broker.test.mjs | 356 ++++++++ packages/bus/tests/business.test.mjs | 28 + packages/bus/tests/credentials.test.mjs | 112 +++ packages/bus/tests/fixtures/crash-writer.mjs | 11 + packages/bus/tests/human.test.mjs | 125 +++ packages/bus/tests/process.test.mjs | 193 +++++ packages/bus/tests/schema.test.mjs | 522 ++++++++++++ packages/bus/tests/store.test.mjs | 159 ++++ packages/bus/tests/transport.test.mjs | 119 +++ 24 files changed, 3646 insertions(+) create mode 100644 packages/bus/README.md create mode 100644 packages/bus/package.json create mode 100644 packages/bus/schema.sql create mode 100644 packages/bus/src/broker.mjs create mode 100644 packages/bus/src/business.mjs create mode 100644 packages/bus/src/client.mjs create mode 100644 packages/bus/src/credentials.mjs create mode 100644 packages/bus/src/human-cli.mjs create mode 100644 packages/bus/src/human.mjs create mode 100644 packages/bus/src/index.mjs create mode 100644 packages/bus/src/process.mjs create mode 100644 packages/bus/src/runtime.mjs create mode 100644 packages/bus/src/server.mjs create mode 100644 packages/bus/src/store.mjs create mode 100644 packages/bus/src/views.mjs create mode 100644 packages/bus/tests/broker.test.mjs create mode 100644 packages/bus/tests/business.test.mjs create mode 100644 packages/bus/tests/credentials.test.mjs create mode 100644 packages/bus/tests/fixtures/crash-writer.mjs create mode 100644 packages/bus/tests/human.test.mjs create mode 100644 packages/bus/tests/process.test.mjs create mode 100644 packages/bus/tests/schema.test.mjs create mode 100644 packages/bus/tests/store.test.mjs create mode 100644 packages/bus/tests/transport.test.mjs diff --git a/packages/bus/README.md b/packages/bus/README.md new file mode 100644 index 00000000..75118426 --- /dev/null +++ b/packages/bus/README.md @@ -0,0 +1,228 @@ +# Bus and broker core β€” slice 1 S2 + +The broker owns `/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//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 ` 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. diff --git a/packages/bus/package.json b/packages/bus/package.json new file mode 100644 index 00000000..08591500 --- /dev/null +++ b/packages/bus/package.json @@ -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" + } +} diff --git a/packages/bus/schema.sql b/packages/bus/schema.sql new file mode 100644 index 00000000..6aef8900 --- /dev/null +++ b/packages/bus/schema.sql @@ -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')); diff --git a/packages/bus/src/broker.mjs b/packages/bus/src/broker.mjs new file mode 100644 index 00000000..8c6347eb --- /dev/null +++ b/packages/bus/src/broker.mjs @@ -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, + ); + } +} diff --git a/packages/bus/src/business.mjs b/packages/bus/src/business.mjs new file mode 100644 index 00000000..1b127881 --- /dev/null +++ b/packages/bus/src/business.mjs @@ -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; +} diff --git a/packages/bus/src/client.mjs b/packages/bus/src/client.mjs new file mode 100644 index 00000000..52d46be3 --- /dev/null +++ b/packages/bus/src/client.mjs @@ -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')); + } + }); + }); + } +} diff --git a/packages/bus/src/credentials.mjs b/packages/bus/src/credentials.mjs new file mode 100644 index 00000000..f132c835 --- /dev/null +++ b/packages/bus/src/credentials.mjs @@ -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(); + } +} diff --git a/packages/bus/src/human-cli.mjs b/packages/bus/src/human-cli.mjs new file mode 100644 index 00000000..cca99883 --- /dev/null +++ b/packages/bus/src/human-cli.mjs @@ -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; +} diff --git a/packages/bus/src/human.mjs b/packages/bus/src/human.mjs new file mode 100644 index 00000000..6d875252 --- /dev/null +++ b/packages/bus/src/human.mjs @@ -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(); +} diff --git a/packages/bus/src/index.mjs b/packages/bus/src/index.mjs new file mode 100644 index 00000000..10cf2b4a --- /dev/null +++ b/packages/bus/src/index.mjs @@ -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'; diff --git a/packages/bus/src/process.mjs b/packages/bus/src/process.mjs new file mode 100644 index 00000000..905d4a15 --- /dev/null +++ b/packages/bus/src/process.mjs @@ -0,0 +1,57 @@ +// Started with fork() by the trusted host; boot refs/capabilities travel over IPC, +// never command arguments, stdout or service-token-bearing environment variables. +import { startBroker } from './runtime.mjs'; +import { BusError } from './broker.mjs'; +let runtime, + booted = false, + closing = false; +async function close(code) { + if (closing) return; + closing = true; + try { + await runtime?.close(); + } catch { + code = 2; + } + process.exitCode = code; + if (process.connected) process.disconnect(); +} +if (!process.send) { + process.stderr.write('trusted-host-required\n'); + process.exitCode = 2; +} else { + const timer = setTimeout(() => close(2), 10000); + process.on('message', async (message) => { + try { + if (message?.op === 'bindLaunch' && runtime) { + try { + process.send({ ok: true, launch: runtime.bindLaunch(message.record) }); + } catch (e) { + process.send({ ok: false, error: e instanceof BusError ? e.code : 'bind-refused' }); + } + return; + } + if (message?.op === 'close') { + clearTimeout(timer); + await close(0); + return; + } + if (message?.op !== 'boot' || booted) throw new BusError('invalid-host-request'); + booted = true; + clearTimeout(timer); + runtime = await startBroker(message.config); + if (closing) { + await runtime.close(); + return; + } + process.send({ ok: true, path: runtime.path, launches: runtime.launches, readers: runtime.readers }); + } catch (e) { + process.send?.({ ok: false, error: e instanceof BusError ? e.code : 'startup-refused' }, () => + close(2), + ); + } + }); + process.on('disconnect', () => close(2)); + process.on('SIGTERM', () => close(0)); + process.on('SIGINT', () => close(0)); +} diff --git a/packages/bus/src/runtime.mjs b/packages/bus/src/runtime.mjs new file mode 100644 index 00000000..158172f3 --- /dev/null +++ b/packages/bus/src/runtime.mjs @@ -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; + } +} diff --git a/packages/bus/src/server.mjs b/packages/bus/src/server.mjs new file mode 100644 index 00000000..47eafc87 --- /dev/null +++ b/packages/bus/src/server.mjs @@ -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; + } + }, + }; +} diff --git a/packages/bus/src/store.mjs b/packages/bus/src/store.mjs new file mode 100644 index 00000000..755e4304 --- /dev/null +++ b/packages/bus/src/store.mjs @@ -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); + } +} diff --git a/packages/bus/src/views.mjs b/packages/bus/src/views.mjs new file mode 100644 index 00000000..d161c2cd --- /dev/null +++ b/packages/bus/src/views.mjs @@ -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 }), + }); +} diff --git a/packages/bus/tests/broker.test.mjs b/packages/bus/tests/broker.test.mjs new file mode 100644 index 00000000..3e908d70 --- /dev/null +++ b/packages/bus/tests/broker.test.mjs @@ -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'); + } +}); diff --git a/packages/bus/tests/business.test.mjs b/packages/bus/tests/business.test.mjs new file mode 100644 index 00000000..462c0bfc --- /dev/null +++ b/packages/bus/tests/business.test.mjs @@ -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/); +}); diff --git a/packages/bus/tests/credentials.test.mjs b/packages/bus/tests/credentials.test.mjs new file mode 100644 index 00000000..1d3a33f4 --- /dev/null +++ b/packages/bus/tests/credentials.test.mjs @@ -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(); +}); diff --git a/packages/bus/tests/fixtures/crash-writer.mjs b/packages/bus/tests/fixtures/crash-writer.mjs new file mode 100644 index 00000000..e2d00592 --- /dev/null +++ b/packages/bus/tests/fixtures/crash-writer.mjs @@ -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(); diff --git a/packages/bus/tests/human.test.mjs b/packages/bus/tests/human.test.mjs new file mode 100644 index 00000000..841f73fd --- /dev/null +++ b/packages/bus/tests/human.test.mjs @@ -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); +}); diff --git a/packages/bus/tests/process.test.mjs b/packages/bus/tests/process.test.mjs new file mode 100644 index 00000000..efdb5af2 --- /dev/null +++ b/packages/bus/tests/process.test.mjs @@ -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); +}); diff --git a/packages/bus/tests/schema.test.mjs b/packages/bus/tests/schema.test.mjs new file mode 100644 index 00000000..4d00d56f --- /dev/null +++ b/packages/bus/tests/schema.test.mjs @@ -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); +}); diff --git a/packages/bus/tests/store.test.mjs b/packages/bus/tests/store.test.mjs new file mode 100644 index 00000000..60723c1f --- /dev/null +++ b/packages/bus/tests/store.test.mjs @@ -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); +}); diff --git a/packages/bus/tests/transport.test.mjs b/packages/bus/tests/transport.test.mjs new file mode 100644 index 00000000..8436de39 --- /dev/null +++ b/packages/bus/tests/transport.test.mjs @@ -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); +});