diff --git a/bot_bottle/backend/macos_container/consolidated_launch.py b/bot_bottle/backend/macos_container/consolidated_launch.py index 5b7efc1..fcb693e 100644 --- a/bot_bottle/backend/macos_container/consolidated_launch.py +++ b/bot_bottle/backend/macos_container/consolidated_launch.py @@ -36,9 +36,12 @@ from dataclasses import dataclass from ...egress import EgressPlan from ...git_gate import GitGatePlan -from ...orchestrator.client import OrchestratorClient +from ...log import info +from ...orchestrator.client import OrchestratorClient, OrchestratorClientError from ...orchestrator.registration import registration_inputs from ..docker.gateway_provision import deprovision_git_gate, provision_git_gate +from . import util as container_mod +from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active from .gateway import GATEWAY_NETWORK from .gateway_provision import AppleGatewayTransport from .infra import MacosInfraService, OrchestratorStartError @@ -89,6 +92,32 @@ def ensure_gateway( ) +def live_source_ips(network: str) -> list[str]: + """Every running agent container's address on `network`. + + The reconciliation input: the orchestrator lives inside the infra + container and cannot enumerate the host's containers, so the host has to + tell it which bottles are actually up. Containers that have not been + assigned an address yet contribute nothing — the reap's grace window, not + this list, is what protects an in-flight launch. + + Raises `EnumerationError` when the live set cannot be determined + authoritatively: either the container listing fails or any individual + inspect fails. Callers must skip reconciliation in that case to avoid + unregistering healthy bottles.""" + ips: list[str] = [] + for agent in enumerate_active(): + name = f"{CONTAINER_NAME_PREFIX}{agent.slug}" + ip = container_mod.inspect_container_network_ip(name, network) + if ip is None: + raise EnumerationError( + f"container inspect {name!r} failed; live set is not authoritative" + ) + if ip: + ips.append(ip) + return ips + + def register_agent( egress_plan: EgressPlan, git_gate_plan: GitGatePlan, @@ -103,6 +132,16 @@ def register_agent( container — it is the attribution key the gateway resolves policy by. Raises on failure; the caller tears down.""" client = OrchestratorClient(endpoint.orchestrator_url) + # Self-heal before registering: a launcher that died hard (SIGKILL, closed + # terminal, host sleep) never ran its teardown callback, leaving an active + # row with no container. vmnet recycles addresses, so such a row can + # collide with this bottle's — and `by_source_ip` fail-closes on ambiguity, + # which would resolve no policy at all and deny every host. Best-effort: a + # reconciliation failure must not block an otherwise-fine launch. + try: + client.reconcile(live_source_ips(endpoint.network)) + except (OrchestratorClientError, EnumerationError) as e: + info(f"registry reconciliation skipped: {e}") inputs = registration_inputs(egress_plan) reg = client.register_bottle( source_ip, image_ref=image_ref, policy=inputs.policy, @@ -136,6 +175,7 @@ __all__ = [ "GatewayEndpoint", "LaunchContext", "ensure_gateway", + "live_source_ips", "register_agent", "teardown_consolidated", "ConsolidatedLaunchError", diff --git a/bot_bottle/backend/macos_container/enumerate.py b/bot_bottle/backend/macos_container/enumerate.py index f8ed062..01ce207 100644 --- a/bot_bottle/backend/macos_container/enumerate.py +++ b/bot_bottle/backend/macos_container/enumerate.py @@ -19,6 +19,10 @@ CONTAINER_NAME_PREFIX = "bot-bottle-" _INFRA_NAMES = frozenset({INFRA_NAME}) +class EnumerationError(RuntimeError): + """container list failed; the resulting live set is not authoritative.""" + + def enumerate_active() -> list[ActiveAgent]: result = subprocess.run( ["container", "list", "--quiet"], @@ -27,7 +31,10 @@ def enumerate_active() -> list[ActiveAgent]: check=False, ) if result.returncode != 0: - return [] + raise EnumerationError( + f"container list failed: " + f"{(result.stderr or '').strip() or ''}" + ) out: list[ActiveAgent] = [] for name in sorted(line.strip() for line in result.stdout.splitlines()): if not name.startswith(CONTAINER_NAME_PREFIX) or name in _INFRA_NAMES: diff --git a/bot_bottle/backend/macos_container/util.py b/bot_bottle/backend/macos_container/util.py index 9c5f3aa..6ea6f2a 100644 --- a/bot_bottle/backend/macos_container/util.py +++ b/bot_bottle/backend/macos_container/util.py @@ -572,6 +572,41 @@ def try_container_ipv4_on_network(name: str, network: str) -> str: return "" +def inspect_container_network_ip(name: str, network: str) -> str | None: + """IP of `name` on `network`, distinguishing inspect failure from "not yet". + + Returns: + - the IP string when the container has one on `network` + - "" when inspect succeeds but no address is assigned yet (in-flight DHCP) + - None when the inspect command itself fails (authoritative list impossible) + """ + result = subprocess.run( + [_CONTAINER, "inspect", name], + capture_output=True, text=True, check=False, + ) + if result.returncode != 0: + return None + try: + data = json.loads(result.stdout or "[]") + except json.JSONDecodeError: + return None + if isinstance(data, list): + data = data[0] if data else {} + if not isinstance(data, dict): + return None + status = data.get("status") + networks = status.get("networks") if isinstance(status, dict) else None + if not isinstance(networks, list): + return "" + for entry in networks: + if not isinstance(entry, dict) or entry.get("network") != network: + continue + raw = entry.get("ipv4Address") + if isinstance(raw, str) and raw: + return raw.split("/", 1)[0] + return "" + + def wait_container_ipv4_on_network( name: str, network: str, *, timeout: float = 15.0, poll: float = 0.25, ) -> str: diff --git a/bot_bottle/egress_addon.py b/bot_bottle/egress_addon.py index b550dc6..30e370d 100644 --- a/bot_bottle/egress_addon.py +++ b/bot_bottle/egress_addon.py @@ -379,6 +379,7 @@ class EgressAddon: env, request_method=flow.request.method, request_headers=req_headers, + deny_reason=config.deny_reason, ) if decision.action == "block": diff --git a/bot_bottle/egress_addon_core.py b/bot_bottle/egress_addon_core.py index d3e2d3a..8b8280e 100644 --- a/bot_bottle/egress_addon_core.py +++ b/bot_bottle/egress_addon_core.py @@ -89,6 +89,14 @@ LOG_FULL = 2 # log block/warn events + full request and response bodies class Config: routes: tuple[Route, ...] log: int = LOG_OFF + # Why this Config is a deny-all, when it is one for a reason *other* than + # the bottle's own policy genuinely not listing the host. A deny-all is + # indistinguishable from "policy loaded, host not allowed" at the decision + # point — both are simply "no matching route" — so without this the + # operator sees `host X is not in the allowlist` and goes hunting for a + # missing route that was never the problem. Empty for a normally-parsed + # policy; `decide` prefers it over the allowlist wording when set. + deny_reason: str = "" @dataclass(frozen=True) @@ -405,16 +413,40 @@ class PolicyResolverLike(typing.Protocol): ... +# Deny-all explanations. Each names the *actual* failure so an operator isn't +# sent looking for a missing egress route when the bottle never had a policy +# to begin with — the failure mode that made a bricked registration read like +# a misconfigured allowlist. +DENY_UNATTRIBUTED = ( + "egress: this request was not attributed to any bottle, so no egress " + "policy applies and every host is denied. Either the bottle's registry " + "row is missing/ambiguous (torn down, or another bottle claimed its " + "source IP), or the request carried no matching identity token — check " + "that the caller's proxy URL includes it. This is not an allowlist problem." +) +DENY_UNPARSEABLE = ( + "egress: this bottle's egress policy could not be parsed, so it is being " + "treated as deny-all. Fix the bottle's egress.routes; every host is denied " + "until it loads." +) +DENY_RESOLVER_ERROR = ( + "egress: the orchestrator could not be reached to resolve this bottle's " + "egress policy, so every host is denied (fail-closed). Check that the " + "control plane is up; this is not an allowlist problem." +) + + def _config_from_policy(policy: "str | None") -> "Config": """Parse a resolved policy blob into a Config, fail-closed: None / empty / unparseable all become a deny-all Config (no routes → every request - blocked).""" + blocked). Each deny-all carries the reason it is one, so the block message + names the real fault instead of blaming the allowlist.""" if not policy: - return Config(routes=()) # unattributed or empty → deny-all + return Config(routes=(), deny_reason=DENY_UNATTRIBUTED) try: return load_config(policy) except ValueError: - return Config(routes=()) # unparseable policy → deny + return Config(routes=(), deny_reason=DENY_UNPARSEABLE) def resolve_client_config( @@ -428,7 +460,7 @@ def resolve_client_config( try: policy = resolver.resolve(client_ip, identity_token) except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught - return Config(routes=()) # orchestrator unreachable/errored → deny + return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR) return _config_from_policy(policy) @@ -457,7 +489,7 @@ def resolve_client_context( client_ip, identity_token, ) except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught - return Config(routes=()), "", {} # orchestrator unreachable/errored → deny + return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR), "", {} return _config_from_policy(policy), (bottle_id or ""), tokens @@ -572,12 +604,16 @@ def decide( *, request_method: str = "GET", request_headers: typing.Mapping[str, str] | None = None, + deny_reason: str = "", ) -> Decision: + """`deny_reason` is `Config.deny_reason`: when the deny-all came from a + missing/unparseable policy rather than the bottle's own allowlist, report + that instead of implying a route is merely absent.""" route = match_route(routes, request_host) if route is None: return Decision( action="block", - reason=( + reason=deny_reason or ( f"egress: host {request_host!r} is not in the " f"bottle's egress.routes allowlist. Declare a " f"route for it or remove the request." @@ -852,6 +888,9 @@ __all__ = [ "is_git_push_request", "is_git_fetch_request", "load_config", + "DENY_UNATTRIBUTED", + "DENY_UNPARSEABLE", + "DENY_RESOLVER_ERROR", "resolve_client_config", "resolve_client_context", "PolicyResolverLike", diff --git a/bot_bottle/orchestrator/client.py b/bot_bottle/orchestrator/client.py index b235313..86c919c 100644 --- a/bot_bottle/orchestrator/client.py +++ b/bot_bottle/orchestrator/client.py @@ -15,6 +15,7 @@ from __future__ import annotations import json import urllib.error import urllib.request +from collections.abc import Iterable from dataclasses import dataclass from ..paths import host_control_plane_token @@ -147,6 +148,20 @@ class OrchestratorClient: raise OrchestratorClientError(f"teardown {bottle_id}: HTTP {status}") return True + def reconcile( + self, live_source_ips: Iterable[str], *, grace_seconds: float | None = None, + ) -> list[str]: + """Drop registry rows for bottles that are no longer running + (`POST /reconcile`), returning the reaped bottle ids. `live_source_ips` + is the caller's enumeration of its live bottles — the orchestrator + can't see the backend from inside the infra container.""" + body: dict[str, object] = {"live_source_ips": list(live_source_ips)} + if grace_seconds is not None: + body["grace_seconds"] = grace_seconds + payload = self._ok("POST", "/reconcile", body) + reaped = payload.get("reaped") + return [r for r in reaped if isinstance(r, str)] if isinstance(reaped, list) else [] + def set_policy(self, bottle_id: str, policy: str) -> bool: """Live-reload a bottle's policy (`PUT /bottles//policy`). False on 404 (unknown bottle).""" diff --git a/bot_bottle/orchestrator/control_plane.py b/bot_bottle/orchestrator/control_plane.py index 389c241..06ad8c7 100644 --- a/bot_bottle/orchestrator/control_plane.py +++ b/bot_bottle/orchestrator/control_plane.py @@ -13,6 +13,9 @@ vsock / unix-socket portability caveats): PUT /bottles//policy -> 200 {"updated": true} | 404 (live reload) body: {"policy"} DELETE /bottles/ -> 200 {"torn_down": true} | 404 (teardown) + POST /reconcile -> 200 {"reaped": [bottle_id, ...]} + body: {"live_source_ips": [...], + ["grace_seconds"]} POST /attribute -> 200 {"bottle_id"} | 403 POST /resolve -> 200 {"bottle_id","policy"} | 403 body: {"source_ip","identity_token"} @@ -141,6 +144,27 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches return 200, {"torn_down": True} return 404, {"error": "no such bottle"} + if method == "POST" and route == "/reconcile": + # Host-driven self-heal: the caller enumerates its live bottles (only + # the host can see the backend) and the orchestrator drops rows for + # every other active bottle. Trusted-caller only — an agent that could + # reach this would be able to unregister its neighbours. + try: + data = _parse_json_object(body) + except ValueError as e: + return 400, {"error": f"invalid JSON: {e}"} + raw_ips = data.get("live_source_ips") + if not isinstance(raw_ips, list): + return 400, {"error": "live_source_ips (list of strings) is required"} + live = [ip for ip in raw_ips if isinstance(ip, str) and ip] + grace = data.get("grace_seconds") + kwargs = ( + {"grace_seconds": float(grace)} + if isinstance(grace, (int, float)) and not isinstance(grace, bool) + else {} + ) + return 200, {"reaped": orch.reconcile(live, **kwargs)} + if method == "POST" and route == "/attribute": try: data = _parse_json_object(body) diff --git a/bot_bottle/orchestrator/registry.py b/bot_bottle/orchestrator/registry.py index 1b9b38b..06c6d74 100644 --- a/bot_bottle/orchestrator/registry.py +++ b/bot_bottle/orchestrator/registry.py @@ -32,6 +32,7 @@ import hmac import secrets import sqlite3 import time +from collections.abc import Iterable from dataclasses import dataclass from pathlib import Path @@ -42,6 +43,12 @@ from ..paths import host_db_path # 256 bits of urandom, URL-safe — unguessable per-bottle identity token. IDENTITY_TOKEN_BYTES = 32 +# How recently a row must have been registered to be exempt from +# `reap_absent`. Covers the window between `container run` and the address +# becoming visible to another launch's enumeration, so reconciliation never +# reaps a bottle that is still coming up. +DEFAULT_REAP_GRACE_SECONDS = 120.0 + def new_identity_token() -> str: """A fresh per-bottle identity token (PRD 0070 attribution defence).""" @@ -225,6 +232,70 @@ class RegistryStore(DbStore): ).fetchall() return [_row_to_record(r) for r in rows] + def reap_absent( + self, + live_source_ips: Iterable[str], + *, + grace_seconds: float = DEFAULT_REAP_GRACE_SECONDS, + now: float | None = None, + ) -> list[BottleRecord]: + """Delete active rows whose source IP is not held by a live bottle. + + A row only ever leaves the registry two ways: an explicit + `teardown_bottle` (the launcher's cleanup callback) or the supersede + sweep in `register`. Neither runs when the launching CLI dies hard — + SIGKILL, a closed terminal, a host sleep/crash — so the row outlives + its container. That orphan is not inert: source IPs are recycled by + the backend's DHCP, and `by_source_ip` fail-closes on ambiguity, so a + leftover row at a reused address can brick the *next* bottle that + lands on it (no policy resolved -> every host denied, reported to the + agent as "not in the allowlist"). Reconciling against the live set at + launch keeps the registry from accumulating those landmines. + + Restores the invariant the data plane needs: **at most one active row + per live address, and none at all for a dead one.** Two cases, because + a dead bottle's address may already have been handed to a live one: + + * no live bottle holds the address — every row there is an orphan; + * a live bottle holds it but several rows claim it — the newest + registration is authoritative and the rest are orphans, the same + rule `register`'s same-IP supersede sweep applies. Without this + second case a recycled address stays ambiguous, which is exactly + the state that resolves no policy. + + `grace_seconds` protects an in-flight launch: registration happens + moments after `container run`, and a concurrent launch's address may + not be visible to the caller's enumeration yet. Rows younger than the + grace window are never reaped, so reconciliation can't race a bottle + that is still coming up. Returns the deleted records.""" + live = {ip for ip in live_source_ips if ip} + cutoff = (time.time() if now is None else now) - grace_seconds + with self._connection() as conn: + rows = conn.execute( + "SELECT * FROM orchestrator_bottles WHERE state = 'active'", + ).fetchall() + by_ip: dict[str, list[BottleRecord]] = {} + for row in rows: + rec = _row_to_record(row) + by_ip.setdefault(rec.source_ip, []).append(rec) + candidates: list[BottleRecord] = [] + for ip, recs in by_ip.items(): + if ip not in live: + candidates.extend(recs) + continue + # Keep the newest claim on a live address; supersede the rest. + recs.sort(key=lambda r: r.created_at) + candidates.extend(recs[:-1]) + doomed = [r for r in candidates if r.created_at <= cutoff] + for rec in doomed: + conn.execute( + "DELETE FROM orchestrator_bottles WHERE bottle_id = ?", + (rec.bottle_id,), + ) + if doomed: + self._chmod() + return doomed + def by_source_ip(self, source_ip: str) -> BottleRecord | None: """Network-layer attribution: the single active bottle at this source IP, or None if unknown or ambiguous (more than one — a @@ -262,4 +333,5 @@ __all__ = [ "new_identity_token", "default_db_path", "IDENTITY_TOKEN_BYTES", + "DEFAULT_REAP_GRACE_SECONDS", ] diff --git a/bot_bottle/orchestrator/service.py b/bot_bottle/orchestrator/service.py index 7701b1c..b449782 100644 --- a/bot_bottle/orchestrator/service.py +++ b/bot_bottle/orchestrator/service.py @@ -13,15 +13,20 @@ Launch lifecycle: and returns the record. If the broker rejects/fails, the registry entry is rolled back so a failed launch leaves no orphan. * `teardown_bottle` sends a signed teardown request, then deregisters. + * `reconcile` sweeps rows whose bottle is no longer running — the + self-heal for the teardown paths that never got to run (a hard-killed + launcher), since an orphan row at a recycled source IP bricks the next + bottle that lands on it. """ from __future__ import annotations import json +from collections.abc import Iterable from datetime import datetime, timezone from .broker import LaunchBroker, LaunchRequest, sign_request -from .registry import BottleRecord, RegistryStore +from .registry import DEFAULT_REAP_GRACE_SECONDS, BottleRecord, RegistryStore from .gateway import Gateway from ..supervise import ( AuditEntry, @@ -117,6 +122,30 @@ class Orchestrator: self._tokens.pop(bottle_id, None) return True + def reconcile( + self, + live_source_ips: Iterable[str], + *, + grace_seconds: float = DEFAULT_REAP_GRACE_SECONDS, + ) -> list[str]: + """Drop registry rows for bottles that are no longer running, and + forget their in-memory egress tokens. Returns the reaped bottle ids. + + The caller supplies the live set because only the host can enumerate + its own containers — the orchestrator runs *inside* the infra + container and has no view of the backend. Deliberately does not + broker a teardown: the container is already gone, so there is nothing + to stop, and a broker error must not stop the sweep from clearing + the row that would otherwise brick the next bottle at that address. + + See `RegistryStore.reap_absent` for why orphans accumulate and why + they are harmful rather than merely untidy.""" + reaped = self.registry.reap_absent( + live_source_ips, grace_seconds=grace_seconds) + for rec in reaped: + self._tokens.pop(rec.bottle_id, None) + return [rec.bottle_id for rec in reaped] + def tokens_for(self, bottle_id: str) -> dict[str, str]: """The bottle's in-memory egress auth tokens (env_name -> value), or empty. The gateway injects these per request; they are never diff --git a/tests/unit/test_egress_multitenant.py b/tests/unit/test_egress_multitenant.py index e35ec95..6d99e0d 100644 --- a/tests/unit/test_egress_multitenant.py +++ b/tests/unit/test_egress_multitenant.py @@ -4,7 +4,14 @@ from __future__ import annotations import unittest -from bot_bottle.egress_addon_core import resolve_client_config, resolve_client_context +from bot_bottle.egress_addon_core import ( + DENY_RESOLVER_ERROR, + DENY_UNATTRIBUTED, + DENY_UNPARSEABLE, + decide, + resolve_client_config, + resolve_client_context, +) from bot_bottle.policy_resolver import PolicyResolveError @@ -108,3 +115,55 @@ class TestResolveClientContext(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class TestDenyReasonNamesTheRealFault(unittest.TestCase): + """A deny-all must not masquerade as a missing allowlist entry. + + Regression: an unregistered bottle resolves no policy, so *every* host is + denied — but the block message said `host X is not in the allowlist`, + which reads as a config problem and sends the operator hunting for a route + that was never missing. The structural reason wins over that wording. + """ + + def _reason(self, resolver: object, host: str = "chatgpt.com") -> str: + cfg = resolve_client_config(resolver, "10.243.0.1") # type: ignore[arg-type] + return decide(cfg.routes, host, "/v1/x", {}, deny_reason=cfg.deny_reason).reason + + def test_unattributed_says_unattributed_not_allowlist(self) -> None: + reason = self._reason(_FakeResolver(result=None)) + self.assertEqual(DENY_UNATTRIBUTED, reason) + # The misleading claim is the one that must be gone: the host was + # never "not in the allowlist" — there was no allowlist at all. + self.assertNotIn("is not in the bottle's egress.routes allowlist", reason) + # Both causes must be named. `/resolve` fail-closes on a missing row + # *and* on a token mismatch, and the message pointing only at the row + # sent us hunting for a deregistered bottle that was registered fine. + self.assertIn("registry row", reason) + self.assertIn("identity token", reason) + + def test_resolver_error_says_orchestrator_unreachable(self) -> None: + self.assertEqual(DENY_RESOLVER_ERROR, self._reason(_FakeResolver(raises=True))) + + def test_unparseable_policy_says_so(self) -> None: + self.assertEqual( + DENY_UNPARSEABLE, self._reason(_FakeResolver(result="routes: notalist\n"))) + + def test_a_real_allowlist_miss_keeps_the_allowlist_wording(self) -> None: + """The message only changes for structural deny-alls — a loaded policy + that genuinely lacks the host still points at the allowlist.""" + reason = self._reason(_FakeResolver(result='routes:\n - host: "api.example.com"\n')) + self.assertIn("is not in the bottle's egress.routes allowlist", reason) + self.assertIn("chatgpt.com", reason) + + def test_allowed_host_is_still_forwarded(self) -> None: + cfg = resolve_client_config( + _FakeResolver(result='routes:\n - host: "api.example.com"\n'), "10.243.0.1") + decision = decide( + cfg.routes, "api.example.com", "/v1/x", {}, deny_reason=cfg.deny_reason) + self.assertEqual("forward", decision.action) + + def test_a_parsed_policy_carries_no_deny_reason(self) -> None: + cfg = resolve_client_config( + _FakeResolver(result='routes:\n - host: "api.example.com"\n'), "10.243.0.1") + self.assertEqual("", cfg.deny_reason) diff --git a/tests/unit/test_macos_consolidated_launch.py b/tests/unit/test_macos_consolidated_launch.py index f3e395d..5e1e7a6 100644 --- a/tests/unit/test_macos_consolidated_launch.py +++ b/tests/unit/test_macos_consolidated_launch.py @@ -87,7 +87,8 @@ class TestRegisterAgent(unittest.TestCase): *, source_ip: str = "192.168.128.9", ): with patch(f"{_MOD}.OrchestratorClient", return_value=client), \ - patch(f"{_MOD}.provision_git_gate", provision or Mock()): + patch(f"{_MOD}.provision_git_gate", provision or Mock()), \ + patch(f"{_MOD}.live_source_ips", return_value=[]): return register_agent( _egress_plan(), _git_plan(), source_ip=source_ip, endpoint=_endpoint(), image_ref="img:1", @@ -134,3 +135,104 @@ class TestTeardown(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class TestLiveSourceIps(unittest.TestCase): + """The reconciliation input: the host enumerates its own bottles because + the orchestrator, inside the infra container, cannot see the backend.""" + + def _agents(self, *slugs: str) -> list[Mock]: + return [Mock(slug=s) for s in slugs] + + def test_maps_slugs_to_container_addresses(self) -> None: + from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips + with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \ + patch(f"{_MOD}.container_mod.inspect_container_network_ip", + side_effect=["10.0.0.1", "10.0.0.2"]) as ip: + got = live_source_ips("net0") + self.assertEqual(["10.0.0.1", "10.0.0.2"], got) + self.assertEqual("bot-bottle-a", ip.call_args_list[0].args[0]) + + def test_containers_without_an_address_are_skipped(self) -> None: + """A container that hasn't been given a DHCP address yet contributes + nothing — the reap's grace window, not this list, protects it.""" + from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips + with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \ + patch(f"{_MOD}.container_mod.inspect_container_network_ip", + side_effect=["", "10.0.0.2"]): + self.assertEqual(["10.0.0.2"], live_source_ips("net0")) + + def test_container_list_failure_raises(self) -> None: + """If container list fails, the live set is not authoritative and + reconciliation must be skipped.""" + from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips + from bot_bottle.backend.macos_container.enumerate import EnumerationError + with patch(f"{_MOD}.enumerate_active", + side_effect=EnumerationError("container list failed")): + with self.assertRaises(EnumerationError): + live_source_ips("net0") + + def test_per_container_inspect_failure_raises(self) -> None: + """If any individual inspect fails, the live set is not authoritative.""" + from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips + from bot_bottle.backend.macos_container.enumerate import EnumerationError + with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \ + patch(f"{_MOD}.container_mod.inspect_container_network_ip", + side_effect=["10.0.0.1", None]): + with self.assertRaises(EnumerationError): + live_source_ips("net0") + + +class TestRegisterAgentReconciles(unittest.TestCase): + """Registration self-heals the registry first: an orphan row at a recycled + address makes attribution ambiguous, which resolves no policy at all and + denies every host for the bottle being launched.""" + + def _register(self, client: Mock) -> None: + with patch(f"{_MOD}.OrchestratorClient", return_value=client), \ + patch(f"{_MOD}.provision_git_gate"), \ + patch(f"{_MOD}.live_source_ips", return_value=["10.0.0.7"]): + register_agent( + _egress_plan(), _git_plan(), + source_ip="10.0.0.7", endpoint=_endpoint(), + ) + + def test_reconciles_before_registering(self) -> None: + client = _client() + calls: list[str] = [] + + def _reconcile(*_args: object, **_kwargs: object) -> list[str]: + calls.append("reconcile") + return [] + + def _register_bottle(*_args: object, **_kwargs: object) -> RegisteredBottle: + calls.append("register") + return RegisteredBottle("b1", "tok") + + client.reconcile.side_effect = _reconcile + client.register_bottle.side_effect = _register_bottle + self._register(client) + self.assertEqual(["reconcile", "register"], calls) + client.reconcile.assert_called_once_with(["10.0.0.7"]) + + def test_a_reconcile_failure_does_not_block_the_launch(self) -> None: + from bot_bottle.orchestrator.client import OrchestratorClientError + client = _client() + client.reconcile.side_effect = OrchestratorClientError("unreachable") + self._register(client) + client.register_bottle.assert_called_once() + + def test_enumeration_error_does_not_block_the_launch(self) -> None: + """A partial container listing must not abort the launch — skip + reconciliation and proceed, just as with an unreachable orchestrator.""" + from bot_bottle.backend.macos_container.enumerate import EnumerationError + client = _client() + with patch(f"{_MOD}.OrchestratorClient", return_value=client), \ + patch(f"{_MOD}.provision_git_gate"), \ + patch(f"{_MOD}.live_source_ips", + side_effect=EnumerationError("container list failed")): + register_agent( + _egress_plan(), _git_plan(), + source_ip="10.0.0.7", endpoint=_endpoint(), + ) + client.register_bottle.assert_called_once() diff --git a/tests/unit/test_macos_container_cleanup.py b/tests/unit/test_macos_container_cleanup.py index 990536e..48b6871 100644 --- a/tests/unit/test_macos_container_cleanup.py +++ b/tests/unit/test_macos_container_cleanup.py @@ -66,8 +66,10 @@ class TestMacosContainerEnumerate(unittest.TestCase): agents = self._enumerate("bot-bottle-mac-infra\nbot-bottle-dev-abc\n") self.assertEqual(["dev-abc"], [a.slug for a in agents]) - def test_empty_when_the_cli_fails(self): - self.assertEqual([], self._enumerate("", returncode=1)) + def test_raises_when_the_cli_fails(self): + from bot_bottle.backend.macos_container.enumerate import EnumerationError + with self.assertRaises(EnumerationError): + self._enumerate("", returncode=1) if __name__ == "__main__": diff --git a/tests/unit/test_macos_container_util.py b/tests/unit/test_macos_container_util.py index b89472f..837d4b4 100644 --- a/tests/unit/test_macos_container_util.py +++ b/tests/unit/test_macos_container_util.py @@ -334,6 +334,50 @@ class TestInspectDigests(unittest.TestCase): self.assertEqual({}, util.container_env("x")) +class TestInspectContainerNetworkIp(unittest.TestCase): + """inspect_container_network_ip must distinguish inspect failure (None) + from 'no DHCP address yet' (""), which is the invariant live_source_ips + relies on to skip reconciliation on partial snapshots.""" + + _NETWORK = "bot-bottle-mac-gateway" + + def _inspect(self, stdout: str, returncode: int = 0) -> str | None: + cp = util.subprocess.CompletedProcess( + args=[], returncode=returncode, stdout=stdout, stderr="", + ) + with patch.object(util.subprocess, "run", return_value=cp): + return util.inspect_container_network_ip("bot-bottle-abc", self._NETWORK) + + def _entry(self, ip: str = "192.168.128.5") -> str: + return ( + f'[{{"status":{{"networks":[' + f'{{"network":"{self._NETWORK}","ipv4Address":"{ip}"}}' + f']}}}}]' + ) + + def test_returns_ip_when_inspect_succeeds(self) -> None: + self.assertEqual("192.168.128.5", self._inspect(self._entry())) + + def test_strips_cidr_prefix(self) -> None: + self.assertEqual("192.168.128.5", self._inspect(self._entry("192.168.128.5/24"))) + + def test_returns_empty_string_when_no_address_assigned_yet(self) -> None: + no_ip = f'[{{"status":{{"networks":[{{"network":"{self._NETWORK}","ipv4Address":""}}]}}}}]' + self.assertEqual("", self._inspect(no_ip)) + + def test_returns_empty_string_when_network_list_absent(self) -> None: + self.assertEqual("", self._inspect('[{"status":{}}]')) + + def test_returns_none_on_nonzero_exit(self) -> None: + self.assertIsNone(self._inspect("", returncode=1)) + + def test_returns_none_on_malformed_json(self) -> None: + self.assertIsNone(self._inspect("not-json")) + + def test_returns_none_on_unexpected_json_shape(self) -> None: + self.assertIsNone(self._inspect("null")) + + class TestWaitContainerIpv4(unittest.TestCase): def test_returns_address_once_dhcp_assigns_it(self): with patch.object(util, "try_container_ipv4_on_network", side_effect=["", "", "192.168.128.4"]), \ diff --git a/tests/unit/test_orchestrator_client.py b/tests/unit/test_orchestrator_client.py index 848ee6f..d6e07bb 100644 --- a/tests/unit/test_orchestrator_client.py +++ b/tests/unit/test_orchestrator_client.py @@ -104,3 +104,32 @@ class TestHealthAndPolicy(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class TestReconcile(unittest.TestCase): + def setUp(self) -> None: + self.c = OrchestratorClient("http://orch:8080") + + def test_posts_live_ips_and_returns_reaped(self) -> None: + with patch(_URLOPEN, return_value=_resp(200, {"reaped": ["b1", "b2"]})) as m: + got = self.c.reconcile(["10.0.0.2", "10.0.0.3"]) + self.assertEqual(["b1", "b2"], got) + sent = json.loads(m.call_args.args[0].data) + self.assertEqual(["10.0.0.2", "10.0.0.3"], sent["live_source_ips"]) + self.assertNotIn("grace_seconds", sent) # omitted -> server default + + def test_grace_seconds_is_forwarded_when_given(self) -> None: + with patch(_URLOPEN, return_value=_resp(200, {"reaped": []})) as m: + self.c.reconcile([], grace_seconds=30) + self.assertEqual(30, json.loads(m.call_args.args[0].data)["grace_seconds"]) + + def test_malformed_reaped_is_tolerated(self) -> None: + with patch(_URLOPEN, return_value=_resp(200, {"reaped": ["ok", 5, None]})): + self.assertEqual(["ok"], self.c.reconcile([])) + with patch(_URLOPEN, return_value=_resp(200, {})): + self.assertEqual([], self.c.reconcile([])) + + def test_error_status_raises(self) -> None: + with patch(_URLOPEN, side_effect=_http_error(500)): + with self.assertRaises(OrchestratorClientError): + self.c.reconcile([]) diff --git a/tests/unit/test_orchestrator_control_plane.py b/tests/unit/test_orchestrator_control_plane.py index b062399..ea75d5c 100644 --- a/tests/unit/test_orchestrator_control_plane.py +++ b/tests/unit/test_orchestrator_control_plane.py @@ -8,11 +8,13 @@ from __future__ import annotations import json import secrets +import sqlite3 import tempfile import threading import unittest import urllib.error import urllib.request +from contextlib import closing from pathlib import Path from unittest.mock import patch @@ -264,6 +266,7 @@ class TestControlPlaneAuth(unittest.TestCase): ("POST", "/bottles", _body({"source_ip": "10.0.0.1"})), ("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})), ("DELETE", "/bottles/x", b""), + ("POST", "/reconcile", _body({"live_source_ips": []})), ("POST", "/resolve", _body({"source_ip": "10.0.0.1", "identity_token": "t"})), ("POST", "/attribute", _body({"source_ip": "10.0.0.1", "identity_token": "t"})), ("GET", "/supervise/proposals", b""), @@ -387,3 +390,56 @@ class TestDispatchSupervise(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class TestReconcileRoute(unittest.TestCase): + """`POST /reconcile` — the host tells the orchestrator which bottles are + actually up, since the orchestrator can't see the backend from inside the + infra container.""" + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.orch = _orchestrator(Path(self._tmp.name) / "r.db") + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _old(self, source_ip: str) -> str: + rec = self.orch.registry.register(source_ip) + with closing(sqlite3.connect(self.orch.registry.db_path)) as conn: + conn.execute( + "UPDATE orchestrator_bottles SET created_at = 0.0 WHERE bottle_id = ?", + (rec.bottle_id,)) + conn.commit() + return rec.bottle_id + + def test_reaps_absent_and_reports_ids(self) -> None: + dead = self._old("10.0.0.1") + alive = self._old("10.0.0.2") + status, payload = dispatch( + self.orch, "POST", "/reconcile", _body({"live_source_ips": ["10.0.0.2"]})) + self.assertEqual(200, status) + self.assertEqual([dead], payload["reaped"]) + self.assertIsNone(self.orch.registry.get(dead)) + self.assertIsNotNone(self.orch.registry.get(alive)) + + def test_missing_live_source_ips_is_400(self) -> None: + status, _ = dispatch(self.orch, "POST", "/reconcile", _body({})) + self.assertEqual(400, status) + + def test_grace_seconds_is_honoured(self) -> None: + """A grace window wide enough to cover the row protects it.""" + self.orch.registry.register("10.0.0.3") + status, payload = dispatch( + self.orch, "POST", "/reconcile", + _body({"live_source_ips": [], "grace_seconds": 3600})) + self.assertEqual(200, status) + self.assertEqual([], payload["reaped"]) + + def test_non_string_entries_are_ignored(self) -> None: + dead = self._old("10.0.0.4") + status, payload = dispatch( + self.orch, "POST", "/reconcile", + _body({"live_source_ips": [None, 7, "10.0.0.9"]})) + self.assertEqual(200, status) + self.assertEqual([dead], payload["reaped"]) diff --git a/tests/unit/test_orchestrator_registry.py b/tests/unit/test_orchestrator_registry.py index 58751ab..3328a71 100644 --- a/tests/unit/test_orchestrator_registry.py +++ b/tests/unit/test_orchestrator_registry.py @@ -4,7 +4,9 @@ from __future__ import annotations import sqlite3 import tempfile +import time import unittest +from contextlib import closing from pathlib import Path from bot_bottle.orchestrator.registry import ( @@ -169,3 +171,81 @@ class TestRegistryStore(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class TestReapAbsent(unittest.TestCase): + """`reap_absent` — the self-heal for rows whose bottle is gone. + + An orphan is not merely untidy: source IPs get recycled, and + `by_source_ip` fail-closes on ambiguity, so a leftover row at a reused + address resolves *no* policy for the next bottle that lands there and + every host it asks for is denied. + """ + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.db = Path(self._tmp.name) / "registry.db" + self.store = RegistryStore(self.db) + self.store.migrate() + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _aged(self, source_ip: str, *, age: float) -> BottleRecord: + """Register a bottle and backdate it past the grace window.""" + rec = self.store.register(source_ip) + with closing(sqlite3.connect(self.db)) as conn: + conn.execute( + "UPDATE orchestrator_bottles SET created_at = ? WHERE bottle_id = ?", + (time.time() - age, rec.bottle_id), + ) + conn.commit() + return rec + + def test_reaps_row_with_no_live_container(self) -> None: + gone = self._aged("10.243.0.9", age=600) + reaped = self.store.reap_absent([]) + self.assertEqual([gone.bottle_id], [r.bottle_id for r in reaped]) + self.assertIsNone(self.store.get(gone.bottle_id)) + + def test_keeps_row_whose_ip_is_live(self) -> None: + alive = self._aged("10.243.0.9", age=600) + self.assertEqual([], self.store.reap_absent(["10.243.0.9"])) + self.assertIsNotNone(self.store.get(alive.bottle_id)) + + def test_grace_window_protects_an_in_flight_launch(self) -> None: + """A bottle registered moments ago is never reaped, even though the + caller's enumeration didn't see its address yet.""" + fresh = self.store.register("10.243.0.10") + self.assertEqual([], self.store.reap_absent([])) + self.assertIsNotNone(self.store.get(fresh.bottle_id)) + + def test_reaping_the_orphan_unbricks_the_reused_address(self) -> None: + """The regression this exists for: an orphan at an address that vmnet + later hands to a new bottle makes `by_source_ip` ambiguous, so the new + bottle resolves no policy at all.""" + orphan = self._aged("10.243.0.11", age=600) + # A new bottle lands on the recycled address. Force the row in directly + # so `register`'s own supersede sweep doesn't mask the ambiguity. + with closing(sqlite3.connect(self.db)) as conn: + conn.execute( + "INSERT INTO orchestrator_bottles " + "(bottle_id, source_ip, identity_token, state, created_at, metadata, policy) " + "VALUES ('newbottle', '10.243.0.11', 'tok-new', 'active', ?, '', 'routes: []')", + (time.time(),), + ) + conn.commit() + self.assertIsNone(self.store.by_source_ip("10.243.0.11")) # bricked + + reaped = self.store.reap_absent(["10.243.0.11"], grace_seconds=60) + self.assertEqual([orphan.bottle_id], [r.bottle_id for r in reaped]) + rec = self.store.by_source_ip("10.243.0.11") + assert rec is not None + self.assertEqual("newbottle", rec.bottle_id) + + def test_ignores_empty_ips_in_the_live_set(self) -> None: + gone = self._aged("10.243.0.12", age=600) + self.assertEqual( + [gone.bottle_id], + [r.bottle_id for r in self.store.reap_absent(["", "10.243.0.99"])], + ) diff --git a/tests/unit/test_orchestrator_service.py b/tests/unit/test_orchestrator_service.py index d99008b..efee477 100644 --- a/tests/unit/test_orchestrator_service.py +++ b/tests/unit/test_orchestrator_service.py @@ -4,8 +4,10 @@ from __future__ import annotations import json import secrets +import sqlite3 import tempfile import unittest +from contextlib import closing from pathlib import Path from unittest.mock import patch @@ -304,3 +306,55 @@ class TestOrchestratorSupervise(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class TestOrchestratorReconcile(unittest.TestCase): + """`reconcile` — drop rows for bottles that are no longer running.""" + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.secret = secrets.token_bytes(16) + self.db = Path(self._tmp.name) / "r.db" + self.store = RegistryStore(self.db) + self.store.migrate() + self.broker = StubBroker(self.secret) + self.orch = Orchestrator(self.store, self.broker, self.secret) + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _age_all(self, seconds: float) -> None: + """Backdate every row past the reap grace window.""" + with closing(sqlite3.connect(self.db)) as conn: + conn.execute( + "UPDATE orchestrator_bottles SET created_at = created_at - ?", (seconds,)) + conn.commit() + + def test_reaps_dead_bottle_and_forgets_its_tokens(self) -> None: + dead = self.orch.launch_bottle("10.243.0.1", tokens={"EGRESS_TOKEN_0": "s3cret"}) + live = self.orch.launch_bottle("10.243.0.2", tokens={"EGRESS_TOKEN_0": "keep"}) + self._age_all(600) + + self.assertEqual([dead.bottle_id], self.orch.reconcile(["10.243.0.2"])) + self.assertIsNone(self.store.get(dead.bottle_id)) + self.assertIsNotNone(self.store.get(live.bottle_id)) + # The in-memory egress credential goes with the row. + self.assertEqual({}, self.orch.tokens_for(dead.bottle_id)) + self.assertEqual({"EGRESS_TOKEN_0": "keep"}, self.orch.tokens_for(live.bottle_id)) + + def test_reconcile_does_not_broker_a_teardown(self) -> None: + """The container is already gone — there is nothing to stop, and a + broker error must not stop the sweep clearing the row.""" + self.orch.launch_bottle("10.243.0.1") + self._age_all(600) + self.broker.launched.clear() + self.orch.reconcile([]) + self.assertEqual([], self.broker.torn_down) + + def test_reconcile_keeps_everything_when_all_are_live(self) -> None: + a = self.orch.launch_bottle("10.243.0.1") + b = self.orch.launch_bottle("10.243.0.2") + self._age_all(600) + self.assertEqual([], self.orch.reconcile(["10.243.0.1", "10.243.0.2"])) + self.assertIsNotNone(self.store.get(a.bottle_id)) + self.assertIsNotNone(self.store.get(b.bottle_id))