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
47 changes: 42 additions & 5 deletions scripts/agent-eval/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -492,9 +492,11 @@ A single supervisor SIGTERM observation can qualify by either the original
already-exited-zero/empty-group proof or this additional shutdown proof:

1. Before cleanup, the supervisor reads a bounded, complete broker log whose last
exchange boundary is a **flushed reply** or completed EOF. It captures request
and call counts. A busy request/call, missing reply, partial/oversized log or
unknown counts cannot qualify.
exchange boundary is a **published reply** or completed EOF. It captures
request and call counts. A busy request/call, missing reply, partial/oversized
log or unknown counts cannot qualify. The worker publishes a request's reply
checkpoint after the request's work completes and immediately before writing
that response (see [Reply checkpoint ordering](#reply-checkpoint-ordering)).
2. After that checkpoint, polling the inherited MCP stdin must show exactly
`POLLHUP`: an empty pipe with no writers. Open input, queued input, non-pipe
input and unknown state cannot qualify. The provider process pidfd must still
Expand All @@ -508,8 +510,9 @@ already-exited-zero/empty-group proof or this additional shutdown proof:
4. Final validation requires exactly the checkpoint's request and call counts.
This fence rejects additional requests prefetched into Python's input buffer,
even if the kernel pipe was already empty. The complete serial ledger must
reconcile every request, completed call and flushed reply, then `stdin_eof`.
A completed tool log before a broken reply pipe is insufficient.
reconcile every request, completed call and reply, then `stdin_eof`, whose
reply count only advances after a successful flush. A completed tool log or a
published reply before a broken reply pipe is insufficient.

Both paths still require worker exit zero, bounded diagnostic EOF without
truncation or supervision errors, no surviving descendants or SIGKILL cleanup,
Expand Down Expand Up @@ -561,6 +564,40 @@ AGENT_EVAL_ORBIT_GRAPH="$PWD/target/debug/orbit-graph" \
python3 -B -m unittest discover -s scripts/agent-eval/tests -v
```

### Reply checkpoint ordering

Earlier broker bytes flushed each MCP response before logging its `reply`
checkpoint. A client that tore down as soon as it held its final response could
have the supervisor observe EOF while the checkpoint was still unpublished; the
null checkpoint correctly refused the drain, so the episode failed
`broker_failed` depending on scheduling (ORB-13847). The worker now logs the
`reply` checkpoint after any wire redaction and immediately before writing the
response, so a delivered response always has a published checkpoint. Only that
write can remain. It is proven by worker exit zero and the `stdin_eof` stop's
reply count. A broken or full reply pipe still fails, as does an expired drain.
The drain never performs tool work.

The two orderings are equivalent for a client. Before, a flushed response could
already sit unread in the pipe at observation. Now the drain may also finish
writing a published response to a pipe that is still open. Telemetry
reconciliation still requires the provider to report every broker call.

Classification rules, request/raw schemas, runner/broker version strings and
`eof-idle-at-term-observation-v2` are unchanged. The broker source hash changes,
and each capture records it (`tool_versions.read`, harness hashes). Captures
from earlier broker bytes keep their post-flush `reply` meaning and their
outcomes. Freeze the new harness hashes before any future use. Corpus locks
that pin the previous broker bytes refuse it intentionally.

`test_broker.py` stops a runtime copy of the worker immediately after its final
response is flushed. The client then reads that response, closes input, sends
TERM, waits for the logged supervisor observation and resumes the worker. With
the earlier ordering, this fails deterministically: the observation has a null
checkpoint and cleanup sends SIGTERM to the worker, whose exit status then
depends on scheduling. With the repair, it passes with the full checkpoint, no
cleanup signals and exit zero. A companion test stops the worker after publication but before the
write and closes the client's reader. The justified drain then still fails.

## Prospective source-reply provenance (runner 5)

Installed-plugin **schema-2** episodes now use runner 5, broker 4 and
Expand Down
22 changes: 16 additions & 6 deletions scripts/agent-eval/eval_broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -1080,19 +1080,26 @@ def serve(self, stdin, stdout):
self.record({"type": "protocol_error", "code": error.code,
"reason": error.message[:200]})
if is_request or not isinstance(message, dict):
self.reply(stdout, request_id, error=(error.code, error.message))
self.reply(stdout, request_id, error=(error.code, error.message),
seq=requests if is_request else None)
if is_request:
replies += 1
self.record({"type": "reply", "seq": requests, "calls": self.calls})
continue
if is_request:
self.reply(stdout, request_id, result=result if result is not None else {})
self.reply(stdout, request_id, result=result if result is not None else {},
seq=requests)
replies += 1
self.record({"type": "reply", "seq": requests, "calls": self.calls})
self.record({"type": "stop", "reason": "stdin_eof", "requests": requests,
"replies": replies, "calls": self.calls, "output_bytes": self.output_bytes})

def reply(self, stdout, request_id, result=None, error=None):
def reply(self, stdout, request_id, result=None, error=None, seq=None):
"""Write one MCP response, publishing a request's reply checkpoint first.

A client may tear down as soon as it holds the response, so the
supervisor's idle checkpoint must already cover it. Only the write
remains after publication. Its flush is proven later by worker exit
zero and the stdin_eof stop's reply count; a failed write never is.
"""
if not (request_id is None or isinstance(request_id, (str, int))) or \
isinstance(request_id, bool):
request_id = None
Expand All @@ -1105,7 +1112,10 @@ def reply(self, stdout, request_id, result=None, error=None):
body, count = self.redactor.value(body)
if count:
self.record({"type": "wire_redaction", "audit_redactions": count})
stdout.write((canonical(body) + "\n").encode())
data = (canonical(body) + "\n").encode()
if seq is not None:
self.record({"type": "reply", "seq": seq, "calls": self.calls})
stdout.write(data)
stdout.flush()


Expand Down
99 changes: 96 additions & 3 deletions scripts/agent-eval/tests/test_broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -393,12 +393,12 @@ def test_empty_kernel_pipe_with_python_buffered_request_fails_count_fence(self):
config, state = support.broker_config(self.work, self.repo, self.snapshot)
runtime = self.work / "buffered-broker.py"
source = (support.TOOL / "eval_broker.py").read_text()
boundary = '\n self.record({"type": "reply", "seq": requests, "calls": self.calls})'
boundary = "\n stdout.flush()"
self.assertEqual(source.count(boundary), 1)
(self.work / "reply_provenance.py").write_bytes(
(support.TOOL / "reply_provenance.py").read_bytes())
runtime.write_text(source.replace(boundary, boundary +
'\n if request_id == 41: os.kill(os.getpid(), signal.SIGSTOP)'))
'\n if request_id == 41: os.kill(os.getpid(), signal.SIGSTOP)'))
client = support.McpClient(argv=[sys.executable, "-B", str(runtime), "--config", str(config)])
self.addCleanup(client.terminate)
client.request("initialize")
Expand Down Expand Up @@ -442,6 +442,94 @@ def test_empty_kernel_pipe_with_python_buffered_request_fails_count_fence(self):
pass
os.close(fd)

def held_reply_client(self, name, boundary, before):
"""A runtime copy whose worker stops itself at one reply() boundary for request 42.

The SIGSTOP is a test probe only. Kernel stop state and the logged
supervisor observation order client teardown, with no sleeps or retries.
"""
config, state = support.broker_config(self.work, self.repo, self.snapshot)
source = (support.TOOL / "eval_broker.py").read_text()
self.assertEqual(source.count(boundary), 1)
probe = "\n if request_id == 42: os.kill(os.getpid(), signal.SIGSTOP)"
runtime = self.work / name
runtime.write_text(source.replace(boundary, probe + boundary if before else boundary + probe))
(self.work / "reply_provenance.py").write_bytes(
(support.TOOL / "reply_provenance.py").read_bytes())
client = support.McpClient(argv=[sys.executable, "-B", str(runtime), "--config", str(config)])
self.addCleanup(client.terminate)
client.request("initialize")
identity = next(e["identity"] for e in support.log_entries(state) if e["type"] == "start")
fd = os.pidfd_open(identity["pid"])

def release():
try:
signal.pidfd_send_signal(fd, signal.SIGCONT)
signal.pidfd_send_signal(fd, signal.SIGTERM)
except ProcessLookupError:
pass
os.close(fd)
self.addCleanup(release) # runs before client.terminate
client.send_raw(json.dumps({"jsonrpc": "2.0", "id": 42, "method": "tools/call",
"params": {"name": "read",
"arguments": {"path": "src/lib.rs"}}}) + "\n")
return client, state, identity["pid"], fd

def teardown_held_worker(self, client, state, pid, fd):
"""Wait for the probe stop, then EOF + supervisor TERM; resume after the observation."""
deadline = time.monotonic() + 10
while True:
with open(f"/proc/{pid}/stat") as stream:
if stream.read().rsplit(")", 1)[1].split()[0] == "T":
break
self.assertLess(time.monotonic(), deadline)
client.process.stdin.close()
os.kill(client.process.pid, signal.SIGTERM)
while not any(e["type"] == "supervisor_signal" for e in support.log_entries(state)):
self.assertLess(time.monotonic(), deadline)
signal.pidfd_send_signal(fd, signal.SIGCONT)

def test_teardown_right_after_final_flushed_reply_keeps_its_checkpoint(self):
# A client may tear down as soon as it holds its final response. Stop
# the worker the instant that response is flushed, before it can log
# anything else: the idle checkpoint must already cover the exchange.
client, state, pid, fd = self.held_reply_client(
"flushed-reply-broker.py", "\n stdout.flush()", before=False)
reply = client.read()
self.assertEqual(reply["id"], 42)
self.assertFalse(reply["result"]["isError"])
self.teardown_held_worker(client, state, pid, fd)
self.assertEqual(client.close(), 0)
log, _ = support.eval_runner.read_broker_log(state / "calls.jsonl")
observation = log["exits"][0]["lifecycle"]["signals"][0]
self.assertFalse(observation["worker_exited_zero"])
self.assertTrue(observation["input_eof"])
self.assertEqual(observation["checkpoint"], {"requests": 2, "calls": 1})
self.assertFalse(log["exits"][0]["lifecycle"]["drain_expired"])
self.assertEqual(log["exits"][0]["cleanup"], {"signals": [], "survivors": []})
self.assertTrue(eval_broker.orderly_broker_exit(log["exits"][0]))
self.assertTrue(support.eval_runner.complete_broker_lifecycle(log))
self.assertEqual(log["lifecycle_events"][-1],
{"type": "stop", "reason": "stdin_eof", "requests": 2, "replies": 2,
"calls": 1, "output_bytes": len(log["calls"][0]["output"].encode())})

def test_published_reply_that_is_never_flushed_does_not_qualify(self):
# Stop after the checkpoint is published but before the response is
# written, then remove its reader: the justified drain must still fail.
client, state, pid, fd = self.held_reply_client(
"unflushed-reply-broker.py", "\n stdout.write(data)", before=True)
client.process.stdout.close()
self.teardown_held_worker(client, state, pid, fd)
self.assertEqual(client.close(), 1)
log, _ = support.eval_runner.read_broker_log(state / "calls.jsonl")
observation = log["exits"][0]["lifecycle"]["signals"][0]
self.assertTrue(observation["input_eof"])
self.assertEqual(observation["checkpoint"], {"requests": 2, "calls": 1})
self.assertNotEqual(log["exits"][0]["exit_code"], 0)
self.assertFalse(any(e["type"] == "stop" for e in log["lifecycle_events"]))
self.assertFalse(eval_broker.orderly_broker_exit(log["exits"][0]))
self.assertFalse(support.eval_runner.complete_broker_lifecycle(log))

def test_worker_stderr_is_drained_bounded_and_exit_is_recorded(self):
config, state = support.broker_config(self.work, self.repo, self.snapshot)
runtime = self.work / "fault-broker.py"
Expand Down Expand Up @@ -469,9 +557,14 @@ def test_flushed_reply_is_not_inferred_from_completed_tool(self):
self.assertEqual(client.close(), 1)
entries = support.log_entries(state)
self.assertEqual(sum(e["type"] == "call" for e in entries), 1)
self.assertEqual(sum(e["type"] == "reply" for e in entries), 1) # initialize only
# The checkpoint is published before the write; only exit zero and the
# stop's reply count prove a flush, and the broken pipe yields neither.
self.assertEqual(sum(e["type"] == "reply" for e in entries), 2)
self.assertFalse(any(e["type"] == "stop" for e in entries))
self.assertNotEqual(entries[-1]["exit_code"], 0)
self.assertFalse(eval_broker.orderly_broker_exit(entries[-1]))
log, _ = support.eval_runner.read_broker_log(state / "calls.jsonl")
self.assertFalse(support.eval_runner.complete_broker_lifecycle(log))

def test_clean_worker_with_descendant_requires_failed_cleanup(self):
config, state = support.broker_config(self.work, self.repo, self.snapshot)
Expand Down
Loading