Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 24 additions & 9 deletions container/invoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,15 @@

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,
validate_runtime_evidence,
)
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,
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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)
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
59 changes: 56 additions & 3 deletions container/native_observer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand All @@ -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):
Expand Down
32 changes: 21 additions & 11 deletions container/native_observer.ts
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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 }
Expand Down Expand Up @@ -173,8 +185,13 @@ function write(record: Record<string, unknown>) {
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
}
Expand Down Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion docs/native-tool-observer.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
4 changes: 2 additions & 2 deletions docs/runtime-evidence-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
4 changes: 2 additions & 2 deletions docs/trusted-checkout-evidence.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.

Expand Down
31 changes: 31 additions & 0 deletions tests/integration/run_runtime_evidence_acceptance.py
Original file line number Diff line number Diff line change
Expand Up @@ -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", {}),
}
Expand Down Expand Up @@ -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),
Expand All @@ -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),
}


Expand Down
34 changes: 34 additions & 0 deletions tests/integration/runtime_evidence_fixture.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { spawnSync } from "node:child_process"
import { appendFileSync } from "node:fs"

const oraclePath = "/workspace/runtime-evidence-oracle.jsonl"
Expand Down Expand Up @@ -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" }
Expand Down
5 changes: 4 additions & 1 deletion tests/test_invoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Loading
Loading