diff --git a/container/invoke.py b/container/invoke.py index 95dadfe..2df9e5b 100644 --- a/container/invoke.py +++ b/container/invoke.py @@ -18,7 +18,7 @@ try: from .evidence_safety import OMITTED, Projection, Sanitizer - from .native_observer import OBSERVATION_PATH, load_runtime_observations + from .native_observer import OBSERVER_STREAM_ENV, RuntimeObservationTransport from .runtime_evidence import ( build_runtime_evidence, unsupported_runtime_evidence, @@ -26,7 +26,7 @@ ) except ImportError: from evidence_safety import OMITTED, Projection, Sanitizer - from native_observer import OBSERVATION_PATH, load_runtime_observations + from native_observer import OBSERVER_STREAM_ENV, RuntimeObservationTransport from runtime_evidence import ( build_runtime_evidence, unsupported_runtime_evidence, @@ -827,8 +827,8 @@ def runtime_sanitizer(env: dict[str, str]) -> Sanitizer: def configure_runtime_observer_env(env: dict[str, str], sanitizer: Sanitizer) -> None: # The same inventory used by the Python pre-sink guard is supplied to the - # in-process observer so no raw dynamic credential value is written to its - # capture file before projection. + # in-process observer so no raw dynamic credential value reaches the + # runner-owned capture stream before projection. env["OPENCODE_EVAL_OBSERVER_CREDENTIALS"] = json.dumps( list(sanitizer.credentials), ensure_ascii=False, @@ -839,6 +839,17 @@ def configure_runtime_observer_env(env: dict[str, str], sanitizer: Sanitizer) -> ) +def finalize_runtime_observer( + transport: RuntimeObservationTransport, + env: dict[str, str], +) -> dict[str, Any]: + """Drain and close the runner-owned observer stream exactly once.""" + try: + return transport.finish() + finally: + env.pop(OBSERVER_STREAM_ENV, None) + + def sanitize_result_for_output( result: dict[str, Any], sanitizer: Sanitizer, @@ -892,9 +903,10 @@ def invoke_opencode( command += ["--agent", agent] command += ["--model", invoked_model, prompt] - # The preflight process can activate the observer. Never mix its records - # with the real invocation. - OBSERVATION_PATH.unlink(missing_ok=True) + # Preflight intentionally receives no observer endpoint. Create a + # one-connection runner-owned stream only for the real invocation. + observer_transport = RuntimeObservationTransport() + env[OBSERVER_STREAM_ENV] = observer_transport.endpoint run_started = time.perf_counter() try: proc = run(command, Path("/workspace"), env, timeout) @@ -905,7 +917,7 @@ def invoke_opencode( events = parse_events(raw_stdout) safe_stdout, _ = sanitizer.json_lines(raw_stdout) sid = session_id(events) - capture = load_runtime_observations() + capture = finalize_runtime_observer(observer_transport, env) runtime_evidence = build_runtime_evidence( capture, sanitizer, @@ -953,13 +965,16 @@ def invoke_opencode( "plugin_preflight": plugin_preflight, } return sanitize_result_for_output(result, sanitizer) + except BaseException: + finalize_runtime_observer(observer_transport, env) + raise run_seconds = time.perf_counter() - run_started events = parse_events(proc.stdout) safe_stdout, _ = sanitizer.json_lines(proc.stdout) safe_stderr, _ = sanitizer.json_lines(proc.stderr) sid = session_id(events) - capture = load_runtime_observations() + capture = finalize_runtime_observer(observer_transport, env) runtime_evidence = build_runtime_evidence( capture, sanitizer, diff --git a/container/native_observer.py b/container/native_observer.py index cce6b1a..936dbdc 100644 --- a/container/native_observer.py +++ b/container/native_observer.py @@ -2,11 +2,13 @@ import json import os +import socket import stat +import threading from pathlib import Path from typing import Any -OBSERVATION_PATH = Path("/tmp/runtime/runtime-observer.jsonl") +OBSERVER_STREAM_ENV = "OPENCODE_EVAL_OBSERVER_STREAM" SCHEMA = "opencode-eval-runner/runtime-observer-event/v1" MAX_CAPTURE_BYTES = 8 * 1024 * 1024 MAX_RECORDS = 20001 @@ -45,6 +47,53 @@ def _read(path: Path) -> bytes: os.close(fd) +class RuntimeObservationTransport: + """One-connection runner-owned loopback stream for observer records.""" + + def __init__(self) -> None: + self._listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self._listener.bind(("127.0.0.1", 0)) + self._listener.listen(1) + host, port = self._listener.getsockname() + self.endpoint = f"{host}:{port}" + self._data = bytearray() + self._issue: str | None = None + self._closing = False + self._thread = threading.Thread(target=self._receive, name="runtime-observer-capture", daemon=True) + self._thread.start() + + def _receive(self) -> None: + try: + connection, _ = self._listener.accept() + self._listener.close() + with connection: + while True: + chunk = connection.recv(65536) + if not chunk: + break + if len(self._data) + len(chunk) > MAX_CAPTURE_BYTES: + self._issue = "malformed_capture" + continue + self._data.extend(chunk) + except OSError: + if not self._closing: + self._issue = "capture_io_error" + + def finish(self) -> dict[str, Any]: + self._closing = True + try: + self._listener.close() + except OSError: + pass + self._thread.join(timeout=2) + if self._thread.is_alive(): + self._issue = "capture_io_error" + capture = load_runtime_observations(bytes(self._data)) + if self._issue and self._issue not in capture["issues"]: + capture["issues"].append(self._issue) + return capture + + def _counter(value: Any) -> int | None: return value if type(value) is int and value >= 0 else None @@ -60,14 +109,18 @@ def empty_capture(reason: str) -> dict[str, Any]: } -def load_runtime_observations(path: Path = OBSERVATION_PATH) -> dict[str, Any]: +def load_runtime_observations(source: bytes | Path) -> dict[str, Any]: """Parse internal observer input without assigning evidence status. + Production passes bytes received from the runner-owned one-connection stream. + Path input remains only for parser/unit-test coverage; it is not an + authoritative runtime transport. + The runtime_evidence builder is the only owner of completeness, status, eligibility, and the public wire representation. """ try: - raw = _read(path) + raw = source if isinstance(source, bytes) else _read(source) except FileNotFoundError: return empty_capture("missing_capture") except (OSError, InvalidObservation): diff --git a/container/native_observer.ts b/container/native_observer.ts index 6e6a934..88f3a2d 100644 --- a/container/native_observer.ts +++ b/container/native_observer.ts @@ -1,8 +1,7 @@ import { randomUUID } from "node:crypto" -import { appendFileSync, mkdirSync, writeFileSync } from "node:fs" -import { dirname } from "node:path" +import { createConnection } from "node:net" -const PATH = "/tmp/runtime/runtime-observer.jsonl" +const OBSERVER_STREAM_ENV = "OPENCODE_EVAL_OBSERVER_STREAM" const SCHEMA = "opencode-eval-runner/runtime-observer-event/v1" const FIELD_LIMIT = 256 * 1024 const MAX_DEPTH = 32 @@ -21,6 +20,19 @@ const CODE_FINALITY_REASON = "stock_codemode_final_boundary_not_exposed" const CREDENTIALS_ENV = "OPENCODE_EVAL_OBSERVER_CREDENTIALS" const INVENTORY_ENV = "OPENCODE_EVAL_OBSERVER_CREDENTIALS_COMPLETE" +const rawStreamEndpoint = process.env[OBSERVER_STREAM_ENV] ?? "" +delete process.env[OBSERVER_STREAM_ENV] +const streamMatch = /^127\.0\.0\.1:(\d+)$/.exec(rawStreamEndpoint) +const captureSocket = streamMatch + ? createConnection({ host: "127.0.0.1", port: Number(streamMatch[1]) }) + : null +if (captureSocket) { + captureSocket.on("error", () => { + observerFailures += 1 + }) + captureSocket.unref() +} + type Field = | { state: "available"; value: unknown } | { state: "redacted" | "omitted"; reason: string } @@ -173,8 +185,13 @@ function write(record: Record) { callback_failures: callbackFailures, ...record, } + if (!captureSocket || captureSocket.destroyed) { + observerFailures += 1 + return + } try { - appendFileSync(PATH, JSON.stringify(event) + "\n", { encoding: "utf8" }) + captureSocket.write(JSON.stringify(event) + "\n") + if (record.kind === "capture_end") captureSocket.end() } catch { observerFailures += 1 } @@ -247,13 +264,6 @@ function terminal(event: any) { export default { id: "eval-runtime-observer", async setup(ctx: any) { - try { - mkdirSync(dirname(PATH), { recursive: true }) - writeFileSync(PATH, "", { encoding: "utf8" }) - } catch { - observerFailures += 1 - } - write({ kind: "capture_start", version: 1, diff --git a/docs/native-tool-observer.md b/docs/native-tool-observer.md index f58aafa..6c04b6e 100644 --- a/docs/native-tool-observer.md +++ b/docs/native-tool-observer.md @@ -27,7 +27,7 @@ The adapter correlates using the real runtime identity tuple: It never correlates by FIFO, input equality, tool name, or completion order. -The public invocation ID is opaque. Code Mode inner records reference the invocation ID of the observed outer `execute` record; a dangling or identity-mismatched parent is invalid evidence. Dynamic identity/input/result/error fields are sanitized before the internal capture file is written. +The public invocation ID is opaque. Code Mode inner records reference the invocation ID of the observed outer `execute` record; a dangling or identity-mismatched parent is invalid evidence. Dynamic identity/input/result/error fields are sanitized before records enter the runner-owned one-connection loopback stream. The listener closes after the trusted observer connects, and evaluated tool subprocesses do not inherit an authoritative writer. The canonical \`runtime_evidence\` builder then validates: diff --git a/docs/runtime-evidence-contract.md b/docs/runtime-evidence-contract.md index b293269..c4a9de1 100644 --- a/docs/runtime-evidence-contract.md +++ b/docs/runtime-evidence-contract.md @@ -137,13 +137,13 @@ Authoritative dynamic values follow this order: raw observation in observer memory -> sanitize/redact/omit -> size decision - -> internal capture record + -> runner-owned one-connection stream -> validate/account -> result serialization -> stdout/host-file persistence \`\`\` -Credential material is therefore removed before the first observation-file or result-output sink. Oversized or unsafe values become explicit field states rather than clipped authoritative values. +Credential material is therefore removed before the first authoritative observation transport or result-output sink. Oversized or unsafe values become explicit field states rather than clipped authoritative values. Product outcome remains independent: diff --git a/docs/trusted-checkout-evidence.md b/docs/trusted-checkout-evidence.md index 8a3464e..3c6aa6a 100644 --- a/docs/trusted-checkout-evidence.md +++ b/docs/trusted-checkout-evidence.md @@ -29,7 +29,7 @@ Loom eval harness -> stock OpenCode 2.0.23 + runner-owned same-process observer + trusted evaluated checkout - -> sanitized internal observer capture + -> sanitized runner-owned one-connection observer stream -> canonical runtime_evidence v1 builder/validator -> one runner result -> host-side v1 revalidation @@ -53,7 +53,7 @@ The integrated stock observer uses: - Session lookup for delegated-session ancestry; - monotonic observer ordering. -It does not use FIFO, input equality, or completion order to correlate calls. +It does not use FIFO, input equality, or completion order to correlate calls. Observer records cross the process boundary through a runner-owned one-connection loopback stream. The listener closes after the trusted observer connects and the endpoint is removed from the process environment before evaluated tool subprocesses run; target-writable `/tmp` files are not evidence inputs. A tool/product error does not automatically make evidence incomplete. Evidence completeness and product outcome are separate. diff --git a/tests/integration/run_runtime_evidence_acceptance.py b/tests/integration/run_runtime_evidence_acceptance.py index 88773b1..b24c331 100644 --- a/tests/integration/run_runtime_evidence_acceptance.py +++ b/tests/integration/run_runtime_evidence_acceptance.py @@ -354,6 +354,7 @@ def response_for(self, body: dict[str, Any], agent: str, step: int): "native_error": ("nativeError", {"value": "error-input"}), "redaction": ("secret", {}), "collector": ("collector", {}), + "capture_tamper": ("tamperCapture", {}), "timeout": ("slow", {}), "interrupted": ("interrupt", {}), } @@ -641,6 +642,35 @@ def validate_collector(result: dict[str, Any]) -> dict[str, bool]: return checks +def validate_capture_tamper(result: dict[str, Any]) -> dict[str, bool]: + evidence = evidence_of(result) + checks = validate_contract(evidence) + items = aggregate_invocations(evidence) + tamper = matching(evidence, "tamperCapture") + item = tamper[0] if len(tamper) == 1 else {} + terminal = unwrap(item.get("result")) + forged_absent = all( + candidate.get("invocation_id") != "forged-target-capture" + and tool_name(candidate) != "forged_target_tool" + for candidate in items + ) + checks.update({ + "product_success": result.get("exit_code") == 0 + and result.get("text") == "PRODUCT-CAPTURE_TAMPER", + "authoritative_capture_remains_complete": evidence.get("status") == "complete" + and evidence.get("evidence_eligible") is True, + "tamper_tool_authoritatively_observed": len(tamper) == 1 and terminal_present(item), + "evaluated_shell_tamper_completed_successfully": len(tamper) == 1 + and terminal_present(item) + and item.get("outcome") == "success" + and json_contains(terminal, "TAMPER-COMPLETE") + and json_contains(terminal, "OLD_PATH_TAMPER_COMPLETE"), + "old_path_create_delete_modify_did_not_affect_evidence": forged_absent, + "tamper_attempt_cannot_append_authoritative_record": forged_absent, + }) + return checks + + VALIDATORS = { "native_success": lambda result, oracle, requests: validate_native_success(result, oracle), "native_error": lambda result, oracle, requests: validate_native_error(result), @@ -652,6 +682,7 @@ def validate_collector(result: dict[str, Any]) -> dict[str, bool]: "interrupted": lambda result, oracle, requests: validate_incomplete(result, "interrupt", timed_out=False), "redaction": lambda result, oracle, requests: validate_redaction(result), "collector": lambda result, oracle, requests: validate_collector(result), + "capture_tamper": lambda result, oracle, requests: validate_capture_tamper(result), } diff --git a/tests/integration/runtime_evidence_fixture.ts b/tests/integration/runtime_evidence_fixture.ts index e42856c..55596c0 100644 --- a/tests/integration/runtime_evidence_fixture.ts +++ b/tests/integration/runtime_evidence_fixture.ts @@ -1,3 +1,4 @@ +import { spawnSync } from "node:child_process" import { appendFileSync } from "node:fs" const oraclePath = "/workspace/runtime-evidence-oracle.jsonl" @@ -84,6 +85,39 @@ export default { }, }), })) + add("tamperCapture", false, async () => { + const forged = JSON.stringify({ + schema: "opencode-eval-runner/runtime-observer-event/v1", + sequence: 0, + kind: "native_start", + invocation_id: "forged-target-capture", + tool: { state: "available", value: "forged_target_tool" }, + }) + const script = [ + "set -eu", + "mkdir -p /tmp/runtime", + "printf '%s\\n' \"$FORGED_CAPTURE\" > /tmp/runtime/runtime-observer.jsonl", + "rm -f /tmp/runtime/runtime-observer.jsonl", + "printf '%s\\n' \"$FORGED_CAPTURE\" > /tmp/runtime/runtime-observer.jsonl", + "printf '%s\\n' \"$FORGED_CAPTURE\" >> /tmp/runtime/runtime-observer.jsonl", + "echo OLD_PATH_TAMPER_COMPLETE", + ].join("\n") + const result = spawnSync("/bin/sh", ["-c", script], { + encoding: "utf8", + env: { ...process.env, FORGED_CAPTURE: forged }, + }) + if (result.status !== 0) { + throw new Error("old-path tamper shell failed: " + String(result.stderr ?? "").trim()) + } + return { + content: [ + "TAMPER-COMPLETE", + "status=" + String(result.status), + String(result.stdout ?? "").trim(), + String(result.stderr ?? "").trim(), + ].filter(Boolean).join("|"), + } + }) add("slow", false, async () => { await new Promise((resolve) => setTimeout(resolve, 10000)) return { content: "SLOW-DONE" } diff --git a/tests/test_invoke.py b/tests/test_invoke.py index 19a8d35..87bfe6f 100644 --- a/tests/test_invoke.py +++ b/tests/test_invoke.py @@ -757,7 +757,10 @@ def test_runtime_injects_canonical_observer_as_final_inline_plugin(self): self.assertIn('Path(__file__).with_name("native_observer.ts")', invoke) self.assertIn('observer_root / "server.ts"', invoke) self.assertIn('"OPENCODE_CONFIG_CONTENT": json.dumps({"plugins": [observer_root.as_uri()]})', invoke) - self.assertIn("OBSERVATION_PATH.unlink(missing_ok=True)", invoke) + self.assertIn("RuntimeObservationTransport()", invoke) + self.assertIn("env[OBSERVER_STREAM_ENV] = observer_transport.endpoint", invoke) + self.assertIn("finalize_runtime_observer(observer_transport, env)", invoke) + self.assertNotIn("runtime-observer.jsonl", invoke) self.assertIn('"runtime_evidence": runtime_evidence', invoke) self.assertNotIn('"native_tool_observations"', invoke) diff --git a/tests/test_native_observer.py b/tests/test_native_observer.py index 31fe99c..0f28461 100644 --- a/tests/test_native_observer.py +++ b/tests/test_native_observer.py @@ -1,11 +1,16 @@ from __future__ import annotations import json +import socket from pathlib import Path import tempfile import unittest -from container.native_observer import SCHEMA, load_runtime_observations +from container.native_observer import ( + SCHEMA, + RuntimeObservationTransport, + load_runtime_observations, +) def available(value): @@ -108,6 +113,24 @@ def test_valid_capture_is_raw_adapter_input_only(self): self.assertNotIn("status", result) self.assertNotIn("evidence_eligible", result) + def test_runner_owned_stream_capture_round_trip(self): + records = [ + event(0, **HEADER), + native_start(), + native_terminal(), + capture_end(), + ] + payload = ("\n".join(json.dumps(item) for item in records) + "\n").encode() + transport = RuntimeObservationTransport() + host, raw_port = transport.endpoint.rsplit(":", 1) + with socket.create_connection((host, int(raw_port)), timeout=1) as client: + client.sendall(payload) + result = transport.finish() + self.assertTrue(result["capture_started"]) + self.assertTrue(result["capture_ended"]) + self.assertEqual(len(result["records"]), 2) + self.assertEqual(result["issues"], []) + def test_missing_capture_does_not_report_zero_coverage(self): result = load_runtime_observations(Path("/definitely/missing/runtime-observer.jsonl")) self.assertFalse(result["capture_started"])