diff --git a/.github/workflows/trust001-preflight.yml b/.github/workflows/trust001-preflight.yml new file mode 100644 index 0000000..5a03daf --- /dev/null +++ b/.github/workflows/trust001-preflight.yml @@ -0,0 +1,77 @@ +name: TRUST-001 provider-free preflight + +on: + pull_request: + paths: + - "Containerfile" + - "runner/trust001/**" + - "container/trust001/**" + - "trust001/**" + - "tests/test_trust001_*.py" + - "tests/trust001/**" + - ".github/workflows/trust001-preflight.yml" + - "docs/experiments/trust-001/checkpoints/**" + +concurrency: + group: trust001-${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + +permissions: + contents: read + +jobs: + unit: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 + + - name: Verify forbidden OpenCode patch tree absent + run: | + test ! -e runtime-patches + test ! -e opencode-patches + + - name: Compile TRUST-001 Python + run: | + python3 -m py_compile runner/trust001/*.py + + - name: Run provider-free TRUST-001 unit tests + run: | + python3 -m unittest discover -s tests -p 'test_trust001_*.py' + + - name: Check dependency-free bridge scripts + run: | + node --check trust001/bridge.mjs + node --check trust001/isolated-host.mjs + node --check trust001/synthetic-loom-peer.mjs + node --check trust001/synthetic-remote-plugin.mjs + python3 -m py_compile tests/trust001/run_stock_bridge_preflight.py + python3 -m py_compile tests/trust001/run_remote_context_preflight.py + + stock-opencode: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 + + - name: Build stock OpenCode candidate image + run: | + docker build -f Containerfile --target opencode --build-arg OPENCODE_VERSION=2.0.23 -t trust001-stock-opencode:preflight . + + - name: Verify exact stock version and labels + run: | + test "$(docker run --rm --entrypoint opencode trust001-stock-opencode:preflight --version)" = "opencode v2.0.23" + test "$(docker image inspect trust001-stock-opencode:preflight --format '{{ index .Config.Labels "io.opencode-eval.stock-opencode-version" }}')" = "2.0.23" + test "$(docker image inspect trust001-stock-opencode:preflight --format '{{ index .Config.Labels "io.opencode-eval.stock-opencode-package" }}')" = "@opencode/cli" + + - name: Exercise stock bridge handshake without inference + run: | + python3 tests/trust001/run_stock_bridge_preflight.py \ + --image trust001-stock-opencode:preflight + + - name: Exercise bounded remote context without inference + run: | + python3 tests/trust001/run_remote_context_preflight.py \ + --image trust001-stock-opencode:preflight + + - name: Record provider-free image identity + run: | + docker image inspect trust001-stock-opencode:preflight --format 'image={{.Id}} config={{.Config.Image}} version={{ index .Config.Labels "io.opencode-eval.stock-opencode-version" }}' diff --git a/Containerfile b/Containerfile index 97fb1f9..e479e29 100644 --- a/Containerfile +++ b/Containerfile @@ -1,5 +1,5 @@ FROM node:24-bookworm-slim@sha256:0e0ff40c39bc087845bfb27465a0df4ea419520094bc35842ff83dd8cbe6f9b6 AS opencode-builder -ARG OPENCODE_VERSION=2.0.18 +ARG OPENCODE_VERSION=2.0.23 RUN npm install --global "@opencode/cli@${OPENCODE_VERSION}" \ && resolved="$(readlink -f "$(command -v opencode)")" \ && test -x "$resolved" \ @@ -53,8 +53,11 @@ ENTRYPOINT ["python3", "/opt/opencode-eval-runner/container/invoke.py"] USER 1000:1000 FROM runtime-base AS opencode +ARG OPENCODE_VERSION=2.0.23 +LABEL io.opencode-eval.stock-opencode-version="${OPENCODE_VERSION}" \ + io.opencode-eval.stock-opencode-package="@opencode/cli" COPY --from=opencode-builder /opencode /usr/local/bin/opencode -RUN opencode --version +RUN test "$(opencode --version)" = "opencode v${OPENCODE_VERSION}" FROM runtime-base AS copilot COPY --from=copilot-builder /opt/copilot/bin/copilot /usr/local/bin/copilot diff --git a/docs/experiments/trust-001/checkpoints/AUTHORIZATION-A.md b/docs/experiments/trust-001/checkpoints/AUTHORIZATION-A.md new file mode 100644 index 0000000..6d974c8 --- /dev/null +++ b/docs/experiments/trust-001/checkpoints/AUTHORIZATION-A.md @@ -0,0 +1,369 @@ +# Authorization A — TRUST-001 stock OpenCode candidate + +**Decision:** GRANTED +**Scope:** prototype construction + explicitly named provider-free local preflights only +**Planning checkpoint:** `fdbc3c9b1c5314f588ffb3cfd34bf0a19502c2aa` +**Candidate branch:** `experiment/trust-001-stock-opencode` +**OpenCode:** stock v2.0.23 only +**Model/provider inference:** NOT AUTHORIZED +**Authorization B:** NOT GRANTED + +This decision authorizes Wave-2 construction under the Gate-1 PASS recorded in PR #43. It does not establish TRUST-001 feasibility. + +## 1. Immutable external constraints + +- OpenCode v2.0.23 remains stock and unmodified. +- The repository's normal OpenCode default pin is authorized to move from 2.0.18 to stock 2.0.23 as part of this candidate branch. +- No OpenCode source patch, fork, custom binary, or upstream PR is permitted. +- Loom production semantics remain unchanged. +- Loom source used by the candidate must be the pinned generation `149406dfa0a01f94491d17054e50a1bc84bb97be`. +- PR #41 patched-runtime work is research/reference only. +- Any runner-owned evidence-safety code reused from PR #41 must be copied/reviewed as runner code and must not carry a patched OpenCode dependency. + +## 2. Candidate construction boundary + +The authorized candidate is: + +> stock OpenCode v2.0.23 with a runner-owned trusted bridge plugin; isolated Loom execution in a separate OCI domain; host-side trusted collector/scope/safety/evidence writer; distinct evidence and capability channels. + +No broader generic remote-PluginHost platform is authorized. + +## 3. Allowed repository mutations + +Construction may change only: + +- `Containerfile`; +- `runner/cli.py`; +- new `runner/trust001/**`; +- `container/invoke.py`; +- new `container/trust001/**`; +- new `trust001/**` bridge / isolated-host sources; +- new `tests/trust001/**`; +- new `tests/test_trust001_*.py`; +- a dedicated provider-free workflow under `.github/workflows/trust001-*.yml`; +- new checkpoint/results documents under `docs/experiments/trust-001/checkpoints/**`; +- this Authorization-A record. + +Existing Gate-1 planning documents are frozen. A change to the reviewed TCB, capability authority, channel model, scope/closure rule, or first-sink policy requires returning to the affected earlier gate. + +No Loom repository mutation is authorized. + +## 4. Stock OpenCode build requirement + +The candidate runtime must package the official upstream stock OpenCode v2.0.23 artifact only. + +The build must: + +- pin `@opencode/cli@2.0.23`; +- verify `opencode --version`; +- record the resulting candidate image digest and image-config digest; +- record the stock OpenCode package/source identity used; +- contain no `runtime-patches/` or equivalent replacement binary. + +The exact built image digest becomes part of the Gate-2 candidate checkpoint. + +## 5. Reviewed trusted boundary + +The candidate may trust only: + +- host runner launcher/orchestrator; +- host collector/scope accountant/safety projector/final writer; +- stock OpenCode v2.0.23 core; +- runner-owned bridge plugin; +- reviewed bridge/correlation/channel implementation; +- selected host kernel + OCI engine enforcement assumptions. + +The following remain untrusted for evidence authority: + +- Loom module/dependencies; +- Loom callbacks/handlers/product state; +- model/product outputs; +- writable workspace; +- Loom subprocesses; +- stock OpenCode shell subprocesses; +- arbitrary external plugins. + +## 6. Capability surface + +The remote Loom broker is limited to the Gate-1 reviewed pinned-Loom surface: + +- immutable location facts; +- legacy plugin storage `get/set/scan`, generation-lifetime and product-only; +- `rpc.register`; +- `agent.transform`; +- `agent.list`; +- `tool.transform`; +- `tool.list`; +- `permission.hook("evaluate")`; +- Session hooks used by pinned Loom; +- `session.get`; +- `session.context`; +- `session.synthetic`; +- `tool.hook("execute.before")`; +- `tool.hook("execute.after")`; +- remote Loom tool execution handlers. + +No generic "call arbitrary PluginHost method" operation is allowed. + +Any newly required capability creates a new candidate checkpoint and requires affected authority review before use. + +## 7. Plugin-source closure + +Before stock OpenCode activation, candidate construction must derive a plugin-source closure manifest covering: + +- effective plugin add/remove operations; +- explicit config documents; +- config roots; +- auto-discovered `plugin/` / `plugins/` paths; +- configured package/local targets; +- resolved symlink targets; +- watched inputs capable of changing the operation set. + +For the tested generation: + +- only the runner bridge may be admitted as evaluated external in-process plugin code; +- every input that can change the operation set must be runner-owned/read-only or otherwise outside evaluated write authority; +- the Loom module must never be imported/evaluated in the OpenCode process; +- a post-activation source-set change is a checkpoint failure, not a hot reload. + +## 8. Mount/process policy + +Authorized prototype defaults: + +- read-only OCI root filesystems; +- non-root execution; +- dropped capabilities; +- no-new-privileges; +- no host PID namespace; +- no container-engine socket inside candidate domains; +- no evidence volume inside evaluated domains; +- bridge/runner source read-only; +- Loom source/dependencies read-only; +- workspace mounted only at the requested normal product mode; +- plugin-source/config closure protected independently of general workspace writability; +- OpenCode and Loom HOME/XDG/state isolated from each other and from host ambient state; +- core dumps disabled or made unavailable/non-persistent. + +Changes that weaken these controls require renewed authority review. + +## 9. Channel model + +Two distinct channels are authorized. + +### Evidence channel + +- direction: trusted OpenCode bridge → host collector; +- only evidence-authoritative candidate observation ingress; +- endpoint outside workspace; +- endpoint inaccessible to Loom; +- no reusable secret in process environment; +- admitted peer bound to the expected bridge/OpenCode generation; +- listener closes/unlinks after admission; +- established descriptor must be close-on-exec; +- no reconnect. + +### Capability channel + +- direction: bridge ↔ isolated Loom; +- product callbacks/results only; +- distinct endpoint/descriptors/material from evidence channel; +- no collector-readable evidence framing or direct evidence-write authority; +- admitted isolated Loom generation bound before product execution; +- listener closes/unlinks after admission; +- no reconnect; +- identity/integrity loss makes affected work unresolved and affected evidence ineligible. + +A shell/Loom child gaining evidence-channel read/write/duplication authority rejects the checkpoint. + +## 10. Request and seal protocol + +Construction must implement only the reviewed lifecycle: + +- trusted request IDs; +- one accepted response per outstanding request; +- reject unknown, unsolicited, duplicate, replayed, stale-generation, post-cancellation, and late responses; +- channel loss never retries product work; +- restart creates a new generation. + +Completeness requires the reviewed two-phase generation close: + +1. stop admitting new case-required work; +2. settle/classify all already-admitted work; +3. assign monotonic evidence sequence; +4. drain accepted observations; +5. emit `seal(generation, finalSequence)`; +6. collector requires a contiguous sequence through `finalSequence`; +7. any accepted post-seal request/event is a completeness failure. + +Construction may implement the mechanism; proof that the chosen stock lifecycle provides a sufficient close/drain boundary remains a provider-free preflight obligation. + +## 11. First-sink policy + +Candidate-controlled raw runtime observation may exist transiently only in bounded trusted memory/IPC. + +Before any candidate-controlled: + +- file write; +- stdout/stderr diagnostic; +- OCI log payload; +- temporary persistence; +- clipping/truncation; +- exported artifact; + +credential protection must already have occurred. + +Diagnostics use fixed/bounded reason codes and non-sensitive identities only. + +Normal stock OpenCode/Loom product stores are not automatically runner evidence stores. If candidate code copies/promotes them, the copy boundary becomes a reviewed evidence sink. + +## 12. Named provider-free preflights authorized + +Only the following candidate executions are authorized before Authorization B: + +### PF-01 — Stock runtime provenance + +Build/start the candidate stock-OpenCode image and verify: + +- v2.0.23; +- no patched runtime; +- exact source/image/config/component provenance; +- no external inference. + +### PF-02 — Plugin-source closure + +Using synthetic project/config fixtures: + +- enumerate effective plugin-source operations; +- prove evaluated-writable workspace mutations cannot add/replace/reload in-process plugin code; +- test symlink/config/watch/restart/cold-start escape attempts. + +### PF-03 — Evidence-channel authority + +Using synthetic bridge observations and a hostile local child: + +- listener admission/closure; +- descriptor close-on-exec; +- no inheritance; +- no reconnect; +- `/proc//fd`, `pidfd_getfd`, ptrace/process-memory attempts where the host permits the test; +- failure/suppression yields incomplete evidence, never false completion. + +### PF-04 — Capability-peer authority + +Using a synthetic isolated plugin peer: + +- expected peer admission; +- unadmitted peer rejection; +- duplicate/replay/stale-generation/late response rejection; +- malformed/oversized/code-bearing message rejection; +- identity/integrity loss makes affected work unresolved/ineligible; +- no evidence-channel authority from capability traffic. + +### PF-05 — Registration/transform fidelity + +Provider-free synthetic plugin fixtures for: + +- tool namespace/definition registration; +- opaque execute-handler identity; +- `tool.list` self-identity parity; +- initial `agent.transform` declarative replay; +- registration order. + +Dynamic agent-registry reload parity remains outside the initial claim unless separately reviewed. + +### PF-06 — Hook/callback fidelity + +Provider-free synthetic fixtures for: + +- `execute.before` mutation/failure; +- handler return/throw; +- `execute.after` mutation; +- permission mutation; +- Session context/retry mutation; +- cancellation + late response; +- channel loss before/after product completion. + +No claim of remote-side-effect rollback is allowed. + +### PF-07 — Legacy storage bounds + +Synthetic generation-lifetime `get/set/scan` tests proving: + +- only the admitted Loom plugin namespace is reachable; +- structural/size limits; +- no arbitrary KV namespace escape; +- no evidence authority. + +### PF-08 — Event ordering / root admission / sealing + +Provider-free stock OpenCode + synthetic tool/plugin fixture proving or falsifying: + +- trusted root Session admission; +- public-event observation path; +- same/different-parent overlapping calls; +- monotonic collector sequence; +- drain and `seal(generation, finalSequence)`; +- sequence gap and post-seal rejection. + +### PF-09 — First-sink confidentiality + +Synthetic credentials at every actual candidate-controlled first sink: + +- literal; +- short; +- JSON escaped; +- nested/deep-key; +- representation-changing echo; +- sensitive-key value; +- boundary-spanning value; +- timeout/error/transport-loss path. + +Plaintext at any candidate-controlled first sink is FAIL. + +### PF-10 — Pinned Loom provider-free activation + +Only after PF-01 through PF-07 pass for the exact candidate checkpoint: + +- execute pinned Loom `149406d` in the isolated Loom domain; +- activate the runner bridge in stock OpenCode; +- no model/provider call; +- verify required registrations/capability use can initialize without importing Loom into OpenCode; +- verify plugin-source closure remains unchanged. + +This preflight does not authorize SCN-01 through SCN-07 or any model-backed `eval:live`. + +## 13. Explicitly not authorized + +Authorization A does not permit: + +- model/provider inference; +- paid/credentialed inference; +- semantic scenarios SCN-01 through SCN-07; +- general adversarial Wave-4 execution beyond the named provider-free preflights; +- runner→Loom model-backed composition; +- OpenCode modification; +- Loom modification; +- merge to main; +- release/publish changes beyond CI-local candidate images; +- production adoption. + +## 14. Review boundary + +Before Authorization B: + +- Gate 2 independent source/authority review must PASS; +- Gate 2C provider-free preflight must PASS; +- exact source/image/component checkpoint must be recorded. + +Review roles: + +- **Gate 2:** independent cold source/authority Reviewer; +- **Gate 2C:** independent cold preflight/adversarial Reviewer. + +The implementation author must not self-promote an `UNPROVEN` or `NOT RUN` item to PASS. + +## 15. Rejection/recovery + +The Gate-1 checkpoint discipline remains authoritative. + +Any demonstrated rejection condition stops the exact candidate checkpoint. Corrections create a new source/image/component checkpoint and return to earlier gates whenever authority, TCB, capability surface, channel model, scope/seal rule, or first-sink boundary changes. diff --git a/runner/trust001/__init__.py b/runner/trust001/__init__.py new file mode 100644 index 0000000..dcaa77e --- /dev/null +++ b/runner/trust001/__init__.py @@ -0,0 +1,5 @@ +"""TRUST-001 candidate primitives. + +Authorization boundary: construction/provider-free preflight only. +These modules do not establish TRUST-001 feasibility. +""" diff --git a/runner/trust001/broker.py b/runner/trust001/broker.py new file mode 100644 index 0000000..ff1c3c9 --- /dev/null +++ b/runner/trust001/broker.py @@ -0,0 +1,75 @@ +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Any, Literal + +from .protocol import ProtocolError, require, validate_capability +from .state import CapabilityLedger + +Peer = Literal["bridge", "loom"] + + +@dataclass +class CapabilityRouter: + generation: str + host_calls: CapabilityLedger = field(init=False) + callbacks: CapabilityLedger = field(init=False) + failed: bool = False + failure: str | None = None + + def __post_init__(self) -> None: + self.host_calls = CapabilityLedger(self.generation) + self.callbacks = CapabilityLedger(self.generation) + + def _reject(self, reason: str) -> None: + self.failed = True + self.failure = self.failure or reason + raise ProtocolError(reason) + + def route(self, origin: Peer, raw: dict[str, Any]) -> tuple[Peer, dict[str, Any]]: + if self.failed: + raise ProtocolError(self.failure or "capability_router_failed") + message = validate_capability(raw) + if message["generation"] != self.generation: + self._reject("stale_generation") + + kind = message["kind"] + if kind == "capability.hello": + self._reject("unexpected_hello_after_admission") + + try: + if origin == "loom": + if kind == "capability.host.request": + self.host_calls.admit(message["request_id"], message["operation"]) + return "bridge", message + if kind == "capability.callback.response": + self.callbacks.settle(message["request_id"]) + return "bridge", message + if kind == "capability.host.cancel": + self.host_calls.cancel(message["request_id"]) + return "bridge", message + self._reject("loom_direction_violation") + + if origin == "bridge": + if kind == "capability.callback.request": + self.callbacks.admit(message["request_id"], message["operation"]) + return "loom", message + if kind == "capability.host.response": + self.host_calls.settle(message["request_id"]) + return "loom", message + if kind == "capability.callback.cancel": + self.callbacks.cancel(message["request_id"]) + return "loom", message + self._reject("bridge_direction_violation") + except ProtocolError as exc: + self._reject(str(exc)) + + self._reject("invalid_origin") + + def close_admission(self) -> None: + self.host_calls.close_admission() + self.callbacks.close_admission() + + @property + def settled(self) -> bool: + return self.host_calls.settled and self.callbacks.settled diff --git a/runner/trust001/channel.py b/runner/trust001/channel.py new file mode 100644 index 0000000..9586433 --- /dev/null +++ b/runner/trust001/channel.py @@ -0,0 +1,112 @@ +from __future__ import annotations + +import os +import socket +import stat +import struct +import tempfile +from dataclasses import dataclass +from pathlib import Path + +from .protocol import ProtocolError, require + + +@dataclass(frozen=True) +class PeerCredentials: + pid: int + uid: int + gid: int + + +def peer_credentials(sock: socket.socket) -> PeerCredentials | None: + if not hasattr(socket, "SO_PEERCRED"): + return None + raw = sock.getsockopt(socket.SOL_SOCKET, socket.SO_PEERCRED, struct.calcsize("3i")) + pid, uid, gid = struct.unpack("3i", raw) + return PeerCredentials(pid=pid, uid=uid, gid=gid) + + +class OneShotUnixListener: + """Private one-accept AF_UNIX listener. + + The socket pathname is unlinked immediately after the first accepted peer. + Python sockets are non-inheritable by default; this class asserts that + property for both listener and accepted descriptor. + """ + + def __init__(self, parent: Path | None = None, *, name: str = "channel.sock") -> None: + root = Path(tempfile.mkdtemp(prefix="trust001-", dir=parent)) + root.chmod(0o700) + self.root = root + self.path = root / name + self.socket = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + require(not self.socket.get_inheritable(), "listener_inheritable") + self.socket.bind(str(self.path)) + self.path.chmod(0o600) + self.socket.listen(1) + self.accepted = False + self.closed = False + + def accept_once( + self, + *, + expected_uid: int | None = None, + expected_gid: int | None = None, + expected_pid: int | None = None, + timeout: float = 5.0, + ) -> tuple[socket.socket, PeerCredentials | None]: + require(not self.accepted and not self.closed, "listener_not_available") + self.socket.settimeout(timeout) + conn, _ = self.socket.accept() + conn.set_inheritable(False) + if conn.get_inheritable(): + conn.close() + raise ProtocolError("accepted_descriptor_inheritable") + credentials = peer_credentials(conn) + if credentials is not None: + if expected_uid is not None and credentials.uid != expected_uid: + conn.close() + raise ProtocolError("unexpected_peer_uid") + if expected_gid is not None and credentials.gid != expected_gid: + conn.close() + raise ProtocolError("unexpected_peer_gid") + if expected_pid is not None and credentials.pid != expected_pid: + conn.close() + raise ProtocolError("unexpected_peer_pid") + elif any(value is not None for value in (expected_uid, expected_gid, expected_pid)): + conn.close() + raise ProtocolError("peer_credentials_unavailable") + self.accepted = True + try: + self.path.unlink() + finally: + self.socket.close() + self.closed = True + return conn, credentials + + def close(self) -> None: + if not self.closed: + self.socket.close() + self.closed = True + self.path.unlink(missing_ok=True) + + def cleanup(self) -> None: + self.close() + try: + self.root.rmdir() + except OSError: + pass + + def __enter__(self) -> "OneShotUnixListener": + return self + + def __exit__(self, *_: object) -> None: + self.cleanup() + + +def private_socket_path(listener: OneShotUnixListener) -> None: + info = listener.root.stat() + require(stat.S_IMODE(info.st_mode) == 0o700, "socket_dir_mode") + info = listener.path.stat() + require(stat.S_ISSOCK(info.st_mode), "not_socket") + require(stat.S_IMODE(info.st_mode) == 0o600, "socket_mode") diff --git a/runner/trust001/protocol.py b/runner/trust001/protocol.py new file mode 100644 index 0000000..57bb583 --- /dev/null +++ b/runner/trust001/protocol.py @@ -0,0 +1,225 @@ +from __future__ import annotations + +import json +import math +import re +import struct +from dataclasses import dataclass +from typing import Any + +WIRE_VERSION = "opencode-eval-runner/trust001-wire/v1" +MAX_FRAME_BYTES = 256 * 1024 +MAX_JSON_DEPTH = 32 +MAX_JSON_NODES = 20_000 + +GENERATION_RE = re.compile(r"[0-9a-f]{64}\Z") +REQUEST_RE = re.compile(r"[0-9a-f]{32}\Z") + +HOST_OPERATIONS = frozenset({ + "location.get", + "storage.get", "storage.set", "storage.scan", + "rpc.register", + "agent.transform.register", "agent.list", + "tool.transform.register", "tool.list", "tool.hook.register", + "permission.hook.register", + "session.hook.register", "session.get", "session.context", "session.synthetic", +}) + +CALLBACK_OPERATIONS = frozenset({ + "rpc.call", + "tool.execute", + "tool.execute.before", + "tool.execute.after", + "permission.evaluate", + "session.context", + "session.retry", + "trust001.preflight.ping", +}) + +CAPABILITY_KINDS = frozenset({ + "capability.hello", + "capability.host.request", + "capability.host.response", + "capability.callback.request", + "capability.callback.response", + "capability.host.cancel", + "capability.callback.cancel", +}) +EVIDENCE_KINDS = frozenset({ + "evidence.hello", + "evidence.observation", + "evidence.close", + "evidence.seal", +}) + + +class ProtocolError(ValueError): + pass + + +def require(ok: bool, message: str = "invalid_protocol") -> None: + if not ok: + raise ProtocolError(message) + + +def _pairs(pairs: list[tuple[str, Any]]) -> dict[str, Any]: + value: dict[str, Any] = {} + for key, item in pairs: + require(key not in value, "duplicate_json_key") + value[key] = item + return value + + +def _bad_number(_: str) -> None: + raise ProtocolError("nonfinite_number") + + +def _bounded_json(value: Any) -> None: + pending = [(value, 0)] + nodes = 0 + while pending: + current, depth = pending.pop() + nodes += 1 + require(nodes <= MAX_JSON_NODES and depth <= MAX_JSON_DEPTH, "json_shape_limit") + if current is None or type(current) in (bool, int, str): + if type(current) is str: + current.encode("utf-8", errors="strict") + continue + if type(current) is float: + require(math.isfinite(current), "nonfinite_number") + continue + if type(current) is list: + pending.extend((item, depth + 1) for item in current) + continue + if type(current) is dict: + require(all(type(key) is str for key in current), "non_string_key") + pending.extend((item, depth + 1) for item in current.values()) + continue + raise ProtocolError("non_json_value") + + +def strict_loads(raw: bytes) -> dict[str, Any]: + require(len(raw) <= MAX_FRAME_BYTES, "frame_limit") + try: + value = json.loads( + raw.decode("utf-8", errors="strict"), + object_pairs_hook=_pairs, + parse_constant=_bad_number, + ) + except (json.JSONDecodeError, UnicodeError) as exc: + raise ProtocolError("invalid_json") from exc + require(type(value) is dict, "frame_not_object") + _bounded_json(value) + return value + + +def encode(value: dict[str, Any]) -> bytes: + _bounded_json(value) + try: + raw = json.dumps( + value, + ensure_ascii=False, + allow_nan=False, + separators=(",", ":"), + ).encode("utf-8") + except (TypeError, ValueError, UnicodeError) as exc: + raise ProtocolError("invalid_json") from exc + require(len(raw) <= MAX_FRAME_BYTES, "frame_limit") + return struct.pack("!I", len(raw)) + raw + + +def _exact(value: dict[str, Any], required: set[str], optional: set[str] = frozenset()) -> None: + require(required <= value.keys(), "missing_field") + require(not (set(value) - required - optional), "unknown_field") + + +def _common(value: dict[str, Any], kinds: frozenset[str]) -> tuple[str, str]: + _exact(value, {"version", "kind", "generation"}, { + "role", "request_id", "operation", "payload", "ok", "error", + "reason", "sequence", "final_sequence", + }) + require(value.get("version") == WIRE_VERSION, "wrong_version") + kind = value.get("kind") + generation = value.get("generation") + require(type(kind) is str and kind in kinds, "wrong_channel_kind") + require(type(generation) is str and GENERATION_RE.fullmatch(generation) is not None, "invalid_generation") + return kind, generation + + +def validate_capability(value: dict[str, Any]) -> dict[str, Any]: + kind, _ = _common(value, CAPABILITY_KINDS) + if kind == "capability.hello": + _exact(value, {"version", "kind", "generation", "role"}) + require(value["role"] in {"bridge", "loom"}, "invalid_role") + return value + + if kind in {"capability.host.request", "capability.callback.request"}: + _exact(value, {"version", "kind", "generation", "request_id", "operation", "payload"}) + require(type(value["request_id"]) is str and REQUEST_RE.fullmatch(value["request_id"]) is not None, + "invalid_request_id") + require(type(value["operation"]) is str, "invalid_operation") + allowed = HOST_OPERATIONS if kind == "capability.host.request" else CALLBACK_OPERATIONS + require(value["operation"] in allowed, "operation_not_admitted") + _bounded_json(value["payload"]) + return value + + if kind in {"capability.host.response", "capability.callback.response"}: + _exact(value, {"version", "kind", "generation", "request_id", "ok"}, {"payload", "error"}) + require(type(value["request_id"]) is str and REQUEST_RE.fullmatch(value["request_id"]) is not None, + "invalid_request_id") + require(type(value["ok"]) is bool, "invalid_response") + require(("payload" in value) != ("error" in value), "invalid_response") + if "error" in value: + require(type(value["error"]) is str and 0 < len(value["error"]) <= 256, "invalid_error") + elif "payload" in value: + _bounded_json(value["payload"]) + return value + + require(kind in {"capability.host.cancel", "capability.callback.cancel"}, "wrong_channel_kind") + _exact(value, {"version", "kind", "generation", "request_id", "reason"}) + require(type(value["request_id"]) is str and REQUEST_RE.fullmatch(value["request_id"]) is not None, + "invalid_request_id") + require(type(value["reason"]) is str and 0 < len(value["reason"]) <= 128, "invalid_reason") + return value + + +def validate_evidence(value: dict[str, Any]) -> dict[str, Any]: + kind, _ = _common(value, EVIDENCE_KINDS) + if kind == "evidence.hello": + _exact(value, {"version", "kind", "generation", "role"}) + require(value["role"] == "bridge", "invalid_role") + elif kind == "evidence.observation": + _exact(value, {"version", "kind", "generation", "sequence", "payload"}) + require(type(value["sequence"]) is int and 1 <= value["sequence"] <= 2**53 - 1, "invalid_sequence") + _bounded_json(value["payload"]) + elif kind == "evidence.close": + _exact(value, {"version", "kind", "generation"}) + elif kind == "evidence.seal": + _exact(value, {"version", "kind", "generation", "final_sequence"}) + require(type(value["final_sequence"]) is int and 0 <= value["final_sequence"] <= 2**53 - 1, + "invalid_final_sequence") + return value + + +@dataclass +class FrameReader: + buffer: bytearray + + def __init__(self) -> None: + self.buffer = bytearray() + + def feed(self, data: bytes, *, channel: str) -> list[dict[str, Any]]: + require(channel in {"capability", "evidence"}, "invalid_channel") + self.buffer.extend(data) + frames: list[dict[str, Any]] = [] + while len(self.buffer) >= 4: + size = struct.unpack("!I", self.buffer[:4])[0] + require(0 < size <= MAX_FRAME_BYTES, "frame_limit") + if len(self.buffer) < 4 + size: + break + raw = bytes(self.buffer[4:4 + size]) + del self.buffer[:4 + size] + value = strict_loads(raw) + frames.append(validate_capability(value) if channel == "capability" else validate_evidence(value)) + require(len(self.buffer) <= MAX_FRAME_BYTES + 4, "frame_buffer_limit") + return frames diff --git a/runner/trust001/runtime.py b/runner/trust001/runtime.py new file mode 100644 index 0000000..e389c12 --- /dev/null +++ b/runner/trust001/runtime.py @@ -0,0 +1,103 @@ +from __future__ import annotations + +import secrets +import socket +from dataclasses import dataclass +from pathlib import Path + +from .channel import OneShotUnixListener, PeerCredentials +from .protocol import require +from .transport import CapabilityRelay, EvidenceCollector + + +@dataclass +class AdmittedChannels: + generation: str + capability: CapabilityRelay + evidence: EvidenceCollector + bridge_capability_peer: PeerCredentials | None + loom_capability_peer: PeerCredentials | None + evidence_peer: PeerCredentials | None + + +class ChannelSet: + """Three physically distinct one-shot endpoints. + + - evidence: visible only to stock OpenCode/bridge + - capability_bridge: visible only to stock OpenCode/bridge + - capability_loom: visible only to isolated Loom + + The two capability sockets are joined only by the trusted host relay. + """ + + def __init__(self, root: Path | None = None, generation: str | None = None) -> None: + self.generation = generation or secrets.token_hex(32) + require(len(self.generation) == 64, "invalid_generation") + self.evidence = OneShotUnixListener(root, name="evidence.sock") + self.capability_bridge = OneShotUnixListener(root, name="bridge.sock") + self.capability_loom = OneShotUnixListener(root, name="loom.sock") + + def paths(self) -> dict[str, Path]: + return { + "evidence": self.evidence.path, + "capability_bridge": self.capability_bridge.path, + "capability_loom": self.capability_loom.path, + } + + def mount_roots(self) -> dict[str, Path]: + return { + "evidence": self.evidence.root, + "capability_bridge": self.capability_bridge.root, + "capability_loom": self.capability_loom.root, + } + + def admit( + self, + *, + bridge_uid: int | None = None, + bridge_gid: int | None = None, + loom_uid: int | None = None, + loom_gid: int | None = None, + timeout: float = 10.0, + ) -> AdmittedChannels: + evidence_sock, evidence_peer = self.evidence.accept_once( + expected_uid=bridge_uid, + expected_gid=bridge_gid, + timeout=timeout, + ) + bridge_sock, bridge_peer = self.capability_bridge.accept_once( + expected_uid=bridge_uid, + expected_gid=bridge_gid, + timeout=timeout, + ) + loom_sock, loom_peer = self.capability_loom.accept_once( + expected_uid=loom_uid, + expected_gid=loom_gid, + timeout=timeout, + ) + try: + capability = CapabilityRelay.admit(self.generation, bridge_sock, loom_sock) + evidence = EvidenceCollector.admit(self.generation, evidence_sock) + except Exception: + evidence_sock.close() + bridge_sock.close() + loom_sock.close() + raise + return AdmittedChannels( + generation=self.generation, + capability=capability, + evidence=evidence, + bridge_capability_peer=bridge_peer, + loom_capability_peer=loom_peer, + evidence_peer=evidence_peer, + ) + + def cleanup(self) -> None: + for listener in (self.evidence, self.capability_bridge, self.capability_loom): + listener.cleanup() + + def __enter__(self) -> "ChannelSet": + return self + + def __exit__(self, *_: object) -> None: + self.cleanup() diff --git a/runner/trust001/source_closure.py b/runner/trust001/source_closure.py new file mode 100644 index 0000000..b40b272 --- /dev/null +++ b/runner/trust001/source_closure.py @@ -0,0 +1,99 @@ +from __future__ import annotations + +import hashlib +import json +import os +import shutil +from pathlib import Path +from typing import Any + +from .protocol import ProtocolError, require + +BRIDGE_ID = "loom" + + +def sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as stream: + for block in iter(lambda: stream.read(1024 * 1024), b""): + digest.update(block) + return digest.hexdigest() + + +def _plain_json_config(source: Path | None) -> dict[str, Any]: + if source is None: + return {"$schema": "https://opencode.ai/config.json"} + try: + value = json.loads(source.read_text(encoding="utf-8")) + except (OSError, UnicodeError, json.JSONDecodeError) as exc: + raise ProtocolError("trust001_config_must_be_strict_json") from exc + require(type(value) is dict, "trust001_config_not_object") + # Candidate OpenCode may not consume product-selected external plugin + # declarations. Loom is isolated and the bridge is runner-owned. + require("plugin" not in value and "plugins" not in value, "external_plugin_declaration") + return value + + +def build_source_closure( + *, + destination: Path, + bridge_source: Path, + product_config: Path | None = None, +) -> dict[str, Any]: + require(bridge_source.is_file() and not bridge_source.is_symlink(), "invalid_bridge_source") + if destination.exists(): + shutil.rmtree(destination) + plugins = destination / "plugins" + plugins.mkdir(parents=True, mode=0o700) + + config = _plain_json_config(product_config) + config_path = destination / "opencode.json" + config_path.write_text( + json.dumps(config, ensure_ascii=False, allow_nan=False, separators=(",", ":")) + "\n", + encoding="utf-8", + ) + bridge_path = plugins / f"{BRIDGE_ID}.js" + shutil.copyfile(bridge_source, bridge_path) + + for path in (config_path, bridge_path): + path.chmod(0o444) + plugins.chmod(0o555) + destination.chmod(0o555) + + return { + "schema": "opencode-eval-runner/trust001-plugin-source-closure/v1", + "project_config_disabled": True, + "external_plugin_declarations": False, + "config": { + "relative_path": "opencode.json", + "sha256": sha256(config_path), + }, + "plugins": [{ + "id": BRIDGE_ID, + "relative_path": f"plugins/{BRIDGE_ID}.js", + "sha256": sha256(bridge_path), + }], + "expected_external_plugin_ids": [BRIDGE_ID], + } + + +def verify_source_closure(root: Path, manifest: dict[str, Any]) -> None: + require(manifest.get("schema") == "opencode-eval-runner/trust001-plugin-source-closure/v1", + "invalid_closure_manifest") + require(manifest.get("project_config_disabled") is True, "project_config_not_disabled") + require(manifest.get("external_plugin_declarations") is False, "external_plugin_declaration") + + config = manifest.get("config") + plugins = manifest.get("plugins") + require(type(config) is dict and type(plugins) is list and len(plugins) == 1, "invalid_closure_manifest") + expected = [ + (root / str(config.get("relative_path")), config.get("sha256")), + (root / str(plugins[0].get("relative_path")), plugins[0].get("sha256")), + ] + for path, digest in expected: + require(path.is_file() and not path.is_symlink(), "closure_path_changed") + require(type(digest) is str and sha256(path) == digest, "closure_digest_changed") + require(path.stat().st_mode & 0o222 == 0, "closure_path_writable") + + require(root.stat().st_mode & 0o222 == 0, "closure_root_writable") + require((root / "plugins").stat().st_mode & 0o222 == 0, "closure_plugins_writable") diff --git a/runner/trust001/state.py b/runner/trust001/state.py new file mode 100644 index 0000000..cae4053 --- /dev/null +++ b/runner/trust001/state.py @@ -0,0 +1,81 @@ +from __future__ import annotations + +import secrets +from dataclasses import dataclass, field +from typing import Any + +from .protocol import ProtocolError, REQUEST_RE, require + + +@dataclass +class CapabilityLedger: + generation: str + outstanding: dict[str, str] = field(default_factory=dict) + terminal: set[str] = field(default_factory=set) + closed: bool = False + + def allocate(self, operation: str) -> str: + require(not self.closed, "generation_closed") + request_id = secrets.token_hex(16) + while request_id in self.outstanding or request_id in self.terminal: + request_id = secrets.token_hex(16) + self.outstanding[request_id] = operation + return request_id + + def admit(self, request_id: str, operation: str) -> None: + """Deterministic helper for provider-free fixtures; production should allocate().""" + require(not self.closed, "generation_closed") + require(REQUEST_RE.fullmatch(request_id) is not None, "invalid_request_id") + require(request_id not in self.outstanding and request_id not in self.terminal, "duplicate_request") + self.outstanding[request_id] = operation + + def settle(self, request_id: str) -> str: + require(request_id in self.outstanding, "unknown_or_late_response") + operation = self.outstanding.pop(request_id) + self.terminal.add(request_id) + return operation + + def cancel(self, request_id: str) -> str: + return self.settle(request_id) + + def close_admission(self) -> None: + self.closed = True + + @property + def settled(self) -> bool: + return self.closed and not self.outstanding + + +@dataclass +class EvidenceLedger: + generation: str + sequence: int = 0 + observations: list[Any] = field(default_factory=list) + sealed: bool = False + invalid: bool = False + failure: str | None = None + + def _fail(self, reason: str) -> None: + self.invalid = True + self.failure = self.failure or reason + raise ProtocolError(reason) + + def observe(self, sequence: int, payload: Any) -> None: + if self.sealed: + self._fail("post_seal_observation") + expected = self.sequence + 1 + if sequence != expected: + self._fail("non_contiguous_sequence") + self.sequence = sequence + self.observations.append(payload) + + def seal(self, final_sequence: int) -> None: + if self.sealed: + self._fail("duplicate_seal") + if final_sequence != self.sequence: + self._fail("seal_sequence_mismatch") + self.sealed = True + + @property + def complete(self) -> bool: + return self.sealed and not self.invalid diff --git a/runner/trust001/transport.py b/runner/trust001/transport.py new file mode 100644 index 0000000..0550dd6 --- /dev/null +++ b/runner/trust001/transport.py @@ -0,0 +1,183 @@ +from __future__ import annotations + +import socket +import threading +from dataclasses import dataclass, field +from queue import Queue +from typing import Any + +from .broker import CapabilityRouter +from .protocol import FrameReader, ProtocolError, encode, require, validate_evidence +from .state import EvidenceLedger + + +class FramedSocket: + def __init__(self, sock: socket.socket, *, channel: str) -> None: + require(channel in {"capability", "evidence"}, "invalid_channel") + self.sock = sock + self.channel = channel + self.reader = FrameReader() + self.pending: list[dict[str, Any]] = [] + + def send(self, value: dict[str, Any]) -> None: + self.sock.sendall(encode(value)) + + def recv(self) -> dict[str, Any]: + if self.pending: + return self.pending.pop(0) + while True: + chunk = self.sock.recv(65536) + if not chunk: + raise ProtocolError("channel_closed") + frames = self.reader.feed(chunk, channel=self.channel) + if frames: + self.pending.extend(frames[1:]) + return frames[0] + + +@dataclass +class CapabilityRelay: + generation: str + bridge: FramedSocket + loom: FramedSocket + router: CapabilityRouter + _lock: threading.Lock = field(default_factory=threading.Lock, repr=False) + + @classmethod + def admit( + cls, + generation: str, + bridge_sock: socket.socket, + loom_sock: socket.socket, + ) -> "CapabilityRelay": + bridge = FramedSocket(bridge_sock, channel="capability") + loom = FramedSocket(loom_sock, channel="capability") + bridge_hello = bridge.recv() + loom_hello = loom.recv() + require( + bridge_hello == { + "version": bridge_hello["version"], + "kind": "capability.hello", + "generation": generation, + "role": "bridge", + }, + "invalid_bridge_hello", + ) + require( + loom_hello == { + "version": loom_hello["version"], + "kind": "capability.hello", + "generation": generation, + "role": "loom", + }, + "invalid_loom_hello", + ) + return cls( + generation=generation, + bridge=bridge, + loom=loom, + router=CapabilityRouter(generation), + ) + + def relay_bridge_once(self) -> dict[str, Any]: + message = self.bridge.recv() + with self._lock: + target, routed = self.router.route("bridge", message) + require(target == "loom", "invalid_route") + self.loom.send(routed) + return routed + + def relay_loom_once(self) -> dict[str, Any]: + message = self.loom.recv() + with self._lock: + target, routed = self.router.route("loom", message) + require(target == "bridge", "invalid_route") + self.bridge.send(routed) + return routed + + def serve(self) -> tuple[threading.Event, Queue[BaseException], list[threading.Thread]]: + """Run both capability directions concurrently until stopped or failed.""" + stop = threading.Event() + errors: Queue[BaseException] = Queue() + + def worker(direction: str) -> None: + relay = self.relay_bridge_once if direction == "bridge" else self.relay_loom_once + while not stop.is_set(): + try: + relay() + except BaseException as exc: + if not stop.is_set(): + clean_close = ( + isinstance(exc, ProtocolError) + and str(exc) == "channel_closed" + and not self.router.host_calls.outstanding + and not self.router.callbacks.outstanding + ) + if not clean_close: + errors.put(exc) + stop.set() + self.close() + return + + threads = [ + threading.Thread(target=worker, args=("bridge",), daemon=True), + threading.Thread(target=worker, args=("loom",), daemon=True), + ] + for thread in threads: + thread.start() + return stop, errors, threads + + def close(self) -> None: + for channel in (self.bridge, self.loom): + try: + channel.sock.shutdown(socket.SHUT_RDWR) + except OSError: + pass + try: + channel.sock.close() + except OSError: + pass + + +@dataclass +class EvidenceCollector: + generation: str + channel: FramedSocket + ledger: EvidenceLedger + + @classmethod + def admit(cls, generation: str, sock: socket.socket) -> "EvidenceCollector": + channel = FramedSocket(sock, channel="evidence") + hello = channel.recv() + require( + hello == { + "version": hello["version"], + "kind": "evidence.hello", + "generation": generation, + "role": "bridge", + }, + "invalid_evidence_hello", + ) + return cls( + generation=generation, + channel=channel, + ledger=EvidenceLedger(generation), + ) + + def request_close(self) -> None: + self.channel.send({ + "version": "opencode-eval-runner/trust001-wire/v1", + "kind": "evidence.close", + "generation": self.generation, + }) + + def receive_once(self) -> dict[str, Any]: + message = validate_evidence(self.channel.recv()) + require(message["generation"] == self.generation, "stale_generation") + if message["kind"] == "evidence.observation": + self.ledger.observe(message["sequence"], message["payload"]) + elif message["kind"] == "evidence.seal": + self.ledger.seal(message["final_sequence"]) + else: + raise ProtocolError("unexpected_evidence_control") + return message diff --git a/tests/test_invoke.py b/tests/test_invoke.py index b1495e7..1ccc93a 100644 --- a/tests/test_invoke.py +++ b/tests/test_invoke.py @@ -706,12 +706,14 @@ def test_workflows_pin_external_actions_by_commit(self): ): self.assertNotIn(mutable, ci + publish) - def test_container_pins_opencode_2_0_18(self): + def test_container_pins_stock_opencode_2_0_23(self): containerfile = (Path(__file__).resolve().parents[1] / "Containerfile").read_text( encoding="utf-8" ) - self.assertIn("ARG OPENCODE_VERSION=2.0.18", containerfile) - self.assertNotIn("ARG OPENCODE_VERSION=2.0.15", containerfile) + self.assertIn("ARG OPENCODE_VERSION=2.0.23", containerfile) + self.assertIn('io.opencode-eval.stock-opencode-version="${OPENCODE_VERSION}"', containerfile) + self.assertIn('io.opencode-eval.stock-opencode-package="@opencode/cli"', containerfile) + self.assertNotIn("ARG OPENCODE_VERSION=2.0.18", containerfile) def test_container_pins_base_images_and_copilot_release_asset(self): containerfile = (Path(__file__).resolve().parents[1] / "Containerfile").read_text( diff --git a/tests/test_trust001_broker.py b/tests/test_trust001_broker.py new file mode 100644 index 0000000..beb315a --- /dev/null +++ b/tests/test_trust001_broker.py @@ -0,0 +1,149 @@ +from __future__ import annotations + +import unittest + +from runner.trust001.broker import CapabilityRouter +from runner.trust001.protocol import ProtocolError, WIRE_VERSION + + +GEN = "a" * 64 +HOST_REQ = "1" * 32 +CALLBACK_REQ = "2" * 32 + + +def message(kind: str, request_id: str, **extra): + return { + "version": WIRE_VERSION, + "kind": kind, + "generation": GEN, + "request_id": request_id, + **extra, + } + + +class BrokerTests(unittest.TestCase): + def test_loom_host_request_round_trip(self): + router = CapabilityRouter(GEN) + target, request = router.route("loom", message( + "capability.host.request", + HOST_REQ, + operation="storage.get", + payload={"key": "project/x"}, + )) + self.assertEqual(target, "bridge") + self.assertEqual(request["operation"], "storage.get") + + target, _ = router.route("bridge", message( + "capability.host.response", + HOST_REQ, + ok=True, + payload={"value": None}, + )) + self.assertEqual(target, "loom") + self.assertNotIn(HOST_REQ, router.host_calls.outstanding) + + def test_bridge_callback_round_trip(self): + router = CapabilityRouter(GEN) + target, request = router.route("bridge", message( + "capability.callback.request", + CALLBACK_REQ, + operation="tool.execute.before", + payload={"tool": "question"}, + )) + self.assertEqual(target, "loom") + self.assertEqual(request["operation"], "tool.execute.before") + + target, _ = router.route("loom", message( + "capability.callback.response", + CALLBACK_REQ, + ok=True, + payload={"input": {}}, + )) + self.assertEqual(target, "bridge") + + def test_request_classes_cannot_cross_direction(self): + cases = [ + ("loom", message( + "capability.callback.request", CALLBACK_REQ, + operation="tool.execute.before", payload={} + )), + ("bridge", message( + "capability.host.request", HOST_REQ, + operation="storage.get", payload={"key": "x"} + )), + ("loom", message( + "capability.host.response", HOST_REQ, + ok=True, payload={} + )), + ("bridge", message( + "capability.callback.response", CALLBACK_REQ, + ok=True, payload={} + )), + ] + for origin, value in cases: + with self.subTest(origin=origin, kind=value["kind"]): + router = CapabilityRouter(GEN) + with self.assertRaises(ProtocolError): + router.route(origin, value) + self.assertTrue(router.failed) + + def test_duplicate_response_fails_router(self): + router = CapabilityRouter(GEN) + router.route("bridge", message( + "capability.callback.request", CALLBACK_REQ, + operation="trust001.preflight.ping", payload={} + )) + response = message( + "capability.callback.response", CALLBACK_REQ, + ok=True, payload={"pong": True} + ) + router.route("loom", response) + with self.assertRaises(ProtocolError): + router.route("loom", response) + self.assertTrue(router.failed) + + def test_stale_generation_fails_router(self): + router = CapabilityRouter(GEN) + stale = message( + "capability.host.request", HOST_REQ, + operation="storage.get", payload={"key": "x"} + ) + stale["generation"] = "b" * 64 + with self.assertRaises(ProtocolError): + router.route("loom", stale) + self.assertTrue(router.failed) + + def test_class_specific_cancel_rejects_late_response(self): + router = CapabilityRouter(GEN) + router.route("bridge", message( + "capability.callback.request", CALLBACK_REQ, + operation="tool.execute", payload={} + )) + router.route("bridge", message( + "capability.callback.cancel", CALLBACK_REQ, + reason="interrupted" + )) + with self.assertRaises(ProtocolError): + router.route("loom", message( + "capability.callback.response", CALLBACK_REQ, + ok=True, payload={} + )) + self.assertTrue(router.failed) + + def test_close_requires_both_namespaces_settled(self): + router = CapabilityRouter(GEN) + router.route("loom", message( + "capability.host.request", HOST_REQ, + operation="storage.get", payload={"key": "x"} + )) + router.close_admission() + self.assertFalse(router.settled) + router.route("bridge", message( + "capability.host.response", HOST_REQ, + ok=True, payload={"value": None} + )) + self.assertTrue(router.settled) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_trust001_channel.py b/tests/test_trust001_channel.py new file mode 100644 index 0000000..4c1dbe1 --- /dev/null +++ b/tests/test_trust001_channel.py @@ -0,0 +1,67 @@ +from __future__ import annotations + +import os +import socket +import threading +import unittest + +from runner.trust001.channel import OneShotUnixListener, private_socket_path +from runner.trust001.protocol import ProtocolError + + +class ChannelTests(unittest.TestCase): + def test_listener_is_private_noninheritable_and_one_shot(self): + with OneShotUnixListener() as listener: + private_socket_path(listener) + self.assertFalse(listener.socket.get_inheritable()) + observed = {} + + def connect(): + client = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + client.connect(str(listener.path)) + observed["client_inheritable"] = client.get_inheritable() + client.sendall(b"x") + client.close() + + thread = threading.Thread(target=connect) + thread.start() + conn, credentials = listener.accept_once( + expected_uid=os.getuid() if hasattr(socket, "SO_PEERCRED") else None, + expected_gid=os.getgid() if hasattr(socket, "SO_PEERCRED") else None, + ) + try: + self.assertFalse(conn.get_inheritable()) + self.assertEqual(conn.recv(1), b"x") + if credentials is not None: + self.assertEqual(credentials.uid, os.getuid()) + self.assertEqual(credentials.gid, os.getgid()) + finally: + conn.close() + thread.join(2) + self.assertFalse(thread.is_alive()) + self.assertFalse(listener.path.exists()) + self.assertFalse(observed["client_inheritable"]) + + def test_wrong_expected_peer_is_rejected(self): + if not hasattr(socket, "SO_PEERCRED"): + self.skipTest("SO_PEERCRED unavailable") + with OneShotUnixListener() as listener: + thread = threading.Thread( + target=lambda: self._connect_and_close(str(listener.path)) + ) + thread.start() + with self.assertRaises(ProtocolError): + listener.accept_once(expected_uid=os.getuid() + 1) + thread.join(2) + + @staticmethod + def _connect_and_close(path: str) -> None: + client = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + try: + client.connect(path) + finally: + client.close() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_trust001_protocol.py b/tests/test_trust001_protocol.py new file mode 100644 index 0000000..982b619 --- /dev/null +++ b/tests/test_trust001_protocol.py @@ -0,0 +1,168 @@ +from __future__ import annotations + +import json +import unittest + +from runner.trust001.protocol import ( + EVIDENCE_KINDS, + FrameReader, + ProtocolError, + WIRE_VERSION, + encode, + validate_capability, + validate_evidence, +) +from runner.trust001.state import CapabilityLedger, EvidenceLedger + + +GEN = "a" * 64 +REQ = "b" * 32 + + +class ProtocolTests(unittest.TestCase): + def test_channel_message_kinds_do_not_overlap(self): + from runner.trust001.protocol import CAPABILITY_KINDS + self.assertFalse(CAPABILITY_KINDS & EVIDENCE_KINDS) + + def test_capability_rejects_evidence_frame(self): + value = { + "version": WIRE_VERSION, + "kind": "evidence.observation", + "generation": GEN, + "sequence": 1, + "payload": {"safe": True}, + } + with self.assertRaises(ProtocolError): + validate_capability(value) + + def test_evidence_rejects_capability_frame(self): + value = { + "version": WIRE_VERSION, + "kind": "capability.callback.response", + "generation": GEN, + "request_id": REQ, + "ok": True, + "payload": {"safe": True}, + } + with self.assertRaises(ProtocolError): + validate_evidence(value) + + def test_frame_reader_handles_fragmented_input(self): + message = { + "version": WIRE_VERSION, + "kind": "evidence.observation", + "generation": GEN, + "sequence": 1, + "payload": {"value": "public"}, + } + raw = encode(message) + reader = FrameReader() + self.assertEqual(reader.feed(raw[:3], channel="evidence"), []) + self.assertEqual(reader.feed(raw[3:9], channel="evidence"), []) + self.assertEqual(reader.feed(raw[9:], channel="evidence"), [message]) + + def test_duplicate_json_keys_rejected(self): + raw = ( + '{"version":"' + WIRE_VERSION + '","kind":"evidence.hello",' + '"generation":"' + GEN + '","role":"bridge","role":"bridge"}' + ).encode() + framed = len(raw).to_bytes(4, "big") + raw + with self.assertRaises(ProtocolError): + FrameReader().feed(framed, channel="evidence") + + def test_host_and_callback_operations_are_separate(self): + host = { + "version": WIRE_VERSION, + "kind": "capability.host.request", + "generation": GEN, + "request_id": REQ, + "operation": "storage.get", + "payload": {"key": "x"}, + } + callback = { + "version": WIRE_VERSION, + "kind": "capability.callback.request", + "generation": GEN, + "request_id": REQ, + "operation": "tool.execute.before", + "payload": {}, + } + self.assertEqual(validate_capability(host), host) + self.assertEqual(validate_capability(callback), callback) + with self.assertRaises(ProtocolError): + validate_capability({**host, "operation": "tool.execute.before"}) + with self.assertRaises(ProtocolError): + validate_capability({**callback, "operation": "storage.get"}) + + def test_capability_response_shape_is_strict(self): + good = { + "version": WIRE_VERSION, + "kind": "capability.callback.response", + "generation": GEN, + "request_id": REQ, + "ok": True, + "payload": {"value": 1}, + } + self.assertEqual(validate_capability(good), good) + for bad in ( + {**good, "error": "also present"}, + {**good, "unknown": 1}, + {**good, "request_id": "caller-picked prose"}, + ): + with self.subTest(bad=bad): + with self.assertRaises(ProtocolError): + validate_capability(bad) + + +class LifecycleTests(unittest.TestCase): + def test_capability_response_is_at_most_once(self): + ledger = CapabilityLedger(GEN) + ledger.admit(REQ, "session.get") + self.assertEqual(ledger.settle(REQ), "session.get") + with self.assertRaises(ProtocolError): + ledger.settle(REQ) + + def test_cancel_closes_request_to_late_response(self): + ledger = CapabilityLedger(GEN) + ledger.admit(REQ, "tool.execute") + self.assertEqual(ledger.cancel(REQ), "tool.execute") + with self.assertRaises(ProtocolError): + ledger.settle(REQ) + + def test_generation_close_blocks_new_requests(self): + ledger = CapabilityLedger(GEN) + ledger.close_admission() + with self.assertRaises(ProtocolError): + ledger.admit(REQ, "session.get") + + def test_seal_requires_contiguous_final_sequence(self): + ledger = EvidenceLedger(GEN) + ledger.observe(1, {"kind": "first"}) + ledger.observe(2, {"kind": "second"}) + ledger.seal(2) + self.assertTrue(ledger.complete) + + def test_sequence_gap_rejects_checkpoint(self): + ledger = EvidenceLedger(GEN) + with self.assertRaises(ProtocolError): + ledger.observe(2, {"kind": "gap"}) + self.assertTrue(ledger.invalid) + self.assertFalse(ledger.complete) + + def test_post_seal_observation_rejects_checkpoint(self): + ledger = EvidenceLedger(GEN) + ledger.observe(1, {"kind": "first"}) + ledger.seal(1) + with self.assertRaises(ProtocolError): + ledger.observe(2, {"kind": "late"}) + self.assertTrue(ledger.invalid) + self.assertFalse(ledger.complete) + + def test_empty_generation_can_seal_at_zero(self): + ledger = EvidenceLedger(GEN) + ledger.seal(0) + self.assertTrue(ledger.complete) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_trust001_runtime.py b/tests/test_trust001_runtime.py new file mode 100644 index 0000000..fd6b023 --- /dev/null +++ b/tests/test_trust001_runtime.py @@ -0,0 +1,67 @@ +from __future__ import annotations + +import os +import socket +import threading +import unittest + +from runner.trust001.protocol import WIRE_VERSION +from runner.trust001.runtime import ChannelSet +from runner.trust001.transport import FramedSocket + + +GEN = "a" * 64 + + +class RuntimeChannelTests(unittest.TestCase): + def test_three_endpoints_are_distinct_and_unlinked_after_admission(self): + with ChannelSet(generation=GEN) as channels: + paths = channels.paths() + roots = channels.mount_roots() + self.assertEqual(len(set(paths.values())), 3) + self.assertEqual(len(set(roots.values())), 3) + + clients = {} + + def connect(name, channel, role): + sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + sock.connect(str(paths[name])) + FramedSocket(sock, channel=channel).send({ + "version": WIRE_VERSION, + "kind": f"{channel}.hello", + "generation": GEN, + "role": role, + }) + clients[name] = sock + + # ChannelSet admits evidence, bridge capability, then Loom capability. + threads = [ + threading.Thread(target=connect, args=("evidence", "evidence", "bridge")), + threading.Thread(target=connect, args=("capability_bridge", "capability", "bridge")), + threading.Thread(target=connect, args=("capability_loom", "capability", "loom")), + ] + for thread in threads: + thread.start() + + admitted = channels.admit( + bridge_uid=os.getuid() if hasattr(socket, "SO_PEERCRED") else None, + bridge_gid=os.getgid() if hasattr(socket, "SO_PEERCRED") else None, + loom_uid=os.getuid() if hasattr(socket, "SO_PEERCRED") else None, + loom_gid=os.getgid() if hasattr(socket, "SO_PEERCRED") else None, + ) + for thread in threads: + thread.join(2) + self.assertFalse(thread.is_alive()) + + self.assertEqual(admitted.generation, GEN) + self.assertTrue(all(not path.exists() for path in paths.values())) + + for sock in clients.values(): + sock.close() + admitted.capability.bridge.sock.close() + admitted.capability.loom.sock.close() + admitted.evidence.channel.sock.close() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_trust001_source_closure.py b/tests/test_trust001_source_closure.py new file mode 100644 index 0000000..5c59330 --- /dev/null +++ b/tests/test_trust001_source_closure.py @@ -0,0 +1,65 @@ +from __future__ import annotations + +import json +import tempfile +import unittest +from pathlib import Path + +from runner.trust001.protocol import ProtocolError +from runner.trust001.source_closure import build_source_closure, verify_source_closure + + +class SourceClosureTests(unittest.TestCase): + def test_builds_read_only_bridge_only_closure(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + bridge = root / "bridge.ts" + bridge.write_text("export default { id: 'trust001-bridge', async setup() {} }\n") + destination = root / "config" + manifest = build_source_closure(destination=destination, bridge_source=bridge) + verify_source_closure(destination, manifest) + self.assertEqual(manifest["expected_external_plugin_ids"], ["loom"]) + self.assertTrue(manifest["project_config_disabled"]) + + def test_product_config_cannot_declare_external_plugins(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + bridge = root / "bridge.ts" + bridge.write_text("export default {}\n") + for key in ("plugin", "plugins"): + config = root / f"{key}.json" + config.write_text(json.dumps({key: ["hostile"]})) + with self.subTest(key=key): + with self.assertRaises(ProtocolError): + build_source_closure( + destination=root / f"out-{key}", + bridge_source=bridge, + product_config=config, + ) + + def test_symlinked_bridge_source_is_rejected(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + real = root / "real.ts" + real.write_text("export default {}\n") + link = root / "link.ts" + link.symlink_to(real) + with self.assertRaises(ProtocolError): + build_source_closure(destination=root / "out", bridge_source=link) + + def test_post_build_mutation_breaks_manifest(self): + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + bridge = root / "bridge.ts" + bridge.write_text("export default {}\n") + destination = root / "config" + manifest = build_source_closure(destination=destination, bridge_source=bridge) + plugin = destination / "plugins" / "loom.js" + plugin.chmod(0o644) + plugin.write_text("export default { id: 'replacement' }\n") + with self.assertRaises(ProtocolError): + verify_source_closure(destination, manifest) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_trust001_transport.py b/tests/test_trust001_transport.py new file mode 100644 index 0000000..eb0b5ec --- /dev/null +++ b/tests/test_trust001_transport.py @@ -0,0 +1,127 @@ +from __future__ import annotations + +import socket +import threading +import unittest + +from runner.trust001.protocol import ProtocolError, WIRE_VERSION, encode +from runner.trust001.transport import CapabilityRelay, EvidenceCollector, FramedSocket + + +GEN = "a" * 64 +REQ = "b" * 32 + + +class TransportTests(unittest.TestCase): + def test_capability_relay_admits_roles_and_round_trips_callback(self): + bridge_client, bridge_server = socket.socketpair() + loom_client, loom_server = socket.socketpair() + try: + FramedSocket(bridge_client, channel="capability").send({ + "version": WIRE_VERSION, + "kind": "capability.hello", + "generation": GEN, + "role": "bridge", + }) + FramedSocket(loom_client, channel="capability").send({ + "version": WIRE_VERSION, + "kind": "capability.hello", + "generation": GEN, + "role": "loom", + }) + relay = CapabilityRelay.admit(GEN, bridge_server, loom_server) + + bridge = FramedSocket(bridge_client, channel="capability") + loom = FramedSocket(loom_client, channel="capability") + request = { + "version": WIRE_VERSION, + "kind": "capability.callback.request", + "generation": GEN, + "request_id": REQ, + "operation": "trust001.preflight.ping", + "payload": {"value": "ping"}, + } + bridge.send(request) + self.assertEqual(relay.relay_bridge_once(), request) + self.assertEqual(loom.recv(), request) + + response = { + "version": WIRE_VERSION, + "kind": "capability.callback.response", + "generation": GEN, + "request_id": REQ, + "ok": True, + "payload": {"value": "pong"}, + } + loom.send(response) + self.assertEqual(relay.relay_loom_once(), response) + self.assertEqual(bridge.recv(), response) + finally: + for sock in (bridge_client, bridge_server, loom_client, loom_server): + sock.close() + + def test_evidence_collector_requires_contiguous_seal(self): + client, server = socket.socketpair() + try: + wire = FramedSocket(client, channel="evidence") + wire.send({ + "version": WIRE_VERSION, + "kind": "evidence.hello", + "generation": GEN, + "role": "bridge", + }) + collector = EvidenceCollector.admit(GEN, server) + wire.send({ + "version": WIRE_VERSION, + "kind": "evidence.observation", + "generation": GEN, + "sequence": 1, + "payload": {"type": "bridge.activated"}, + }) + collector.receive_once() + wire.send({ + "version": WIRE_VERSION, + "kind": "evidence.seal", + "generation": GEN, + "final_sequence": 1, + }) + collector.receive_once() + self.assertTrue(collector.ledger.complete) + finally: + client.close() + server.close() + + def test_evidence_post_seal_is_rejected(self): + client, server = socket.socketpair() + try: + wire = FramedSocket(client, channel="evidence") + wire.send({ + "version": WIRE_VERSION, + "kind": "evidence.hello", + "generation": GEN, + "role": "bridge", + }) + collector = EvidenceCollector.admit(GEN, server) + wire.send({ + "version": WIRE_VERSION, + "kind": "evidence.seal", + "generation": GEN, + "final_sequence": 0, + }) + collector.receive_once() + wire.send({ + "version": WIRE_VERSION, + "kind": "evidence.observation", + "generation": GEN, + "sequence": 1, + "payload": {"late": True}, + }) + with self.assertRaises(ProtocolError): + collector.receive_once() + finally: + client.close() + server.close() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/trust001/run_remote_context_preflight.py b/tests/trust001/run_remote_context_preflight.py new file mode 100644 index 0000000..4458f61 --- /dev/null +++ b/tests/trust001/run_remote_context_preflight.py @@ -0,0 +1,324 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import argparse +import json +import os +import queue +import selectors +import subprocess +import sys +import tempfile +import threading +import time +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from runner.trust001.runtime import ChannelSet +from runner.trust001.source_closure import build_source_closure, verify_source_closure + + +def container_script() -> str: + return r""" +import json +import os +import sys +sys.path.insert(0, "/opt/opencode-eval-runner") +from container.invoke import _start_preflight_server, _standalone_json_request, _stop_preflight_server + +env = dict(os.environ) +server, base, auth = _start_preflight_server(env, 30) +try: + created = _standalone_json_request( + base, + "/api/session", + method="POST", + payload={ + "title": "TRUST-001 remote-context preflight", + "location": {"directory": "/workspace"}, + }, + timeout=10, + authorization=auth, + ) + session = created.get("data", created) + _standalone_json_request( + base, + f"/api/session/{session['id']}/prompt", + method="POST", + payload={"text": "activate plugins only", "resume": False}, + timeout=10, + authorization=auth, + ) + inventory = _standalone_json_request( + base, + "/api/plugin?location%5Bdirectory%5D=%2Fworkspace", + timeout=10, + authorization=auth, + ) + plugins = inventory.get("data", inventory) + tools = _standalone_json_request( + base, + "/experimental/tool/ids?location%5Bdirectory%5D=%2Fworkspace", + timeout=10, + authorization=auth, + ) + print(json.dumps({"plugins": plugins, "tools": tools}, separators=(",", ":")), flush=True) + print("TRUST001_CONTAINER_READY", flush=True) + sys.stdin.readline() +finally: + _stop_preflight_server(server) +""" + + +def wait_for_marker(stream, marker: str, timeout: float) -> list[str]: + selector = selectors.DefaultSelector() + selector.register(stream, selectors.EVENT_READ) + deadline = time.monotonic() + timeout + lines: list[str] = [] + try: + while time.monotonic() < deadline: + events = selector.select(max(0.1, deadline - time.monotonic())) + if not events: + continue + line = stream.readline() + if line == "": + raise RuntimeError(f"{marker} stream closed before readiness") + lines.append(line.rstrip("\n")) + if marker in line: + return lines + finally: + selector.close() + raise RuntimeError(f"{marker} readiness timeout") + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--image", required=True) + args = parser.parse_args() + + generation = "b" * 64 + with tempfile.TemporaryDirectory(prefix="trust001-remote-context-") as tmp: + root = Path(tmp) + workspace = root / "workspace" + workspace.mkdir() + hostile = workspace / ".opencode" / "plugins" + hostile.mkdir(parents=True) + (hostile / "hostile.mjs").write_text( + "throw new Error('workspace plugin escape loaded');\n" + "export default { id: 'hostile', async setup() {} };\n", + encoding="utf-8", + ) + + product_config = root / "product-config.json" + product_config.write_text( + json.dumps({ + "$schema": "https://opencode.ai/config.json", + "model": "openai/preflight-no-inference", + }) + "\n", + encoding="utf-8", + ) + closure = root / "trusted-config" + manifest = build_source_closure( + destination=closure, + bridge_source=ROOT / "trust001" / "bridge.mjs", + product_config=product_config, + ) + verify_source_closure(closure, manifest) + + with ChannelSet(generation=generation) as channels: + paths = channels.paths() + roots = channels.mount_roots() + + isolated_env = dict(os.environ) + isolated_env.update({ + "TRUST001_GENERATION": generation, + "TRUST001_CAPABILITY_SOCKET": str(paths["capability_loom"]), + "TRUST001_LOOM_ENTRYPOINT": str(ROOT / "trust001" / "synthetic-remote-plugin.mjs"), + }) + isolated = subprocess.Popen( + ["node", str(ROOT / "trust001" / "isolated-host.mjs")], + cwd=ROOT, + env=isolated_env, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + + result_box: queue.Queue[object] = queue.Queue() + close_request = threading.Event() + + def run_channels() -> None: + try: + admitted = channels.admit( + bridge_uid=os.getuid(), + bridge_gid=os.getgid(), + loom_uid=os.getuid(), + loom_gid=os.getgid(), + timeout=30, + ) + stop, errors, threads = admitted.capability.serve() + if not close_request.wait(30): + raise RuntimeError("trusted close was not requested") + admitted.evidence.request_close() + while not admitted.evidence.ledger.sealed: + admitted.evidence.receive_once() + stop.wait(10) + for thread in threads: + thread.join(2) + if not errors.empty(): + raise errors.get_nowait() + result_box.put({ + "complete": admitted.evidence.ledger.complete, + "observations": admitted.evidence.ledger.observations, + "router_failed": admitted.capability.router.failed, + "host_outstanding": len(admitted.capability.router.host_calls.outstanding), + "callback_outstanding": len(admitted.capability.router.callbacks.outstanding), + }) + admitted.evidence.channel.sock.close() + except BaseException as exc: + result_box.put(exc) + + channel_thread = threading.Thread(target=run_channels, daemon=True) + channel_thread.start() + + uid = str(os.getuid()) + gid = str(os.getgid()) + command = [ + "docker", "run", "--rm", "-i", + "--network", "none", + "--read-only", + "--tmpfs", "/tmp:rw,exec,nosuid,nodev,size=256m", + "--cap-drop", "ALL", + "--security-opt", "no-new-privileges", + "--user", f"{uid}:{gid}", + "--workdir", "/workspace", + "--volume", f"{workspace}:/workspace:rw", + "--volume", f"{closure}:/trusted-opencode-config:ro", + "--volume", f"{roots['evidence']}:/run/trust001/evidence:rw", + "--volume", f"{roots['capability_bridge']}:/run/trust001/capability:rw", + "--env", "HOME=/tmp/home", + "--env", "XDG_DATA_HOME=/tmp/data", + "--env", "XDG_CACHE_HOME=/tmp/cache", + "--env", "XDG_STATE_HOME=/tmp/state", + "--env", "XDG_CONFIG_HOME=/tmp/config", + "--env", "OPENCODE_CONFIG_DIR=/trusted-opencode-config", + "--env", "OPENCODE_DISABLE_PROJECT_CONFIG=1", + "--env", "OPENCODE_DISABLE_AUTOUPDATE=1", + "--env", f"TRUST001_GENERATION={generation}", + "--env", "TRUST001_EVIDENCE_SOCKET=/run/trust001/evidence/evidence.sock", + "--env", "TRUST001_CAPABILITY_SOCKET=/run/trust001/capability/bridge.sock", + "--entrypoint", "python3", + args.image, + "-c", container_script(), + ] + proc = subprocess.Popen( + command, + cwd=ROOT, + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + bufsize=1, + ) + if proc.stdout is None or proc.stdin is None or isolated.stdout is None: + raise RuntimeError("preflight process pipes unavailable") + + try: + container_lines = wait_for_marker(proc.stdout, "TRUST001_CONTAINER_READY", 30) + isolated_lines = wait_for_marker(isolated.stdout, "TRUST001_READY", 30) + except RuntimeError as exc: + container_code = proc.poll() + isolated_code = isolated.poll() + if proc.poll() is None: + proc.terminate() + proc.communicate(timeout=5) + if isolated.poll() is None: + isolated.terminate() + isolated_out, _ = isolated.communicate(timeout=5) + stages = [ + line for line in isolated_out.splitlines() + if line.startswith("TRUST001_") + ] + raise RuntimeError( + f"provider-free readiness failed: container_exit={container_code} " + f"isolated_exit={isolated_code} markers={stages!r}" + ) from exc + close_request.set() + + channel_thread.join(15) + if channel_thread.is_alive(): + isolated.terminate() + isolated.communicate(timeout=5) + proc.terminate() + proc.communicate(timeout=5) + raise RuntimeError("channel service did not terminate") + outcome = result_box.get_nowait() + if isinstance(outcome, BaseException): + isolated.terminate() + isolated.communicate(timeout=5) + proc.terminate() + proc.communicate(timeout=5) + raise outcome + + proc.stdin.write("\n") + proc.stdin.flush() + remaining_out, remaining_err = proc.communicate(timeout=10) + container_output = "\n".join(container_lines) + "\n" + remaining_out + + fenced = False + try: + remaining_isolated_out, isolated_err = isolated.communicate(timeout=2) + except subprocess.TimeoutExpired: + fenced = True + isolated.terminate() + remaining_isolated_out, isolated_err = isolated.communicate(timeout=5) + isolated_out = "\n".join(isolated_lines) + "\n" + remaining_isolated_out + + if proc.returncode != 0: + raise RuntimeError(f"remote-context OpenCode preflight failed: exit={proc.returncode}") + if isolated.returncode not in (0, -15): + raise RuntimeError("isolated remote context failed") + if "TRUST001_READY" not in isolated_out: + raise RuntimeError("isolated remote context did not become ready") + + json_line = next((line for line in container_lines if line.startswith("{")), None) + if json_line is None: + raise RuntimeError("container preflight payload missing") + payload = json.loads(json_line) + plugins = payload["plugins"] + ids = {item.get("id") for item in plugins if isinstance(item, dict)} + if "loom" not in ids or "hostile" in ids: + raise RuntimeError(f"unexpected plugin inventory: {sorted(str(x) for x in ids)}") + if not outcome["complete"] or outcome["router_failed"]: + raise RuntimeError("remote-context generation did not close cleanly") + if outcome["host_outstanding"] or outcome["callback_outstanding"]: + raise RuntimeError("remote-context generation closed with outstanding requests") + + kinds = [ + item.get("type") + for item in outcome["observations"] + if isinstance(item, dict) + ] + if "bridge.activated" not in kinds: + raise RuntimeError("bridge activation observation missing") + + print(json.dumps({ + "schema": "trust001-remote-context-preflight/v1", + "passed": True, + "generation": generation, + "plugin_ids": sorted(str(x) for x in ids if x), + "isolated_ready": True, + "isolated_fenced_after_seal": fenced, + "evidence_complete": True, + "router_failed": False, + "project_config_disabled": manifest["project_config_disabled"], + "provider_inference": False, + "network": "none", + }, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/trust001/run_stock_bridge_preflight.py b/tests/trust001/run_stock_bridge_preflight.py new file mode 100644 index 0000000..3c4116d --- /dev/null +++ b/tests/trust001/run_stock_bridge_preflight.py @@ -0,0 +1,242 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import argparse +import json +import os +import queue +import subprocess +import sys +import tempfile +import threading +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from runner.trust001.runtime import ChannelSet +from runner.trust001.source_closure import build_source_closure, verify_source_closure + + +def container_script() -> str: + return r""" +import json +import os +import sys +sys.path.insert(0, "/opt/opencode-eval-runner") +from container.invoke import _start_preflight_server, _standalone_json_request, _stop_preflight_server + +env = dict(os.environ) +server, base, auth = _start_preflight_server(env, 30) +try: + created = _standalone_json_request( + base, + "/api/session", + method="POST", + payload={ + "title": "TRUST-001 provider-free activation", + "location": {"directory": "/workspace"}, + }, + timeout=10, + authorization=auth, + ) + session = created.get("data", created) + session_id = session["id"] + _standalone_json_request( + base, + f"/api/session/{session_id}/prompt", + method="POST", + payload={"text": "activate plugins only", "resume": False}, + timeout=10, + authorization=auth, + ) + inventory = _standalone_json_request( + base, + "/api/plugin?location%5Bdirectory%5D=%2Fworkspace", + timeout=10, + authorization=auth, + ) + plugins = inventory.get("data", inventory) + print(json.dumps({"plugins": plugins}, separators=(",", ":"))) +finally: + _stop_preflight_server(server) +""" + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--image", required=True) + args = parser.parse_args() + + generation = "a" * 64 + result_queue: queue.Queue[object] = queue.Queue() + + with tempfile.TemporaryDirectory(prefix="trust001-stock-preflight-") as tmp: + root = Path(tmp) + workspace = root / "workspace" + workspace.mkdir() + hostile = workspace / ".opencode" / "plugins" + hostile.mkdir(parents=True) + (hostile / "hostile.mjs").write_text( + "throw new Error('workspace plugin escape loaded');\n" + "export default { id: 'hostile', async setup() {} };\n", + encoding="utf-8", + ) + + product_config = root / "product-config.json" + product_config.write_text( + json.dumps({ + "$schema": "https://opencode.ai/config.json", + "model": "openai/preflight-no-inference", + }) + "\n", + encoding="utf-8", + ) + closure = root / "trusted-config" + manifest = build_source_closure( + destination=closure, + bridge_source=ROOT / "trust001" / "bridge.mjs", + product_config=product_config, + ) + verify_source_closure(closure, manifest) + + with ChannelSet(generation=generation) as channels: + paths = channels.paths() + roots = channels.mount_roots() + + peer_env = dict(os.environ) + peer_env.update({ + "TRUST001_GENERATION": generation, + "TRUST001_CAPABILITY_SOCKET": str(paths["capability_loom"]), + }) + peer = subprocess.Popen( + ["node", str(ROOT / "trust001" / "synthetic-loom-peer.mjs")], + cwd=ROOT, + env=peer_env, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + + result_box: queue.Queue[object] = queue.Queue() + + def admit_and_relay() -> None: + try: + admitted = channels.admit( + bridge_uid=os.getuid(), + bridge_gid=os.getgid(), + loom_uid=os.getuid(), + loom_gid=os.getgid(), + timeout=30, + ) + admitted.capability.relay_bridge_once() + admitted.capability.relay_loom_once() + saw_preflight = False + while not saw_preflight: + message = admitted.evidence.receive_once() + payload = message.get("payload") if isinstance(message, dict) else None + if isinstance(payload, dict) and payload.get("type") == "capability.preflight": + saw_preflight = True + admitted.evidence.request_close() + while not admitted.evidence.ledger.sealed: + admitted.evidence.receive_once() + result_box.put({ + "complete": admitted.evidence.ledger.complete, + "observations": admitted.evidence.ledger.observations, + }) + admitted.capability.bridge.sock.close() + admitted.capability.loom.sock.close() + admitted.evidence.channel.sock.close() + except BaseException as exc: + result_box.put(exc) + + relay = threading.Thread(target=admit_and_relay, daemon=True) + relay.start() + + uid = str(os.getuid()) + gid = str(os.getgid()) + command = [ + "docker", "run", "--rm", + "--network", "none", + "--read-only", + "--tmpfs", "/tmp:rw,exec,nosuid,nodev,size=256m", + "--cap-drop", "ALL", + "--security-opt", "no-new-privileges", + "--user", f"{uid}:{gid}", + "--workdir", "/workspace", + "--volume", f"{workspace}:/workspace:rw", + "--volume", f"{closure}:/trusted-opencode-config:ro", + "--volume", f"{roots['evidence']}:/run/trust001/evidence:rw", + "--volume", f"{roots['capability_bridge']}:/run/trust001/capability:rw", + "--env", "HOME=/tmp/home", + "--env", "XDG_DATA_HOME=/tmp/data", + "--env", "XDG_CACHE_HOME=/tmp/cache", + "--env", "XDG_STATE_HOME=/tmp/state", + "--env", "XDG_CONFIG_HOME=/tmp/config", + "--env", "OPENCODE_CONFIG_DIR=/trusted-opencode-config", + "--env", "OPENCODE_DISABLE_PROJECT_CONFIG=1", + "--env", "OPENCODE_DISABLE_AUTOUPDATE=1", + "--env", f"TRUST001_GENERATION={generation}", + "--env", "TRUST001_EVIDENCE_SOCKET=/run/trust001/evidence/evidence.sock", + "--env", "TRUST001_CAPABILITY_SOCKET=/run/trust001/capability/bridge.sock", + "--env", "TRUST001_PREFLIGHT=1", + "--entrypoint", "python3", + args.image, + "-c", container_script(), + ] + proc = subprocess.run( + command, + cwd=ROOT, + capture_output=True, + text=True, + timeout=60, + check=False, + ) + + peer_out, peer_err = peer.communicate(timeout=10) + relay.join(10) + if relay.is_alive(): + raise RuntimeError("channel relay did not terminate") + outcome = result_box.get_nowait() + if isinstance(outcome, BaseException): + raise outcome + if proc.returncode != 0: + raise RuntimeError( + f"stock OpenCode bridge preflight failed: exit={proc.returncode} " + f"stdout={proc.stdout!r} stderr={proc.stderr!r} peer={peer_err!r}" + ) + if peer.returncode != 0: + raise RuntimeError(f"synthetic Loom peer failed: {peer_err!r}") + + payload = json.loads(proc.stdout.strip().splitlines()[-1]) + plugins = payload["plugins"] + ids = {item.get("id") for item in plugins if isinstance(item, dict)} + if "loom" not in ids: + raise RuntimeError(f"bridge not active: {sorted(str(x) for x in ids)}") + if "hostile" in ids: + raise RuntimeError("workspace plugin escaped source closure") + if not outcome["complete"]: + raise RuntimeError("evidence generation did not seal completely") + kinds = [ + item.get("type") + for item in outcome["observations"] + if isinstance(item, dict) + ] + if "bridge.activated" not in kinds or "capability.preflight" not in kinds: + raise RuntimeError(f"missing bridge observations: {kinds!r}") + + print(json.dumps({ + "schema": "trust001-stock-bridge-preflight/v1", + "passed": True, + "generation": generation, + "plugin_ids": sorted(str(x) for x in ids if x), + "observation_types": kinds, + "project_config_disabled": manifest["project_config_disabled"], + "external_plugin_declarations": manifest["external_plugin_declarations"], + "provider_inference": False, + "network": "none", + }, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/trust001/bridge.mjs b/trust001/bridge.mjs new file mode 100644 index 0000000..bd45f0e --- /dev/null +++ b/trust001/bridge.mjs @@ -0,0 +1,527 @@ +import net from "node:net" +import crypto from "node:crypto" + +const VERSION = "opencode-eval-runner/trust001-wire/v1" +const generation = process.env.TRUST001_GENERATION ?? "" +const evidencePath = process.env.TRUST001_EVIDENCE_SOCKET ?? "" +const capabilityPath = process.env.TRUST001_CAPABILITY_SOCKET ?? "" +const preflight = process.env.TRUST001_PREFLIGHT === "1" + +function requireValue(ok, message) { + if (!ok) throw new Error(message) +} + +requireValue(/^[0-9a-f]{64}$/.test(generation), "invalid TRUST001_GENERATION") +requireValue(evidencePath.startsWith("/"), "invalid TRUST001_EVIDENCE_SOCKET") +requireValue(capabilityPath.startsWith("/"), "invalid TRUST001_CAPABILITY_SOCKET") + +function requestID() { + return crypto.randomBytes(16).toString("hex") +} + +function handlerID(value) { + requireValue(typeof value === "string" && /^[0-9a-f]{32}$/.test(value), "invalid_handler_id") + return value +} + +function frame(value) { + const raw = Buffer.from(JSON.stringify(value), "utf8") + requireValue(raw.length > 0 && raw.length <= 256 * 1024, "frame_limit") + const prefix = Buffer.allocUnsafe(4) + prefix.writeUInt32BE(raw.length, 0) + return Buffer.concat([prefix, raw]) +} + +class Reader { + constructor(socket) { + this.socket = socket + this.buffer = Buffer.alloc(0) + this.pending = [] + } + + next() { + if (this.pending.length) return Promise.resolve(this.pending.shift()) + return new Promise((resolve, reject) => { + const onError = () => cleanup(() => reject(new Error("channel_error"))) + const onClose = () => cleanup(() => reject(new Error("channel_closed"))) + const onData = (chunk) => { + this.buffer = Buffer.concat([this.buffer, chunk]) + while (this.buffer.length >= 4) { + const size = this.buffer.readUInt32BE(0) + if (size < 1 || size > 256 * 1024) { + cleanup(() => reject(new Error("frame_limit"))) + return + } + if (this.buffer.length < 4 + size) return + const raw = this.buffer.subarray(4, 4 + size) + this.buffer = this.buffer.subarray(4 + size) + let value + try { + value = JSON.parse(raw.toString("utf8")) + } catch { + cleanup(() => reject(new Error("invalid_json"))) + return + } + this.pending.push(value) + } + if (this.pending.length) cleanup(() => resolve(this.pending.shift())) + } + const cleanup = (finish) => { + this.socket.off("data", onData) + this.socket.off("error", onError) + this.socket.off("close", onClose) + finish() + } + this.socket.on("data", onData) + this.socket.once("error", onError) + this.socket.once("close", onClose) + }) + } +} + +function connect(path) { + return new Promise((resolve, reject) => { + const socket = net.createConnection({ path }) + socket.once("connect", () => resolve(socket)) + socket.once("error", reject) + }) +} + +function send(socket, value) { + socket.write(frame(value)) +} + +function fixedError(error, fallback) { + if (error instanceof Error && error.message && error.message.length <= 256) return error.message + return fallback +} + +function plainContext(tool) { + return { + sessionID: tool.sessionID, + agent: tool.agent, + messageID: tool.messageID, + id: tool.id, + } +} + +export default { + id: "loom", + + async setup(ctx) { + const evidence = await connect(evidencePath) + const capability = await connect(capabilityPath) + const capabilityReader = new Reader(capability) + const evidenceReader = new Reader(evidence) + + send(evidence, { + version: VERSION, + kind: "evidence.hello", + generation, + role: "bridge", + }) + send(capability, { + version: VERSION, + kind: "capability.hello", + generation, + role: "bridge", + }) + + let sequence = 0 + let closing = false + let closed = false + let capabilityFailure = null + const eventController = new AbortController() + const pendingCallbacks = new Map() + const activeHostRequests = new Set() + const cancelledHostRequests = new Set() + const remoteToolHandlers = new Map() + const registrations = [] + + const observe = (payload) => { + if (closed) throw new Error("evidence_generation_closed") + sequence += 1 + send(evidence, { + version: VERSION, + kind: "evidence.observation", + generation, + sequence, + payload, + }) + } + + const callback = (operation, payload, signal) => { + requireValue(!closed && !closing && !capabilityFailure, "capability_unavailable") + const id = requestID() + return new Promise((resolve, reject) => { + const abort = () => { + if (!pendingCallbacks.has(id)) return + pendingCallbacks.delete(id) + send(capability, { + version: VERSION, + kind: "capability.callback.cancel", + generation, + request_id: id, + reason: "interrupted", + }) + reject(new Error("remote_callback_cancelled")) + } + pendingCallbacks.set(id, { resolve, reject, abort, signal }) + signal?.addEventListener("abort", abort, { once: true }) + send(capability, { + version: VERSION, + kind: "capability.callback.request", + generation, + request_id: id, + operation, + payload, + }) + }) + } + + const respondHost = (message, ok, payload, error) => { + if (cancelledHostRequests.has(message.request_id)) return + send(capability, { + version: VERSION, + kind: "capability.host.response", + generation, + request_id: message.request_id, + ok, + ...(ok ? { payload } : { error }), + }) + } + + const registerAgentTransform = async (payload) => { + const operations = payload?.operations + requireValue(Array.isArray(operations) && operations.length <= 16, "invalid_agent_transform") + for (const op of operations) { + requireValue(op && op.kind === "default" && typeof op.id === "string", "unsupported_agent_transform") + } + const registration = await ctx.agent.transform((editor) => { + for (const op of operations) editor.default(op.id) + }) + registrations.push(registration) + return { registered: true } + } + + const effectiveToolID = (definition) => { + const name = definition.name.replace(/[^a-zA-Z0-9_-]/g, "_") + const namespace = definition.options?.namespace + return namespace === undefined ? name : namespace.replaceAll(".", "_") + "_" + name + } + + const registerToolTransform = async (payload) => { + const namespaces = payload?.namespaces ?? [] + const tools = payload?.tools ?? [] + requireValue(Array.isArray(namespaces) && namespaces.length <= 8, "invalid_tool_transform") + requireValue(Array.isArray(tools) && tools.length <= 128, "invalid_tool_transform") + const registration = await ctx.tool.transform((editor) => { + for (const namespace of namespaces) { + requireValue( + namespace && typeof namespace.name === "string" && typeof namespace.description === "string", + "invalid_tool_namespace", + ) + editor.namespace(namespace) + } + for (const item of tools) { + requireValue(item && typeof item.definition === "object", "invalid_tool_definition") + const remoteHandler = handlerID(item.handler_id) + const definition = item.definition + requireValue(typeof definition.name === "string" && typeof definition.description === "string", + "invalid_tool_definition") + requireValue(definition.input && typeof definition.input === "object", "invalid_tool_definition") + const execute = async (input, tool) => { + const result = await callback( + "tool.execute", + { + handler_id: remoteHandler, + input, + context: plainContext(tool), + }, + tool.signal, + ) + requireValue(result && typeof result === "object", "invalid_remote_tool_result") + return result + } + const effectiveID = effectiveToolID(definition) + requireValue(!remoteToolHandlers.has(effectiveID), "duplicate_remote_tool_id") + remoteToolHandlers.set(effectiveID, remoteHandler) + editor.add({ ...definition, execute }) + } + }) + registrations.push(registration) + return { registered: true } + } + + const listRemoteTools = async () => { + const tools = await ctx.tool.list() + const actual = new Set(tools.map((tool) => tool.id)) + return { + tools: Array.from(remoteToolHandlers.entries()).flatMap(([id, remoteHandler]) => + actual.has(id) ? [{ id, handler_id: remoteHandler }] : [] + ), + } + } + + const registerToolHook = async (payload) => { + const name = payload?.name + const remoteHandler = handlerID(payload?.handler_id) + requireValue(name === "execute.before" || name === "execute.after", "unsupported_tool_hook") + const registration = await ctx.tool.hook(name, async (event) => { + await callback(name === "execute.before" ? "tool.execute.before" : "tool.execute.after", { + handler_id: remoteHandler, + event, + }) + }) + registrations.push(registration) + return { registered: true } + } + + const registerPermissionHook = async (payload) => { + requireValue(payload?.name === "evaluate", "unsupported_permission_hook") + const remoteHandler = handlerID(payload?.handler_id) + const registration = await ctx.permission.hook("evaluate", async (event) => { + const result = await callback("permission.evaluate", { handler_id: remoteHandler, event }) + const next = result?.event + requireValue(next && typeof next === "object", "invalid_permission_callback") + requireValue(["allow", "deny", "ask"].includes(next.effect), "invalid_permission_effect") + event.effect = next.effect + if (next.message === undefined) delete event.message + else { + requireValue(typeof next.message === "string", "invalid_permission_message") + event.message = next.message + } + }) + registrations.push(registration) + return { registered: true } + } + + const registerSessionHook = async (payload) => { + const name = payload?.name + const remoteHandler = handlerID(payload?.handler_id) + requireValue(name === "context" || name === "retry", "unsupported_session_hook") + const registration = await ctx.session.hook(name, async (event) => { + const result = await callback(name === "context" ? "session.context" : "session.retry", { + handler_id: remoteHandler, + event, + }) + const next = result?.event + requireValue(next && typeof next === "object", "invalid_session_callback") + if (name === "context") { + requireValue(Array.isArray(next.system), "invalid_session_context") + event.system = next.system + } else { + const decision = next.decision + requireValue( + decision && typeof decision === "object" && typeof decision.retry === "boolean" && + (decision.retry === false || typeof decision.delay === "number"), + "invalid_retry_decision", + ) + event.decision = decision + } + }) + registrations.push(registration) + return { registered: true } + } + + const registerRpc = async (payload) => { + const definition = payload?.definition + const handlers = payload?.handlers + requireValue(definition && typeof definition === "object" && typeof definition.id === "string", + "invalid_rpc_definition") + requireValue(handlers && typeof handlers === "object", "invalid_rpc_handlers") + const remoteHandlers = Object.fromEntries( + Object.entries(handlers).map(([method, value]) => { + const remoteHandler = handlerID(value) + return [method, async (input) => callback("rpc.call", { + handler_id: remoteHandler, + method, + input, + })] + }), + ) + const registration = await ctx.rpc.register(definition, remoteHandlers) + registrations.push(registration) + return { registered: true } + } + + const dispatchHost = async (message) => { + requireValue(message.generation === generation, "stale_generation") + const payload = message.payload ?? {} + switch (message.operation) { + case "location.get": + return { location: ctx.location } + case "storage.get": { + requireValue(typeof payload.key === "string" && payload.key.length <= 4096, "invalid_storage_key") + const value = await ctx.storage.get(payload.key) + return value === undefined ? { present: false } : { present: true, value } + } + case "storage.set": + requireValue(typeof payload.key === "string" && payload.key.length <= 4096, "invalid_storage_key") + await ctx.storage.set(payload.key, payload.value) + return { stored: true } + case "storage.scan": { + requireValue(typeof payload.prefix === "string" && payload.prefix.length <= 4096, "invalid_storage_prefix") + const limit = payload.limit === undefined ? undefined : payload.limit + requireValue(limit === undefined || Number.isInteger(limit) && limit >= 1 && limit <= 1000, + "invalid_storage_limit") + requireValue(payload.after === undefined || typeof payload.after === "string", "invalid_storage_after") + return await ctx.storage.scan({ + prefix: payload.prefix, + ...(limit === undefined ? {} : { limit }), + ...(payload.after === undefined ? {} : { after: payload.after }), + }) + } + case "agent.list": + return await ctx.agent.list(payload.input) + case "agent.transform.register": + return await registerAgentTransform(payload) + case "tool.transform.register": + return await registerToolTransform(payload) + case "tool.list": + return await listRemoteTools() + case "tool.hook.register": + return await registerToolHook(payload) + case "permission.hook.register": + return await registerPermissionHook(payload) + case "session.hook.register": + return await registerSessionHook(payload) + case "session.get": + return { value: await ctx.session.get(payload.input) } + case "session.context": + return { value: await ctx.session.context(payload.input) } + case "session.synthetic": + return { value: await ctx.session.synthetic(payload.input) } + case "rpc.register": + return await registerRpc(payload) + default: + throw new Error("operation_not_admitted") + } + } + + const capabilityLoop = (async () => { + try { + while (!closed) { + const message = await capabilityReader.next() + requireValue(message?.version === VERSION, "wrong_version") + requireValue(message?.generation === generation, "stale_generation") + if (message.kind === "capability.callback.response") { + const pending = pendingCallbacks.get(message.request_id) + requireValue(pending, "unknown_or_late_response") + pendingCallbacks.delete(message.request_id) + pending.signal?.removeEventListener("abort", pending.abort) + if (message.ok === true) pending.resolve(message.payload) + else pending.reject(new Error(fixedError(message.error, "remote_callback_failed"))) + continue + } + if (message.kind === "capability.host.cancel") { + cancelledHostRequests.add(message.request_id) + continue + } + requireValue(message.kind === "capability.host.request", "capability_direction_violation") + if (closing) { + capabilityFailure = capabilityFailure ?? "post_close_request" + respondHost(message, false, undefined, "generation_closing") + continue + } + const id = message.request_id + requireValue(typeof id === "string" && /^[0-9a-f]{32}$/.test(id), "invalid_request_id") + requireValue(!activeHostRequests.has(id) && !cancelledHostRequests.has(id), "duplicate_request") + activeHostRequests.add(id) + void dispatchHost(message).then( + (result) => respondHost(message, true, result), + (error) => respondHost(message, false, undefined, fixedError(error, "host_operation_failed")), + ).finally(() => activeHostRequests.delete(id)) + } + } catch (error) { + if (!closed) { + capabilityFailure = fixedError(error, "capability_channel_failed") + for (const pending of pendingCallbacks.values()) { + pending.reject(new Error("capability_channel_failed")) + } + pendingCallbacks.clear() + } + } + })() + + observe({ + type: "bridge.activated", + implementation: "trust001-bridge", + plugin: "loom", + opencode: ctx.app?.version ?? null, + location: { + directory: ctx.location?.directory ?? null, + workspaceID: ctx.location?.workspaceID ?? null, + }, + }) + + const eventTask = (async () => { + try { + for await (const event of ctx.event.subscribe({ signal: eventController.signal })) { + observe({ type: "opencode.event", event }) + } + } catch { + if (!eventController.signal.aborted) capabilityFailure = capabilityFailure ?? "event_stream_failed" + } + })() + + const closeGeneration = async () => { + if (closed || closing) return + closing = true + eventController.abort() + await eventTask.catch(() => undefined) + const eligible = + !capabilityFailure && + pendingCallbacks.size === 0 && + activeHostRequests.size === 0 + if (eligible) { + send(evidence, { + version: VERSION, + kind: "evidence.seal", + generation, + final_sequence: sequence, + }) + } else { + capabilityFailure = capabilityFailure ?? "close_with_outstanding_work" + } + closed = true + capability.end() + evidence.end() + } + + const evidenceControlTask = (async () => { + try { + const message = await evidenceReader.next() + requireValue(message?.version === VERSION, "wrong_evidence_control_version") + requireValue(message?.generation === generation, "stale_evidence_control_generation") + requireValue(message?.kind === "evidence.close", "unexpected_evidence_control") + await closeGeneration() + } catch { + if (!closed) { + capabilityFailure = capabilityFailure ?? "evidence_control_failed" + closed = true + capability.destroy() + evidence.destroy() + } + } + })() + + if (preflight) { + const response = await callback("trust001.preflight.ping", { value: "ping" }) + requireValue(response?.value === "pong", "invalid_preflight_response") + observe({ type: "capability.preflight", result: "pong" }) + } + + return async () => { + if (!closed) { + eventController.abort() + await eventTask.catch(() => undefined) + closed = true + capability.end() + evidence.end() + } + await capabilityLoop.catch(() => undefined) + await evidenceControlTask.catch(() => undefined) + } + }, +} diff --git a/trust001/isolated-host.mjs b/trust001/isolated-host.mjs new file mode 100644 index 0000000..91533b6 --- /dev/null +++ b/trust001/isolated-host.mjs @@ -0,0 +1,392 @@ +import net from "node:net" +import crypto from "node:crypto" +import { pathToFileURL } from "node:url" + +const VERSION = "opencode-eval-runner/trust001-wire/v1" +const generation = process.env.TRUST001_GENERATION ?? "" +const capabilityPath = process.env.TRUST001_CAPABILITY_SOCKET ?? "" +const loomEntrypoint = process.env.TRUST001_LOOM_ENTRYPOINT ?? "" + +function requireValue(ok, code) { + if (!ok) throw new Error(code) +} + +requireValue(/^[0-9a-f]{64}$/.test(generation), "invalid_generation") +requireValue(capabilityPath.startsWith("/"), "invalid_capability_socket") +requireValue(loomEntrypoint.startsWith("/"), "invalid_loom_entrypoint") + +function requestID() { + return crypto.randomBytes(16).toString("hex") +} + +function handlerID() { + return crypto.randomBytes(16).toString("hex") +} + +function frame(value) { + const raw = Buffer.from(JSON.stringify(value), "utf8") + requireValue(raw.length > 0 && raw.length <= 256 * 1024, "frame_limit") + const prefix = Buffer.allocUnsafe(4) + prefix.writeUInt32BE(raw.length, 0) + return Buffer.concat([prefix, raw]) +} + +class Reader { + constructor(socket) { + this.socket = socket + this.buffer = Buffer.alloc(0) + this.pending = [] + } + + next() { + if (this.pending.length) return Promise.resolve(this.pending.shift()) + return new Promise((resolve, reject) => { + const onError = () => cleanup(() => reject(new Error("channel_error"))) + const onClose = () => cleanup(() => reject(new Error("channel_closed"))) + const onData = (chunk) => { + this.buffer = Buffer.concat([this.buffer, chunk]) + while (this.buffer.length >= 4) { + const size = this.buffer.readUInt32BE(0) + if (size < 1 || size > 256 * 1024) { + cleanup(() => reject(new Error("frame_limit"))) + return + } + if (this.buffer.length < 4 + size) return + const raw = this.buffer.subarray(4, 4 + size) + this.buffer = this.buffer.subarray(4 + size) + let value + try { + value = JSON.parse(raw.toString("utf8")) + } catch { + cleanup(() => reject(new Error("invalid_json"))) + return + } + this.pending.push(value) + } + if (this.pending.length) cleanup(() => resolve(this.pending.shift())) + } + const cleanup = (finish) => { + this.socket.off("data", onData) + this.socket.off("error", onError) + this.socket.off("close", onClose) + finish() + } + this.socket.on("data", onData) + this.socket.once("error", onError) + this.socket.once("close", onClose) + }) + } +} + +function connect(path) { + return new Promise((resolve, reject) => { + const socket = net.createConnection({ path }) + socket.once("connect", () => resolve(socket)) + socket.once("error", reject) + }) +} + +function send(socket, value) { + socket.write(frame(value)) +} + +function fixedError(error) { + if (error instanceof Error && error.message) return error.message.slice(0, 256) + return "remote_handler_failed" +} + +function clone(value) { + if (value === undefined) return undefined + return JSON.parse(JSON.stringify(value)) +} + +async function main() { + const socket = await connect(capabilityPath) + const reader = new Reader(socket) + const pendingHost = new Map() + const callbacks = new Map() + const activeCallbacks = new Map() + let closed = false + let loopFailure = null + + send(socket, { + version: VERSION, + kind: "capability.hello", + generation, + role: "loom", + }) + + const callHost = (operation, payload) => { + requireValue(!closed && !loopFailure, "capability_unavailable") + const id = requestID() + return new Promise((resolve, reject) => { + pendingHost.set(id, { resolve, reject }) + send(socket, { + version: VERSION, + kind: "capability.host.request", + generation, + request_id: id, + operation, + payload, + }) + }) + } + + const remember = (handler) => { + requireValue(typeof handler === "function", "invalid_handler") + const id = handlerID() + callbacks.set(id, handler) + return id + } + + const callbackResult = async (message) => { + const operation = message.operation + if (operation === "trust001.preflight.ping") { + return { value: "pong" } + } + + const id = message.payload?.handler_id + requireValue(typeof id === "string" && callbacks.has(id), "unknown_handler") + const handler = callbacks.get(id) + const controller = new AbortController() + activeCallbacks.set(message.request_id, controller) + try { + if (operation === "tool.execute") { + const ctx = { + ...message.payload.context, + signal: controller.signal, + progress: async () => undefined, + } + return await handler(message.payload.input, ctx) + } + if (operation === "rpc.call") { + return await handler(message.payload.input, { + signal: controller.signal, + error: (type, text, data) => ({ type, message: text, data }), + }) + } + if ([ + "tool.execute.before", + "tool.execute.after", + "permission.evaluate", + "session.context", + "session.retry", + ].includes(operation)) { + const event = message.payload.event + await handler(event) + return { event } + } + throw new Error("operation_not_admitted") + } finally { + activeCallbacks.delete(message.request_id) + } + } + + const loop = (async () => { + try { + while (!closed) { + const message = await reader.next() + requireValue(message?.version === VERSION, "wrong_version") + requireValue(message?.generation === generation, "stale_generation") + + if (message.kind === "capability.host.response") { + const pending = pendingHost.get(message.request_id) + requireValue(pending, "unknown_or_late_host_response") + pendingHost.delete(message.request_id) + if (message.ok === true) pending.resolve(message.payload) + else pending.reject(new Error(typeof message.error === "string" ? message.error : "host_operation_failed")) + continue + } + + if (message.kind === "capability.callback.cancel") { + const controller = activeCallbacks.get(message.request_id) + if (controller) controller.abort() + continue + } + + requireValue(message.kind === "capability.callback.request", "capability_direction_violation") + void callbackResult(message).then( + (payload) => send(socket, { + version: VERSION, + kind: "capability.callback.response", + generation, + request_id: message.request_id, + ok: true, + payload: payload === undefined ? null : payload, + }), + (error) => send(socket, { + version: VERSION, + kind: "capability.callback.response", + generation, + request_id: message.request_id, + ok: false, + error: fixedError(error), + }), + ) + } + } catch (error) { + if (!closed && error instanceof Error && error.message === "channel_closed" && pendingHost.size === 0) { + closed = true + return + } + if (!closed) { + loopFailure = fixedError(error) + for (const pending of pendingHost.values()) pending.reject(new Error("capability_channel_failed")) + pendingHost.clear() + } + } + })() + + const locationReply = await callHost("location.get", {}) + const location = locationReply?.location + requireValue(location && typeof location.directory === "string", "invalid_location") + + const storage = { + async get(key) { + const reply = await callHost("storage.get", { key }) + return reply?.present ? reply.value : undefined + }, + async set(key, value) { + await callHost("storage.set", { key, value }) + }, + async scan(input) { + return await callHost("storage.scan", input) + }, + } + + const agent = { + async list(input) { + return await callHost("agent.list", { input }) + }, + async transform(transform) { + const snapshot = await agent.list() + const agents = clone(snapshot?.data ?? []) + const operations = [] + const editor = { + list: () => agents, + get: (id) => agents.find((item) => item?.id === id), + default: (id) => operations.push({ kind: "default", id }), + update: () => { throw new Error("unsupported_agent_transform_update") }, + remove: () => { throw new Error("unsupported_agent_transform_remove") }, + } + const result = transform(editor) + requireValue(!result || typeof result.then !== "function", "async_agent_transform_unsupported") + await callHost("agent.transform.register", { operations }) + return { dispose: async () => undefined } + }, + } + + const toolHandlers = new Map() + const tool = { + async transform(transform) { + const namespaces = [] + const tools = [] + const editor = { + list: () => [], + get: () => undefined, + namespace: (namespace) => namespaces.push(clone(namespace)), + add: (definition) => { + requireValue(definition && typeof definition.execute === "function", "invalid_tool_definition") + const id = remember(definition.execute) + toolHandlers.set(id, definition.execute) + const { execute: _execute, ...serializable } = definition + tools.push({ definition: clone(serializable), handler_id: id }) + }, + update: () => { throw new Error("unsupported_tool_transform_update") }, + remove: () => { throw new Error("unsupported_tool_transform_remove") }, + } + const result = transform(editor) + requireValue(!result || typeof result.then !== "function", "async_tool_transform_unsupported") + await callHost("tool.transform.register", { namespaces, tools }) + return { dispose: async () => undefined } + }, + async list() { + const reply = await callHost("tool.list", {}) + return (reply?.tools ?? []).map((item) => ({ + id: item.id, + execute: toolHandlers.get(item.handler_id), + })) + }, + async hook(name, handler) { + const id = remember(handler) + await callHost("tool.hook.register", { name, handler_id: id }) + return { dispose: async () => undefined } + }, + } + + const permission = { + async hook(name, handler) { + const id = remember(handler) + await callHost("permission.hook.register", { name, handler_id: id }) + return { dispose: async () => undefined } + }, + } + + const session = { + async get(input) { + return (await callHost("session.get", { input }))?.value + }, + async context(input) { + return (await callHost("session.context", { input }))?.value + }, + async synthetic(input) { + return (await callHost("session.synthetic", { input }))?.value + }, + async hook(name, handler) { + const id = remember(handler) + await callHost("session.hook.register", { name, handler_id: id }) + return { dispose: async () => undefined } + }, + } + + const rpc = { + async register(definition, handlers) { + const ids = Object.fromEntries( + Object.entries(handlers).map(([name, handler]) => [name, remember(handler)]), + ) + await callHost("rpc.register", { + definition: clone(definition), + handlers: ids, + }) + return { + dispose: async () => undefined, + events: { emit: async () => { throw new Error("rpc_event_emit_not_admitted") } }, + } + }, + } + + const ctx = { + app: {}, + location, + options: {}, + storage, + rpc, + agent, + tool, + permission, + session, + } + + let cleanup + try { + const module = await import(pathToFileURL(loomEntrypoint).href) + const plugin = module.default + requireValue(plugin && plugin.id === "loom" && typeof plugin.setup === "function", "invalid_loom_plugin") + cleanup = await plugin.setup(ctx) + process.stdout.write("TRUST001_READY\n") + await loop + } finally { + closed = true + if (typeof cleanup === "function") { + try { await cleanup() } catch {} + } + socket.end() + } + + requireValue(!loopFailure, "capability_channel_failed") +} + +main().catch(() => { + process.stderr.write("TRUST001_FAILED\n") + process.exitCode = 2 +}) diff --git a/trust001/synthetic-loom-peer.mjs b/trust001/synthetic-loom-peer.mjs new file mode 100644 index 0000000..0628464 --- /dev/null +++ b/trust001/synthetic-loom-peer.mjs @@ -0,0 +1,67 @@ +import net from "node:net" + +const VERSION = "opencode-eval-runner/trust001-wire/v1" +const generation = process.env.TRUST001_GENERATION ?? "" +const socketPath = process.env.TRUST001_CAPABILITY_SOCKET ?? "" + +if (!/^[0-9a-f]{64}$/.test(generation)) throw new Error("invalid generation") +if (!socketPath.startsWith("/")) throw new Error("invalid socket path") + +function frame(value) { + const raw = Buffer.from(JSON.stringify(value), "utf8") + if (raw.length < 1 || raw.length > 256 * 1024) throw new Error("frame_limit") + const prefix = Buffer.allocUnsafe(4) + prefix.writeUInt32BE(raw.length, 0) + return Buffer.concat([prefix, raw]) +} + +function send(socket, value) { + socket.write(frame(value)) +} + +let buffer = Buffer.alloc(0) +const socket = net.createConnection({ path: socketPath }) + +socket.on("connect", () => { + send(socket, { + version: VERSION, + kind: "capability.hello", + generation, + role: "loom", + }) +}) + +socket.on("data", (chunk) => { + buffer = Buffer.concat([buffer, chunk]) + while (buffer.length >= 4) { + const size = buffer.readUInt32BE(0) + if (size < 1 || size > 256 * 1024) throw new Error("frame_limit") + if (buffer.length < 4 + size) return + const message = JSON.parse(buffer.subarray(4, 4 + size).toString("utf8")) + buffer = buffer.subarray(4 + size) + + if ( + message.version !== VERSION || + message.generation !== generation || + message.kind !== "capability.callback.request" || + message.operation !== "trust001.preflight.ping" || + !/^[0-9a-f]{32}$/.test(message.request_id) + ) { + throw new Error("unexpected capability request") + } + + send(socket, { + version: VERSION, + kind: "capability.callback.response", + generation, + request_id: message.request_id, + ok: true, + payload: { value: "pong" }, + }) + } +}) + +socket.on("error", (error) => { + process.stderr.write("synthetic peer failed: " + error.message + "\n") + process.exitCode = 2 +}) diff --git a/trust001/synthetic-remote-plugin.mjs b/trust001/synthetic-remote-plugin.mjs new file mode 100644 index 0000000..1ecb1f2 --- /dev/null +++ b/trust001/synthetic-remote-plugin.mjs @@ -0,0 +1,88 @@ +export default { + id: "loom", + async setup(ctx) { + if (!ctx.location || typeof ctx.location.directory !== "string") { + throw new Error("missing location") + } + + await ctx.storage.set("trust001/synthetic", { ready: true }) + const stored = await ctx.storage.get("trust001/synthetic") + if (!stored || stored.ready !== true) throw new Error("storage round trip failed") + const page = await ctx.storage.scan({ prefix: "trust001/", limit: 10 }) + if (!Array.isArray(page.entries) || !page.entries.some((entry) => entry.key === "trust001/synthetic")) { + throw new Error("storage scan failed") + } + process.stdout.write("TRUST001_SYNTHETIC_STAGE:storage\n") + + await ctx.rpc.register({ + id: "trust001.synthetic", + methods: { + ping: { + input: { + type: "object", + properties: { value: { type: "string" } }, + required: ["value"], + additionalProperties: false, + }, + output: { + type: "object", + properties: { value: { type: "string" } }, + required: ["value"], + additionalProperties: false, + }, + }, + }, + events: {}, + }, { + ping: async (input) => ({ value: input.value }), + }) + process.stdout.write("TRUST001_SYNTHETIC_STAGE:rpc\n") + + await ctx.agent.transform((editor) => { + if (editor.get("general")) editor.default("general") + }) + process.stdout.write("TRUST001_SYNTHETIC_STAGE:agent\n") + + let expectedRoster + await ctx.tool.transform((editor) => { + editor.namespace({ name: "loom", description: "TRUST-001 synthetic namespace" }) + const execute = async () => ({ content: "synthetic-roster" }) + expectedRoster = execute + editor.add({ + name: "roster", + description: "Synthetic roster tool", + input: { + type: "object", + properties: {}, + additionalProperties: false, + }, + options: { namespace: "loom", codemode: false }, + execute, + }) + }) + process.stdout.write("TRUST001_SYNTHETIC_STAGE:tool-transform\n") + + const registrations = await ctx.tool.list() + const roster = registrations.filter((entry) => entry.id === "loom_roster") + if (roster.length !== 1 || roster[0].execute !== expectedRoster) { + throw new Error("tool identity failed") + } + process.stdout.write("TRUST001_SYNTHETIC_STAGE:tool-list\n") + + await ctx.permission.hook("evaluate", async (event) => { + if (event.action === "trust001.synthetic.deny") { + event.effect = "deny" + event.message = "synthetic-deny" + } + }) + await ctx.session.hook("context", async (event) => { + event.system.push({ type: "text", text: "trust001-synthetic-context" }) + }) + await ctx.session.hook("retry", async (event) => { + if (event.attempt >= 2) event.decision = { retry: false } + }) + await ctx.tool.hook("execute.before", async () => undefined) + await ctx.tool.hook("execute.after", async () => undefined) + process.stdout.write("TRUST001_SYNTHETIC_STAGE:hooks\n") + }, +}