diff --git a/Dockerfile.gateway b/Dockerfile.gateway index a9f6d6d..8849d21 100644 --- a/Dockerfile.gateway +++ b/Dockerfile.gateway @@ -26,8 +26,9 @@ # /git-gate-entrypoint.sh docker-cp'd at start time # /git-gate/creds/* docker-cp'd at start time # /git/* bare repos, populated at runtime -# /run/supervise/bot-bottle.db bind-mounted at run time # /home/mitmproxy/.mitmproxy/ mitmproxy CA dir +# (No bot-bottle.db mount: the data plane reaches the supervise queue over the +# control-plane RPC and never opens the DB — PRD 0070 / issue #469.) # # Exposed ports inside the container: # 9099 egress (mitmproxy, agent-facing HTTPS proxy) diff --git a/bot_bottle/backend/firecracker/infra_vm.py b/bot_bottle/backend/firecracker/infra_vm.py index 8ccf6e9..00d962f 100644 --- a/bot_bottle/backend/firecracker/infra_vm.py +++ b/bot_bottle/backend/firecracker/infra_vm.py @@ -493,10 +493,11 @@ BOT_BOTTLE_ROOT=/var/lib/bot-bottle python3 -m bot_bottle.orchestrator \\ # Gateway data plane, multi-tenant: each request resolves source-IP -> # policy against the local control plane. The VM backend reaches git over # git-http (9420), so the git:// daemon (git-gate, needs a per-bottle -# entrypoint the consolidated model doesn't use) is left out. +# entrypoint the consolidated model doesn't use) is left out. No +# SUPERVISE_DB_PATH: the data plane reaches the supervise queue over the +# control-plane RPC and never opens bot-bottle.db (PRD 0070 / #469). BOT_BOTTLE_GATEWAY_DAEMONS=egress,git-http,supervise \\ BOT_BOTTLE_ORCHESTRATOR_URL=http://127.0.0.1:{CONTROL_PLANE_PORT} \\ -SUPERVISE_DB_PATH=/var/lib/bot-bottle/db/bot-bottle.db \\ python3 -m bot_bottle.gateway_init & # Reap as PID 1; children are backgrounded, so `wait` blocks. diff --git a/bot_bottle/backend/macos_container/infra.py b/bot_bottle/backend/macos_container/infra.py index 1574bdb..cb92540 100644 --- a/bot_bottle/backend/macos_container/infra.py +++ b/bot_bottle/backend/macos_container/infra.py @@ -100,10 +100,12 @@ def _init_script(port: int) -> str: f"( cd {_SRC_IN_CONTAINER} && BOT_BOTTLE_ROOT={_DB_ROOT_IN_CONTAINER} " f"python3 -m bot_bottle.orchestrator --host 0.0.0.0 --port {port} " "--broker stub ) &\n" - # Gateway data plane, multi-tenant against the local control plane. + # Gateway data plane, multi-tenant against the local control plane. No + # SUPERVISE_DB_PATH: the data plane reaches the supervise queue over the + # control-plane RPC and never opens bot-bottle.db (PRD 0070 / #469). f"( cd /app && BOT_BOTTLE_GATEWAY_DAEMONS={_GATEWAY_DAEMONS} " f"BOT_BOTTLE_ORCHESTRATOR_URL=http://127.0.0.1:{port} " - f"SUPERVISE_DB_PATH={_DB_PATH_IN_CONTAINER} python3 -m bot_bottle.gateway_init ) &\n" + f"python3 -m bot_bottle.gateway_init ) &\n" "while : ; do wait ; done\n" ) diff --git a/bot_bottle/egress_addon.py b/bot_bottle/egress_addon.py index 04dd3fa..8766623 100644 --- a/bot_bottle/egress_addon.py +++ b/bot_bottle/egress_addon.py @@ -40,8 +40,13 @@ from bot_bottle.egress_addon_core import ( scan_inbound, scan_outbound, ) -from bot_bottle import supervise as _sv -from bot_bottle.policy_resolver import PolicyResolver +from bot_bottle.policy_resolver import PolicyResolveError, PolicyResolver +from bot_bottle.supervise_types import ( + STATUS_APPROVED, + STATUS_MODIFIED, + STATUSES, + TOOL_EGRESS_TOKEN_ALLOW, +) INTROSPECT_HOST = "_egress.local" @@ -314,8 +319,15 @@ class EgressAddon: flow.request.headers.pop("Proxy-Authorization", None) flow.request.headers.pop(IDENTITY_HEADER, None) conn = flow.client_conn - if not token and conn is not None: - token = self._conn_tokens.get(getattr(conn, "id", ""), "") + conn_id = getattr(conn, "id", "") if conn is not None else "" + if not token and conn_id: + token = self._conn_tokens.get(conn_id, "") + # Remember the token per connection so a later token-block on this flow + # can attribute its supervise proposal by (source_ip, identity_token) + # over RPC — plain-HTTP requests carry it here, HTTPS tunnels captured + # it at CONNECT (http_connect); both land in `_conn_tokens`. + if token and conn_id: + self._conn_tokens[conn_id] = token return token def http_connect(self, flow: http.HTTPFlow) -> None: @@ -592,42 +604,48 @@ class EgressAddon: redact_tokens(request_path, env=env), result, ) - proposal = _sv.Proposal.new( - bottle_slug=slug, - tool=_sv.TOOL_EGRESS_TOKEN_ALLOW, - proposed_file=payload, - justification=_TOKEN_ALLOW_JUSTIFICATION, - current_file_hash=_sv.sha256_hex(payload), - ) + # Attribute the proposal by (source_ip, identity_token) over the control + # plane — the data plane no longer opens bot-bottle.db (PRD 0070 / #469). + # source_ip + token come from this bottle's connection; `slug` is kept + # only for its own DLP safelist. + conn = flow.client_conn + source_ip = conn.peername[0] if conn is not None and conn.peername else "" + token = self._conn_tokens.get(getattr(conn, "id", ""), "") if conn is not None else "" try: - _sv.write_proposal(proposal) - except OSError as e: + proposal_id = self._resolver.propose_supervise( + source_ip, token, + tool=TOOL_EGRESS_TOKEN_ALLOW, + proposed_file=payload, + justification=_TOKEN_ALLOW_JUSTIFICATION, + ) + except PolicyResolveError as e: sys.stderr.write( f"egress: could not queue token-allow proposal: {e}; " "blocking request\n" ) + proposal_id = None + if proposal_id is None: self._block(flow, f"egress DLP: {result.reason}", ctx=self._req_ctx(flow)) return False sys.stderr.write(json.dumps({ "event": "egress_token_supervise", "reason": f"egress DLP: {result.reason}", - "proposal": proposal.id, + "proposal": proposal_id, **self._req_ctx(flow), }) + "\n") - response = await self._await_token_response(proposal.id, slug) - _sv.archive_proposal(slug, proposal.id) + response = await self._await_token_response(proposal_id, source_ip, token) - if response is not None and response.status in ( - _sv.STATUS_APPROVED, _sv.STATUS_MODIFIED, + if response is not None and response.get("status") in ( + STATUS_APPROVED, STATUS_MODIFIED, ): self._safe_tokens_for(slug).add(result.matched) if self._flow_log(flow) >= LOG_BLOCKS: sys.stderr.write(json.dumps({ "event": "egress_token_allowed", "reason": f"egress DLP: {result.reason}", - "proposal": proposal.id, + "proposal": proposal_id, **self._req_ctx(flow), }) + "\n") return True @@ -645,19 +663,23 @@ class EgressAddon: async def _await_token_response( self, proposal_id: str, - slug: str, - ) -> "_sv.Response | None": - """Poll the DB for the operator's response without blocking the - proxy event loop. Returns the Response, or None on timeout.""" + source_ip: str, + token: str, + ) -> "dict[str, object] | None": + """Poll the control plane for the operator's decision without blocking + the proxy event loop. Returns the terminal `{status, ...}` payload once + decided, or None on timeout. A transient orchestrator error (or a + `pending`/`unknown` status) is retried until the deadline, then fails + closed — the caller blocks the request.""" loop = asyncio.get_running_loop() deadline = loop.time() + self._token_allow_timeout while True: try: - return _sv.read_response(slug, proposal_id) - except (OSError, ValueError, KeyError): - # Not written yet, or a partial/malformed write — retry until - # the deadline, then fail closed. - pass + result = self._resolver.poll_supervise(source_ip, token, proposal_id) + except PolicyResolveError: + result = None + if result is not None and result.get("status") in STATUSES: + return result if loop.time() >= deadline: return None await asyncio.sleep(TOKEN_ALLOW_POLL_INTERVAL_SECONDS) diff --git a/bot_bottle/git_gate_render.py b/bot_bottle/git_gate_render.py index 2a2378c..04c038c 100644 --- a/bot_bottle/git_gate_render.py +++ b/bot_bottle/git_gate_render.py @@ -272,21 +272,27 @@ supervise_gitleaks_allow() { proposal_id=$( PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" GITLEAKS_ALLOW_REF="$ref" python3 - "$report_file" <<'PY' -import datetime -import hashlib import json import os import sys from pathlib import Path +# Queue over the control-plane RPC — the git-gate data plane no longer opens +# bot-bottle.db (PRD 0070 / issue #469). Attribution is by (source_ip, +# identity_token), resolved server-side, so the proposal lands under the +# calling bottle exactly as a direct write once did. try: - import supervise as _sv + from bot_bottle.policy_resolver import PolicyResolver, PolicyResolveError + from bot_bottle.supervise_types import TOOL_GITLEAKS_ALLOW except ImportError: - from bot_bottle import supervise as _sv + from policy_resolver import PolicyResolver, PolicyResolveError + from supervise_types import TOOL_GITLEAKS_ALLOW report_path = Path(sys.argv[1]) -slug = os.environ.get("SUPERVISE_BOTTLE_SLUG", "") -if not slug: +source_ip = os.environ.get("SUPERVISE_SOURCE_IP", "") +token = os.environ.get("SUPERVISE_IDENTITY_TOKEN", "") +orch_url = os.environ.get("BOT_BOTTLE_ORCHESTRATOR_URL", "") +if not source_ip or not orch_url: sys.exit(2) try: @@ -323,19 +329,21 @@ for i, finding in enumerate(raw, 1): ]) payload = "\n".join(lines).rstrip() + "\n" -proposal = _sv.Proposal.new( - bottle_slug=slug, - tool=_sv.TOOL_GITLEAKS_ALLOW, - proposed_file=payload, - justification=( - "git-gate found gitleaks findings hidden by # gitleaks:allow; " - "approve only for dummy test fixtures or confirmed false positives" - ), - current_file_hash=hashlib.sha256(payload.encode("utf-8")).hexdigest(), - now=datetime.datetime.now(datetime.timezone.utc), -) -_sv.write_proposal(proposal) -print(proposal.id) +try: + proposal_id = PolicyResolver(orch_url).propose_supervise( + source_ip, token, + tool=TOOL_GITLEAKS_ALLOW, + proposed_file=payload, + justification=( + "git-gate found gitleaks findings hidden by # gitleaks:allow; " + "approve only for dummy test fixtures or confirmed false positives" + ), + ) +except PolicyResolveError: + sys.exit(4) +if not proposal_id: + sys.exit(2) +print(proposal_id) PY ) rc=$? @@ -348,7 +356,6 @@ PY return 1 fi - slug=${SUPERVISE_BOTTLE_SLUG:-} timeout=${SUPERVISE_GITLEAKS_ALLOW_TIMEOUT_SECONDS:-300} case "$timeout" in ''|*[!0-9]*) @@ -360,20 +367,30 @@ PY echo "git-gate: approve with './cli.py supervise' to continue this push" >&2 waited=0 while [ "$waited" -lt "$timeout" ]; do - status=$(PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" python3 - "$slug" "$proposal_id" <<'PY' + status=$(PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" python3 - "$proposal_id" <<'PY' +import os import sys +# Non-blocking poll over the control plane. A decided proposal is archived +# server-side on read, so no separate archive step is needed here. try: - import supervise as _sv + from bot_bottle.policy_resolver import PolicyResolver, PolicyResolveError except ImportError: - from bot_bottle import supervise as _sv + from policy_resolver import PolicyResolver, PolicyResolveError -slug = sys.argv[1] +source_ip = os.environ.get("SUPERVISE_SOURCE_IP", "") +token = os.environ.get("SUPERVISE_IDENTITY_TOKEN", "") +orch_url = os.environ.get("BOT_BOTTLE_ORCHESTRATOR_URL", "") try: - response = _sv.read_response(slug, sys.argv[2]) -except FileNotFoundError: + result = PolicyResolver(orch_url).poll_supervise(source_ip, token, sys.argv[1]) +except PolicyResolveError: sys.exit(2) -print(response.status) +if result is None: + sys.exit(2) +status = result.get("status", "") +if status in ("pending", "unknown"): + sys.exit(2) +print(status) PY ) rc=$? @@ -385,16 +402,6 @@ PY if [ -n "$status" ]; then case "$status" in approved|modified) - PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" python3 - "$slug" "$proposal_id" <<'PY' || true -import sys - -try: - import supervise as _sv -except ImportError: - from bot_bottle import supervise as _sv - -_sv.archive_proposal(sys.argv[1], sys.argv[2]) -PY echo "git-gate: supervisor approved # gitleaks:allow for $ref" >&2 return 0 ;; diff --git a/bot_bottle/git_http_backend.py b/bot_bottle/git_http_backend.py index 22e538e..21f924f 100644 --- a/bot_bottle/git_http_backend.py +++ b/bot_bottle/git_http_backend.py @@ -167,11 +167,14 @@ class GitHttpHandler(BaseHTTPRequestHandler): "SERVER_PORT": str(self.server.server_port), # type: ignore "SERVER_PROTOCOL": self.request_version, }) - # Attribute the gitleaks-allow supervise proposal (written by + # Attribute the gitleaks-allow supervise proposal (queued by # receive-pack's pre-receive hook, a child of the CGI we spawn below) to - # the calling bottle. The namespaced root is `/`, so its - # final component is the bottle id — the same per-bottle key egress uses. - env["SUPERVISE_BOTTLE_SLUG"] = sandbox_root.name + # the calling bottle. The hook queues over the control-plane RPC (PRD + # 0070 / issue #469), which re-resolves the bottle from these — the same + # (source_ip, identity_token) pair that selected the namespaced root + # above — so it can only ever queue its own proposals. + env["SUPERVISE_SOURCE_IP"] = self.client_address[0] + env["SUPERVISE_IDENTITY_TOKEN"] = self.headers.get(IDENTITY_HEADER, "") for header, variable in ( ("accept", "HTTP_ACCEPT"), ("content-encoding", "HTTP_CONTENT_ENCODING"), diff --git a/bot_bottle/orchestrator/control_plane.py b/bot_bottle/orchestrator/control_plane.py index 1625cec..87d7714 100644 --- a/bot_bottle/orchestrator/control_plane.py +++ b/bot_bottle/orchestrator/control_plane.py @@ -27,6 +27,18 @@ vsock / unix-socket portability caveats): POST /supervise/respond -> 200 {"responded": true} | 409 (operator) body: {"proposal_id","bottle_slug", "decision", ["notes"],["final_file"]} + POST /supervise/propose -> 201 {"proposal_id"} | 403 (agent) + body: {"source_ip","identity_token", + "tool","proposed_file","justification"} + POST /supervise/poll -> 200 {"status", ["notes"],["final_file"]} | 403 + body: {"source_ip","identity_token", + "proposal_id"} + +The `/supervise/propose` + `/supervise/poll` pair is the **agent** half of the +supervise flow: the data plane (supervise / egress / git-gate) queues a +proposal and polls for its response over RPC instead of opening `bot-bottle.db` +directly. Both attribute the caller by `(source_ip, identity_token)` exactly +like `/resolve`, so a bottle can only ever queue or read its own proposals. `POST /bottles` / `DELETE` drive the full launch lifecycle: they mint (or tear down) the bottle in the registry AND broker the backend-native launch @@ -51,6 +63,7 @@ import typing from urllib.parse import urlsplit from ..paths import CONTROL_PLANE_TOKEN_ENV +from ..supervise_types import TOOLS from .service import Orchestrator # JSON body payload type (parsed request / rendered response). @@ -235,6 +248,56 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches return 200, {"responded": True} return 409, {"error": err} + if method == "POST" and route == "/supervise/propose": + # Agent half: queue a proposal, attributed to the caller resolved from + # (source_ip, identity_token) — never a caller-supplied slug — so the + # data plane can't forge attribution. Fail-closed 403 when unattributed. + try: + data = _parse_json_object(body) + except ValueError as e: + return 400, {"error": f"invalid JSON: {e}"} + source_ip = data.get("source_ip") + token = data.get("identity_token") + tool = data.get("tool") + proposed_file = data.get("proposed_file") + justification = data.get("justification") + if not isinstance(source_ip, str) or not source_ip: + return 400, {"error": "source_ip (string) is required"} + if not isinstance(tool, str) or tool not in TOOLS: + return 400, {"error": f"tool (string) must be one of {TOOLS}"} + if not isinstance(proposed_file, str) or not proposed_file: + return 400, {"error": "proposed_file (string) is required"} + if not isinstance(justification, str) or not justification: + return 400, {"error": "justification (string) is required"} + rec = orch.resolve(source_ip, token if isinstance(token, str) else "") + if rec is None: + return 403, {"error": "unattributed"} + proposal_id = orch.supervise_queue_proposal( + rec.bottle_id, tool=tool, proposed_file=proposed_file, + justification=justification, + ) + return 201, {"proposal_id": proposal_id} + + if method == "POST" and route == "/supervise/poll": + # Agent half: non-blocking read of the caller's own proposal decision. + # Attributed like /propose, and scoped to the resolved bottle id, so a + # guessed proposal_id can never read another bottle's response. + try: + data = _parse_json_object(body) + except ValueError as e: + return 400, {"error": f"invalid JSON: {e}"} + source_ip = data.get("source_ip") + token = data.get("identity_token") + proposal_id = data.get("proposal_id") + if not isinstance(source_ip, str) or not source_ip: + return 400, {"error": "source_ip (string) is required"} + if not isinstance(proposal_id, str) or not proposal_id: + return 400, {"error": "proposal_id (string) is required"} + rec = orch.resolve(source_ip, token if isinstance(token, str) else "") + if rec is None: + return 403, {"error": "unattributed"} + return 200, orch.supervise_poll_response(rec.bottle_id, proposal_id) + if method == "POST" and route == "/resolve": # The per-request lookup the multi-tenant gateway makes: returns the # bottle's policy. Requires a matching (source_ip, identity_token) diff --git a/bot_bottle/orchestrator/gateway.py b/bot_bottle/orchestrator/gateway.py index 741e20e..16bf1e6 100644 --- a/bot_bottle/orchestrator/gateway.py +++ b/bot_bottle/orchestrator/gateway.py @@ -26,15 +26,8 @@ from ..docker_cmd import run_docker from ..paths import ( CONTROL_PLANE_TOKEN_ENV, host_control_plane_token, - host_db_path, host_gateway_ca_dir, ) -from ..supervise import DB_PATH_IN_CONTAINER - -# The host DB dir is bind-mounted here so the gateway's supervise daemon -# writes its queued proposals into the ONE host DB (the same file the -# orchestrator container opens and the operator reaches over HTTP). -_SUPERVISE_DB_DIR_IN_CONTAINER = os.path.dirname(DB_PATH_IN_CONTAINER) # The gateway's mitmproxy writes its CA a beat after the container starts, so # reads poll for it rather than assuming it's there on a fresh launch. @@ -75,14 +68,6 @@ GATEWAY_DOCKERFILE = "Dockerfile.gateway" _REPO_ROOT = Path(__file__).resolve().parents[2] -def _host_db_dir() -> str: - """The host DB directory (created if missing), for the gateway's - supervise-DB bind-mount.""" - db_dir = host_db_path().parent - db_dir.mkdir(parents=True, exist_ok=True) - return str(db_dir) - - def rotate_gateway_ca(ca_dir: Path | None = None) -> list[Path]: """Delete the persisted mitmproxy CA so the next gateway start mints a fresh one — the explicit, deliberate CA-rollover path (issue #450). @@ -255,11 +240,10 @@ class DockerGateway(Gateway): # container recreation AND docker volume pruning (agents trust it) # — see host_gateway_ca_dir / issue #450. "--volume", f"{host_gateway_ca_dir()}:{MITMPROXY_HOME}", - # Share the one host DB: the supervise daemon queues proposals - # into the same file the orchestrator (and the operator, over - # HTTP) reads — no second, disconnected DB in the container. - "--volume", f"{_host_db_dir()}:{_SUPERVISE_DB_DIR_IN_CONTAINER}", - "--env", f"SUPERVISE_DB_PATH={DB_PATH_IN_CONTAINER}", + # No DB mount: the data plane (egress / supervise / git-gate) reaches + # the supervise queue over the control-plane RPC and never opens + # bot-bottle.db, so the gateway container gets no file handle on it + # (PRD 0070 / issue #469). ] for port in self._host_port_bindings: argv += ["--publish", f"0.0.0.0:{port}:{port}"] diff --git a/bot_bottle/orchestrator/lifecycle.py b/bot_bottle/orchestrator/lifecycle.py index cc3d82a..9560c07 100644 --- a/bot_bottle/orchestrator/lifecycle.py +++ b/bot_bottle/orchestrator/lifecycle.py @@ -29,14 +29,12 @@ from ..paths import ( host_control_plane_token, host_gateway_ca_dir, ) -from ..supervise import DB_PATH_IN_CONTAINER from .gateway import ( GATEWAY_DOCKERFILE, GATEWAY_IMAGE, GATEWAY_NETWORK, GatewayError, MITMPROXY_HOME, - _host_db_dir, ) DEFAULT_PORT = 8099 @@ -69,12 +67,12 @@ _INFRA_DAEMONS = "egress,git-http,supervise,orchestrator" # container. Separate from /app so the gateway's baked scripts # (egress_addon.py, egress-entrypoint.sh) are not overlaid. _SRC_IN_CONTAINER = "/bot-bottle-src" -# Bot-bottle host-root bind-mount inside the container (DB + state). +# Bot-bottle host-root bind-mount inside the container (DB + state). The +# orchestrator control plane opens bot-bottle.db under here (via BOT_BOTTLE_ROOT +# -> host_db_path()); it is the ONLY process in the container with a file +# handle on it (PRD 0070 / issue #469). _ROOT_IN_CONTAINER = "/bot-bottle-root" -# The supervise daemon writes proposals into the host DB directory. -_SUPERVISE_DB_DIR_IN_CONTAINER = os.path.dirname(DB_PATH_IN_CONTAINER) - _HEALTH_POLL_SECONDS = 0.25 _HEALTH_REQUEST_TIMEOUT_SECONDS = 1.0 @@ -202,9 +200,6 @@ class OrchestratorService: # recreation AND docker volume pruning (issue #450): every agent # trusts this one CA, so a fresh one would break all running bottles. "--volume", f"{host_gateway_ca_dir()}:{MITMPROXY_HOME}", - # Shared supervise DB (same file the operator reads over HTTP). - "--volume", f"{_host_db_dir()}:{_SUPERVISE_DB_DIR_IN_CONTAINER}", - "--env", f"SUPERVISE_DB_PATH={DB_PATH_IN_CONTAINER}", # Live control-plane source, mounted to a path that does not # overlay the gateway's baked /app scripts. "--volume", f"{self._repo_root}:{_SRC_IN_CONTAINER}:ro", diff --git a/bot_bottle/orchestrator/service.py b/bot_bottle/orchestrator/service.py index f1d2963..1f042f3 100644 --- a/bot_bottle/orchestrator/service.py +++ b/bot_bottle/orchestrator/service.py @@ -31,16 +31,23 @@ from .gateway import Gateway from ..supervise import ( AuditEntry, COMPONENT_FOR_TOOL, + POLL_STATUS_PENDING, + POLL_STATUS_UNKNOWN, + Proposal, Response, STATUS_APPROVED, STATUS_MODIFIED, STATUS_REJECTED, TOOL_EGRESS_ALLOW, TOOL_EGRESS_BLOCK, + archive_proposal, list_all_pending_proposals, read_proposal, + read_response, render_diff, + sha256_hex, write_audit_entry, + write_proposal, write_response, ) @@ -188,6 +195,60 @@ class Orchestrator: # The orchestrator owns the single DB *and* the live policy, so operator # decisions are applied here, server-side, and reached over HTTP by the # host TUI (no direct-DB access, one path for every backend). + # + # The *agent* half of the queue (queue a proposal, poll for its response) + # is served here too (PRD 0070): the data plane no longer opens the DB, so + # supervise_server / egress / git-gate reach the queue over the control + # plane. Both are keyed on the caller's orchestrator-resolved bottle id — + # never a caller-supplied slug — so a bottle can only queue or read its own + # proposals even if the data plane is compromised. + + def supervise_queue_proposal( + self, bottle_id: str, *, tool: str, proposed_file: str, justification: str, + ) -> str: + """Queue an agent-side proposal under `bottle_id` and return its id. + + `bottle_id` is the caller resolved from `(source_ip, identity_token)` by + the control plane, so the proposal is attributed to the calling bottle, + not to anything the (possibly hostile) data plane asserts. The + proposed-file self-hash is computed here so the agent never supplies it.""" + proposal = Proposal.new( + bottle_slug=bottle_id, + tool=tool, + proposed_file=proposed_file, + justification=justification, + current_file_hash=sha256_hex(proposed_file), + ) + write_proposal(proposal) + return proposal.id + + def supervise_poll_response(self, bottle_id: str, proposal_id: str) -> dict[str, object]: + """Non-blocking poll of one of `bottle_id`'s queued proposals. + + Returns `{"status": ...}`: a terminal decision (`approved`/`modified`/ + `rejected`, with `notes` + `final_file`) once the operator has responded, + `pending` while it's still queued undecided, or `unknown` when no such + queued proposal exists for this bottle (wrong id, or already resolved and + read). A decided proposal is archived here — the same terminal step the + data plane used to run after reading the response — so `pending` + proposals stay visible to the operator until decided *and* polled. + + Reads are scoped to `bottle_id` (the queue key), so a caller can never + read another bottle's proposal even with a guessed id.""" + try: + response = read_response(bottle_id, proposal_id) + except FileNotFoundError: + try: + read_proposal(bottle_id, proposal_id) + except FileNotFoundError: + return {"status": POLL_STATUS_UNKNOWN} + return {"status": POLL_STATUS_PENDING} + archive_proposal(bottle_id, proposal_id) + return { + "status": response.status, + "notes": response.notes, + "final_file": response.final_file, + } def supervise_pending(self) -> list[dict[str, object]]: """All pending proposals across bottles, FIFO, as JSON dicts diff --git a/bot_bottle/policy_resolver.py b/bot_bottle/policy_resolver.py index 06ed563..5a5dd15 100644 --- a/bot_bottle/policy_resolver.py +++ b/bot_bottle/policy_resolver.py @@ -64,28 +64,80 @@ class PolicyResolver: self._base = base_url.rstrip("/") self._timeout = timeout + def _post_json(self, path: str, payload: dict[str, object]) -> dict[str, object] | None: + """POST `payload` to control-plane `path`, returning the JSON object body + — or None when the orchestrator answers `403` (unattributed / fail + closed). Raises `PolicyResolveError` on an unreachable / unexpected-status + / malformed response so every caller can fail closed. Shared by + `_post_resolve` and the supervise propose/poll RPCs.""" + body = json.dumps(payload).encode() + req = urllib.request.Request( + f"{self._base}{path}", data=body, method="POST", + headers={"Content-Type": "application/json", **_control_auth_headers()}, + ) + try: + with urllib.request.urlopen(req, timeout=self._timeout) as resp: + data = json.loads(resp.read()) + except urllib.error.HTTPError as e: + if e.code == 403: + return None # unattributed → fail closed (caller denies) + raise PolicyResolveError(f"{path} returned HTTP {e.code}") from e + except (urllib.error.URLError, TimeoutError, OSError, ValueError) as e: + raise PolicyResolveError(f"{path} unreachable or malformed: {e}") from e + return data if isinstance(data, dict) else {} + def _post_resolve(self, source_ip: str, identity_token: str) -> dict[str, object] | None: """The orchestrator's `/resolve` payload for this client, or None if unattributed (a clean `403`). Raises `PolicyResolveError` on an unreachable / unexpected-status / malformed response so every caller can fail closed. Shared by `resolve` and `resolve_bottle_id`.""" - body = json.dumps( - {"source_ip": source_ip, "identity_token": identity_token} - ).encode() - req = urllib.request.Request( - f"{self._base}/resolve", data=body, method="POST", - headers={"Content-Type": "application/json", **_control_auth_headers()}, + return self._post_json( + "/resolve", {"source_ip": source_ip, "identity_token": identity_token}, ) - try: - with urllib.request.urlopen(req, timeout=self._timeout) as resp: - payload = json.loads(resp.read()) - except urllib.error.HTTPError as e: - if e.code == 403: - return None # unattributed → fail closed (caller denies) - raise PolicyResolveError(f"/resolve returned HTTP {e.code}") from e - except (urllib.error.URLError, TimeoutError, OSError, ValueError) as e: - raise PolicyResolveError(f"/resolve unreachable or malformed: {e}") from e - return payload if isinstance(payload, dict) else {} + + def propose_supervise( + self, + source_ip: str, + identity_token: str, + *, + tool: str, + proposed_file: str, + justification: str, + ) -> str | None: + """Queue a supervise proposal on the control plane and return its + `proposal_id`. The orchestrator attributes the proposal to the calling + bottle by `(source_ip, identity_token)` — exactly like `/resolve` — so a + bottle can only ever queue its *own* proposals. Returns None when the + pair is unattributed (a clean `403`); raises `PolicyResolveError` if the + orchestrator can't be reached, so the data-plane caller fails closed + (blocks / refuses) rather than silently dropping the proposal.""" + payload = self._post_json("/supervise/propose", { + "source_ip": source_ip, + "identity_token": identity_token, + "tool": tool, + "proposed_file": proposed_file, + "justification": justification, + }) + if payload is None: + return None + proposal_id = payload.get("proposal_id") + return proposal_id if isinstance(proposal_id, str) and proposal_id else None + + def poll_supervise( + self, source_ip: str, identity_token: str, proposal_id: str, + ) -> dict[str, object] | None: + """Poll a queued proposal for the operator's decision, **non-blocking**. + Attributed by `(source_ip, identity_token)` so a bottle can only read its + *own* proposal's response. Returns `{"status": ...}` where status is one + of the terminal decisions (`approved`/`modified`/`rejected`, carrying + `notes` + `final_file`), `pending` (queued, no decision yet), or + `unknown` (no such queued proposal for this bottle). Returns None when + unattributed; raises `PolicyResolveError` if unreachable.""" + return self._post_json("/supervise/poll", { + "source_ip": source_ip, + "identity_token": identity_token, + "proposal_id": proposal_id, + }) def resolve(self, source_ip: str, identity_token: str = "") -> str | None: """The calling bottle's policy blob, or None if unattributed. Always diff --git a/bot_bottle/supervise.py b/bot_bottle/supervise.py index a1bc0aa..827bd30 100644 --- a/bot_bottle/supervise.py +++ b/bot_bottle/supervise.py @@ -40,6 +40,8 @@ from pathlib import Path from .supervise_types import ( ACTION_OPERATOR_EDIT, AuditEntry, + POLL_STATUS_PENDING, + POLL_STATUS_UNKNOWN, Proposal, Response, STATUSES, @@ -257,6 +259,8 @@ __all__ = [ "STATUS_APPROVED", "STATUS_MODIFIED", "STATUS_REJECTED", + "POLL_STATUS_PENDING", + "POLL_STATUS_UNKNOWN", "SUPERVISE_HOSTNAME", "SUPERVISE_PORT", "Supervise", diff --git a/bot_bottle/supervise_server.py b/bot_bottle/supervise_server.py index 64c4069..7a90d21 100644 --- a/bot_bottle/supervise_server.py +++ b/bot_bottle/supervise_server.py @@ -7,9 +7,13 @@ config changes when stuck. The tools are `egress-allow`, Each queued proposal tool call: 1. Validates the proposed file syntactically. - 2. Writes a Proposal to the host SQLite database. - 3. Blocks polling for a matching Response row, up to a short grace - window (`SUPERVISE_RESPONSE_TIMEOUT_SECONDS`, default 30s). + 2. Queues a Proposal over the control-plane RPC (`POST /supervise/propose`), + which attributes it to the calling bottle by (source_ip, identity_token) + and writes it to the one orchestrator-owned SQLite database. This daemon + never opens `bot-bottle.db` itself (PRD 0070 / issue #469). + 3. Polls the control plane (`POST /supervise/poll`) for a matching Response, + up to a short grace window (`SUPERVISE_RESPONSE_TIMEOUT_SECONDS`, + default 30s). 4. On a decision within the window, returns the operator's `{status, notes}`. On timeout, returns `status: pending` **with the proposal id** and leaves the proposal queued — the flow is @@ -24,8 +28,8 @@ without holding an HTTP request open. One shared server fronts every bottle (PRD 0070) and attributes each proposal to the calling bottle by source IP, resolved from the orchestrator — an unattributed or unreachable source fails closed. BOT_BOTTLE_ORCHESTRATOR_URL -is mandatory: there is no fixed-slug single-tenant fallback. SUPERVISE_DB_PATH -points at the bind-mounted host database. +is mandatory: there is no fixed-slug single-tenant fallback, and the queue is +reached only over that control-plane RPC (no direct database handle). Speaks MCP over HTTP+JSON-RPC. Methods handled: @@ -51,7 +55,7 @@ import socketserver import sys import time import typing -from dataclasses import dataclass, replace +from dataclasses import dataclass from bot_bottle.constants import IDENTITY_HEADER from bot_bottle.egress_addon_core import ( @@ -312,7 +316,9 @@ def validate_proposed_file(tool: str, content: str) -> None: @dataclass(frozen=True) class ServerConfig: - bottle_slug: str + # No bottle_slug: the calling bottle is attributed per request by the + # control plane from (source_ip, identity_token), so nothing on this daemon + # is keyed by a single slug (PRD 0070 / issue #469). response_timeout_seconds: float = DEFAULT_RESPONSE_TIMEOUT_SECONDS @@ -331,14 +337,22 @@ def handle_tools_list(_params: dict[str, object]) -> dict[str, object]: def handle_tools_call( params: dict[str, object], config: ServerConfig, + *, + resolver: "PolicyResolver", + source_ip: str, + identity_token: str, ) -> dict[str, object]: - """Validates the proposal, writes it to the queue, blocks waiting - for a Response, returns the result wrapped in MCP `content`. + """Validates the proposal, queues it over the control-plane RPC, polls for + a Response through the grace window, and returns the result wrapped in MCP + `content`. - `list-egress-routes` never reaches here — the handler answers it from - the calling bottle's resolved policy before dispatching (see - `MCPHandler._dispatch`); this path is the queued, operator-approved - `egress-allow` / `egress-block` tools.""" + The queue lives behind the orchestrator now — `propose_supervise` attributes + the proposal to the calling bottle by `(source_ip, identity_token)` and + `poll_supervise` reads only that bottle's own responses, so this daemon + never opens `bot-bottle.db`. `list-egress-routes` never reaches here — the + handler answers it from the calling bottle's resolved policy before + dispatching (see `MCPHandler._dispatch`); this path is the queued, + operator-approved `egress-allow` / `egress-block` tools.""" name = params.get("name") if not isinstance(name, str): raise _RpcClientError(ERR_INVALID_PARAMS, "tools/call missing 'name'") @@ -366,59 +380,100 @@ def handle_tools_call( else: raise _RpcClientError(ERR_INVALID_PARAMS, f"unknown tool {name!r}") - proposal = _sv.Proposal.new( - bottle_slug=config.bottle_slug, - tool=name, - proposed_file=proposed_file, - justification=justification, - current_file_hash=_sv.sha256_hex(proposed_file), + proposal_id = _queue_proposal( + resolver, source_ip, identity_token, + tool=name, proposed_file=proposed_file, justification=justification, ) - try: - _sv.write_proposal(proposal) - except OSError as e: - raise _RpcInternalError(f"failed to write proposal to queue: {e}") from e sys.stderr.write( - f"supervise: queued proposal {proposal.id} ({name}) " - f"for bottle {config.bottle_slug}; waiting for operator...\n" + f"supervise: queued proposal {proposal_id} ({name}); waiting for operator...\n" ) sys.stderr.flush() - deadline = time.monotonic() + config.response_timeout_seconds - try: - response = _sv.wait_for_response( - config.bottle_slug, - proposal.id, - poll_interval=MIN_RESPONSE_POLL_INTERVAL_SECONDS, - deadline=deadline, - ) - except TimeoutError: - text = format_pending_response_text(proposal.id, config.response_timeout_seconds) - return { - "content": [{"type": "text", "text": text}], - "isError": False, - } - try: - _sv.archive_proposal(config.bottle_slug, proposal.id) - except OSError as e: - raise _RpcInternalError(f"failed to archive proposal: {e}") from e - text = format_response_text(response) + deadline = time.monotonic() + config.response_timeout_seconds + while time.monotonic() < deadline: + response = _poll_response(resolver, source_ip, identity_token, proposal_id) + if response is not None: + return { + "content": [{"type": "text", "text": format_response_text(response)}], + "isError": response.status == _sv.STATUS_REJECTED, + } + time.sleep(MIN_RESPONSE_POLL_INTERVAL_SECONDS) + + text = format_pending_response_text(proposal_id, config.response_timeout_seconds) return { "content": [{"type": "text", "text": text}], - "isError": response.status == _sv.STATUS_REJECTED, + "isError": False, } +def _queue_proposal( + resolver: "PolicyResolver", + source_ip: str, + identity_token: str, + *, + tool: str, + proposed_file: str, + justification: str, +) -> str: + """Queue a proposal over the control plane and return its id, mapping the + resolver's fail-closed outcomes to internal errors so `do_POST` returns a + clean -32603 rather than leaking the cause to the agent.""" + try: + proposal_id = resolver.propose_supervise( + source_ip, identity_token, + tool=tool, proposed_file=proposed_file, justification=justification, + ) + except PolicyResolveError as e: + raise _RpcInternalError(f"orchestrator unreachable, cannot queue proposal: {e}") from e + if proposal_id is None: + raise _RpcInternalError("request source is not attributed to a bottle") + return proposal_id + + +def _poll_response( + resolver: "PolicyResolver", source_ip: str, identity_token: str, proposal_id: str, +) -> "_sv.Response | None": + """One non-blocking control-plane poll for `proposal_id`'s decision. Returns + the operator's Response once decided, or None while it is still pending + (or `unknown` — treated as still-waiting so a transient blip degrades to the + pending-timeout path rather than a spurious error mid grace window).""" + try: + result = resolver.poll_supervise(source_ip, identity_token, proposal_id) + except PolicyResolveError as e: + raise _RpcInternalError(f"orchestrator unreachable, cannot poll proposal: {e}") from e + if result is None: + raise _RpcInternalError("request source is not attributed to a bottle") + if result.get("status") in _sv.STATUSES: + return _response_from_poll(proposal_id, result) + return None + + +def _response_from_poll(proposal_id: str, result: "dict[str, object]") -> "_sv.Response": + """Build a Response from a decided `/supervise/poll` payload.""" + final_file = result.get("final_file") + return _sv.Response( + proposal_id=proposal_id, + status=str(result.get("status")), + notes=str(result.get("notes") or ""), + final_file=final_file if isinstance(final_file, str) else None, + ) + + def handle_check_proposal( params: dict[str, object], - config: ServerConfig, + *, + resolver: "PolicyResolver", + source_ip: str, + identity_token: str, ) -> dict[str, object]: """Non-blocking poll of a queued proposal's decision, by id. Never creates a Proposal (so `check-proposal` isn't in `TOOLS`); it only - reads the queue. Resolution order mirrors the synchronous path's terminal - step — a decided proposal is archived here exactly as `handle_tools_call` - archives it after `wait_for_response`, so `pending` proposals stay visible - to the operator until they're both decided *and* polled.""" + reads the queue over the control-plane RPC, attributed by + `(source_ip, identity_token)` so it can only read *this* bottle's own + proposals. The orchestrator archives a decided proposal on read — the same + terminal step the synchronous path triggers — so `pending` proposals stay + visible to the operator until they're both decided *and* polled.""" args_raw = params.get("arguments", {}) if not isinstance(args_raw, dict): raise _RpcClientError(ERR_INVALID_PARAMS, "tools/call 'arguments' must be an object") @@ -431,25 +486,24 @@ def handle_check_proposal( proposal_id = proposal_id.strip() try: - response = _sv.read_response(config.bottle_slug, proposal_id) - except FileNotFoundError: - # No decision yet — distinguish "still queued" from "unknown id". - try: - _sv.read_proposal(config.bottle_slug, proposal_id) - except FileNotFoundError: - return { - "content": [{"type": "text", "text": format_unknown_proposal_text(proposal_id)}], - "isError": True, - } + result = resolver.poll_supervise(source_ip, identity_token, proposal_id) + except PolicyResolveError as e: + raise _RpcInternalError(f"orchestrator unreachable, cannot poll proposal: {e}") from e + if result is None: + raise _RpcInternalError("request source is not attributed to a bottle") + status = result.get("status") + + if status == _sv.POLL_STATUS_UNKNOWN: + return { + "content": [{"type": "text", "text": format_unknown_proposal_text(proposal_id)}], + "isError": True, + } + if status == _sv.POLL_STATUS_PENDING: return { "content": [{"type": "text", "text": format_still_pending_text(proposal_id)}], "isError": False, } - - try: - _sv.archive_proposal(config.bottle_slug, proposal_id) - except OSError as e: - raise _RpcInternalError(f"failed to archive proposal: {e}") from e + response = _response_from_poll(proposal_id, result) return { "content": [{"type": "text", "text": format_response_text(response)}], "isError": response.status == _sv.STATUS_REJECTED, @@ -591,16 +645,34 @@ class MCPHandler(http.server.BaseHTTPRequestHandler): # — silently dropping base routes like api.anthropic.com on approval. if req.params.get("name") == _sv.TOOL_LIST_EGRESS_ROUTES: return self._resolved_routes_payload() + resolver = self._resolver_or_fail() + source_ip = self.client_address[0] + token = self._identity_token() # `check-proposal` is a non-blocking read of the calling bottle's - # own queue — attributed by source IP like a proposal, but it - # never queues or blocks. + # own queue — attributed by (source_ip, identity_token) like a + # proposal, but it never queues or blocks. if req.params.get("name") == _sv.TOOL_CHECK_PROPOSAL: - return handle_check_proposal(req.params, self._attributed_config(config)) - # Attribute the proposal to the source-IP-resolved bottle, so the one - # shared server queues each bottle's proposal under its own slug. - return handle_tools_call(req.params, self._attributed_config(config)) + return handle_check_proposal( + req.params, resolver=resolver, + source_ip=source_ip, identity_token=token, + ) + # The control plane attributes the proposal to the source-IP + token + # resolved bottle, so the one shared queue holds each bottle's + # proposal under its own id — no slug is asserted by this daemon. + return handle_tools_call( + req.params, config, resolver=resolver, + source_ip=source_ip, identity_token=token, + ) raise _RpcClientError(ERR_METHOD_NOT_FOUND, f"method not found: {method}") + def _identity_token(self) -> str: + """The agent's per-bottle identity token from the request header (the + MCP client sends it via `mcp add --header`), or empty. The control plane + requires the (source_ip, token) pair, so a missing/wrong token + fail-closes there.""" + headers = getattr(self, "headers", None) + return headers.get(IDENTITY_HEADER, "") if headers is not None else "" + def _resolver_or_fail(self) -> "PolicyResolver": """This server's policy resolver. A server started without one is a misconfiguration, not a tenancy mode — fail closed rather than @@ -612,40 +684,18 @@ class MCPHandler(http.server.BaseHTTPRequestHandler): def _resolved_routes_payload(self) -> dict[str, object]: """The calling bottle's live egress routes as the `list-egress-routes` - JSON payload, resolved by (source_ip, identity token). Fail-closed like - `_attributed_config`: an unattributed source or an unreachable - orchestrator yields an empty route list (never another bottle's), - courtesy of `resolve_client_context`.""" + JSON payload, resolved by (source_ip, identity token). Fail-closed: an + unattributed source or an unreachable orchestrator yields an empty route + list (never another bottle's), courtesy of `resolve_client_context`.""" resolver = self._resolver_or_fail() - headers = getattr(self, "headers", None) - token = headers.get(IDENTITY_HEADER, "") if headers is not None else "" conf, _slug, _tokens = resolve_client_context( - resolver, self.client_address[0], token, + resolver, self.client_address[0], self._identity_token(), ) body = json.dumps( {"routes": [route_to_yaml_dict(r) for r in conf.routes]}, indent=2, ) return {"content": [{"type": "text", "text": body}], "isError": False} - def _attributed_config(self, config: ServerConfig) -> ServerConfig: - """The ServerConfig with `bottle_slug` bound to *this request's* bottle: - the bottle id attributed from the source IP — **fail-closed**, an - unattributed or unreachable source raises so no proposal is queued under - the wrong (or empty) slug.""" - resolver = self._resolver_or_fail() - # The agent's MCP client sends the identity token as a request header - # (provisioned via `mcp add --header`); the orchestrator requires the - # (source_ip, token) pair, so a missing/wrong token fail-closes below. - headers = getattr(self, "headers", None) - token = headers.get(IDENTITY_HEADER, "") if headers is not None else "" - try: - bottle_id = resolver.resolve_bottle_id(self.client_address[0], token) - except PolicyResolveError as e: - raise _RpcInternalError(f"orchestrator unreachable, cannot attribute: {e}") from e - if not bottle_id: - raise _RpcInternalError("request source is not attributed to a bottle") - return replace(config, bottle_slug=bottle_id) - def _write_jsonrpc(self, body: bytes) -> None: self.send_response(200) self.send_header("Content-Type", "application/json") @@ -668,10 +718,10 @@ class MCPHandler(http.server.BaseHTTPRequestHandler): class MCPServer(socketserver.ThreadingMixIn, http.server.HTTPServer): allow_reuse_address = True daemon_threads = True - config: ServerConfig = ServerConfig(bottle_slug="") - # Set by `serve`; every proposal is attributed to the source-IP-resolved - # bottle. The class default is a placeholder — a server without a resolver - # fails closed per request (see `_resolver_or_fail`). + config: ServerConfig = ServerConfig() + # Set by `serve`; every proposal is attributed to the source-IP + token + # resolved bottle by the control plane. A server without a resolver fails + # closed per request (see `_resolver_or_fail`). policy_resolver: "PolicyResolver | None" = None @@ -686,12 +736,9 @@ def serve( response_timeout_seconds: float = DEFAULT_RESPONSE_TIMEOUT_SECONDS, ) -> typing.NoReturn: server = MCPServer((bind, port), MCPHandler) - # bottle_slug is a placeholder: every request's proposal is attributed to - # the source-IP-resolved bottle (see MCPHandler._attributed_config). - server.config = ServerConfig( - bottle_slug="", - response_timeout_seconds=response_timeout_seconds, - ) + # Every request's proposal is attributed to the source-IP + token resolved + # bottle by the control plane (see MCPHandler._dispatch). + server.config = ServerConfig(response_timeout_seconds=response_timeout_seconds) server.policy_resolver = resolver sys.stderr.write( f"supervise listening on {bind}:{port}; multi-tenant; " diff --git a/bot_bottle/supervise_types.py b/bot_bottle/supervise_types.py index f88203a..6bb65e5 100644 --- a/bot_bottle/supervise_types.py +++ b/bot_bottle/supervise_types.py @@ -37,6 +37,15 @@ STATUS_MODIFIED = "modified" STATUS_REJECTED = "rejected" STATUSES: tuple[str, ...] = (STATUS_APPROVED, STATUS_MODIFIED, STATUS_REJECTED) +# Non-terminal markers the control-plane supervise-poll RPC returns when a +# proposal has no operator decision yet (`PENDING`) or no queued proposal with +# that id exists for the calling bottle (`UNKNOWN`). They are never a +# `Response.status` — only the poll wire contract (see +# `Orchestrator.supervise_poll_response`) — so they are deliberately outside +# `STATUSES`. +POLL_STATUS_PENDING = "pending" +POLL_STATUS_UNKNOWN = "unknown" + ACTION_OPERATOR_EDIT = "operator-edit" @@ -157,6 +166,8 @@ __all__ = [ "STATUS_APPROVED", "STATUS_MODIFIED", "STATUS_REJECTED", + "POLL_STATUS_PENDING", + "POLL_STATUS_UNKNOWN", "TOOLS", "TOOL_EGRESS_ALLOW", "TOOL_EGRESS_BLOCK", diff --git a/docs/prds/0070-per-host-orchestrator.md b/docs/prds/0070-per-host-orchestrator.md index 40e7fb3..0b7c579 100644 --- a/docs/prds/0070-per-host-orchestrator.md +++ b/docs/prds/0070-per-host-orchestrator.md @@ -299,11 +299,14 @@ and integrate a console against. the data plane and the console reach state through the control-plane RPC, never a direct file handle.** No agent-facing component gets the file, so none can forge attribution. (This supersedes the earlier `ro`-mount idea.) - - *Transitional caveat:* today the per-bottle **supervise gateway - rw-bind-mounts `bot-bottle.db`** to write proposals — exactly the - pattern the orchestrator removes (supervise consolidates into the - orchestrator; gateway writes become RPC calls). Until that lands, don't - put the attribution registry behind a data-plane-writable mount. + This rule is now **in force** across the data plane (issue #469): the + supervise daemon, the egress addon, and the git-gate pre-receive hook queue + proposals and poll for their responses over the agent-side supervise RPC + (`POST /supervise/propose` + `POST /supervise/poll`, attributed by + `(source_ip, identity_token)` exactly like `/resolve`), so a bottle can only + ever queue or read its own proposals. The gateway/data-plane containers no + longer bind-mount the DB directory or carry `SUPERVISE_DB_PATH`; the + orchestrator is the sole opener of the one file. Implementation note for the VM slices: SQLite **WAL** over a guest share (virtiofs/9p) is finicky (the `-shm`/`-wal` files need real mmap/locking), diff --git a/tests/unit/test_egress_addon_request_flow.py b/tests/unit/test_egress_addon_request_flow.py index fd90c76..3de75b9 100644 --- a/tests/unit/test_egress_addon_request_flow.py +++ b/tests/unit/test_egress_addon_request_flow.py @@ -209,6 +209,7 @@ from bot_bottle.egress_addon_core import ( # noqa: E402 Route, route_to_yaml_dict, ) +from bot_bottle.policy_resolver import PolicyResolveError # noqa: E402 # --------------------------------------------------------------------------- @@ -266,7 +267,42 @@ def _config_to_policy(config: Config) -> str: }) + "\n" -class _StaticResolver: +class _SuperviseRpcFake: + """Mixin adding the control-plane supervise RPCs to a fake resolver, now + that the egress data plane queues/polls proposals over RPC instead of + opening the DB (issue #469). `supervise_status` is the operator's eventual + decision (None models a timeout — poll stays `pending`); `propose_error` + models an unreachable orchestrator. `propose_calls` records what was queued + so a test can assert the proposal was attributed by the caller's source IP.""" + + supervise_status: "str | None" = None + propose_error: bool = False + + @property + def propose_calls(self) -> list: + if not hasattr(self, "_propose_calls"): + self._propose_calls: list = [] + return self._propose_calls + + def propose_supervise( + self, source_ip, identity_token, *, tool, proposed_file, justification, + ): + del proposed_file, justification + self.propose_calls.append( + {"source_ip": source_ip, "identity_token": identity_token, "tool": tool} + ) + if self.propose_error: + raise PolicyResolveError("orchestrator down") + return "prop-1" + + def poll_supervise(self, source_ip, identity_token, proposal_id): + del source_ip, identity_token, proposal_id + if self.supervise_status is None: + return {"status": "pending"} + return {"status": self.supervise_status, "notes": "", "final_file": None} + + +class _StaticResolver(_SuperviseRpcFake): """Fake orchestrator resolver that serves one Config (+ optional bottle id and per-bottle tokens) for every client — the host-test stand-in for a bottle's policy now that egress is resolver-only.""" @@ -323,7 +359,7 @@ def _with_client_ip(flow: _Flow, ip: str) -> _Flow: return flow -class _CtxResolver: +class _CtxResolver(_SuperviseRpcFake): """Fake orchestrator resolver: maps source IP -> bottle id, and grants the same allow-list to any attributed bottle (unattributed -> deny).""" @@ -516,66 +552,33 @@ class TestOutboundDlpPolicy(unittest.TestCase): # --------------------------------------------------------------------------- -def _fake_sv(response_status: str | None) -> types.SimpleNamespace: - """Stand-in for the `supervise` module the adapter queues proposals to. - - `response_status` of None models a timeout (read_response never returns a - decision); a status string models the operator's eventual answer.""" - def _new_proposal(**_kw: Any) -> Any: - return types.SimpleNamespace(id="prop-1") - - def _sha256_hex(_payload: Any) -> str: - return "hash" - - def _noop(*_args: Any) -> None: - return None - - def _read_response(_slug: Any, _pid: Any) -> Any: - if response_status is None: - raise OSError("not written yet") # forces poll -> timeout - return types.SimpleNamespace(status=response_status) - - ns = types.SimpleNamespace() - ns.STATUS_APPROVED = "approved" - ns.STATUS_MODIFIED = "modified" - ns.TOOL_EGRESS_TOKEN_ALLOW = "egress_token_allow" - ns.Proposal = types.SimpleNamespace(new=_new_proposal) - ns.sha256_hex = _sha256_hex - ns.write_proposal = _noop - ns.archive_proposal = _noop - ns.read_response = _read_response - return ns - - class TestSuperviseBranch(unittest.TestCase): - def _supervised_addon(self) -> EgressAddon: + def _supervised_addon(self, status: str | None) -> EgressAddon: addon = _addon(Config(routes=(Route(host="api.example.com"),)), slug="test-bottle") addon._token_allow_timeout = 0.05 + cast(Any, addon._resolver).supervise_status = status return addon def test_operator_approval_allows_token_and_forwards(self) -> None: - addon = self._supervised_addon() + addon = self._supervised_addon("approved") flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")) - with patch.object(_ea_mod, "_sv", _fake_sv("approved")): - _run_request(addon, flow) + _run_request(addon, flow) self.assertIsNone(flow.response) # forwarded after approval # Approval lands in the calling bottle's safelist (keyed by slug). self.assertIn(_OPENAI_KEY, addon._safe_tokens_for("test-bottle")) def test_operator_rejection_blocks(self) -> None: - addon = self._supervised_addon() + addon = self._supervised_addon("rejected") flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")) - with patch.object(_ea_mod, "_sv", _fake_sv("rejected")): - _run_request(addon, flow) + _run_request(addon, flow) assert flow.response is not None self.assertEqual(403, flow.response.status_code) self.assertIn("rejected", flow.response.get_text()) def test_supervise_timeout_blocks(self) -> None: - addon = self._supervised_addon() + addon = self._supervised_addon(None) # poll stays pending -> timeout flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")) - with patch.object(_ea_mod, "_sv", _fake_sv(None)): - _run_request(addon, flow) + _run_request(addon, flow) assert flow.response is not None self.assertEqual(403, flow.response.status_code) self.assertIn("timed out", flow.response.get_text()) @@ -772,19 +775,14 @@ class TestRedactSurfaces(unittest.TestCase): class TestSuperviseWriteFailure(unittest.TestCase): - def test_write_proposal_oserror_blocks(self) -> None: + def test_propose_rpc_error_blocks(self) -> None: + # An unreachable orchestrator (propose RPC raises) fails closed: the + # request is blocked rather than forwarded unsupervised. addon = _addon(Config(routes=(Route(host="api.example.com"),)), slug="test-bottle") addon._token_allow_timeout = 0.05 + cast(Any, addon._resolver).propose_error = True flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")) - - fake = _fake_sv("approved") - - def _raise(_p: Any) -> None: - raise OSError("disk full") - - fake.write_proposal = _raise - with patch.object(_ea_mod, "_sv", fake): - _run_request(addon, flow) + _run_request(addon, flow) assert flow.response is not None self.assertEqual(403, flow.response.status_code) @@ -844,22 +842,23 @@ class TestSuperviseMultiTenant(unittest.TestCase): """Consolidated gateway: supervise proposals + the DLP safelist are keyed per bottle, resolved by source IP (PRD 0070).""" - def _consolidated_addon(self) -> EgressAddon: + def _consolidated_addon(self, status: str | None = "approved") -> EgressAddon: # Static config is empty; the resolver supplies each bottle's config. addon = _addon(Config(routes=())) - addon._resolver = cast(Any, _CtxResolver({"10.0.0.1": "bottle-a", "10.0.0.2": "bottle-b"})) + resolver = _CtxResolver({"10.0.0.1": "bottle-a", "10.0.0.2": "bottle-b"}) + resolver.supervise_status = status + addon._resolver = cast(Any, resolver) addon._token_allow_timeout = 0.05 return addon def test_approval_is_scoped_to_the_calling_bottle(self) -> None: - addon = self._consolidated_addon() + addon = self._consolidated_addon("approved") # bottle-a (10.0.0.1) sends the token; the operator approves. flow = _with_client_ip( _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")), "10.0.0.1", ) - with patch.object(_ea_mod, "_sv", _fake_sv("approved")): - _run_request(addon, flow) + _run_request(addon, flow) self.assertIsNone(flow.response) # forwarded after approval # The approval lands ONLY in bottle-a's safelist — never bottle-b's. # A global set here would be the cross-tenant leak this slice closes. @@ -867,22 +866,19 @@ class TestSuperviseMultiTenant(unittest.TestCase): self.assertNotIn(_OPENAI_KEY, addon._safe_tokens_for("bottle-b")) def test_proposal_is_attributed_to_the_source_ip_bottle(self) -> None: - addon = self._consolidated_addon() - seen: list[str] = [] - fake = _fake_sv("approved") - - def _capture(**kw: Any) -> Any: - seen.append(kw["bottle_slug"]) - return types.SimpleNamespace(id="p") - - fake.Proposal = types.SimpleNamespace(new=_capture) + addon = self._consolidated_addon("approved") flow = _with_client_ip( _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")), "10.0.0.2", ) - with patch.object(_ea_mod, "_sv", fake): - _run_request(addon, flow) - self.assertEqual(["bottle-b"], seen) # proposal keyed by the resolved bottle + _run_request(addon, flow) + # The addon forwards the caller's source IP to the control plane, which + # attributes the proposal server-side by (source_ip, identity_token) — + # the addon never asserts a slug. The resolved bottle keys the safelist. + calls = cast(Any, addon._resolver).propose_calls + self.assertEqual(["10.0.0.2"], [c["source_ip"] for c in calls]) + self.assertIn(_OPENAI_KEY, addon._safe_tokens_for("bottle-b")) + self.assertNotIn(_OPENAI_KEY, addon._safe_tokens_for("bottle-a")) def test_auth_token_injected_from_resolved_tokens(self) -> None: # The bottle's upstream token comes from /resolve (in-memory on the @@ -925,16 +921,17 @@ class TestSuperviseMultiTenant(unittest.TestCase): self.assertIsNotNone(flow.response) # blocked — token unset def test_unattributed_source_ip_cannot_supervise(self) -> None: - addon = self._consolidated_addon() + addon = self._consolidated_addon("approved") # 10.9.9.9 is not in the resolver map -> deny-all config, empty slug. flow = _with_client_ip( _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")), "10.9.9.9", ) - with patch.object(_ea_mod, "_sv", _fake_sv("approved")): - _run_request(addon, flow) + _run_request(addon, flow) self.assertIsNotNone(flow.response) # blocked (no route, no supervise) self.assertNotIn(_OPENAI_KEY, addon._safe_tokens_for("")) + # Never even reached the queue — no route means no supervise proposal. + self.assertEqual([], cast(Any, addon._resolver).propose_calls) class TestMultiTenantInboundDlp(unittest.TestCase): diff --git a/tests/unit/test_git_gate.py b/tests/unit/test_git_gate.py index cca39c5..fd899f0 100644 --- a/tests/unit/test_git_gate.py +++ b/tests/unit/test_git_gate.py @@ -213,22 +213,25 @@ class TestHookRender(unittest.TestCase): # the suppressed findings for human approval. self.assertIn("--ignore-gitleaks-allow", hook) self.assertIn("--report-format=json", hook) - self.assertIn("tool=_sv.TOOL_GITLEAKS_ALLOW", hook) - self.assertIn("_sv.write_proposal", hook) - self.assertIn("_sv.read_response", hook) - self.assertIn("SUPERVISE_BOTTLE_SLUG", hook) + # The hook queues + polls over the control-plane RPC — it no longer + # opens the DB directly (PRD 0070 / issue #469). + self.assertIn("tool=TOOL_GITLEAKS_ALLOW", hook) + self.assertIn("propose_supervise", hook) + self.assertIn("poll_supervise", hook) + self.assertIn("SUPERVISE_SOURCE_IP", hook) + self.assertIn("SUPERVISE_IDENTITY_TOKEN", hook) self.assertIn("supervisor approved # gitleaks:allow", hook) self.assertIn("supervisor rejected # gitleaks:allow", hook) def test_inline_gitleaks_allow_python_imports_work_in_gateway_layout(self): hook = git_gate_render_hook() - # The gateway image copies supervise.py flat under /app, while - # host-side tests import it through the bot_bottle package. - # Hooks execute from the bare repo directory, so the embedded - # Python must include /app and support both import layouts. + # The gateway image copies the package modules flat under /app, while + # host-side tests import them through the bot_bottle package. Hooks + # execute from the bare repo directory, so the embedded Python must + # include /app and support both import layouts. self.assertIn('PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}"', hook) - self.assertIn("import supervise as _sv", hook) - self.assertIn("from bot_bottle import supervise as _sv", hook) + self.assertIn("from bot_bottle.policy_resolver import PolicyResolver", hook) + self.assertIn("from policy_resolver import PolicyResolver", hook) def test_inline_gitleaks_allow_fails_closed_without_supervisor(self): hook = git_gate_render_hook() diff --git a/tests/unit/test_git_http_backend.py b/tests/unit/test_git_http_backend.py index afe264c..405f0b0 100644 --- a/tests/unit/test_git_http_backend.py +++ b/tests/unit/test_git_http_backend.py @@ -112,11 +112,13 @@ class TestGitHttpBackend(unittest.TestCase): ).strip() self.assertEqual(head, cloned) - def test_consolidated_push_stamps_bottle_slug_for_the_hook(self): - # In consolidated mode the backend attributes the push by source IP and - # stamps SUPERVISE_BOTTLE_SLUG= into the CGI env, so the - # gitleaks-allow pre-receive hook queues its proposal under the right - # bottle. The hook here just records what it received. + def test_consolidated_push_stamps_supervise_attribution_for_the_hook(self): + # In consolidated mode the backend attributes the push by (source IP, + # identity token) and stamps SUPERVISE_SOURCE_IP + SUPERVISE_IDENTITY_TOKEN + # into the CGI env, so the gitleaks-allow pre-receive hook can queue its + # proposal over the control-plane RPC (which re-resolves the bottle from + # exactly that pair — PRD 0070 / issue #469). The hook here just records + # the source IP it received. from http.server import ThreadingHTTPServer bottle_id = "bottleab12" @@ -130,10 +132,10 @@ class TestGitHttpBackend(unittest.TestCase): ["git", "-C", str(bare), "config", "http.receivepack", "true"], check=True, ) - capture = root / "slug-capture" + capture = root / "source-ip-capture" hook = bare / "hooks" / "pre-receive" hook.write_text( - f"#!/bin/sh\nprintf '%s' \"${{SUPERVISE_BOTTLE_SLUG:-UNSET}}\" > " + f"#!/bin/sh\nprintf '%s' \"${{SUPERVISE_SOURCE_IP:-UNSET}}\" > " f"{capture}\ncat >/dev/null\nexit 0\n" ) hook.chmod(0o755) @@ -166,7 +168,7 @@ class TestGitHttpBackend(unittest.TestCase): ["git", "push", url, "HEAD:refs/heads/main"], cwd=work, check=True, capture_output=True, text=True, timeout=5, ) - self.assertEqual(bottle_id, capture.read_text()) + self.assertEqual("127.0.0.1", capture.read_text()) def test_post_forwards_git_cgi_headers(self): from http.server import ThreadingHTTPServer diff --git a/tests/unit/test_orchestrator_control_plane.py b/tests/unit/test_orchestrator_control_plane.py index 9ebd8b6..fe445b2 100644 --- a/tests/unit/test_orchestrator_control_plane.py +++ b/tests/unit/test_orchestrator_control_plane.py @@ -426,6 +426,112 @@ class TestDispatchSupervise(unittest.TestCase): self.assertIn("no such proposal", str(payload["error"])) +class TestDispatchSuperviseAgentRpc(unittest.TestCase): + """The agent half — `/supervise/propose` + `/supervise/poll` — attributed by + (source_ip, identity_token) like /resolve, so a bottle can only ever queue + or read its own proposals (PRD 0070 / issue #469).""" + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + root = Path(self._tmp.name) + db = root / "db" / "bot-bottle.db" + db.parent.mkdir(parents=True) + self._env = patch.dict("os.environ", { + "BOT_BOTTLE_ROOT": str(root), + "SUPERVISE_DB_PATH": str(db), + }) + self._env.start() + self.store = RegistryStore(db) + self.store.migrate() + StoreManager(db).migrate() + secret = secrets.token_bytes(16) + self.orch = Orchestrator(self.store, StubBroker(secret), secret) + + def tearDown(self) -> None: + self._env.stop() + self._tmp.cleanup() + + def _register(self, source_ip: str = "10.243.0.9", slug: str = "demo"): + return self.store.register( + source_ip, metadata=json.dumps({"slug": slug}), policy="routes: []\n") + + def _propose(self, rec, proposed: str = "routes:\n - host: g.com\n"): + return dispatch(self.orch, "POST", "/supervise/propose", _body({ + "source_ip": rec.source_ip, "identity_token": rec.identity_token, + "tool": TOOL_EGRESS_ALLOW, "proposed_file": proposed, "justification": "need it", + })) + + def _poll(self, rec, proposal_id: str): + return dispatch(self.orch, "POST", "/supervise/poll", _body({ + "source_ip": rec.source_ip, "identity_token": rec.identity_token, + "proposal_id": proposal_id, + })) + + def test_propose_queues_under_the_resolved_bottle(self) -> None: + rec = self._register() + status, payload = self._propose(rec) + self.assertEqual(201, status) + pid = payload["proposal_id"] + assert isinstance(pid, str) and pid + _, listing = dispatch(self.orch, "GET", "/supervise/proposals", b"") + self.assertEqual(pid, listing["proposals"][0]["id"]) + # Queued under the orchestrator-resolved bottle id, never a caller slug. + self.assertEqual(rec.bottle_id, listing["proposals"][0]["bottle_slug"]) + + def test_propose_unattributed_is_403(self) -> None: + status, payload = dispatch(self.orch, "POST", "/supervise/propose", _body({ + "source_ip": "10.9.9.9", "identity_token": "wrong", + "tool": TOOL_EGRESS_ALLOW, "proposed_file": "x\n", "justification": "j", + })) + self.assertEqual(403, status) + self.assertIn("unattributed", str(payload["error"])) + + def test_propose_rejects_unknown_tool(self) -> None: + rec = self._register() + status, _ = dispatch(self.orch, "POST", "/supervise/propose", _body({ + "source_ip": rec.source_ip, "identity_token": rec.identity_token, + "tool": "not-a-tool", "proposed_file": "x\n", "justification": "j", + })) + self.assertEqual(400, status) + + def test_poll_pending_then_decided_then_archived(self) -> None: + rec = self._register() + _, proposed = self._propose(rec) + pid = proposed["proposal_id"] + assert isinstance(pid, str) + + status, poll = self._poll(rec, pid) + self.assertEqual(200, status) + self.assertEqual("pending", poll["status"]) + + # Operator decides server-side. + dispatch(self.orch, "POST", "/supervise/respond", _body({ + "proposal_id": pid, "bottle_slug": rec.bottle_id, + "decision": "approve", "notes": "ok", + })) + _, decided = self._poll(rec, pid) + self.assertEqual("approved", decided["status"]) + self.assertEqual("ok", decided["notes"]) + + # The decided poll archived it: gone from pending, and a re-poll is + # 'unknown' rather than replaying the decision forever. + _, listing = dispatch(self.orch, "GET", "/supervise/proposals", b"") + self.assertEqual([], listing["proposals"]) + _, again = self._poll(rec, pid) + self.assertEqual("unknown", again["status"]) + + def test_poll_cannot_read_another_bottles_proposal(self) -> None: + rec_a = self._register("10.0.0.1", "a") + rec_b = self._register("10.0.0.2", "b") + _, proposed = self._propose(rec_a) + pid = proposed["proposal_id"] + assert isinstance(pid, str) + # b polls a's proposal id: scoped to b's own queue → never a's response. + status, poll = self._poll(rec_b, pid) + self.assertEqual(200, status) + self.assertEqual("unknown", poll["status"]) + + if __name__ == "__main__": unittest.main() diff --git a/tests/unit/test_orchestrator_gateway.py b/tests/unit/test_orchestrator_gateway.py index a413e66..40a14c8 100644 --- a/tests/unit/test_orchestrator_gateway.py +++ b/tests/unit/test_orchestrator_gateway.py @@ -124,13 +124,11 @@ class TestDockerGateway(unittest.TestCase): src = ca_mounts[0].rsplit(":", 1)[0] self.assertTrue(src.endswith("/" + GATEWAY_CA_DIRNAME), src) self.assertTrue(Path(src).is_absolute(), src) - # Shares the ONE host DB: the supervise daemon queues into the same - # file the orchestrator + operator (over HTTP) use. - self.assertTrue(any( - a.startswith("SUPERVISE_DB_PATH=") and a.endswith("/run/supervise/bot-bottle.db") - for a in runs[0])) - self.assertTrue(any( - a.endswith(":/run/supervise") for a in runs[0])) + # No DB handle on the data plane: the supervise queue is reached over + # the control-plane RPC, so the gateway container carries neither the + # DB bind-mount nor SUPERVISE_DB_PATH (PRD 0070 / issue #469). + self.assertFalse(any(a.startswith("SUPERVISE_DB_PATH=") for a in runs[0])) + self.assertFalse(any(a.endswith(":/run/supervise") for a in runs[0])) # Data plane resolves policy against the orchestrator control plane. self.assertIn(f"BOT_BOTTLE_ORCHESTRATOR_URL={_ORCH_URL}", runs[0]) diff --git a/tests/unit/test_orchestrator_service.py b/tests/unit/test_orchestrator_service.py index 61ed48d..cd50bc6 100644 --- a/tests/unit/test_orchestrator_service.py +++ b/tests/unit/test_orchestrator_service.py @@ -323,6 +323,43 @@ class TestOrchestratorSupervise(unittest.TestCase): self.assertFalse(ok) self.assertIn("no longer registered", err) + # --- agent half: queue + poll (issue #469) ----------------------------- + + def test_queue_proposal_then_poll_pending(self) -> None: + bottle_id = self._register("demo", "routes: []\n") + pid = self.orch.supervise_queue_proposal( + bottle_id, tool=TOOL_EGRESS_ALLOW, + proposed_file="routes:\n - host: google.com\n", justification="need it") + # Visible to the operator, keyed by the bottle id. + pending = self.orch.supervise_pending() + self.assertEqual([pid], [p["id"] for p in pending]) + self.assertEqual( + {"status": "pending"}, self.orch.supervise_poll_response(bottle_id, pid)) + + def test_poll_returns_decision_and_archives(self) -> None: + bottle_id = self._register("demo", "routes: []\n") + pid = self.orch.supervise_queue_proposal( + bottle_id, tool=TOOL_EGRESS_ALLOW, + proposed_file="routes:\n - host: google.com\n", justification="need it") + self.orch.supervise_respond( + pid, bottle_slug=bottle_id, decision="approve", notes="ok") + decided = self.orch.supervise_poll_response(bottle_id, pid) + self.assertEqual("approved", decided["status"]) + self.assertEqual("ok", decided["notes"]) + # Archived on read → gone from pending, and a re-poll is 'unknown'. + self.assertEqual([], self.orch.supervise_pending()) + self.assertEqual( + {"status": "unknown"}, self.orch.supervise_poll_response(bottle_id, pid)) + + def test_poll_unknown_for_other_bottle(self) -> None: + bottle_id = self._register("demo", "routes: []\n") + pid = self.orch.supervise_queue_proposal( + bottle_id, tool=TOOL_EGRESS_ALLOW, + proposed_file="routes:\n - host: google.com\n", justification="j") + # A different bottle id can't read demo's proposal (scoped by queue key). + self.assertEqual( + {"status": "unknown"}, self.orch.supervise_poll_response("other-bottle", pid)) + if __name__ == "__main__": unittest.main() diff --git a/tests/unit/test_policy_resolver.py b/tests/unit/test_policy_resolver.py index c06b7a4..6ab5d1c 100644 --- a/tests/unit/test_policy_resolver.py +++ b/tests/unit/test_policy_resolver.py @@ -106,6 +106,58 @@ class TestPolicyResolver(unittest.TestCase): with self.assertRaises(PolicyResolveError): self.r.resolve_policy_and_bottle_id("10.243.0.1") + # --- supervise agent RPCs (issue #469) --------------------------------- + + def test_propose_supervise_returns_id_and_posts_payload(self) -> None: + with patch(_URLOPEN, return_value=_resp({"proposal_id": "p-7"})) as m: + pid = self.r.propose_supervise( + "10.243.0.7", "the-token", + tool="egress-allow", proposed_file="routes:\n", justification="j", + ) + self.assertEqual("p-7", pid) + req = m.call_args.args[0] + self.assertTrue(req.full_url.endswith("/supervise/propose")) + sent = json.loads(req.data) + self.assertEqual("10.243.0.7", sent["source_ip"]) + self.assertEqual("the-token", sent["identity_token"]) + self.assertEqual("egress-allow", sent["tool"]) + self.assertEqual("routes:\n", sent["proposed_file"]) + + def test_propose_supervise_unattributed_is_none(self) -> None: + with patch(_URLOPEN, side_effect=_http_error(403)): + self.assertIsNone(self.r.propose_supervise( + "10.9.9.9", "t", tool="egress-allow", proposed_file="x", justification="j")) + + def test_propose_supervise_missing_id_is_none(self) -> None: + with patch(_URLOPEN, return_value=_resp({})): + self.assertIsNone(self.r.propose_supervise( + "10.243.0.1", "t", tool="egress-allow", proposed_file="x", justification="j")) + + def test_propose_supervise_unreachable_raises(self) -> None: + with patch(_URLOPEN, side_effect=urllib.error.URLError("refused")): + with self.assertRaises(PolicyResolveError): + self.r.propose_supervise( + "10.243.0.1", "t", tool="egress-allow", proposed_file="x", justification="j") + + def test_poll_supervise_returns_status(self) -> None: + with patch(_URLOPEN, return_value=_resp( + {"status": "approved", "notes": "ok", "final_file": None}) + ) as m: + result = self.r.poll_supervise("10.243.0.7", "tok", "p-7") + self.assertEqual("approved", result["status"]) + req = m.call_args.args[0] + self.assertTrue(req.full_url.endswith("/supervise/poll")) + self.assertEqual("p-7", json.loads(req.data)["proposal_id"]) + + def test_poll_supervise_unattributed_is_none(self) -> None: + with patch(_URLOPEN, side_effect=_http_error(403)): + self.assertIsNone(self.r.poll_supervise("10.9.9.9", "t", "p-7")) + + def test_poll_supervise_unreachable_raises(self) -> None: + with patch(_URLOPEN, side_effect=urllib.error.URLError("refused")): + with self.assertRaises(PolicyResolveError): + self.r.poll_supervise("10.243.0.1", "t", "p-7") + if __name__ == "__main__": unittest.main() diff --git a/tests/unit/test_supervise_server.py b/tests/unit/test_supervise_server.py index 115b7e0..725943f 100644 --- a/tests/unit/test_supervise_server.py +++ b/tests/unit/test_supervise_server.py @@ -1,4 +1,11 @@ -"""Unit: supervise daemon MCP server (PRD 0013).""" +"""Unit: supervise daemon MCP server (PRD 0013, PRD 0070). + +The daemon no longer opens bot-bottle.db: it queues proposals and polls for +their responses over the control-plane RPC (issue #469). These tests drive the +handlers with `_FakeSuperviseResolver`, an in-process stand-in for +`PolicyResolver.propose_supervise` / `poll_supervise` backed by the real queue +store — so the operator-response and archive contracts are still exercised +end-to-end, just through the RPC seam instead of a direct file handle.""" import http.client import json @@ -44,6 +51,81 @@ from bot_bottle.supervise_server import ( validate_proposed_file, ) +# Fixed caller identity for the handler tests. The control plane attributes by +# (source_ip, identity_token); the fake resolver ignores them and answers for a +# fixed bottle, since attribution itself is covered by the orchestrator tests. +_SRC = "10.0.0.7" +_TOK = "tok" + + +class _FakeSuperviseResolver: + """Stand-in for `PolicyResolver`'s supervise RPCs, backed by the real queue + store (as the orchestrator is). `bottle_id=None` models an unattributed + caller (a clean 403 → None); `raises=True` models an unreachable + orchestrator (`PolicyResolveError`).""" + + def __init__( + self, bottle_id: str | None = "dev", raises: bool = False, policy: str = "", + ) -> None: + self.bottle_id = bottle_id + self.raises = raises + self._policy = policy + + def propose_supervise( + self, source_ip, identity_token, *, tool, proposed_file, justification, + ): + del source_ip, identity_token + if self.raises: + raise supervise_server.PolicyResolveError("orchestrator down") + if self.bottle_id is None: + return None + proposal = _sv.Proposal.new( + bottle_slug=self.bottle_id, tool=tool, proposed_file=proposed_file, + justification=justification, current_file_hash=_sv.sha256_hex(proposed_file), + ) + _sv.write_proposal(proposal) + return proposal.id + + def poll_supervise(self, source_ip, identity_token, proposal_id): + del source_ip, identity_token + if self.raises: + raise supervise_server.PolicyResolveError("orchestrator down") + if self.bottle_id is None: + return None + try: + response = _sv.read_response(self.bottle_id, proposal_id) + except FileNotFoundError: + try: + _sv.read_proposal(self.bottle_id, proposal_id) + except FileNotFoundError: + return {"status": _sv.POLL_STATUS_UNKNOWN} + return {"status": _sv.POLL_STATUS_PENDING} + _sv.archive_proposal(self.bottle_id, proposal_id) + return { + "status": response.status, "notes": response.notes, + "final_file": response.final_file, + } + + # Used by list-egress-routes (`_resolved_routes_payload`), unchanged path. + def resolve_policy_and_bottle_id(self, source_ip, identity_token=""): + del source_ip, identity_token + if self.raises: + raise supervise_server.PolicyResolveError("orchestrator down") + return self._policy, self.bottle_id, {} + + +def _tools_call(resolver, params, config=None): + return handle_tools_call( + params, config or ServerConfig(), + resolver=resolver, source_ip=_SRC, identity_token=_TOK, + ) + + +def _check(resolver, params): + return handle_check_proposal( + params, resolver=resolver, source_ip=_SRC, identity_token=_TOK, + ) + # --- Validation ------------------------------------------------------------ @@ -111,32 +193,35 @@ class TestRpcErrorTaxonomy(unittest.TestCase): validate_proposed_file(_sv.TOOL_EGRESS_ALLOW, "routes: nope\n") def test_unknown_tool_in_tools_call_is_client_error(self): - config = ServerConfig(bottle_slug="dev") with self.assertRaises(_RpcClientError) as cm: - handle_tools_call({"name": "no-such-tool", "arguments": {}}, config) + _tools_call(_FakeSuperviseResolver(), {"name": "no-such-tool", "arguments": {}}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) -class TestRpcInternalErrorOnIoFailure(unittest.TestCase): - def test_write_proposal_os_error_raises_internal(self): - config = ServerConfig( - bottle_slug="dev", - ) - with patch.object(_sv, "write_proposal", side_effect=OSError("disk full")), \ - self.assertRaises(_RpcInternalError) as cm: - handle_tools_call( - { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": { - "routes_yaml": "routes:\n - host: example.com\n", - "justification": "x", - }, - }, - config, - ) +class TestRpcInternalErrorOnRpcFailure(unittest.TestCase): + """A queue RPC that can't reach the orchestrator (or returns unattributed) + surfaces as ERR_INTERNAL — the daemon fails closed rather than leaking the + cause to the agent.""" + + _ARGS = { + "name": _sv.TOOL_EGRESS_ALLOW, + "arguments": { + "routes_yaml": "routes:\n - host: example.com\n", + "justification": "x", + }, + } + + def test_unreachable_orchestrator_raises_internal(self): + with self.assertRaises(_RpcInternalError) as cm: + _tools_call(_FakeSuperviseResolver(raises=True), self._ARGS) self.assertEqual(ERR_INTERNAL, cm.exception.code) self.assertIsNotNone(cm.exception.__cause__) + def test_unattributed_source_raises_internal(self): + with self.assertRaises(_RpcInternalError) as cm: + _tools_call(_FakeSuperviseResolver(bottle_id=None), self._ARGS) + self.assertEqual(ERR_INTERNAL, cm.exception.code) + # --- JSON-RPC parsing ------------------------------------------------------ @@ -265,8 +350,8 @@ class TestHandleToolsList(unittest.TestCase): class TestHandleToolsCall(unittest.TestCase): def setUp(self): self._tmp = tempfile.TemporaryDirectory(prefix="supervise-server-test.") - self._home_patch = self._patch_home(Path(self._tmp.name)) - self.config = ServerConfig(bottle_slug="dev") + self._home_patch = use_bottle_root(Path(self._tmp.name) / ".bot-bottle") + self.resolver = _FakeSuperviseResolver("dev") _qs.QueueStore("dev").migrate() _as.AuditStore().migrate() @@ -274,12 +359,9 @@ class TestHandleToolsCall(unittest.TestCase): self._home_patch() self._tmp.cleanup() - def _patch_home(self, fake_home: Path): - return use_bottle_root(fake_home / ".bot-bottle") - def _respond_when_proposal_appears(self, status: str, notes: str = "") -> threading.Thread: """Background thread: poll the queue for a fresh proposal, write a - matching response. Returns the thread so the test can join it.""" + matching response — the operator half, out of band.""" def runner(): for _ in range(200): pending = _sv.list_pending_proposals("dev") @@ -298,16 +380,13 @@ class TestHandleToolsCall(unittest.TestCase): def test_call_round_trips_through_queue(self): responder = self._respond_when_proposal_appears(_sv.STATUS_APPROVED, notes="lgtm") try: - result = handle_tools_call( - { - "name": _sv.TOOL_EGRESS_BLOCK, - "arguments": { - "routes_yaml": "routes:\n - host: example.com\n", - "justification": "need example.com", - }, + result = _tools_call(self.resolver, { + "name": _sv.TOOL_EGRESS_BLOCK, + "arguments": { + "routes_yaml": "routes:\n - host: example.com\n", + "justification": "need example.com", }, - self.config, - ) + }) finally: responder.join() self.assertFalse(result["isError"]) # type: ignore[index] @@ -318,16 +397,13 @@ class TestHandleToolsCall(unittest.TestCase): def test_allow_round_trips_through_queue(self): responder = self._respond_when_proposal_appears(_sv.STATUS_APPROVED, notes="ok") try: - result = handle_tools_call( - { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": { - "routes_yaml": "routes:\n - host: example.com\n", - "justification": "need example.com", - }, + result = _tools_call(self.resolver, { + "name": _sv.TOOL_EGRESS_ALLOW, + "arguments": { + "routes_yaml": "routes:\n - host: example.com\n", + "justification": "need example.com", }, - self.config, - ) + }) finally: responder.join() self.assertFalse(result["isError"]) # type: ignore[index] @@ -338,94 +414,70 @@ class TestHandleToolsCall(unittest.TestCase): def test_rejected_response_sets_isError(self): responder = self._respond_when_proposal_appears(_sv.STATUS_REJECTED, notes="nope") try: - result = handle_tools_call( - { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": { - "routes_yaml": "routes:\n - host: example.com\n", - "justification": "needed for tests", - }, + result = _tools_call(self.resolver, { + "name": _sv.TOOL_EGRESS_ALLOW, + "arguments": { + "routes_yaml": "routes:\n - host: example.com\n", + "justification": "needed for tests", }, - self.config, - ) + }) finally: responder.join() self.assertTrue(result["isError"]) # type: ignore[index] def test_invalid_tool_name_raises(self): with self.assertRaises(_RpcError) as cm: - handle_tools_call( - {"name": "not-a-tool", "arguments": {}}, - self.config, - ) + _tools_call(self.resolver, {"name": "not-a-tool", "arguments": {}}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) def test_missing_justification_raises(self): with self.assertRaises(_RpcError): - handle_tools_call( - { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": {"routes_yaml": "routes:\n - host: example.com\n"}, - }, - self.config, - ) + _tools_call(self.resolver, { + "name": _sv.TOOL_EGRESS_ALLOW, + "arguments": {"routes_yaml": "routes:\n - host: example.com\n"}, + }) def test_missing_name_raises(self): with self.assertRaises(_RpcError) as cm: - handle_tools_call({"arguments": {}}, self.config) + _tools_call(self.resolver, {"arguments": {}}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) def test_arguments_must_be_object(self): with self.assertRaises(_RpcError) as cm: - handle_tools_call( - { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": [], - }, - self.config, - ) + _tools_call(self.resolver, {"name": _sv.TOOL_EGRESS_ALLOW, "arguments": []}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) self.assertIn("must be an object", cm.exception.message) def test_capability_block_call_raises_unknown_tool(self): with self.assertRaises(_RpcError) as cm: - handle_tools_call( - { - "name": "capability-block", - "arguments": { - "dockerfile": "FROM python:3.13\n", - "justification": "need git", - }, + _tools_call(self.resolver, { + "name": "capability-block", + "arguments": { + "dockerfile": "FROM python:3.13\n", + "justification": "need git", }, - self.config, - ) + }) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) self.assertIn("unknown tool", cm.exception.message) def test_archives_proposal_after_response(self): responder = self._respond_when_proposal_appears(_sv.STATUS_APPROVED) try: - handle_tools_call( - { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": { - "routes_yaml": "routes:\n - host: example.com\n", - "justification": "x", - }, + _tools_call(self.resolver, { + "name": _sv.TOOL_EGRESS_ALLOW, + "arguments": { + "routes_yaml": "routes:\n - host: example.com\n", + "justification": "x", }, - self.config, - ) + }) finally: responder.join() - # No pending proposals left after archive. + # No pending proposals left after the decided poll archives it. self.assertEqual([], _sv.list_pending_proposals("dev")) def test_pending_response_times_out_without_archive(self): - config = ServerConfig( - bottle_slug="dev", - response_timeout_seconds=0.05, - ) - result = handle_tools_call( + result = _tools_call( + self.resolver, { "name": _sv.TOOL_EGRESS_ALLOW, "arguments": { @@ -433,9 +485,8 @@ class TestHandleToolsCall(unittest.TestCase): "justification": "need egress", }, }, - config, + ServerConfig(response_timeout_seconds=0.05), ) - self.assertFalse(result["isError"]) # type: ignore[index] text = result["content"][0]["text"] # type: ignore[index] self.assertIn("status: pending", text) @@ -510,7 +561,7 @@ class TestHttpEndToEnd(unittest.TestCase): self.port = s.getsockname()[1] s.close() self.server = MCPServer(("127.0.0.1", self.port), MCPHandler) - self.server.config = ServerConfig(bottle_slug="dev") + self.server.config = ServerConfig() self.thread = threading.Thread( target=self.server.serve_forever, daemon=True, ) @@ -552,23 +603,21 @@ class TestHttpEndToEnd(unittest.TestCase): ) self.assertEqual(ERR_METHOD_NOT_FOUND, result["error"]["code"]) # type: ignore[index] - def test_internal_error_returns_err_internal_over_http(self): - with patch.object( - supervise_server._sv, "write_proposal", - side_effect=OSError("disk full"), - ): - result = self._post_jsonrpc({ - "jsonrpc": "2.0", - "id": 99, - "method": "tools/call", - "params": { - "name": _sv.TOOL_EGRESS_ALLOW, - "arguments": { - "routes_yaml": "routes:\n - host: example.com\n", - "justification": "x", - }, + def test_no_resolver_fails_closed_over_http(self): + # The test server has no policy_resolver wired, so a proposal tools/call + # fails closed with ERR_INTERNAL rather than queuing anything. + result = self._post_jsonrpc({ + "jsonrpc": "2.0", + "id": 99, + "method": "tools/call", + "params": { + "name": _sv.TOOL_EGRESS_ALLOW, + "arguments": { + "routes_yaml": "routes:\n - host: example.com\n", + "justification": "x", }, - }) + }, + }) self.assertIn("error", result) self.assertEqual(ERR_INTERNAL, result["error"]["code"]) # type: ignore[index] @@ -583,72 +632,23 @@ class TestHttpEndToEnd(unittest.TestCase): conn.close() -class _FakeResolver: - def __init__( - self, - bottle_id: str | None = None, - raises: bool = False, - policy: str = "", - ) -> None: - self._bottle_id = bottle_id - self._raises = raises - self._policy = policy - self.calls: list[str] = [] - - def resolve_bottle_id(self, source_ip: str, identity_token: str = "") -> str | None: - del identity_token - self.calls.append(source_ip) - if self._raises: - # Raise the exact class supervise_server catches (it imports - # policy_resolver flat inside the bundle, package-side in tests). - raise supervise_server.PolicyResolveError("orchestrator down") - return self._bottle_id - - def resolve_policy_and_bottle_id( - self, source_ip: str, identity_token: str = "", - ) -> "tuple[str, str | None, dict[str, str]]": - del identity_token - self.calls.append(source_ip) - if self._raises: - raise supervise_server.PolicyResolveError("orchestrator down") - return self._policy, self._bottle_id, {} - - def _handler(resolver: object) -> MCPHandler: """A bare MCPHandler wired with a server (carrying the resolver) and a - client address, enough to exercise `_attributed_config` off-socket.""" + client address, enough to exercise the resolver-backed paths off-socket.""" h: MCPHandler = MCPHandler.__new__(MCPHandler) h.server = types.SimpleNamespace(policy_resolver=resolver) # type: ignore[assignment] h.client_address = ("10.0.0.7", 4321) + h.headers = {} # type: ignore[assignment] return h -class TestAttributedConfig(unittest.TestCase): - """Each proposal is attributed to the calling bottle by source IP (PRD - 0070); a server without a resolver fails closed rather than queuing under an - unattributed slug.""" +class TestResolverFailClosed(unittest.TestCase): + """A dispatch without a resolver is a misconfiguration, not a tenancy mode + — fail closed rather than queue (or list) anything (PRD 0070).""" - def test_missing_resolver_fails_closed(self) -> None: + def test_missing_resolver_raises(self) -> None: with self.assertRaises(_RpcInternalError): - _handler(None)._attributed_config(ServerConfig(bottle_slug="dev")) - - def test_consolidated_binds_source_ip_bottle(self) -> None: - r = _FakeResolver(bottle_id="bottle-x") - cfg = _handler(r)._attributed_config(ServerConfig(bottle_slug="ignored")) - self.assertEqual("bottle-x", cfg.bottle_slug) # resolved slug wins - self.assertEqual(["10.0.0.7"], r.calls) - - def test_unattributed_source_fails_closed(self) -> None: - with self.assertRaises(_RpcInternalError): - _handler(_FakeResolver(bottle_id=None))._attributed_config( - ServerConfig(bottle_slug="x") - ) - - def test_resolver_error_fails_closed(self) -> None: - with self.assertRaises(_RpcInternalError): - _handler(_FakeResolver(raises=True))._attributed_config( - ServerConfig(bottle_slug="x") - ) + _handler(None)._resolver_or_fail() class TestResolvedRoutesPayload(unittest.TestCase): @@ -664,7 +664,7 @@ class TestResolvedRoutesPayload(unittest.TestCase): " - host: www.google.com\n" ) payload = _handler( - _FakeResolver(bottle_id="b1", policy=policy) + _FakeSuperviseResolver(bottle_id="b1", policy=policy) )._resolved_routes_payload() assert payload is not None self.assertFalse(payload["isError"]) # type: ignore[index] @@ -676,7 +676,7 @@ class TestResolvedRoutesPayload(unittest.TestCase): # resolve_client_context swallows resolver errors → deny-all (empty), # never another bottle's routes. payload = _handler( - _FakeResolver(raises=True) + _FakeSuperviseResolver(raises=True) )._resolved_routes_payload() assert payload is not None data = json.loads(payload["content"][0]["text"]) # type: ignore[index] @@ -698,7 +698,7 @@ class TestNonBlockingSupervise(unittest.TestCase): def setUp(self): self._tmp = tempfile.TemporaryDirectory(prefix="supervise-nonblock-test.") self._home_patch = use_bottle_root(Path(self._tmp.name) / ".bot-bottle") - self.config = ServerConfig(bottle_slug="dev") + self.resolver = _FakeSuperviseResolver("dev") _qs.QueueStore("dev").migrate() _as.AuditStore().migrate() @@ -717,9 +717,6 @@ class TestNonBlockingSupervise(unittest.TestCase): _sv.write_proposal(p) return p - def _check(self, proposal_id: str) -> dict[str, object]: - return handle_check_proposal({"arguments": {"proposal_id": proposal_id}}, self.config) - # --- pending response carries the id --- def test_pending_text_includes_id_and_pointer(self): @@ -730,12 +727,13 @@ class TestNonBlockingSupervise(unittest.TestCase): def test_tools_call_timeout_returns_pending_with_id_and_stays_queued(self): # No responder → the grace window expires → pending, not blocked forever. - result = handle_tools_call( + result = _tools_call( + self.resolver, { "name": _sv.TOOL_EGRESS_ALLOW, "arguments": {"routes_yaml": self._ROUTES, "justification": "x"}, }, - ServerConfig(bottle_slug="dev", response_timeout_seconds=0.05), + ServerConfig(response_timeout_seconds=0.05), ) self.assertFalse(result["isError"]) # type: ignore[index] text = result["content"][0]["text"] # type: ignore[index] @@ -749,7 +747,7 @@ class TestNonBlockingSupervise(unittest.TestCase): def test_check_returns_approved_and_archives(self): p = self._seed_proposal() _sv.write_response("dev", _sv.Response(proposal_id=p.id, status=_sv.STATUS_APPROVED, notes="ok")) - result = self._check(p.id) + result = _check(self.resolver, {"arguments": {"proposal_id": p.id}}) self.assertFalse(result["isError"]) text = result["content"][0]["text"] # type: ignore[index] self.assertIn("status: approved", text) @@ -760,13 +758,13 @@ class TestNonBlockingSupervise(unittest.TestCase): def test_check_rejected_sets_isError(self): p = self._seed_proposal() _sv.write_response("dev", _sv.Response(proposal_id=p.id, status=_sv.STATUS_REJECTED, notes="no")) - result = self._check(p.id) + result = _check(self.resolver, {"arguments": {"proposal_id": p.id}}) self.assertTrue(result["isError"]) self.assertIn("status: rejected", result["content"][0]["text"]) # type: ignore[index] def test_check_pending_when_no_decision_yet(self): p = self._seed_proposal() - result = self._check(p.id) + result = _check(self.resolver, {"arguments": {"proposal_id": p.id}}) self.assertFalse(result["isError"]) text = result["content"][0]["text"] # type: ignore[index] self.assertIn("status: pending", text) @@ -774,40 +772,41 @@ class TestNonBlockingSupervise(unittest.TestCase): self.assertEqual(1, len(_sv.list_pending_proposals("dev"))) # not archived def test_check_unknown_id_is_error(self): - result = self._check("no-such-proposal") + result = _check(self.resolver, {"arguments": {"proposal_id": "no-such-proposal"}}) self.assertTrue(result["isError"]) self.assertIn("status: unknown", result["content"][0]["text"]) # type: ignore[index] def test_check_missing_id_raises(self): with self.assertRaises(_RpcClientError) as cm: - handle_check_proposal({"arguments": {}}, self.config) + _check(self.resolver, {"arguments": {}}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) def test_check_empty_id_raises(self): with self.assertRaises(_RpcClientError) as cm: - handle_check_proposal({"arguments": {"proposal_id": " "}}, self.config) + _check(self.resolver, {"arguments": {"proposal_id": " "}}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) def test_check_arguments_must_be_object(self): with self.assertRaises(_RpcClientError) as cm: - handle_check_proposal({"arguments": []}, self.config) + _check(self.resolver, {"arguments": []}) self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code) def test_full_nonblocking_round_trip(self): # 1. tools/call times out → pending with id - result = handle_tools_call( + result = _tools_call( + self.resolver, { "name": _sv.TOOL_EGRESS_ALLOW, "arguments": {"routes_yaml": self._ROUTES, "justification": "x"}, }, - ServerConfig(bottle_slug="dev", response_timeout_seconds=0.05), + ServerConfig(response_timeout_seconds=0.05), ) pid = _sv.list_pending_proposals("dev")[0].id self.assertIn(pid, result["content"][0]["text"]) # type: ignore[index] # 2. operator decides out-of-band _sv.write_response("dev", _sv.Response(proposal_id=pid, status=_sv.STATUS_APPROVED, notes="ok")) # 3. agent resumes by polling — no re-proposing - poll = self._check(pid) + poll = _check(self.resolver, {"arguments": {"proposal_id": pid}}) self.assertFalse(poll["isError"]) self.assertIn("status: approved", poll["content"][0]["text"]) # type: ignore[index] self.assertEqual([], _sv.list_pending_proposals("dev")) # resolved + archived