fix(lease): ignore benign observer idle replies
This commit is contained in:
@@ -22,13 +22,23 @@ from typing import Final
|
|||||||
MAX_FRAME: Final = 64 * 1024
|
MAX_FRAME: Final = 64 * 1024
|
||||||
BROKER_TIMEOUT_SECONDS: Final = 1.5
|
BROKER_TIMEOUT_SECONDS: Final = 1.5
|
||||||
MAX_TRANSCRIPT_BYTES: Final = 4 * 1024 * 1024
|
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]:
|
def read_json(stream: object) -> dict[str, object]:
|
||||||
raw = getattr(stream, "buffer", stream).read(MAX_FRAME + 1)
|
raw = getattr(stream, "buffer", stream).read(MAX_FRAME + 1)
|
||||||
if not isinstance(raw, bytes) or len(raw) > MAX_FRAME:
|
if not isinstance(raw, bytes) or len(raw) > MAX_FRAME:
|
||||||
raise ValueError("invalid observer input")
|
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):
|
if not isinstance(value, dict):
|
||||||
raise ValueError("invalid observer input")
|
raise ValueError("invalid observer input")
|
||||||
return value
|
return value
|
||||||
@@ -100,9 +110,9 @@ def observer_request(socket_path: Path, request: dict[str, object]) -> dict[str,
|
|||||||
if not chunk:
|
if not chunk:
|
||||||
break
|
break
|
||||||
response.extend(chunk)
|
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")
|
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):
|
if not isinstance(value, dict):
|
||||||
raise ValueError("invalid observer reply")
|
raise ValueError("invalid observer reply")
|
||||||
return value
|
return value
|
||||||
@@ -133,10 +143,18 @@ def main(argv: Sequence[str] | None = None, *, environ: Mapping[str, str] | None
|
|||||||
"runtime": arguments.runtime,
|
"runtime": arguments.runtime,
|
||||||
"latest_assistant_message": message,
|
"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)
|
print(f"Mosaic receipt observer refused: {error}", file=sys.stderr)
|
||||||
return 2
|
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__":
|
if __name__ == "__main__":
|
||||||
|
|||||||
@@ -25,7 +25,7 @@
|
|||||||
"lint": "eslint src",
|
"lint": "eslint src",
|
||||||
"typecheck": "tsc --noEmit",
|
"typecheck": "tsc --noEmit",
|
||||||
"test": "vitest run --passWithNoTests && pnpm run test:framework-shell",
|
"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": {
|
"dependencies": {
|
||||||
"@mosaicstack/brain": "workspace:*",
|
"@mosaicstack/brain": "workspace:*",
|
||||||
|
|||||||
@@ -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()
|
||||||
Reference in New Issue
Block a user