fix(orchestrator): reap registry rows whose bottle is no longer running

A registry row only ever left the registry two ways: an explicit
teardown_bottle (the launcher's cleanup callback) or the same-IP supersede
sweep in register(). Neither runs when the launching CLI dies hard, 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
resolves no policy at all for the next bottle that lands there — and a
bottle with no policy denies every host, which surfaces to the agent as
"host X is not in the allowlist" for hosts that were never the problem.

Add reap_absent/reconcile and call it from the macOS launch path before
registering, so each launch self-heals the registry. Restores the invariant
the data plane needs: at most one active row per live address, and none for
a dead one. The second half matters as much as the first — when several rows
claim a *live* address the newest wins and the rest are swept, otherwise a
recycled address stays ambiguous, which is exactly the bricked state.

The host supplies the live set because the orchestrator runs inside the
infra container and cannot see the backend. A grace window exempts rows
younger than it, so reconciliation cannot race a bottle still coming up, and
a reconcile failure is logged rather than blocking an otherwise-fine launch.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-20 22:29:46 -04:00
committed by codex
parent 5e01c28016
commit e4d53fd360
10 changed files with 460 additions and 2 deletions
@@ -36,9 +36,12 @@ from dataclasses import dataclass
from ...egress import EgressPlan from ...egress import EgressPlan
from ...git_gate import GitGatePlan 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 ...orchestrator.registration import registration_inputs
from ..docker.gateway_provision import deprovision_git_gate, provision_git_gate 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 import GATEWAY_NETWORK
from .gateway_provision import AppleGatewayTransport from .gateway_provision import AppleGatewayTransport
from .infra import MacosInfraService, OrchestratorStartError 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( def register_agent(
egress_plan: EgressPlan, egress_plan: EgressPlan,
git_gate_plan: GitGatePlan, git_gate_plan: GitGatePlan,
@@ -103,6 +125,16 @@ def register_agent(
container — it is the attribution key the gateway resolves policy by. container — it is the attribution key the gateway resolves policy by.
Raises on failure; the caller tears down.""" Raises on failure; the caller tears down."""
client = OrchestratorClient(endpoint.orchestrator_url) 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) inputs = registration_inputs(egress_plan)
reg = client.register_bottle( reg = client.register_bottle(
source_ip, image_ref=image_ref, policy=inputs.policy, source_ip, image_ref=image_ref, policy=inputs.policy,
@@ -136,6 +168,7 @@ __all__ = [
"GatewayEndpoint", "GatewayEndpoint",
"LaunchContext", "LaunchContext",
"ensure_gateway", "ensure_gateway",
"live_source_ips",
"register_agent", "register_agent",
"teardown_consolidated", "teardown_consolidated",
"ConsolidatedLaunchError", "ConsolidatedLaunchError",
+15
View File
@@ -15,6 +15,7 @@ from __future__ import annotations
import json import json
import urllib.error import urllib.error
import urllib.request import urllib.request
from collections.abc import Iterable
from dataclasses import dataclass from dataclasses import dataclass
from ..paths import host_control_plane_token from ..paths import host_control_plane_token
@@ -147,6 +148,20 @@ class OrchestratorClient:
raise OrchestratorClientError(f"teardown {bottle_id}: HTTP {status}") raise OrchestratorClientError(f"teardown {bottle_id}: HTTP {status}")
return True 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: def set_policy(self, bottle_id: str, policy: str) -> bool:
"""Live-reload a bottle's policy (`PUT /bottles/<id>/policy`). False on """Live-reload a bottle's policy (`PUT /bottles/<id>/policy`). False on
404 (unknown bottle).""" 404 (unknown bottle)."""
+24
View File
@@ -13,6 +13,9 @@ vsock / unix-socket portability caveats):
PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload) PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload)
body: {"policy"} body: {"policy"}
DELETE /bottles/<bottle_id> -> 200 {"torn_down": true} | 404 (teardown) DELETE /bottles/<bottle_id> -> 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 /attribute -> 200 {"bottle_id"} | 403
POST /resolve -> 200 {"bottle_id","policy"} | 403 POST /resolve -> 200 {"bottle_id","policy"} | 403
body: {"source_ip","identity_token"} 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 200, {"torn_down": True}
return 404, {"error": "no such bottle"} 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": if method == "POST" and route == "/attribute":
try: try:
data = _parse_json_object(body) data = _parse_json_object(body)
+72
View File
@@ -32,6 +32,7 @@ import hmac
import secrets import secrets
import sqlite3 import sqlite3
import time import time
from collections.abc import Iterable
from dataclasses import dataclass from dataclasses import dataclass
from pathlib import Path 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. # 256 bits of urandom, URL-safe — unguessable per-bottle identity token.
IDENTITY_TOKEN_BYTES = 32 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: def new_identity_token() -> str:
"""A fresh per-bottle identity token (PRD 0070 attribution defence).""" """A fresh per-bottle identity token (PRD 0070 attribution defence)."""
@@ -225,6 +232,70 @@ class RegistryStore(DbStore):
).fetchall() ).fetchall()
return [_row_to_record(r) for r in rows] 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: def by_source_ip(self, source_ip: str) -> BottleRecord | None:
"""Network-layer attribution: the single active bottle at this source """Network-layer attribution: the single active bottle at this source
IP, or None if unknown or ambiguous (more than one — a IP, or None if unknown or ambiguous (more than one — a
@@ -262,4 +333,5 @@ __all__ = [
"new_identity_token", "new_identity_token",
"default_db_path", "default_db_path",
"IDENTITY_TOKEN_BYTES", "IDENTITY_TOKEN_BYTES",
"DEFAULT_REAP_GRACE_SECONDS",
] ]
+30 -1
View File
@@ -13,15 +13,20 @@ Launch lifecycle:
and returns the record. If the broker rejects/fails, the registry entry and returns the record. If the broker rejects/fails, the registry entry
is rolled back so a failed launch leaves no orphan. is rolled back so a failed launch leaves no orphan.
* `teardown_bottle` sends a signed teardown request, then deregisters. * `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 from __future__ import annotations
import json import json
from collections.abc import Iterable
from datetime import datetime, timezone from datetime import datetime, timezone
from .broker import LaunchBroker, LaunchRequest, sign_request 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 .gateway import Gateway
from ..supervise import ( from ..supervise import (
AuditEntry, AuditEntry,
@@ -117,6 +122,30 @@ class Orchestrator:
self._tokens.pop(bottle_id, None) self._tokens.pop(bottle_id, None)
return True 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]: def tokens_for(self, bottle_id: str) -> dict[str, str]:
"""The bottle's in-memory egress auth tokens (env_name -> value), or """The bottle's in-memory egress auth tokens (env_name -> value), or
empty. The gateway injects these per request; they are never empty. The gateway injects these per request; they are never
@@ -134,3 +134,69 @@ class TestTeardown(unittest.TestCase):
if __name__ == "__main__": if __name__ == "__main__":
unittest.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()
+29
View File
@@ -104,3 +104,32 @@ class TestHealthAndPolicy(unittest.TestCase):
if __name__ == "__main__": if __name__ == "__main__":
unittest.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([])
@@ -8,11 +8,13 @@ from __future__ import annotations
import json import json
import secrets import secrets
import sqlite3
import tempfile import tempfile
import threading import threading
import unittest import unittest
import urllib.error import urllib.error
import urllib.request import urllib.request
from contextlib import closing
from pathlib import Path from pathlib import Path
from unittest.mock import patch from unittest.mock import patch
@@ -264,6 +266,7 @@ class TestControlPlaneAuth(unittest.TestCase):
("POST", "/bottles", _body({"source_ip": "10.0.0.1"})), ("POST", "/bottles", _body({"source_ip": "10.0.0.1"})),
("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})), ("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})),
("DELETE", "/bottles/x", b""), ("DELETE", "/bottles/x", b""),
("POST", "/reconcile", _body({"live_source_ips": []})),
("POST", "/resolve", _body({"source_ip": "10.0.0.1", "identity_token": "t"})), ("POST", "/resolve", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("POST", "/attribute", _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""), ("GET", "/supervise/proposals", b""),
@@ -387,3 +390,56 @@ class TestDispatchSupervise(unittest.TestCase):
if __name__ == "__main__": if __name__ == "__main__":
unittest.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"])
+80
View File
@@ -4,7 +4,9 @@ from __future__ import annotations
import sqlite3 import sqlite3
import tempfile import tempfile
import time
import unittest import unittest
from contextlib import closing
from pathlib import Path from pathlib import Path
from bot_bottle.orchestrator.registry import ( from bot_bottle.orchestrator.registry import (
@@ -169,3 +171,81 @@ class TestRegistryStore(unittest.TestCase):
if __name__ == "__main__": if __name__ == "__main__":
unittest.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"])],
)
+54
View File
@@ -4,8 +4,10 @@ from __future__ import annotations
import json import json
import secrets import secrets
import sqlite3
import tempfile import tempfile
import unittest import unittest
from contextlib import closing
from pathlib import Path from pathlib import Path
from unittest.mock import patch from unittest.mock import patch
@@ -304,3 +306,55 @@ class TestOrchestratorSupervise(unittest.TestCase):
if __name__ == "__main__": if __name__ == "__main__":
unittest.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))