From 9c3f054f204980bb43c85b7b7a1ea371f219ade7 Mon Sep 17 00:00:00 2001 From: Jason Woltje Date: Fri, 7 Aug 2026 15:15:25 -0500 Subject: [PATCH] fix(lease): ignore benign observer idle replies --- .../lease-broker/receipt-observer-client.py | 28 ++- packages/mosaic/package.json | 2 +- .../receipt_observer_client_unittest.py | 199 ++++++++++++++++++ 3 files changed, 223 insertions(+), 6 deletions(-) create mode 100644 packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py diff --git a/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py b/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py index 02f22d3f..2289d606 100644 --- a/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py +++ b/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py @@ -22,13 +22,23 @@ from typing import Final MAX_FRAME: Final = 64 * 1024 BROKER_TIMEOUT_SECONDS: Final = 1.5 MAX_TRANSCRIPT_BYTES: Final = 4 * 1024 * 1024 +BENIGN_OBSERVATION_UNAVAILABLE_CODE: Final = "OBSERVATION_UNAVAILABLE" + + +def reject_duplicate_json_keys(pairs: list[tuple[str, object]]) -> dict[str, object]: + value: dict[str, object] = {} + for key, item in pairs: + if key in value: + raise ValueError("duplicate observer JSON key") + value[key] = item + return value def read_json(stream: object) -> dict[str, object]: raw = getattr(stream, "buffer", stream).read(MAX_FRAME + 1) if not isinstance(raw, bytes) or len(raw) > MAX_FRAME: raise ValueError("invalid observer input") - value = json.loads(raw) + value = json.loads(raw, object_pairs_hook=reject_duplicate_json_keys) if not isinstance(value, dict): raise ValueError("invalid observer input") return value @@ -100,9 +110,9 @@ def observer_request(socket_path: Path, request: dict[str, object]) -> dict[str, if not chunk: break response.extend(chunk) - if len(response) > MAX_FRAME or not response.endswith(b"\n"): + if len(response) > MAX_FRAME or response.count(b"\n") != 1 or not response.endswith(b"\n"): raise ValueError("invalid observer reply") - value = json.loads(response) + value = json.loads(response[:-1], object_pairs_hook=reject_duplicate_json_keys) if not isinstance(value, dict): raise ValueError("invalid observer reply") return value @@ -133,10 +143,18 @@ def main(argv: Sequence[str] | None = None, *, environ: Mapping[str, str] | None "runtime": arguments.runtime, "latest_assistant_message": message, }) - except (KeyError, OSError, ValueError, json.JSONDecodeError) as error: + except (KeyError, OSError, RecursionError, ValueError, json.JSONDecodeError) as error: print(f"Mosaic receipt observer refused: {error}", file=sys.stderr) return 2 - return 0 if reply == {"ok": True} else 2 + if set(reply) == {"ok"} and reply.get("ok") is True: + return 0 + if ( + set(reply) == {"ok", "code"} + and reply.get("ok") is False + and reply.get("code") == BENIGN_OBSERVATION_UNAVAILABLE_CODE + ): + return 0 + return 2 if __name__ == "__main__": diff --git a/packages/mosaic/package.json b/packages/mosaic/package.json index 9b22b470..cc843cec 100644 --- a/packages/mosaic/package.json +++ b/packages/mosaic/package.json @@ -25,7 +25,7 @@ "lint": "eslint src", "typecheck": "tsc --noEmit", "test": "vitest run --passWithNoTests && pnpm run test:framework-shell", - "test:framework-shell": "bash framework/tools/quality/scripts/check-test-enumeration.sh && bash framework/tools/quality/scripts/test-check-test-enumeration.sh && python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/promotion_binding_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/lease-broker/context_recovery_unittest.py && python3 src/lease-broker/recovery_runtime_unittest.py && python3 src/lease-broker/recovery_b1_adversarial_unittest.py && python3 src/lease-broker/invariant_r_unittest.py && python3 src/lease-broker/framework_skill_portability_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 src/mutator-gate/version_coupling_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh && bash framework/tools/qa/test-deps-preflight.sh && bash framework/tools/git/test-pr-review-gitea-comment.sh && bash framework/tools/git/test-pr-review-repo-host-override.sh && bash framework/tools/git/test-ci-queue-wait-branch-absent.sh && bash framework/tools/git/test-ci-queue-wait-tristate.sh && bash framework/tools/git/test-ci-queue-wait-github-checks.sh && bash framework/tools/git/test-pr-merge-queue-branch.sh && bash framework/tools/git/test-pr-merge-head-pin.sh && bash framework/tools/git/test-pr-merge-message-field.sh && bash framework/tools/git/test-git-credential-mosaic.sh && bash framework/tools/git/test-gitea-token-identity.sh && bash framework/tools/woodpecker/test-terminal-green-contract.sh && bash framework/tools/_scripts/test-install-ordering-guard.sh && bash framework/tools/tmux/agent-send.test.sh && bash framework/tools/wake/test-wake-store-ack.sh && bash framework/tools/wake/test-wake-store-enqueue-race.sh && bash framework/tools/wake/test-wake-digest-hmac.sh && bash framework/tools/wake/test-wake-digest-quarantine.sh && bash framework/tools/wake/test-wake-detector.sh && bash framework/tools/wake/test-wake-fn-oracle.sh && bash framework/tools/wake/test-wake-reconcile.sh && bash framework/tools/wake/test-wake-beacon.sh && bash framework/tools/wake/test-wake-preimage.sh && bash framework/tools/wake/test-wake-install.sh" + "test:framework-shell": "bash framework/tools/quality/scripts/check-test-enumeration.sh && bash framework/tools/quality/scripts/test-check-test-enumeration.sh && python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/promotion_binding_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/lease-broker/context_recovery_unittest.py && python3 src/lease-broker/recovery_runtime_unittest.py && python3 src/lease-broker/recovery_b1_adversarial_unittest.py && python3 src/lease-broker/receipt_observer_client_unittest.py && python3 src/lease-broker/invariant_r_unittest.py && python3 src/lease-broker/framework_skill_portability_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 src/mutator-gate/version_coupling_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh && bash framework/tools/qa/test-deps-preflight.sh && bash framework/tools/git/test-pr-review-gitea-comment.sh && bash framework/tools/git/test-pr-review-repo-host-override.sh && bash framework/tools/git/test-ci-queue-wait-branch-absent.sh && bash framework/tools/git/test-ci-queue-wait-tristate.sh && bash framework/tools/git/test-ci-queue-wait-github-checks.sh && bash framework/tools/git/test-pr-merge-queue-branch.sh && bash framework/tools/git/test-pr-merge-head-pin.sh && bash framework/tools/git/test-pr-merge-message-field.sh && bash framework/tools/git/test-git-credential-mosaic.sh && bash framework/tools/git/test-gitea-token-identity.sh && bash framework/tools/woodpecker/test-terminal-green-contract.sh && bash framework/tools/_scripts/test-install-ordering-guard.sh && bash framework/tools/tmux/agent-send.test.sh && bash framework/tools/wake/test-wake-store-ack.sh && bash framework/tools/wake/test-wake-store-enqueue-race.sh && bash framework/tools/wake/test-wake-digest-hmac.sh && bash framework/tools/wake/test-wake-digest-quarantine.sh && bash framework/tools/wake/test-wake-detector.sh && bash framework/tools/wake/test-wake-fn-oracle.sh && bash framework/tools/wake/test-wake-reconcile.sh && bash framework/tools/wake/test-wake-beacon.sh && bash framework/tools/wake/test-wake-preimage.sh && bash framework/tools/wake/test-wake-install.sh" }, "dependencies": { "@mosaicstack/brain": "workspace:*", diff --git a/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py b/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py new file mode 100644 index 00000000..bbe60018 --- /dev/null +++ b/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py @@ -0,0 +1,199 @@ +#!/usr/bin/env python3 +"""Exit-semantics tests for the receipt observer Stop-hook client.""" + +from __future__ import annotations + +import importlib.util +import io +import json +import unittest +from contextlib import redirect_stderr +from pathlib import Path +from unittest import mock + + +TOOLS = Path(__file__).parents[2] / "framework/tools/lease-broker" +CLIENT_PATH = TOOLS / "receipt-observer-client.py" + + +def load_client(): + spec = importlib.util.spec_from_file_location("receipt_observer_client_test", CLIENT_PATH) + if spec is None or spec.loader is None: + raise RuntimeError("unable to load receipt-observer-client.py") + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +CLIENT = load_client() +VALID_INPUT = json.dumps({"latest_assistant_message": "ordinary turn"}).encode() +DEEPLY_NESTED_JSON = b"[" * 2_000 + b"0" + b"]" * 2_000 +ENVIRONMENT = { + "MOSAIC_RECEIPT_OBSERVER_SOCKET": "/unused/observer.sock", + "MOSAIC_LEASE_SESSION_ID": "a" * 64, + "MOSAIC_RUNTIME_GENERATION": "1", +} + + +class FakeObserverSocket: + def __init__(self, response: bytes) -> None: + self.response = response + + def __enter__(self): + return self + + def __exit__(self, *_args: object) -> None: + return None + + def settimeout(self, _timeout: float) -> None: + return None + + def connect(self, _path: str) -> None: + return None + + def sendall(self, _payload: bytes) -> None: + return None + + def shutdown(self, _how: int) -> None: + return None + + def recv(self, _size: int) -> bytes: + response, self.response = self.response, b"" + return response + + +class ReceiptObserverClientExitSemanticsTest(unittest.TestCase): + def run_client( + self, + *, + input_bytes: bytes = VALID_INPUT, + reply: dict[str, object] | None = None, + transport_error: OSError | None = None, + ) -> tuple[int, str, mock.Mock]: + request = mock.Mock(return_value=reply) + if transport_error is not None: + request.side_effect = transport_error + stderr = io.StringIO() + with ( + mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(input_bytes)), + mock.patch.object(CLIENT, "observer_request", request), + redirect_stderr(stderr), + ): + result = CLIENT.main(["--runtime", "pi"], environ=ENVIRONMENT) + return result, stderr.getvalue(), request + + def test_nothing_pending_observation_refusal_is_benign(self) -> None: + result, stderr, request = self.run_client( + reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE"} + ) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + request.assert_called_once() + + def test_success_reply_remains_successful(self) -> None: + result, stderr, _request = self.run_client(reply={"ok": True}) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + + def test_pending_cycle_auth_failure_stays_fail_closed(self) -> None: + result, _stderr, _request = self.run_client( + reply={"ok": False, "code": "ANCESTRY_MISMATCH"} + ) + + self.assertEqual(result, 2) + + def test_transport_failure_stays_fail_closed(self) -> None: + result, stderr, _request = self.run_client( + transport_error=ConnectionRefusedError("observer unavailable") + ) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr) + + def test_parse_failure_stays_fail_closed(self) -> None: + result, stderr, request = self.run_client(input_bytes=b"{") + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr) + request.assert_not_called() + + def test_malformed_wire_replies_stay_fail_closed(self) -> None: + for response in ( + b"not-json\n", + b'{"ok":true}', + b"{}\n{}\n", + b'{"ok":false,"code":"OBSERVATION_UNAVAILABLE"}\n\n', + b'{"ok":true,"ok":false,"code":"OBSERVATION_UNAVAILABLE"}\n', + DEEPLY_NESTED_JSON + b"\n", + b"x" * (CLIENT.MAX_FRAME + 1), + ): + with self.subTest(response=response): + stderr = io.StringIO() + with ( + mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(VALID_INPUT)), + mock.patch.object( + CLIENT.socket, + "socket", + return_value=FakeObserverSocket(response), + ), + redirect_stderr(stderr), + ): + result = CLIENT.main(["--runtime", "pi"], environ=ENVIRONMENT) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr.getvalue()) + + def test_oversized_input_stays_fail_closed(self) -> None: + result, stderr, request = self.run_client(input_bytes=b"x" * (CLIENT.MAX_FRAME + 1)) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr) + request.assert_not_called() + + def test_deeply_nested_input_stays_fail_closed(self) -> None: + result, stderr, request = self.run_client(input_bytes=DEEPLY_NESTED_JSON) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr) + request.assert_not_called() + + def test_json_recursion_failure_stays_fail_closed(self) -> None: + request = mock.Mock() + stderr = io.StringIO() + with ( + mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(VALID_INPUT)), + mock.patch.object(CLIENT, "observer_request", request), + mock.patch.object( + CLIENT.json, + "loads", + side_effect=RecursionError("maximum JSON nesting exceeded"), + ), + redirect_stderr(stderr), + ): + result = CLIENT.main(["--runtime", "pi"], environ=ENVIRONMENT) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr.getvalue()) + request.assert_not_called() + + def test_observation_unavailable_with_unexpected_fields_stays_fail_closed(self) -> None: + result, _stderr, _request = self.run_client( + reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE", "unexpected": True} + ) + + self.assertEqual(result, 2) + + def test_non_boolean_ok_values_stay_fail_closed(self) -> None: + for reply in ( + {"ok": 1}, + {"ok": 0, "code": "OBSERVATION_UNAVAILABLE"}, + ): + with self.subTest(reply=reply): + result, _stderr, _request = self.run_client(reply=reply) + self.assertEqual(result, 2) + + +if __name__ == "__main__": + unittest.main()