diff --git a/bot_bottle/backend/macos_container/consolidated_launch.py b/bot_bottle/backend/macos_container/consolidated_launch.py index 5b7efc1..0d7fdfd 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, enumerate_active from .gateway import GATEWAY_NETWORK from .gateway_provision import AppleGatewayTransport from .infra import MacosInfraService, OrchestratorStartError @@ -89,6 +92,25 @@ 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 read is non-fatal) — the + reap's grace window, not this list, is what protects an in-flight + launch.""" + ips: list[str] = [] + for agent in enumerate_active(): + ip = container_mod.try_container_ipv4_on_network( + f"{CONTAINER_NAME_PREFIX}{agent.slug}", network, + ) + if ip: + ips.append(ip) + return ips + + def register_agent( egress_plan: EgressPlan, git_gate_plan: GitGatePlan, @@ -103,6 +125,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 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 +168,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 46dcb72..d80311d 100644 --- a/bot_bottle/backend/macos_container/enumerate.py +++ b/bot_bottle/backend/macos_container/enumerate.py @@ -8,7 +8,10 @@ from ...bottle_state import read_metadata from .. import ActiveAgent from .infra import INFRA_NAME -_PREFIX = "bot-bottle-" +# The name every agent container carries: `bot-bottle-`. Exported +# because reconciliation has to map an enumerated slug back to a container +# name to read its address. +CONTAINER_NAME_PREFIX = "bot-bottle-" # The shared per-host infra container carries the same prefix as agent # containers but is infrastructure, not a bottle — one control plane + gateway # serves every agent, so listing it as an agent would invent one per host. @@ -26,9 +29,9 @@ def enumerate_active() -> list[ActiveAgent]: return [] out: list[ActiveAgent] = [] for name in sorted(line.strip() for line in result.stdout.splitlines()): - if not name.startswith(_PREFIX) or name in _INFRA_NAMES: + if not name.startswith(CONTAINER_NAME_PREFIX) or name in _INFRA_NAMES: continue - slug = name[len(_PREFIX):] + slug = name[len(CONTAINER_NAME_PREFIX):] metadata = read_metadata(slug) out.append(ActiveAgent( backend_name="macos-container", 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_macos_consolidated_launch.py b/tests/unit/test_macos_consolidated_launch.py index f3e395d..aa67469 100644 --- a/tests/unit/test_macos_consolidated_launch.py +++ b/tests/unit/test_macos_consolidated_launch.py @@ -134,3 +134,69 @@ 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.try_container_ipv4_on_network", + 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.try_container_ipv4_on_network", + side_effect=["", "10.0.0.2"]): + self.assertEqual(["10.0.0.2"], 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() 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))