Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 76037f36a6 | |||
| 4b17e6d683 | |||
| 7a48ea2b0c | |||
| ec953ceda7 | |||
| ed0f95f445 | |||
| 794e4e662d | |||
| f2fe1f9b2d |
@@ -166,9 +166,17 @@ def register_agent(
|
||||
# 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.
|
||||
#
|
||||
# The orchestrator now pulls the live set from the broker itself, so this is
|
||||
# a bare trigger. Until macOS moves onto the host controller (the follow-up
|
||||
# that routes launch through the broker), the orchestrator's broker is the
|
||||
# stub, whose live set is the bottles it recorded launches for — so the sweep
|
||||
# only under-reaps (leaves an orphan a cycle longer), never reaps a healthy
|
||||
# bottle. The Apple-container enumeration below (`live_source_ips`) becomes
|
||||
# that host controller's `list_live` there.
|
||||
try:
|
||||
client.reconcile(live_source_ips(endpoint.network))
|
||||
except (OrchestratorClientError, EnumerationError) as e:
|
||||
client.reconcile()
|
||||
except OrchestratorClientError as e:
|
||||
info(f"registry reconciliation skipped: {e}")
|
||||
reg = provision_bottle(
|
||||
client, source_ip, egress_plan, git_gate_plan, MacosGatewayTransport(),
|
||||
|
||||
@@ -17,7 +17,10 @@ from pathlib import Path
|
||||
|
||||
from .. import log
|
||||
from .store.store_manager import StoreManager
|
||||
from .broker import LaunchBroker, StubBroker
|
||||
from ..paths import LAUNCH_BROKER_KEY_ENV
|
||||
from .broker import StubBroker, SubmitBroker
|
||||
from .broker_client import BrokerClient
|
||||
from .host_server import DEFAULT_PORT, broker_secret
|
||||
from .server import make_server
|
||||
from .docker_broker import DockerBroker
|
||||
from .store.registry_store import RegistryStore, default_db_path
|
||||
@@ -34,8 +37,13 @@ def main(argv: list[str] | None = None) -> int:
|
||||
help=f"registry DB path (default: {default_db_path()})",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--broker", choices=("stub", "docker"), default="stub",
|
||||
help="launch broker: 'stub' records requests; 'docker' runs containers",
|
||||
"--broker", choices=("stub", "docker", "http"), default="stub",
|
||||
help="launch broker: 'stub' records requests; 'docker' runs containers "
|
||||
"in-process; 'http' relays signed requests to a host control server",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--host-controller-url", default=f"http://127.0.0.1:{DEFAULT_PORT}",
|
||||
help="host control server URL (used only with --broker http)",
|
||||
)
|
||||
args = parser.parse_args(argv)
|
||||
|
||||
@@ -47,11 +55,27 @@ def main(argv: list[str] | None = None) -> int:
|
||||
# operator reaches it over HTTP (never a second, disconnected DB).
|
||||
StoreManager(registry.db_path).migrate()
|
||||
|
||||
# An ephemeral signing secret ties the orchestrator (signer) to its
|
||||
# broker (verifier). 'stub' records launches instead of starting
|
||||
# anything; 'docker' runs real containers (firecracker drops in later).
|
||||
secret = secrets.token_bytes(32)
|
||||
broker: LaunchBroker = DockerBroker(secret) if args.broker == "docker" else StubBroker(secret)
|
||||
# A signing secret ties the orchestrator (signer) to its broker (verifier).
|
||||
# 'stub' records launches instead of starting anything; 'docker' runs real
|
||||
# containers in-process; 'http' relays signed requests to a separate host
|
||||
# control server, which verifies and launches. For 'stub'/'docker' the secret
|
||||
# is ephemeral (signer and verifier share this process). For 'http' it must be
|
||||
# the SAME key the host controller holds — and this process is the *guest*
|
||||
# (signer), so it must be given that key by injection, NOT mint its own
|
||||
# process-local one (which would diverge from the host's and 401 every launch).
|
||||
broker: SubmitBroker
|
||||
if args.broker == "http":
|
||||
secret = broker_secret() # env-injected only; no host-file fallback here
|
||||
if secret is None:
|
||||
parser.error(
|
||||
f"--broker http requires the launch-broker key injected as "
|
||||
f"${LAUNCH_BROKER_KEY_ENV} (the host controller owns/mints it); the "
|
||||
"orchestrator must not mint its own or it would diverge from the host's"
|
||||
)
|
||||
broker = BrokerClient(args.host_controller_url)
|
||||
else:
|
||||
secret = secrets.token_bytes(32)
|
||||
broker = DockerBroker(secret) if args.broker == "docker" else StubBroker(secret)
|
||||
orchestrator = OrchestratorCore(registry, broker, secret)
|
||||
|
||||
server = make_server(orchestrator, host=args.host, port=args.port)
|
||||
|
||||
@@ -29,29 +29,71 @@ import json
|
||||
import secrets
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from typing import Protocol
|
||||
|
||||
_JWT_HEADER = {"alg": "HS256", "typ": "JWT"}
|
||||
_ALLOWED_OPS = ("launch", "teardown")
|
||||
|
||||
# The closed op vocabulary, split by kind (PRD "ids + static flags" rule):
|
||||
# * mutation ops act on ONE bottle — they carry its id (+ static launch flags);
|
||||
# * query ops enumerate host state — they carry NO ids or flags at all.
|
||||
# Each new host-privileged op is added here deliberately; `verify_request`
|
||||
# enforces the per-kind shape so the privileged surface can't drift.
|
||||
_MUTATION_OPS = ("launch", "teardown")
|
||||
_QUERY_OPS = ("list_live",)
|
||||
_ALLOWED_OPS = _MUTATION_OPS + _QUERY_OPS
|
||||
|
||||
# Every claim key a signed request may carry — the request fields plus the two
|
||||
# per-signature envelope claims. A strict allow-list (open question 1, resolved
|
||||
# "yes"): an unknown key is rejected outright, so the schema can't silently widen
|
||||
# as ops are added.
|
||||
_ALLOWED_CLAIMS = frozenset(
|
||||
{"op", "bottle_id", "source_ip", "image_ref", "slot", "jti", "iat"}
|
||||
)
|
||||
|
||||
|
||||
class BrokerAuthError(Exception):
|
||||
"""A broker request failed provenance or schema verification —
|
||||
bad/absent signature, malformed token, or a payload that doesn't match
|
||||
the fixed launch-request shape. Fail-closed: the broker must not act."""
|
||||
the fixed launch-request shape. Fail-closed: the broker must not act.
|
||||
|
||||
A **definite** negative: nothing was launched, so a caller may safely roll
|
||||
back as if the op never happened."""
|
||||
|
||||
|
||||
class BrokerUnavailableError(Exception):
|
||||
"""A brokered request could not be carried to a verdict: the broker (or the
|
||||
wire to it) was unreachable, timed out, or dropped the response.
|
||||
|
||||
Crucially **ambiguous** — unlike `BrokerAuthError`, the op MAY already have
|
||||
taken effect on the backend before the response was lost, so a caller must
|
||||
NOT assume it did nothing (e.g. must not roll a registry row back as if no
|
||||
launch happened, which would orphan a running container). Only the in-process
|
||||
brokers never raise this; the out-of-process `BrokerClient` does."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class LaunchRequest:
|
||||
"""The structured, un-coercible launch/teardown request. Ids + flags
|
||||
only — `image_ref` is a content-addressed id, never a path/argv."""
|
||||
class BrokerRequest:
|
||||
"""A structured, un-coercible broker request. Ids + static flags only —
|
||||
`image_ref` is a content-addressed id, never a path/argv.
|
||||
|
||||
A **mutation** op (`launch`/`teardown`) names the bottle it acts on in
|
||||
`bottle_id` (plus launch's static flags). A **query** op (`list_live`) acts
|
||||
on no single bottle, so it carries nothing but its `op` name — `bottle_id`
|
||||
and every flag stay empty, and `verify_request` rejects a query that smuggles
|
||||
any in. `LaunchRequest` is kept as an alias for the launch/teardown callers."""
|
||||
|
||||
op: str # one of _ALLOWED_OPS
|
||||
bottle_id: str
|
||||
bottle_id: str = ""
|
||||
source_ip: str = ""
|
||||
image_ref: str = ""
|
||||
slot: int | None = None
|
||||
|
||||
|
||||
# Historical name — the request type predates the query ops. Kept so the launch
|
||||
# path (`OrchestratorCore`, `DockerBroker`) reads unchanged.
|
||||
LaunchRequest = BrokerRequest
|
||||
|
||||
|
||||
# --- minimal JWS/JWT (HS256), stdlib only ----------------------------------
|
||||
|
||||
def _b64url(data: bytes) -> str:
|
||||
@@ -103,26 +145,55 @@ def verify_request(token: str, secret: bytes) -> LaunchRequest:
|
||||
raise BrokerAuthError("unexpected header/alg")
|
||||
if not isinstance(claims, dict):
|
||||
raise BrokerAuthError("claims are not an object")
|
||||
# Strict schema (open question 1): a signed request may carry ONLY the known
|
||||
# claim keys, so the privileged surface can't widen by smuggling an extra
|
||||
# claim past the fixed field checks below.
|
||||
unknown = set(claims) - _ALLOWED_CLAIMS
|
||||
if unknown:
|
||||
raise BrokerAuthError(f"unknown claim(s): {', '.join(sorted(unknown))}")
|
||||
|
||||
op = claims.get("op")
|
||||
bottle_id = claims.get("bottle_id")
|
||||
if not isinstance(op, str) or op not in _ALLOWED_OPS:
|
||||
raise BrokerAuthError("request does not match the broker-request schema")
|
||||
bottle_id = claims.get("bottle_id", "")
|
||||
source_ip = claims.get("source_ip", "")
|
||||
image_ref = claims.get("image_ref", "")
|
||||
slot = claims.get("slot")
|
||||
if (
|
||||
not isinstance(op, str) or op not in _ALLOWED_OPS
|
||||
or not isinstance(bottle_id, str) or not bottle_id
|
||||
not isinstance(bottle_id, str)
|
||||
or not isinstance(source_ip, str) or not isinstance(image_ref, str)
|
||||
or not (slot is None or isinstance(slot, int))
|
||||
):
|
||||
raise BrokerAuthError("request does not match the launch-request schema")
|
||||
return LaunchRequest(
|
||||
raise BrokerAuthError("request does not match the broker-request schema")
|
||||
# Per-kind shape: a mutation names exactly one bottle; a query names none and
|
||||
# carries no flags (the "no arguments" rule for `list_live`). Enforcing both
|
||||
# halves keeps `bottle_id`/flags from being a coercible field on a query op.
|
||||
if op in _MUTATION_OPS:
|
||||
if not bottle_id:
|
||||
raise BrokerAuthError(f"{op} requires a bottle_id")
|
||||
elif bottle_id or source_ip or image_ref or slot is not None:
|
||||
raise BrokerAuthError(f"{op} is a query op and takes no ids or flags")
|
||||
return BrokerRequest(
|
||||
op=op, bottle_id=bottle_id, source_ip=source_ip, image_ref=image_ref, slot=slot
|
||||
)
|
||||
|
||||
|
||||
# --- the broker itself ------------------------------------------------------
|
||||
|
||||
class SubmitBroker(Protocol):
|
||||
"""The broker surface `OrchestratorCore` depends on. `submit` verifies a
|
||||
signed *mutation* token and performs its op, returning the verified request;
|
||||
`list_live` verifies a signed `list_live` token and returns the source IPs of
|
||||
the bottles the backend currently has running (reconcile's live set). Both the
|
||||
in-process `LaunchBroker` and the out-of-process `BrokerClient` (which relays
|
||||
the token to the host control server) satisfy it structurally, so the core is
|
||||
unchanged whether the backend is local or a real host service."""
|
||||
|
||||
def submit(self, token: str) -> LaunchRequest: ...
|
||||
|
||||
def list_live(self, token: str) -> list[str]: ...
|
||||
|
||||
|
||||
class LaunchBroker(abc.ABC):
|
||||
"""Verifies a signed request came from the orchestrator, then performs
|
||||
the backend-native launch/teardown. Subclasses implement `_launch` /
|
||||
@@ -132,15 +203,35 @@ class LaunchBroker(abc.ABC):
|
||||
self._secret = secret
|
||||
|
||||
def submit(self, token: str) -> LaunchRequest:
|
||||
"""Verify `token` and perform its op. Returns the verified request;
|
||||
raises `BrokerAuthError` if provenance/shape fails."""
|
||||
"""Verify `token` and perform its mutation op. Returns the verified
|
||||
request; raises `BrokerAuthError` if provenance/shape fails, or if the
|
||||
token is a non-mutation op (a `list_live` token routed here by mistake)."""
|
||||
req = verify_request(token, self._secret)
|
||||
if req.op == "launch":
|
||||
self._launch(req)
|
||||
else:
|
||||
elif req.op == "teardown":
|
||||
self._teardown(req)
|
||||
else:
|
||||
raise BrokerAuthError(f"{req.op} is not a submit op")
|
||||
return req
|
||||
|
||||
def list_live(self, token: str) -> list[str]:
|
||||
"""Verify a `list_live` `token` and return the source IPs of the bottles
|
||||
the backend currently has running. Raises `BrokerAuthError` on bad
|
||||
provenance/shape or a non-`list_live` op; a backend enumeration failure
|
||||
is converted to `BrokerUnavailableError` — the single "live set could not
|
||||
be determined" signal reconcile catches to skip the sweep rather than
|
||||
reaping healthy rows against a partial set (the out-of-process
|
||||
`BrokerClient` raises the same on an enumeration 502 / no response)."""
|
||||
req = verify_request(token, self._secret)
|
||||
if req.op != "list_live":
|
||||
raise BrokerAuthError(f"{req.op} is not a list_live op")
|
||||
try:
|
||||
return list(self._list_live())
|
||||
except Exception as e: # noqa: BLE001 — any enumeration failure means the
|
||||
# live set is unknown; reconcile must skip, never reap against it.
|
||||
raise BrokerUnavailableError(f"list_live enumeration failed: {e}") from e
|
||||
|
||||
@abc.abstractmethod
|
||||
def _launch(self, req: LaunchRequest) -> None:
|
||||
...
|
||||
@@ -149,15 +240,31 @@ class LaunchBroker(abc.ABC):
|
||||
def _teardown(self, req: LaunchRequest) -> None:
|
||||
...
|
||||
|
||||
@abc.abstractmethod
|
||||
def _list_live(self) -> list[str]:
|
||||
"""Every running bottle's source IP, backend-native. Must raise (not
|
||||
return a partial list) if the live set can't be determined
|
||||
authoritatively, so reconcile skips rather than reaping healthy rows."""
|
||||
...
|
||||
|
||||
|
||||
class StubBroker(LaunchBroker):
|
||||
"""Dev-harness broker: records verified requests without launching
|
||||
anything. Exercises the full sign -> verify -> act contract in-process."""
|
||||
anything. Exercises the full sign -> verify -> act contract in-process.
|
||||
|
||||
Its live set (`_list_live`) is, by default, every bottle it recorded a
|
||||
launch for and no teardown since — the in-process analogue of enumerating
|
||||
the backend. Tests that need to simulate a bottle dying out from under the
|
||||
registry (the case reconcile exists for) set `live_source_ips` explicitly to
|
||||
override that derived set."""
|
||||
|
||||
def __init__(self, secret: bytes) -> None:
|
||||
super().__init__(secret)
|
||||
self.launched: list[LaunchRequest] = []
|
||||
self.torn_down: list[LaunchRequest] = []
|
||||
# When not None, the exact live set `_list_live` reports — lets a test
|
||||
# say "only these IPs are still up" regardless of what was launched.
|
||||
self.live_source_ips: list[str] | None = None
|
||||
|
||||
def _launch(self, req: LaunchRequest) -> None:
|
||||
self.launched.append(req)
|
||||
@@ -165,10 +272,20 @@ class StubBroker(LaunchBroker):
|
||||
def _teardown(self, req: LaunchRequest) -> None:
|
||||
self.torn_down.append(req)
|
||||
|
||||
def _list_live(self) -> list[str]:
|
||||
if self.live_source_ips is not None:
|
||||
return list(self.live_source_ips)
|
||||
torn = {r.bottle_id for r in self.torn_down}
|
||||
return [r.source_ip for r in self.launched
|
||||
if r.bottle_id not in torn and r.source_ip]
|
||||
|
||||
|
||||
__all__ = [
|
||||
"BrokerAuthError",
|
||||
"BrokerUnavailableError",
|
||||
"BrokerRequest",
|
||||
"LaunchRequest",
|
||||
"SubmitBroker",
|
||||
"LaunchBroker",
|
||||
"StubBroker",
|
||||
"sign_request",
|
||||
|
||||
@@ -0,0 +1,168 @@
|
||||
"""Orchestrator-side broker transport (issue #468, chunk 1).
|
||||
|
||||
The signer's half of the launch-broker transport gap. `BrokerClient` satisfies
|
||||
the exact `submit(token)` contract `OrchestratorCore` already depends on (see
|
||||
`broker.SubmitBroker`), but instead of verifying and launching in-process it POSTs
|
||||
the signed token to the host control server over HTTP (stdlib `urllib`, like
|
||||
`orchestrator/client.py`). Because it is drop-in for that interface, wiring a real
|
||||
out-of-process backend does not change the core: it still signs a request and
|
||||
calls `submit()`; only the wire is new.
|
||||
|
||||
A provenance/schema rejection from the host controller (HTTP 401) is re-raised as
|
||||
the same `BrokerAuthError` the in-process broker raises, so the launch path's
|
||||
rollback-on-failure (`OrchestratorCore.launch_bottle`) behaves identically whether
|
||||
the broker is local or remote.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
from .broker import BrokerAuthError, BrokerUnavailableError, LaunchRequest
|
||||
|
||||
DEFAULT_TIMEOUT_SECONDS = 5.0
|
||||
|
||||
|
||||
class BrokerClientError(RuntimeError):
|
||||
"""The host control server *responded*, but with an unexpected status other
|
||||
than the fail-closed 401 (which surfaces as `BrokerAuthError`) — e.g. a 502
|
||||
backend failure or a malformed body. A definite negative: the host processed
|
||||
the request and it did not launch. (A *no-response* failure — unreachable /
|
||||
timeout / dropped — is the ambiguous `BrokerUnavailableError` instead.)"""
|
||||
|
||||
|
||||
class BrokerClient:
|
||||
"""Drop-in `submit(token)` that relays a signed request to the host control
|
||||
server. Holds no secret — provenance rides entirely in the signed token, so a
|
||||
caller that can reach this client still cannot forge a launch."""
|
||||
|
||||
def __init__(self, base_url: str, *, timeout: float = DEFAULT_TIMEOUT_SECONDS) -> None:
|
||||
self._base = base_url.rstrip("/")
|
||||
self._timeout = timeout
|
||||
|
||||
def submit(self, token: str) -> LaunchRequest:
|
||||
"""POST the signed token to the host controller and return the request it
|
||||
verified and acted on.
|
||||
|
||||
Raises `BrokerAuthError` on a fail-closed 401 (bad provenance/schema —
|
||||
the same exception the in-process broker raises); `BrokerClientError` if
|
||||
the host *responds* with any other non-success status or a malformed
|
||||
body (a definite negative); or `BrokerUnavailableError` if no response is
|
||||
obtained (unreachable / timeout / dropped) — the **ambiguous** case, where
|
||||
the host may already have acted, so the caller must not roll back."""
|
||||
data = json.dumps({"token": token}).encode()
|
||||
req = urllib.request.Request(
|
||||
f"{self._base}/broker", data=data, method="POST",
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=self._timeout) as resp:
|
||||
return _request_from(_json_object(resp.read()))
|
||||
except urllib.error.HTTPError as e:
|
||||
detail = _error_detail(e)
|
||||
if e.code == 401:
|
||||
raise BrokerAuthError(
|
||||
detail or "host controller rejected the request"
|
||||
) from e
|
||||
raise BrokerClientError(
|
||||
f"POST /broker: HTTP {e.code} {detail}".rstrip()
|
||||
) from e
|
||||
except (urllib.error.URLError, TimeoutError, OSError) as e:
|
||||
# No usable response — unreachable, timed out, or the connection
|
||||
# dropped mid-exchange. Ambiguous: the request may already have
|
||||
# launched the bottle, so this is NOT a definite failure.
|
||||
raise BrokerUnavailableError(f"POST /broker: {e}") from e
|
||||
|
||||
def list_live(self, token: str) -> list[str]:
|
||||
"""POST the signed `list_live` token to the host controller and return
|
||||
the source IPs of the bottles it enumerated as running.
|
||||
|
||||
Raises `BrokerAuthError` on a fail-closed 401 (bad provenance/schema);
|
||||
`BrokerClientError` on any other non-success status or a malformed body
|
||||
(a backend-enumeration 502 included); or `BrokerUnavailableError` if no
|
||||
response is obtained. A query has no backend side effect, so — unlike
|
||||
`submit` — every one of these is a *definite* "no live set"; reconcile
|
||||
catches all three and skips the sweep rather than reaping against an
|
||||
unknown or partial set."""
|
||||
data = json.dumps({"token": token}).encode()
|
||||
req = urllib.request.Request(
|
||||
f"{self._base}/broker/live", data=data, method="POST",
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=self._timeout) as resp:
|
||||
return _source_ips_from(_json_object(resp.read()))
|
||||
except urllib.error.HTTPError as e:
|
||||
detail = _error_detail(e)
|
||||
if e.code == 401:
|
||||
raise BrokerAuthError(
|
||||
detail or "host controller rejected the request"
|
||||
) from e
|
||||
raise BrokerClientError(
|
||||
f"POST /broker/live: HTTP {e.code} {detail}".rstrip()
|
||||
) from e
|
||||
except (urllib.error.URLError, TimeoutError, OSError) as e:
|
||||
raise BrokerUnavailableError(f"POST /broker/live: {e}") from e
|
||||
|
||||
|
||||
def _json_object(raw: bytes) -> dict[str, object]:
|
||||
"""Parse a JSON object, tolerating an empty or malformed body (→ {}), like
|
||||
the orchestrator client — a bad body becomes a clean 'missing field' error
|
||||
downstream rather than an opaque JSON crash."""
|
||||
if not raw:
|
||||
return {}
|
||||
try:
|
||||
obj = json.loads(raw)
|
||||
except ValueError:
|
||||
return {}
|
||||
return obj if isinstance(obj, dict) else {}
|
||||
|
||||
|
||||
def _error_detail(e: urllib.error.HTTPError) -> str:
|
||||
"""The `error` string from a structured error response, best-effort — an
|
||||
error body may be absent or unreadable, in which case there is no detail."""
|
||||
try:
|
||||
detail = _json_object(e.read()).get("error", "")
|
||||
except Exception: # noqa: BLE001 — the error body is advisory only
|
||||
return ""
|
||||
return detail if isinstance(detail, str) else ""
|
||||
|
||||
|
||||
def _source_ips_from(payload: dict[str, object]) -> list[str]:
|
||||
"""The `source_ips` list from a `/broker/live` response — every string
|
||||
entry, ignoring any non-string the host controller should never send. A
|
||||
missing/!list field is a malformed response (the query is meaningless
|
||||
without it), so it fails rather than silently reconciling against []."""
|
||||
raw = payload.get("source_ips")
|
||||
if not isinstance(raw, list):
|
||||
raise BrokerClientError("host controller response missing source_ips")
|
||||
return [ip for ip in raw if isinstance(ip, str) and ip]
|
||||
|
||||
|
||||
def _request_from(payload: dict[str, object]) -> LaunchRequest:
|
||||
"""Reconstruct the verified `LaunchRequest` the controller echoed, so the
|
||||
returned value matches the in-process broker's (which returns the request it
|
||||
acted on). A missing op/bottle_id means a malformed response."""
|
||||
op = payload.get("op")
|
||||
bottle_id = payload.get("bottle_id")
|
||||
if not isinstance(op, str) or not isinstance(bottle_id, str) or not bottle_id:
|
||||
raise BrokerClientError("host controller response missing op/bottle_id")
|
||||
source_ip = payload.get("source_ip")
|
||||
image_ref = payload.get("image_ref")
|
||||
slot = payload.get("slot")
|
||||
return LaunchRequest(
|
||||
op=op,
|
||||
bottle_id=bottle_id,
|
||||
source_ip=source_ip if isinstance(source_ip, str) else "",
|
||||
image_ref=image_ref if isinstance(image_ref, str) else "",
|
||||
slot=slot if isinstance(slot, int) and not isinstance(slot, bool) else None,
|
||||
)
|
||||
|
||||
|
||||
__all__ = [
|
||||
"BrokerClient",
|
||||
"BrokerClientError",
|
||||
"DEFAULT_TIMEOUT_SECONDS",
|
||||
]
|
||||
@@ -15,7 +15,6 @@ from __future__ import annotations
|
||||
import json
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
from collections.abc import Iterable
|
||||
from dataclasses import dataclass
|
||||
|
||||
from ..log import debug
|
||||
@@ -194,14 +193,13 @@ 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 orchestrator container/VM."""
|
||||
body: dict[str, object] = {"live_source_ips": list(live_source_ips)}
|
||||
def reconcile(self, *, grace_seconds: float | None = None) -> list[str]:
|
||||
"""Trigger the orchestrator's self-heal sweep (`POST /reconcile`),
|
||||
returning the reaped bottle ids. The caller no longer enumerates the
|
||||
live set: the orchestrator pulls it from the host controller
|
||||
(`list_live`) itself, so this is a bare trigger with an optional
|
||||
`grace_seconds`."""
|
||||
body: dict[str, object] = {}
|
||||
if grace_seconds is not None:
|
||||
body["grace_seconds"] = grace_seconds
|
||||
payload = self._ok("POST", "/reconcile", body)
|
||||
|
||||
@@ -71,6 +71,31 @@ class DockerBroker(LaunchBroker):
|
||||
f"docker rm failed for {req.bottle_id}: {proc.stderr.strip()}"
|
||||
)
|
||||
|
||||
def _list_live(self) -> list[str]:
|
||||
"""Every running bottle container's network IP, from `docker ps` filtered
|
||||
to this broker's label. Raises `DockerBrokerError` if the listing or any
|
||||
inspect fails, so reconcile treats the live set as unknown and skips the
|
||||
sweep rather than reaping healthy rows against a partial enumeration."""
|
||||
listing = self._docker([
|
||||
"docker", "ps", "--filter", f"label={BOTTLE_ID_LABEL}",
|
||||
"--format", "{{.ID}}",
|
||||
])
|
||||
if listing.returncode != 0:
|
||||
raise DockerBrokerError(f"docker ps failed: {listing.stderr.strip()}")
|
||||
ips: list[str] = []
|
||||
for cid in (line for line in listing.stdout.splitlines() if line):
|
||||
inspected = self._docker([
|
||||
"docker", "inspect", "--format",
|
||||
"{{range .NetworkSettings.Networks}}{{.IPAddress}} {{end}}", cid,
|
||||
])
|
||||
if inspected.returncode != 0:
|
||||
raise DockerBrokerError(
|
||||
f"docker inspect {cid} failed; live set is not authoritative: "
|
||||
f"{inspected.stderr.strip()}"
|
||||
)
|
||||
ips.extend(ip for ip in inspected.stdout.split() if ip)
|
||||
return ips
|
||||
|
||||
|
||||
__all__ = [
|
||||
"DockerBroker",
|
||||
|
||||
@@ -0,0 +1,316 @@
|
||||
"""Host control server (issue #468) — the launch broker as a real host service.
|
||||
|
||||
Chunk 1 of the host-control-server stack closes the **transport** gap the PRD
|
||||
opens with: today `LaunchBroker.submit(token)` is an in-process method call from
|
||||
`OrchestratorCore`, and a real host service needs it reachable over the wire.
|
||||
This module is that service — the single privileged host component — reached over
|
||||
**HTTP** (the universal transport 0070 chose), mirroring the orchestrator control
|
||||
plane's shape (`orchestrator/server.py`): a pure `dispatch()` for socket-free
|
||||
testing, wrapped by a thin stdlib `http.server` adapter.
|
||||
|
||||
GET /health -> 200 {"status": "ok"}
|
||||
POST /broker -> 200 {"op", "bottle_id", "source_ip", "image_ref", "slot"}
|
||||
400 (bad body) | 401 (bad provenance/schema) | 502 (backend)
|
||||
body: {"token": "<signed launch/teardown JWT>"}
|
||||
POST /broker/live -> 200 {"source_ips": [...]}
|
||||
400 (bad body) | 401 (bad provenance/schema) | 502 (backend)
|
||||
body: {"token": "<signed list_live JWT>"}
|
||||
|
||||
The op vocabulary grows one signed op at a time (PRD gap 3); each verb reaches
|
||||
the backend only through `verify_request`, so nothing free-form ever crosses the
|
||||
wire. `/broker/live` is the first query op — reconcile's live-bottle enumeration,
|
||||
which the orchestrator (blind to the backend from inside its container) now pulls
|
||||
from the host controller instead of being handed by the CLI.
|
||||
|
||||
Only the **signed token** crosses the wire; the server holds the shared HS256
|
||||
secret and a real `LaunchBroker` (e.g. `DockerBroker`) and runs the existing
|
||||
`verify_request` + `_launch`/`_teardown` path behind the endpoint, so nothing
|
||||
free-form ever reaches it. Provenance/schema failures are fail-closed 401s that
|
||||
never touch the backend (`LaunchBroker.submit` verifies before acting), and a
|
||||
backend launch failure is a 502 the caller must surface — neither takes the
|
||||
controller down.
|
||||
|
||||
The signed launch token *is* the endpoint's authentication (its provenance is the
|
||||
whole point of the JWS), so `/broker` needs no separate caller credential; the
|
||||
host controller's own lifecycle endpoints, which do, arrive with the `host`-role
|
||||
tokens of the separate `HOST_CONTROLLER` trust domain in a later chunk.
|
||||
|
||||
The shared signing secret is the durable **launch-broker `TrustDomain` key**
|
||||
(#468/#476): a host-canonical key file minted 0600 on first use, provisioned to
|
||||
the orchestrator (signer) and this server (verifier). A backend launcher injects
|
||||
it via `$BOT_BOTTLE_LAUNCH_BROKER_KEY`; a host-side dev-harness process reads the
|
||||
key file directly. Durability is the point — a restarted orchestrator re-verifies
|
||||
against the same key, so re-adoption works.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import http.server
|
||||
import json
|
||||
import os
|
||||
import socketserver
|
||||
import sys
|
||||
import typing
|
||||
from urllib.parse import urlsplit
|
||||
|
||||
from .. import log
|
||||
from ..paths import LAUNCH_BROKER_KEY_ENV
|
||||
from ..trust_domain import LAUNCH_BROKER
|
||||
from .broker import BrokerAuthError, LaunchBroker
|
||||
from .docker_broker import DockerBroker
|
||||
|
||||
# JSON body payload type (parsed request / rendered response).
|
||||
Json = dict[str, object]
|
||||
|
||||
# Default host-controller port. Distinct from the orchestrator control plane
|
||||
# (8099) — a separate privileged component listening on its own socket.
|
||||
DEFAULT_PORT = 8091
|
||||
|
||||
# Cap on the request body. A signed broker request is tiny, so rejecting anything
|
||||
# larger *before reading it* keeps a caller that can merely reach the socket (no
|
||||
# signed token needed) from exhausting memory or a handler thread with a huge
|
||||
# Content-Length — the signed token, not mere reachability, is the authority.
|
||||
MAX_BODY_BYTES = 64 * 1024
|
||||
|
||||
# Per-request socket timeout, bounding how long a stalled / slow-loris caller can
|
||||
# hold a handler thread on this privileged listener.
|
||||
REQUEST_TIMEOUT_SECONDS = 15
|
||||
|
||||
|
||||
def _parse_json_object(body: bytes) -> Json:
|
||||
"""Parse a JSON object body. Raises ValueError for non-objects / bad JSON."""
|
||||
if not body:
|
||||
return {}
|
||||
obj = json.loads(body) # raises json.JSONDecodeError (a ValueError)
|
||||
if not isinstance(obj, dict):
|
||||
raise ValueError("request body must be a JSON object")
|
||||
return obj
|
||||
|
||||
|
||||
def broker_secret(
|
||||
environ: typing.Mapping[str, str] | None = None, *, allow_host_file: bool = False,
|
||||
) -> bytes | None:
|
||||
"""The shared launch-broker HS256 secret, as this process should use it.
|
||||
|
||||
Always prefers the key injected into this process's env
|
||||
(`$BOT_BOTTLE_LAUNCH_BROKER_KEY`). `allow_host_file` decides the fallback when
|
||||
it is absent, and the distinction is a security boundary:
|
||||
|
||||
- **Host-side** processes — the host controller and the host dev-harness — pass
|
||||
``allow_host_file=True`` to read (minting on first use) the durable host key
|
||||
file (``bot_bottle_root()/launch-broker-key``) they legitimately own.
|
||||
- The **guest orchestrator** (``--broker http``) keeps the default ``False``.
|
||||
It runs inside a container/VM whose ``bot_bottle_root()`` is process-local,
|
||||
so minting a file there would silently create a key UNRELATED to the host
|
||||
controller's — startup would succeed but every launch would be rejected 401.
|
||||
It must instead be *given* the key by its launcher, and fail closed (None)
|
||||
if it wasn't, rather than diverge.
|
||||
|
||||
None when no key is available (a guest with no injection, or an unwritable
|
||||
host root)."""
|
||||
key = LAUNCH_BROKER.key_from_env(environ)
|
||||
if not key and allow_host_file:
|
||||
try:
|
||||
key = LAUNCH_BROKER.signing_key() # host-canonical, minted on first use
|
||||
except OSError:
|
||||
return None
|
||||
return key.encode("utf-8") if key else None
|
||||
|
||||
|
||||
def dispatch( # pylint: disable=too-many-return-statements
|
||||
broker: LaunchBroker, method: str, path: str, body: bytes,
|
||||
) -> tuple[int, Json]:
|
||||
"""Route one host-control request to a (status, payload) pair. Pure — the
|
||||
only side effect is the broker's own backend launch — so routing is testable
|
||||
without a socket.
|
||||
|
||||
Total by design: a provenance/schema failure becomes 401 and a backend launch
|
||||
failure becomes 502 rather than raising, so one bad request can neither act
|
||||
on the backend nor take the controller down for the next caller."""
|
||||
route = urlsplit(path).path.rstrip("/") or "/"
|
||||
|
||||
if method == "GET" and route == "/health":
|
||||
return 200, {"status": "ok"}
|
||||
|
||||
if method == "POST" and route == "/broker":
|
||||
try:
|
||||
data = _parse_json_object(body)
|
||||
except ValueError as e:
|
||||
return 400, {"error": f"invalid JSON: {e}"}
|
||||
token = data.get("token")
|
||||
if not isinstance(token, str) or not token:
|
||||
return 400, {"error": "token (string) is required"}
|
||||
try:
|
||||
req = broker.submit(token)
|
||||
except BrokerAuthError as e:
|
||||
# Fail-closed: bad signature, malformed token, or off-schema payload.
|
||||
# `submit` verifies before acting, so nothing was launched.
|
||||
return 401, {"error": f"broker auth failed: {e}"}
|
||||
except Exception as e: # noqa: BLE001 — a backend launch failure (docker
|
||||
# down, image gone) is operational, not a control-plane bug; the
|
||||
# caller must see it as a distinct 502, and the server must stay up.
|
||||
return 502, {"error": f"backend launch failed: {e}"}
|
||||
return 200, {
|
||||
"op": req.op,
|
||||
"bottle_id": req.bottle_id,
|
||||
"source_ip": req.source_ip,
|
||||
"image_ref": req.image_ref,
|
||||
"slot": req.slot,
|
||||
}
|
||||
|
||||
if method == "POST" and route == "/broker/live":
|
||||
try:
|
||||
data = _parse_json_object(body)
|
||||
except ValueError as e:
|
||||
return 400, {"error": f"invalid JSON: {e}"}
|
||||
token = data.get("token")
|
||||
if not isinstance(token, str) or not token:
|
||||
return 400, {"error": "token (string) is required"}
|
||||
try:
|
||||
source_ips = broker.list_live(token)
|
||||
except BrokerAuthError as e:
|
||||
# Fail-closed: bad signature, malformed token, or a non-list_live op.
|
||||
return 401, {"error": f"broker auth failed: {e}"}
|
||||
except Exception as e: # noqa: BLE001 — a backend enumeration failure
|
||||
# (docker down, inspect failed) is operational, not a control-plane
|
||||
# bug; surface it as a 502 so the caller skips reconcile rather than
|
||||
# reaping against a partial set, and keep the controller up.
|
||||
return 502, {"error": f"backend enumeration failed: {e}"}
|
||||
return 200, {"source_ips": source_ips}
|
||||
|
||||
return 404, {"error": "not found"}
|
||||
|
||||
|
||||
class Handler(http.server.BaseHTTPRequestHandler):
|
||||
"""Thin stdlib adapter: read the body, call `dispatch`, write JSON."""
|
||||
|
||||
# Socket timeout per request (applied by StreamRequestHandler.setup) so a
|
||||
# stalled caller can't pin a handler thread on this privileged listener.
|
||||
timeout = REQUEST_TIMEOUT_SECONDS
|
||||
|
||||
# Quiet by default; opt back into stdlib access logging with
|
||||
# BOT_BOTTLE_HOST_CONTROLLER_DEBUG (the controller has its own logging).
|
||||
def log_message(self, format: str, *args: typing.Any) -> None: # noqa: A002
|
||||
if os.environ.get("BOT_BOTTLE_HOST_CONTROLLER_DEBUG"):
|
||||
super().log_message(format, *args)
|
||||
|
||||
def _serve(self, method: str) -> None:
|
||||
"""Read the request body (bounded), dispatch it, and write the JSON
|
||||
reply. A dispatch that raises (it shouldn't — dispatch is total) still
|
||||
returns a 500 rather than dropping the connection."""
|
||||
server = self.server
|
||||
assert isinstance(server, HostControlServer)
|
||||
try:
|
||||
length = int(self.headers.get("Content-Length") or 0)
|
||||
except ValueError:
|
||||
self._reply(400, {"error": "invalid Content-Length"})
|
||||
return
|
||||
if length < 0 or length > MAX_BODY_BYTES:
|
||||
# Reject before reading: nothing legitimate is this big, so an
|
||||
# oversized declared length is a bug or a resource-exhaustion attempt.
|
||||
self._reply(413, {"error": "request body too large"})
|
||||
return
|
||||
body = self.rfile.read(length) if length > 0 else b""
|
||||
try:
|
||||
status, payload = dispatch(server.broker, method, self.path, body)
|
||||
except Exception as e: # noqa: BLE001 — the controller must stay up
|
||||
sys.stderr.write(f"host controller: {method} {self.path} failed: {e!r}\n")
|
||||
sys.stderr.flush()
|
||||
status, payload = 500, {"error": f"internal error: {e}"}
|
||||
self._reply(status, payload)
|
||||
|
||||
def _reply(self, status: int, payload: typing.Mapping[str, object]) -> None:
|
||||
"""Write one JSON response with an explicit Content-Length."""
|
||||
data = json.dumps(payload).encode()
|
||||
self.send_response(status)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(data)))
|
||||
self.end_headers()
|
||||
self.wfile.write(data)
|
||||
|
||||
def do_GET(self) -> None:
|
||||
self._serve("GET")
|
||||
|
||||
def do_POST(self) -> None:
|
||||
self._serve("POST")
|
||||
|
||||
|
||||
class HostControlServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
|
||||
"""Threading HTTP server that carries the launch broker for its handlers.
|
||||
|
||||
The broker holds the shared signing secret and performs the backend-native
|
||||
launch/teardown; the server itself keeps no secret of its own — provenance
|
||||
rides entirely in each request's signed token."""
|
||||
|
||||
daemon_threads = True
|
||||
allow_reuse_address = True
|
||||
|
||||
def __init__(self, address: tuple[str, int], broker: LaunchBroker) -> None:
|
||||
self.broker = broker
|
||||
super().__init__(address, Handler)
|
||||
|
||||
|
||||
def make_host_server(
|
||||
broker: LaunchBroker, host: str = "127.0.0.1", port: int = DEFAULT_PORT
|
||||
) -> HostControlServer:
|
||||
"""Build (but do not start) a host control server. `port=0` binds an
|
||||
ephemeral port — read `server.server_address` for the actual one."""
|
||||
return HostControlServer((host, port), broker)
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
"""Run the host control server as a plain process (dev-harness).
|
||||
|
||||
python -m bot_bottle.orchestrator.host_server [--host H] [--port P]
|
||||
|
||||
Fail-closed: without the launch-broker key the server can verify no request's
|
||||
provenance, so it refuses to start rather than run a launcher that accepts
|
||||
unsigned input. As the host-side owner of the key, it may mint/read the host
|
||||
key file (`allow_host_file=True`)."""
|
||||
parser = argparse.ArgumentParser(prog="bot_bottle.orchestrator.host_server")
|
||||
parser.add_argument("--host", default="127.0.0.1", help="bind address")
|
||||
parser.add_argument("--port", type=int, default=DEFAULT_PORT, help="bind port (0 = ephemeral)")
|
||||
args = parser.parse_args(argv)
|
||||
|
||||
secret = broker_secret(allow_host_file=True)
|
||||
if secret is None:
|
||||
sys.stderr.write(
|
||||
f"host controller: refusing to start without the launch-broker key "
|
||||
f"(${LAUNCH_BROKER_KEY_ENV}, or a writable host root to mint it) — it "
|
||||
"could verify no request's provenance and would relay unsigned "
|
||||
"launches to the backend\n"
|
||||
)
|
||||
sys.stderr.flush()
|
||||
return 2
|
||||
|
||||
broker = DockerBroker(secret)
|
||||
server = make_host_server(broker, host=args.host, port=args.port)
|
||||
bound_host, bound_port = server.server_address[0], server.server_address[1]
|
||||
log.info(
|
||||
"host control server listening",
|
||||
context={"host": bound_host, "port": bound_port},
|
||||
)
|
||||
try:
|
||||
server.serve_forever()
|
||||
except KeyboardInterrupt:
|
||||
log.info("host controller shutting down")
|
||||
finally:
|
||||
server.server_close()
|
||||
return 0
|
||||
|
||||
|
||||
__all__ = [
|
||||
"dispatch",
|
||||
"Handler",
|
||||
"HostControlServer",
|
||||
"make_host_server",
|
||||
"broker_secret",
|
||||
"main",
|
||||
"Json",
|
||||
"DEFAULT_PORT",
|
||||
]
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -18,8 +18,8 @@ vsock / unix-socket portability caveats):
|
||||
body: {"env_var_secret"}
|
||||
DELETE /bottles/<bottle_id> -> 200 {"torn_down": true} | 404 (teardown)
|
||||
POST /reconcile -> 200 {"reaped": [bottle_id, ...]}
|
||||
body: {"live_source_ips": [...],
|
||||
["grace_seconds"]}
|
||||
body: {["grace_seconds"]} (live set is
|
||||
pulled from the host controller, not sent)
|
||||
POST /attribute -> 200 {"bottle_id"} | 403
|
||||
POST /resolve -> 200 {"bottle_id","policy"} | 403
|
||||
body: {"source_ip","identity_token"}
|
||||
@@ -207,20 +207,15 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
|
||||
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.
|
||||
# Host-driven self-heal trigger: the orchestrator pulls its own live set
|
||||
# from the host controller (`list_live`) and drops rows for every other
|
||||
# active bottle — the caller no longer supplies the live IPs. Trusted-
|
||||
# caller only: an agent that could reach this would trigger a sweep that
|
||||
# unregisters its neighbours. Body is optional (just `grace_seconds`).
|
||||
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"}
|
||||
if any(not isinstance(ip, str) or not ip for ip in raw_ips):
|
||||
return 400, {"error": "live_source_ips must contain non-empty strings"}
|
||||
live = raw_ips
|
||||
grace = data.get("grace_seconds")
|
||||
kwargs: dict[str, float] = {}
|
||||
if grace is not None:
|
||||
@@ -230,7 +225,7 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
|
||||
if not math.isfinite(parsed_grace) or parsed_grace < 0:
|
||||
return 400, {"error": "grace_seconds must be a non-negative finite number"}
|
||||
kwargs["grace_seconds"] = parsed_grace
|
||||
return 200, {"reaped": orch.reconcile(live, **kwargs)}
|
||||
return 200, {"reaped": orch.reconcile(**kwargs)}
|
||||
|
||||
if method == "POST" and route == "/attribute":
|
||||
try:
|
||||
|
||||
@@ -22,10 +22,17 @@ Launch lifecycle:
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections.abc import Iterable
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from .broker import LaunchBroker, LaunchRequest, sign_request
|
||||
from .. import log
|
||||
from .broker import (
|
||||
BrokerAuthError,
|
||||
BrokerUnavailableError,
|
||||
LaunchRequest,
|
||||
SubmitBroker,
|
||||
sign_request,
|
||||
)
|
||||
from .broker_client import BrokerClientError
|
||||
from .store.registry_store import DEFAULT_REAP_GRACE_SECONDS, BottleRecord, RegistryStore
|
||||
from .supervisor import (
|
||||
AuditEntry,
|
||||
@@ -62,7 +69,7 @@ class OrchestratorCore:
|
||||
def __init__(
|
||||
self,
|
||||
registry: RegistryStore,
|
||||
broker: LaunchBroker,
|
||||
broker: SubmitBroker,
|
||||
sign_secret: bytes,
|
||||
supervisor: Supervisor | None = None,
|
||||
) -> None:
|
||||
@@ -111,14 +118,23 @@ class OrchestratorCore:
|
||||
image_ref=image_ref,
|
||||
slot=slot,
|
||||
)
|
||||
launched = False
|
||||
try:
|
||||
self._broker.submit(sign_request(req, self._secret))
|
||||
launched = True
|
||||
finally:
|
||||
if not launched:
|
||||
self.registry.deregister(rec.bottle_id)
|
||||
self._tokens.pop(rec.bottle_id, None)
|
||||
except BrokerUnavailableError:
|
||||
# Ambiguous delivery failure (timeout / dropped response): the broker
|
||||
# may already have launched the bottle before the response was lost.
|
||||
# Do NOT deregister — that would orphan a running container with no
|
||||
# registry row (reconcile reaps rows, never containers). Keep the row
|
||||
# so reconcile reaps it iff the bottle is not actually live; surface
|
||||
# the error so the caller knows the launch is unconfirmed.
|
||||
raise
|
||||
except Exception:
|
||||
# A definite failure — a fail-closed rejection, a backend launch
|
||||
# error, or the host reporting it did not launch: nothing is running,
|
||||
# so roll the registry entry back to leave no orphan.
|
||||
self.registry.deregister(rec.bottle_id)
|
||||
self._tokens.pop(rec.bottle_id, None)
|
||||
raise
|
||||
return rec
|
||||
|
||||
def teardown_bottle(self, bottle_id: str) -> bool:
|
||||
@@ -136,22 +152,36 @@ class OrchestratorCore:
|
||||
|
||||
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.
|
||||
The live set is pulled from the **broker** (`list_live`), not passed in:
|
||||
only the host can enumerate its own containers, and the host controller
|
||||
is now that host component, so the orchestrator — blind to the backend
|
||||
from inside its container — asks it over the same signed seam it launches
|
||||
through. (Before the host control server this was a `live_source_ips`
|
||||
argument the CLI handed in; the tell that the orchestrator couldn't see
|
||||
the backend goes away with it — PRD gap 3.)
|
||||
|
||||
Fail-safe: if the broker can't return an authoritative live set
|
||||
(unreachable, timed out, enumeration failed), the sweep is **skipped**,
|
||||
never run against an empty/partial set — reaping a healthy bottle is far
|
||||
worse than leaving an orphan one more cycle. Deliberately does not broker
|
||||
a teardown: a reaped container is already gone, so there is nothing to
|
||||
stop.
|
||||
|
||||
See `RegistryStore.reap_absent` for why orphans accumulate and why
|
||||
they are harmful rather than merely untidy."""
|
||||
token = sign_request(LaunchRequest(op="list_live"), self._secret)
|
||||
try:
|
||||
live_source_ips = self._broker.list_live(token)
|
||||
except (BrokerUnavailableError, BrokerAuthError, BrokerClientError) as e:
|
||||
log.info("reconcile skipped: live set unavailable",
|
||||
context={"error": str(e)})
|
||||
return []
|
||||
reaped = self.registry.reap_absent(
|
||||
live_source_ips, grace_seconds=grace_seconds)
|
||||
for rec in reaped:
|
||||
|
||||
@@ -12,12 +12,22 @@ reattachment path reads ENV_VAR_SECRET from the running agent container via
|
||||
``POST /bottles/<id>/reprovision_gateway``; the orchestrator decrypts the
|
||||
stored rows and re-populates ``_tokens``.
|
||||
|
||||
Encryption scheme: HMAC-SHA256 used as a PRF in CTR mode (stdlib-only,
|
||||
no external deps). Each value is encrypted independently. The output blob is
|
||||
``nonce (16 bytes) || ciphertext`` encoded as URL-safe base64 (no padding).
|
||||
Encryption scheme: HMAC-SHA256 used as a PRF in CTR mode, **authenticated**
|
||||
encrypt-then-MAC (stdlib-only, no external deps). Each value is encrypted
|
||||
independently. The output blob is ``nonce (16 bytes) || ciphertext || tag
|
||||
(32 bytes)`` encoded as URL-safe base64 (no padding).
|
||||
|
||||
keystream_block_i = HMAC-SHA256(key, nonce || i.to_bytes(4, "big"))
|
||||
ciphertext_i = plaintext_i XOR keystream_block_i[:len(plaintext_i)]
|
||||
mac_key = HMAC-SHA256(key, "bottled-secret-mac-v1")
|
||||
tag = HMAC-SHA256(mac_key, nonce || ciphertext)
|
||||
|
||||
The tag is what makes a **wrong key deterministically detectable**: without it,
|
||||
CTR decryption with the wrong key yields garbage that only fails when it isn't
|
||||
valid UTF-8 (so ``reprovision`` would sometimes "succeed" with a wrong
|
||||
ENV_VAR_SECRET and inject garbage egress tokens). The MAC key is derived from
|
||||
the ENV_VAR_SECRET by a domain-separated HMAC so the same key never both
|
||||
generates the keystream and signs the tag with the same message shape.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -29,6 +39,7 @@ import secrets
|
||||
|
||||
_KEY_BYTES = 32 # 256-bit key from ENV_VAR_SECRET
|
||||
_NONCE_BYTES = 16 # 128-bit random nonce per encrypt call
|
||||
_TAG_BYTES = 32 # HMAC-SHA256 authentication tag
|
||||
_BLOCK = 32 # HMAC-SHA256 output width == one keystream block
|
||||
|
||||
# Env-var name the agent container receives at startup.
|
||||
@@ -50,45 +61,58 @@ def _keystream(key: bytes, nonce: bytes, block_index: int) -> bytes:
|
||||
).digest()
|
||||
|
||||
|
||||
def _tag(key: bytes, nonce: bytes, ciphertext: bytes) -> bytes:
|
||||
"""The authentication tag over ``nonce || ciphertext``, keyed by a MAC
|
||||
subkey domain-separated from the keystream key."""
|
||||
mac_key = hmac.new(key, b"bottled-secret-mac-v1", hashlib.sha256).digest()
|
||||
return hmac.new(mac_key, nonce + ciphertext, hashlib.sha256).digest()
|
||||
|
||||
|
||||
def _ctr(key: bytes, nonce: bytes, data: bytes) -> bytes:
|
||||
"""CTR keystream XOR — its own inverse, so it both encrypts and decrypts."""
|
||||
out = bytearray()
|
||||
for i in range(0, len(data), _BLOCK):
|
||||
chunk = data[i : i + _BLOCK]
|
||||
ks = _keystream(key, nonce, i)[: len(chunk)]
|
||||
out.extend(b ^ k for b, k in zip(chunk, ks))
|
||||
return bytes(out)
|
||||
|
||||
|
||||
def encrypt_value(secret_b64: str, plaintext: str) -> str:
|
||||
"""Encrypt a single string value with *secret_b64* (the ENV_VAR_SECRET).
|
||||
|
||||
Returns a URL-safe base64 blob ``nonce || ciphertext`` suitable for
|
||||
Returns a URL-safe base64 blob ``nonce || ciphertext || tag`` suitable for
|
||||
the ``bottled_agent_secrets.value`` column."""
|
||||
key = _b64dec(secret_b64)
|
||||
pt = plaintext.encode()
|
||||
nonce = secrets.token_bytes(_NONCE_BYTES)
|
||||
ct = bytearray()
|
||||
for i in range(0, len(pt), _BLOCK):
|
||||
chunk = pt[i : i + _BLOCK]
|
||||
ks = _keystream(key, nonce, i)[: len(chunk)]
|
||||
ct.extend(p ^ k for p, k in zip(chunk, ks))
|
||||
return base64.urlsafe_b64encode(nonce + bytes(ct)).rstrip(b"=").decode()
|
||||
ct = _ctr(key, nonce, plaintext.encode())
|
||||
tag = _tag(key, nonce, ct)
|
||||
return base64.urlsafe_b64encode(nonce + ct + tag).rstrip(b"=").decode()
|
||||
|
||||
|
||||
def decrypt_value(secret_b64: str, blob_b64: str) -> str:
|
||||
"""Decrypt a blob produced by :func:`encrypt_value`.
|
||||
|
||||
Returns the original plaintext string. Raises ``ValueError`` for malformed
|
||||
input or a key mismatch (wrong key produces garbage, not an error, unless
|
||||
the plaintext is non-UTF-8 — treat all such failures as wrong key)."""
|
||||
input, a **wrong key**, or a tampered ciphertext — all caught by the
|
||||
authentication tag before any plaintext is returned, so a wrong
|
||||
ENV_VAR_SECRET is rejected deterministically (never a garbage token)."""
|
||||
key = _b64dec(secret_b64)
|
||||
try:
|
||||
blob = _b64dec(blob_b64)
|
||||
except Exception as exc:
|
||||
raise ValueError(f"invalid ciphertext blob: {exc}") from exc
|
||||
if len(blob) < _NONCE_BYTES:
|
||||
if len(blob) < _NONCE_BYTES + _TAG_BYTES:
|
||||
raise ValueError("ciphertext blob too short")
|
||||
nonce, ciphertext = blob[:_NONCE_BYTES], blob[_NONCE_BYTES:]
|
||||
pt = bytearray()
|
||||
for i in range(0, len(ciphertext), _BLOCK):
|
||||
chunk = ciphertext[i : i + _BLOCK]
|
||||
ks = _keystream(key, nonce, i)[: len(chunk)]
|
||||
pt.extend(c ^ k for c, k in zip(chunk, ks))
|
||||
nonce = blob[:_NONCE_BYTES]
|
||||
tag = blob[-_TAG_BYTES:]
|
||||
ciphertext = blob[_NONCE_BYTES:-_TAG_BYTES]
|
||||
if not hmac.compare_digest(tag, _tag(key, nonce, ciphertext)):
|
||||
raise ValueError("ciphertext failed authentication (wrong key or tampered)")
|
||||
try:
|
||||
return bytes(pt).decode()
|
||||
except UnicodeDecodeError as exc:
|
||||
raise ValueError(f"decryption produced non-UTF-8 output (wrong key?): {exc}") from exc
|
||||
return _ctr(key, nonce, ciphertext).decode()
|
||||
except UnicodeDecodeError as exc: # pragma: no cover - authenticated, so unreachable
|
||||
raise ValueError(f"decryption produced non-UTF-8 output: {exc}") from exc
|
||||
|
||||
|
||||
__all__ = ["ENV_VAR_SECRET_NAME", "new_env_var_secret", "encrypt_value", "decrypt_value"]
|
||||
|
||||
@@ -36,6 +36,13 @@ ROLE_GATEWAY = "gateway"
|
||||
ROLE_CLI = "cli"
|
||||
ROLES: frozenset[str] = frozenset({ROLE_GATEWAY, ROLE_CLI})
|
||||
|
||||
# The host controller's own lifecycle role (#468). Deliberately OUTSIDE `ROLES`:
|
||||
# it belongs to a separate trust domain (`HOST_CONTROLLER`) signed by a key the
|
||||
# orchestrator never holds, so the orchestrator's control-plane key can neither
|
||||
# mint nor accept it — the orchestrator must not be able to forge the credential
|
||||
# used to start and stop it.
|
||||
ROLE_HOST = "host"
|
||||
|
||||
_ALG = "HS256"
|
||||
|
||||
|
||||
@@ -103,4 +110,4 @@ def verify(token: str, secret: str, *, roles: frozenset[str] = ROLES) -> str | N
|
||||
return role if isinstance(role, str) and role in roles else None
|
||||
|
||||
|
||||
__all__ = ["ROLE_GATEWAY", "ROLE_CLI", "ROLES", "mint", "verify"]
|
||||
__all__ = ["ROLE_GATEWAY", "ROLE_CLI", "ROLE_HOST", "ROLES", "mint", "verify"]
|
||||
|
||||
@@ -47,6 +47,22 @@ ORCHESTRATOR_TOKEN_ENV = "BOT_BOTTLE_ORCHESTRATOR_TOKEN"
|
||||
# cannot forge a higher-privilege `cli` token (issue #469 review).
|
||||
ORCHESTRATOR_AUTH_JWT_ENV = "BOT_BOTTLE_ORCHESTRATOR_AUTH_JWT"
|
||||
|
||||
# The durable launch-broker signing key: the HS256 secret the orchestrator
|
||||
# (signer) and the host control server (verifier) share to sign/verify launch
|
||||
# requests (#468). A host-canonical key file (minted 0600 on first use) so it
|
||||
# survives orchestrator restarts — re-adoption re-verifies against the same key —
|
||||
# instead of the ephemeral per-process secret of the in-process broker.
|
||||
LAUNCH_BROKER_KEY_FILENAME = "launch-broker-key"
|
||||
LAUNCH_BROKER_KEY_ENV = "BOT_BOTTLE_LAUNCH_BROKER_KEY"
|
||||
# The host controller's OWN key, for its lifecycle endpoints (the direct
|
||||
# cli -> host controller path that starts/stops the orchestrator). Separate from
|
||||
# the launch-broker key and never held by the orchestrator: the controller starts
|
||||
# and stops the orchestrator, so the orchestrator must not be able to mint the
|
||||
# credentials used to drive it (#468/#476).
|
||||
HOST_CONTROLLER_KEY_FILENAME = "host-controller-key"
|
||||
HOST_CONTROLLER_KEY_ENV = "BOT_BOTTLE_HOST_CONTROLLER_KEY"
|
||||
HOST_CONTROLLER_AUTH_JWT_ENV = "BOT_BOTTLE_HOST_CONTROLLER_AUTH_JWT"
|
||||
|
||||
# The host directory holding the gateway's persistent mitmproxy CA. Bind-mounted
|
||||
# into the infra/gateway container at mitmproxy's confdir so the self-generated
|
||||
# CA survives container recreation — every agent installs this one CA to trust
|
||||
@@ -142,6 +158,11 @@ __all__ = [
|
||||
"ORCHESTRATOR_TOKEN_FILENAME",
|
||||
"ORCHESTRATOR_TOKEN_ENV",
|
||||
"ORCHESTRATOR_AUTH_JWT_ENV",
|
||||
"LAUNCH_BROKER_KEY_FILENAME",
|
||||
"LAUNCH_BROKER_KEY_ENV",
|
||||
"HOST_CONTROLLER_KEY_FILENAME",
|
||||
"HOST_CONTROLLER_KEY_ENV",
|
||||
"HOST_CONTROLLER_AUTH_JWT_ENV",
|
||||
"GATEWAY_CA_DIRNAME",
|
||||
"bot_bottle_root",
|
||||
"host_db_path",
|
||||
|
||||
@@ -29,8 +29,13 @@ from collections.abc import Mapping
|
||||
from dataclasses import dataclass
|
||||
|
||||
from . import orchestrator_auth
|
||||
from .orchestrator_auth import ROLE_GATEWAY
|
||||
from .orchestrator_auth import ROLE_GATEWAY, ROLE_HOST
|
||||
from .paths import (
|
||||
HOST_CONTROLLER_AUTH_JWT_ENV,
|
||||
HOST_CONTROLLER_KEY_ENV,
|
||||
HOST_CONTROLLER_KEY_FILENAME,
|
||||
LAUNCH_BROKER_KEY_ENV,
|
||||
LAUNCH_BROKER_KEY_FILENAME,
|
||||
ORCHESTRATOR_AUTH_JWT_ENV,
|
||||
ORCHESTRATOR_TOKEN_ENV,
|
||||
ORCHESTRATOR_TOKEN_FILENAME,
|
||||
@@ -99,6 +104,40 @@ CONTROL_PLANE = TrustDomain(
|
||||
)
|
||||
|
||||
|
||||
# The launch-broker domain (#468): durable key material for the broker's own
|
||||
# signed launch requests (`broker.py`'s HS256 launch JWT), shared by the
|
||||
# orchestrator (signer) and the host control server (verifier). Unlike
|
||||
# `CONTROL_PLANE` it mints no role tokens — the broker's provenance is the launch
|
||||
# JWT, not a role token — so its `roles` set is empty and it is used only as a
|
||||
# provider of durable, host-canonical key material (`signing_key` / `key_from_env`).
|
||||
# The durability is the point: the key survives orchestrator restarts, so a
|
||||
# restarted orchestrator re-verifies against the same key instead of the
|
||||
# ephemeral per-process secret the in-process broker used.
|
||||
LAUNCH_BROKER = TrustDomain(
|
||||
name="launch-broker",
|
||||
key_filename=LAUNCH_BROKER_KEY_FILENAME,
|
||||
roles=frozenset(),
|
||||
key_env=LAUNCH_BROKER_KEY_ENV,
|
||||
token_env="",
|
||||
)
|
||||
|
||||
# The host controller's own domain (#468) — the SECOND domain #476 reserves. Its
|
||||
# key, which the orchestrator never holds, signs the `host`-role tokens the CLI
|
||||
# presents on the host controller's lifecycle endpoints (start / restart / status
|
||||
# of the orchestrator itself). Keeping it separate from `CONTROL_PLANE` is the
|
||||
# whole point: the host controller starts and stops the orchestrator, so the
|
||||
# orchestrator must not be able to mint the credentials used to drive it. (The
|
||||
# lifecycle endpoints themselves arrive in a later chunk; the domain is
|
||||
# established here alongside the durable launch-broker key.)
|
||||
HOST_CONTROLLER = TrustDomain(
|
||||
name="host-controller",
|
||||
key_filename=HOST_CONTROLLER_KEY_FILENAME,
|
||||
roles=frozenset({ROLE_HOST}),
|
||||
key_env=HOST_CONTROLLER_KEY_ENV,
|
||||
token_env=HOST_CONTROLLER_AUTH_JWT_ENV,
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ControlPlaneProvisioning:
|
||||
"""The one seam every backend launcher uses to provision control-plane auth,
|
||||
@@ -132,9 +171,56 @@ class ControlPlaneProvisioning:
|
||||
return self.domain.mint(ROLE_GATEWAY)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class LaunchBrokerProvisioning:
|
||||
"""The seam that provisions the host-side launch broker's durable keys (#468),
|
||||
the counterpart to `ControlPlaneProvisioning`. Both the orchestrator (signer)
|
||||
and the host control server (verifier) receive the SAME launch-broker key
|
||||
(carry it in `broker_domain.key_env`); the host controller ALSO receives its
|
||||
own lifecycle key (`controller_domain.key_env`) the orchestrator never holds.
|
||||
|
||||
Fail-closed like the control-plane seam: minting returns "" only if the host
|
||||
root is unwritable, and an empty launch-broker key would leave the verifier
|
||||
unable to authenticate any launch — so we raise rather than hand back a key
|
||||
that would make the host controller reject (or, if a caller defaulted it,
|
||||
accept) unsigned input."""
|
||||
|
||||
broker_domain: TrustDomain = LAUNCH_BROKER
|
||||
controller_domain: TrustDomain = HOST_CONTROLLER
|
||||
|
||||
def broker_key(self) -> str:
|
||||
"""The durable launch-broker key both the orchestrator and the host
|
||||
control server must receive (in `broker_domain.key_env`). Raises rather
|
||||
than return ""."""
|
||||
key = self.broker_domain.signing_key()
|
||||
if not key:
|
||||
raise ProvisioningError(
|
||||
f"refusing to provision the {self.broker_domain.name} broker "
|
||||
"without a signing key: the host controller could then verify no "
|
||||
"launch request's provenance"
|
||||
)
|
||||
return key
|
||||
|
||||
def controller_key(self) -> str:
|
||||
"""The host controller's own lifecycle key — provisioned ONLY to the host
|
||||
controller (in `controller_domain.key_env`), never to the orchestrator, so
|
||||
the orchestrator cannot mint the `host`-role tokens that start and stop
|
||||
it. Raises rather than return ""."""
|
||||
key = self.controller_domain.signing_key()
|
||||
if not key:
|
||||
raise ProvisioningError(
|
||||
f"refusing to provision the {self.controller_domain.name} without "
|
||||
"a signing key: its lifecycle endpoints would authenticate no one"
|
||||
)
|
||||
return key
|
||||
|
||||
|
||||
__all__ = [
|
||||
"ProvisioningError",
|
||||
"TrustDomain",
|
||||
"CONTROL_PLANE",
|
||||
"LAUNCH_BROKER",
|
||||
"HOST_CONTROLLER",
|
||||
"ControlPlaneProvisioning",
|
||||
"LaunchBrokerProvisioning",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,273 @@
|
||||
# PRD prd-new: Host control server
|
||||
|
||||
- **Status:** Draft
|
||||
- **Author:** Claude
|
||||
- **Created:** 2026-07-26
|
||||
- **Issue:** #468
|
||||
|
||||
## Summary
|
||||
|
||||
Promote the in-process launch broker into a standalone **host control
|
||||
server**: the single privileged component on the host. Both the CLI and the
|
||||
orchestrator drive it over HTTP; it brokers agent launches, owns the
|
||||
orchestrator's own lifecycle, and is the sole writer of host-durable state (the
|
||||
tamper-evident audit record). This closes the three gaps between today's
|
||||
well-formed broker *contract* ([`orchestrator/broker.py`](../../bot_bottle/orchestrator/broker.py))
|
||||
and a real out-of-process service — transport, durable provisioned secret,
|
||||
and a disciplined op vocabulary — and splits host state by
|
||||
owner and lifetime. The prize: **the CLI no longer needs the Docker socket**,
|
||||
which is what finally lets a dedicated Gitea runner user drop the
|
||||
root-equivalent `docker` group (PRD 0070, "Relationship to other work").
|
||||
|
||||
## Problem
|
||||
|
||||
Container launches run directly from a short-lived CLI process against the
|
||||
Docker socket. That socket is root-equivalent, so every host that launches
|
||||
bottles hands root to whoever invokes the CLI — including a CI runner user we
|
||||
want to keep unprivileged. PRD 0070 already argues for replacing the fat socket
|
||||
with a **thin, structured, auditable** launch broker, and the contract for that
|
||||
broker exists and is tested in-process. But it is *only* in-process:
|
||||
`LaunchBroker.submit(token)` is a method call from
|
||||
`OrchestratorCore.launch_bottle` ([`service.py:116`](../../bot_bottle/orchestrator/service.py)),
|
||||
and `DockerBroker` is on no production path — every backend starts the
|
||||
orchestrator with `--broker stub` ([`__main__.py:54`](../../bot_bottle/orchestrator/__main__.py)).
|
||||
|
||||
Three gaps stand between that scaffold and a host service:
|
||||
|
||||
1. **No transport.** `submit` is an in-process call. A real service needs a
|
||||
`BrokerClient` that POSTs the signed token and a host-side HTTP server that
|
||||
verifies and acts.
|
||||
2. **The signing secret is ephemeral and self-generated.**
|
||||
[`__main__.py:53`](../../bot_bottle/orchestrator/__main__.py) does
|
||||
`secrets.token_bytes(32)` and hands the *same value* to signer and verifier —
|
||||
viable only because they share a process. A separate daemon needs the secret
|
||||
provisioned out of band and durable across orchestrator restarts.
|
||||
3. **The op vocabulary is `launch` / `teardown` only.** Everything else
|
||||
host-privileged still lives in the CLI, so the schema has to grow — carefully,
|
||||
since PRD 0070's security argument rests on "structured requests only, static
|
||||
flags + ids."
|
||||
|
||||
Separately, host state has no clear owner. `OrchestratorCore.reconcile` takes
|
||||
`live_source_ips` as a parameter *only because the orchestrator cannot see the
|
||||
backend* ([`service.py:137`](../../bot_bottle/orchestrator/service.py)); the
|
||||
egress traffic log is written to the container's stderr; and there is no durable,
|
||||
tamper-evident home for the audit record that survives orchestrator destruction.
|
||||
|
||||
## Goals / Success Criteria
|
||||
|
||||
- A standalone host control server that the CLI and orchestrator reach over
|
||||
**HTTP**, with three entry paths working end to end:
|
||||
- `web console -(iroh)-> orchestrator -(http)-> host controller -> launch`
|
||||
- `cli -(http)-> orchestrator -(http)-> host controller -> launch`
|
||||
- `cli -(http)-> host controller` — start / restart / status of the
|
||||
orchestrator **itself** (the bootstrap/recovery path #391 targets).
|
||||
- The launch op is expressed as a **signed JWT of static flags + ids only**,
|
||||
verified against a closed schema.
|
||||
- The signing secret is **provisioned out of band and durable** across
|
||||
orchestrator restarts (a `TrustDomain` per #476, with a key the orchestrator
|
||||
never holds for the host controller's *own* endpoints).
|
||||
- Host-privileged operations move off the CLI to the control server; **the CLI
|
||||
no longer opens the Docker socket** for bottle operations.
|
||||
- `Orchestrator.reconcile` no longer takes `live_source_ips` — live-bottle
|
||||
enumeration becomes an internal control-server call.
|
||||
- Host-durable state lands as an **append-only, hash-chained JSONL** audit log
|
||||
owned solely by the host controller; operational state stays SQLite owned
|
||||
solely by the orchestrator.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- **Removing standing privilege.** This converts on-demand privilege (a CLI the
|
||||
user invokes) into standing privilege (a daemon under launchd/systemd). The
|
||||
win is that the privilege is *narrower* (structured requests vs. a raw socket),
|
||||
not that it disappears. "Always running" is an accepted new property.
|
||||
- **Asymmetric signing.** We stay HS256 — see Design / "Signing stays
|
||||
symmetric."
|
||||
- **Integrity against a live compromised orchestrator.** Host-location of the
|
||||
audit log does not buy this: the orchestrator makes the decisions being audited
|
||||
and can forge or omit entries wherever the file lives. An off-box copy is the
|
||||
answer, tracked separately.
|
||||
- **A single unified DB for all state.** Impossible over a guest-kernel share
|
||||
(SQLite locking is not coherent); state is split by owner and lifetime instead.
|
||||
- **The generic `SecretProvider` (#355)** and **remote terminal design (#478)** —
|
||||
both ride the same door but are their own work.
|
||||
|
||||
## Design
|
||||
|
||||
### Topology
|
||||
|
||||
The host controller is the sole privileged component. The orchestrator becomes a
|
||||
client of it for launches, and the CLI becomes a client of it for *both* bottle
|
||||
operations (indirectly, through the orchestrator) and orchestrator lifecycle
|
||||
(directly, for bootstrap/recovery — startup can't route through the thing being
|
||||
started).
|
||||
|
||||
```
|
||||
web console ─(iroh)─▶ orchestrator ─┐
|
||||
├─(http, signed JWT)─▶ host controller ─▶ launch
|
||||
cli ────────(http)──▶ orchestrator ─┘
|
||||
cli ────────(http, bearer)──────────────────────────────▶ host controller (orchestrator lifecycle)
|
||||
```
|
||||
|
||||
### Transport: `BrokerClient` + host server
|
||||
|
||||
`LaunchBroker.submit(token)` keeps its exact signature and semantics; only the
|
||||
*wire* changes. A new `BrokerClient` implements the same submit contract by
|
||||
POSTing the signed token to the host controller (stdlib `urllib`, like the
|
||||
existing [`orchestrator/client.py`](../../bot_bottle/orchestrator/client.py)),
|
||||
and the host controller's launch handler is the existing `verify_request` +
|
||||
`_launch`/`_teardown` path, now reached over HTTP instead of a method call. The
|
||||
in-process `StubBroker` stays for the dev-harness and tests; `DockerBroker`'s
|
||||
`_launch`/`_teardown` bodies move behind the server unchanged. Because the client
|
||||
satisfies the same interface `OrchestratorCore` already depends on, the core does
|
||||
not change to gain a real backend.
|
||||
|
||||
### Signing stays symmetric (HS256)
|
||||
|
||||
PRD 0070 nominally specifies asymmetric; the code is HS256 and we keep it.
|
||||
Asymmetric matters when the verifier is *less* privileged than the signer — here
|
||||
it is the reverse: the host controller (verifier) is strictly more privileged
|
||||
than the orchestrator (signer), and a controller that could forge orchestrator
|
||||
requests gains nothing, since it is already the component that launches. Staying
|
||||
symmetric also honors the no-runtime-deps policy (stdlib has no Ed25519). This
|
||||
matches the reasoning already inlined in `broker.py`'s module docstring.
|
||||
|
||||
### Replay protection is out of scope (tracked in #494)
|
||||
|
||||
Once the launch token travels over a wire, a captured token could be replayed —
|
||||
`sign_request` already emits `jti`/`iat` but `verify_request` reads neither, so
|
||||
there is no expiry window or `jti` cache today. Enforcing that (an `iat` window +
|
||||
a self-trimming `jti` cache) is a pure in-process change that lands independently
|
||||
of this work, and it is deferred to **#494** rather than gating the MVP of the
|
||||
host control server. Nothing here depends on it; it can merge before or after.
|
||||
|
||||
### Op vocabulary and the "ids + static flags" rule (gap 3)
|
||||
|
||||
Each op moved off the CLI widens the privileged surface, so growth is governed by
|
||||
one explicit rule, enforced in `verify_request`'s schema check:
|
||||
|
||||
> A broker op carries **only ids and enumerated static flags** — a bottle id, a
|
||||
> pool slot, a **content-addressed** image ref chosen from a fixed set, an op
|
||||
> name from a closed vocabulary. Never a free-form path, argv, command, or
|
||||
> caller-supplied filesystem location. If an operation cannot be expressed that
|
||||
> way, it does not become a broker op.
|
||||
|
||||
Operations that fit and move off the CLI (all today in
|
||||
`backend/*/consolidated_launch.py`, driven by a short-lived CLI process):
|
||||
|
||||
| Op | What it does | Fits the rule because |
|
||||
|---|---|---|
|
||||
| `launch` / `teardown` | existing | ids + slot + image ref |
|
||||
| `orchestrator.ensure_running` | start the infra container | no arguments |
|
||||
| `orchestrator.{start,restart,status}` | lifecycle (the #391 path) | no arguments |
|
||||
| `list_live` | enumerate running bottles for reconcile | no arguments; returns ids/IPs |
|
||||
| `allocate_ip` | `next_free_ip` over `_network_container_ips` | no arguments; returns an IP |
|
||||
| `provision_git_gate` | `cp`/`exec` a per-bottle deploy key into the gateway | bottle id + key handle, no path |
|
||||
| `reprovision` | `docker exec printenv <ENV_VAR_SECRET>` on a live agent | bottle id + secret *name* |
|
||||
|
||||
Image **builds** stay with the orchestrator for v1 (PRD 0070 §Memory: builds run
|
||||
control-plane-side; a dedicated slim build unit is later, #468-adjacent), so no
|
||||
`build` broker op is added here.
|
||||
|
||||
With `list_live` as an internal control-server call, `Orchestrator.reconcile`'s
|
||||
`live_source_ips` parameter goes away — the tell PRD 0070 called out that the
|
||||
orchestrator couldn't see the backend disappears with it.
|
||||
|
||||
### Secret provisioning (gap 2)
|
||||
|
||||
The shared HS256 secret becomes a durable, out-of-band artifact via the
|
||||
**`TrustDomain`** seam (#476,
|
||||
[`trust_domain.py`](../../bot_bottle/trust_domain.py)):
|
||||
|
||||
- The **launch-broker secret** is a `TrustDomain` whose key
|
||||
(`host_signing_key(<file>)`, minted 0600 on first use, durable under
|
||||
`bot_bottle_root()`) is provisioned to the orchestrator (signer) and the host
|
||||
controller (verifier). Durability across orchestrator restarts is what makes
|
||||
re-adoption work — a restart re-verifies against the same key.
|
||||
- The **host controller's own lifecycle endpoints** (the direct `cli -> host
|
||||
controller` path) get a **separate** `TrustDomain` key the orchestrator never
|
||||
holds — exactly the second domain #476's PRD reserves. The orchestrator must
|
||||
not be able to mint the credentials used to start and stop it.
|
||||
|
||||
This reuses the seam #476 landed rather than re-deriving provisioning per
|
||||
backend (the PR #471 bug class).
|
||||
|
||||
### One daemon, structurally separate handlers (open decision 1)
|
||||
|
||||
The audit writer and the broker live in **one daemon** for install simplicity,
|
||||
but with **no shared parsing** and **different credentials per handler**:
|
||||
|
||||
- the **launch** handler requires the signed launch **JWT** (provenance +
|
||||
un-coercible schema);
|
||||
- the **audit-append** handler takes a plain **bearer token** and writes to the
|
||||
JSONL log.
|
||||
|
||||
This does not defend against orchestrator compromise (it holds both creds) — it
|
||||
stops a bug in the boring audit path from reaching the privileged launch path.
|
||||
The launcher stays small enough to audit line-by-line, per PRD 0070.
|
||||
|
||||
### State ownership: split by owner and lifetime
|
||||
|
||||
A single mounted DB is impossible — SQLite locking is not coherent across guest
|
||||
kernels over a share, which is why the macOS backend already uses a container-only
|
||||
volume (`INFRA_DB_VOLUME`). So state splits three ways (depends on #469, which
|
||||
gets `bot-bottle.db` off the data plane first):
|
||||
|
||||
| Owner | State | Home | Shape |
|
||||
|---|---|---|---|
|
||||
| **Orchestrator** | `orchestrator_bottles` registry; `bottled_agent_secrets` (encrypted egress tokens); `supervise_proposals` / `supervise_responses` | volume nothing else mounts (generalizing the macOS design) | **SQLite** — mutable, transactional, queried |
|
||||
| **Host controller** | supervise audit entries; egress traffic log (today → container stderr); host-side config | host filesystem, survives orchestrator/volume destruction | **JSONL** — append-only |
|
||||
| **Gateway** | none | — | after #469 the data plane holds no DB state |
|
||||
|
||||
The historical record is **JSONL, not SQLite**, because it is append-only, never
|
||||
updated, never transactionally queried: `O_APPEND` writes are atomic, there is no
|
||||
locking protocol to get wrong, hash-chaining for tamper-evidence is cheap, and it
|
||||
survives container-runtime volume pruning (the #450 lesson) and stays readable
|
||||
without the orchestrator running. Both halves of "the audit record" — supervise
|
||||
decisions and the egress traffic log — land in the one place.
|
||||
|
||||
The orchestrator is **sole mounter and sole writer** of its SQLite volume; the
|
||||
host controller is **sole writer** of the JSONL log, over the authenticated
|
||||
audit-append channel.
|
||||
|
||||
## Implementation chunks
|
||||
|
||||
Ordered, each independently mergeable:
|
||||
|
||||
1. **`BrokerClient` + host launch server** over HTTP, reusing `verify_request`
|
||||
and the existing `DockerBroker` bodies. Wire `OrchestratorCore` to a
|
||||
`BrokerClient` behind a flag; keep `StubBroker` for the dev-harness. Closes
|
||||
gap 1.
|
||||
2. **Durable secret via `TrustDomain`** — provision the launch-broker key to
|
||||
signer + verifier; add the host controller's own lifecycle `TrustDomain`.
|
||||
Closes gap 2.
|
||||
3. **Grow the op vocabulary** one op at a time (`list_live` first — it also
|
||||
removes `reconcile`'s `live_source_ips`), each behind the ids + static-flags
|
||||
rule. Closes gap 3.
|
||||
4. **JSONL audit log** — the host-controller-owned, hash-chained historical
|
||||
record with the plain-bearer audit-append handler; redirect the egress traffic
|
||||
log into it.
|
||||
5. **Drop the Docker socket from the CLI** once every host-privileged op it used
|
||||
is a broker op — the payoff that unblocks the unprivileged Gitea runner user.
|
||||
|
||||
## Open questions
|
||||
|
||||
1. **Schema-width rule enforcement.** The "ids + static flags" rule is stated;
|
||||
should `verify_request` reject unknown claim keys outright (strict schema) to
|
||||
keep the surface from drifting? Leaning yes.
|
||||
2. **Audit-append back-pressure.** What the audit handler does if the JSONL sink
|
||||
is unavailable (fail-closed vs. buffer) — resolve before shipping chunk 5.
|
||||
|
||||
## References
|
||||
|
||||
- **PRD 0070** — the contract, the launch broker, and the state tiers this
|
||||
implements.
|
||||
- **#469** — get `bot-bottle.db` off the data plane (lands underneath this).
|
||||
- **#476** ([`prd-new-control-plane-auth-provisioning`](prd-new-control-plane-auth-provisioning.md))
|
||||
— the `TrustDomain` seam this plugs the host controller's key into.
|
||||
- **#391** — backend-agnostic orchestrator restart (the bootstrap path).
|
||||
- **#494** — enforce broker replay protection (`iat` window + `jti` cache); split
|
||||
out of this PRD as an independent in-process change.
|
||||
- **#386** — prebuilt images from the Gitea OCI registry (the fixed image set the
|
||||
broker validates against).
|
||||
- **#355** — generic `SecretProvider`.
|
||||
- **#478** — remote terminal design.
|
||||
@@ -190,8 +190,7 @@ class TestRegisterAgentReconciles(unittest.TestCase):
|
||||
|
||||
def _register(self, client: Mock) -> None:
|
||||
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
|
||||
patch(f"{_UTIL}.provision_git_gate"), \
|
||||
patch(f"{_MOD}.live_source_ips", return_value=["10.0.0.7"]):
|
||||
patch(f"{_UTIL}.provision_git_gate"):
|
||||
register_agent(
|
||||
_egress_plan(), _git_plan(),
|
||||
source_ip="10.0.0.7", endpoint=_endpoint(),
|
||||
@@ -213,7 +212,9 @@ class TestRegisterAgentReconciles(unittest.TestCase):
|
||||
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"])
|
||||
# A bare trigger now — the orchestrator enumerates the live set itself
|
||||
# (from the host controller), so the CLI passes no IPs.
|
||||
client.reconcile.assert_called_once_with()
|
||||
|
||||
def test_a_reconcile_failure_does_not_block_the_launch(self) -> None:
|
||||
from bot_bottle.orchestrator.client import OrchestratorClientError
|
||||
@@ -221,18 +222,3 @@ class TestRegisterAgentReconciles(unittest.TestCase):
|
||||
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"{_UTIL}.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()
|
||||
|
||||
@@ -2,6 +2,10 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import secrets
|
||||
import unittest
|
||||
|
||||
@@ -14,6 +18,21 @@ from bot_bottle.orchestrator.broker import (
|
||||
)
|
||||
|
||||
|
||||
def _sign_claims(claims: dict[str, object], secret: bytes) -> str:
|
||||
"""Sign an arbitrary claims dict (bypassing `sign_request`'s fixed fields)
|
||||
so a test can craft off-schema payloads — e.g. an extra claim key, or a
|
||||
query op carrying ids — with a *valid* signature, isolating the schema check
|
||||
from the provenance check."""
|
||||
def b64(data: bytes) -> str:
|
||||
return base64.urlsafe_b64encode(data).rstrip(b"=").decode("ascii")
|
||||
header = b64(json.dumps({"alg": "HS256", "typ": "JWT"}, sort_keys=True,
|
||||
separators=(",", ":")).encode())
|
||||
payload = b64(json.dumps(claims, sort_keys=True, separators=(",", ":")).encode())
|
||||
signing_input = f"{header}.{payload}"
|
||||
sig = b64(hmac.new(secret, signing_input.encode("ascii"), hashlib.sha256).digest())
|
||||
return f"{signing_input}.{sig}"
|
||||
|
||||
|
||||
class TestSignVerify(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.secret = secrets.token_bytes(16)
|
||||
@@ -54,6 +73,41 @@ class TestSignVerify(unittest.TestCase):
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
verify_request(token, self.secret)
|
||||
|
||||
def test_list_live_round_trip(self) -> None:
|
||||
req = LaunchRequest(op="list_live")
|
||||
got = verify_request(sign_request(req, self.secret), self.secret)
|
||||
self.assertEqual("list_live", got.op)
|
||||
self.assertEqual("", got.bottle_id)
|
||||
|
||||
def test_mutation_without_bottle_id_rejected(self) -> None:
|
||||
# A launch/teardown must name its bottle even with a valid signature.
|
||||
for op in ("launch", "teardown"):
|
||||
token = _sign_claims(
|
||||
{"op": op, "bottle_id": "", "jti": "x", "iat": 1}, self.secret)
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
verify_request(token, self.secret)
|
||||
|
||||
def test_query_op_carrying_ids_or_flags_rejected(self) -> None:
|
||||
# list_live is "no arguments" — a signed one that smuggles a bottle_id /
|
||||
# source_ip / image_ref / slot is off-schema and refused.
|
||||
for extra in (
|
||||
{"bottle_id": "b1"}, {"source_ip": "10.0.0.1"},
|
||||
{"image_ref": "img"}, {"slot": 2},
|
||||
):
|
||||
claims: dict[str, object] = {"op": "list_live", "jti": "x", "iat": 1, **extra}
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
verify_request(_sign_claims(claims, self.secret), self.secret)
|
||||
|
||||
def test_unknown_claim_key_rejected(self) -> None:
|
||||
# Strict schema (open question 1): a well-signed token with an extra claim
|
||||
# key is refused, so the surface can't widen past the fixed fields.
|
||||
token = _sign_claims(
|
||||
{"op": "launch", "bottle_id": "b1", "cmd": "rm -rf /", "jti": "x", "iat": 1},
|
||||
self.secret,
|
||||
)
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
verify_request(token, self.secret)
|
||||
|
||||
|
||||
class TestStubBroker(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
@@ -81,6 +135,37 @@ class TestStubBroker(unittest.TestCase):
|
||||
self.broker.submit(forged)
|
||||
self.assertEqual([], self.broker.launched) # nothing acted on
|
||||
|
||||
def test_list_live_derives_from_launches(self) -> None:
|
||||
# Default live set: launched, minus torn-down, by source IP.
|
||||
self.broker.submit(
|
||||
sign_request(LaunchRequest(op="launch", bottle_id="b1", source_ip="10.0.0.1"),
|
||||
self.secret))
|
||||
self.broker.submit(
|
||||
sign_request(LaunchRequest(op="launch", bottle_id="b2", source_ip="10.0.0.2"),
|
||||
self.secret))
|
||||
self.broker.submit(
|
||||
sign_request(LaunchRequest(op="teardown", bottle_id="b1", source_ip="10.0.0.1"),
|
||||
self.secret))
|
||||
live = self.broker.list_live(
|
||||
sign_request(LaunchRequest(op="list_live"), self.secret))
|
||||
self.assertEqual(["10.0.0.2"], live)
|
||||
|
||||
def test_list_live_override(self) -> None:
|
||||
self.broker.live_source_ips = ["10.0.0.9"]
|
||||
live = self.broker.list_live(
|
||||
sign_request(LaunchRequest(op="list_live"), self.secret))
|
||||
self.assertEqual(["10.0.0.9"], live)
|
||||
|
||||
def test_submit_rejects_a_list_live_token(self) -> None:
|
||||
# The verbs don't cross: a query token routed to submit is fail-closed.
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
self.broker.submit(sign_request(LaunchRequest(op="list_live"), self.secret))
|
||||
|
||||
def test_list_live_rejects_a_mutation_token(self) -> None:
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
self.broker.list_live(
|
||||
sign_request(LaunchRequest(op="launch", bottle_id="b1"), self.secret))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
"""Unit: orchestrator-side broker client (issue #468, chunk 1). HTTP mocked."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import io
|
||||
import json
|
||||
import unittest
|
||||
import urllib.error
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from bot_bottle.orchestrator.broker import (
|
||||
BrokerAuthError,
|
||||
BrokerUnavailableError,
|
||||
LaunchRequest,
|
||||
)
|
||||
from bot_bottle.orchestrator.broker_client import BrokerClient, BrokerClientError
|
||||
|
||||
_URLOPEN = "bot_bottle.orchestrator.broker_client.urllib.request.urlopen"
|
||||
|
||||
|
||||
def _resp(payload: object) -> MagicMock:
|
||||
m = MagicMock()
|
||||
m.__enter__.return_value.read.return_value = json.dumps(payload).encode()
|
||||
return m
|
||||
|
||||
|
||||
def _http_error(code: int, payload: object = None) -> urllib.error.HTTPError:
|
||||
body = json.dumps(payload).encode() if payload is not None else b""
|
||||
return urllib.error.HTTPError(
|
||||
"http://host/broker", code, "err", {}, io.BytesIO(body)) # type: ignore[arg-type]
|
||||
|
||||
|
||||
class TestSubmit(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.c = BrokerClient("http://host:8091")
|
||||
|
||||
def test_returns_the_verified_request(self) -> None:
|
||||
echo = {
|
||||
"op": "launch", "bottle_id": "b1", "source_ip": "10.0.0.1",
|
||||
"image_ref": "img", "slot": 3,
|
||||
}
|
||||
with patch(_URLOPEN, return_value=_resp(echo)):
|
||||
got = self.c.submit("tok")
|
||||
self.assertEqual(
|
||||
LaunchRequest(op="launch", bottle_id="b1", source_ip="10.0.0.1",
|
||||
image_ref="img", slot=3),
|
||||
got,
|
||||
)
|
||||
|
||||
def test_posts_token_to_broker_endpoint(self) -> None:
|
||||
with patch(_URLOPEN, return_value=_resp({"op": "teardown", "bottle_id": "b1"})) as m:
|
||||
self.c.submit("signed-token")
|
||||
request = m.call_args.args[0]
|
||||
self.assertEqual("POST", request.get_method())
|
||||
self.assertTrue(request.full_url.endswith("/broker"))
|
||||
self.assertEqual({"token": "signed-token"}, json.loads(request.data))
|
||||
|
||||
def test_401_raises_broker_auth_error(self) -> None:
|
||||
# A fail-closed provenance/schema rejection surfaces as the SAME exception
|
||||
# the in-process broker raises, so the launch path's rollback is identical.
|
||||
with patch(_URLOPEN, side_effect=_http_error(401, {"error": "bad signature"})):
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
self.c.submit("forged")
|
||||
|
||||
def test_502_is_a_definite_client_error(self) -> None:
|
||||
# The host responded — it processed the request and did not launch, so a
|
||||
# definite BrokerClientError (the caller may safely roll back).
|
||||
with patch(_URLOPEN, side_effect=_http_error(502, {"error": "docker down"})):
|
||||
with self.assertRaises(BrokerClientError):
|
||||
self.c.submit("tok")
|
||||
|
||||
def test_unreachable_is_ambiguous_unavailable(self) -> None:
|
||||
# No response at all — the request may already have launched, so the
|
||||
# AMBIGUOUS BrokerUnavailableError (the caller must NOT roll back).
|
||||
with patch(_URLOPEN, side_effect=urllib.error.URLError("refused")):
|
||||
with self.assertRaises(BrokerUnavailableError):
|
||||
self.c.submit("tok")
|
||||
|
||||
def test_timeout_is_ambiguous_unavailable(self) -> None:
|
||||
# A dropped/late response after the request was sent is the exact orphan
|
||||
# risk: the host may have launched. Must be ambiguous, not a definite fail.
|
||||
with patch(_URLOPEN, side_effect=TimeoutError("read timed out")):
|
||||
with self.assertRaises(BrokerUnavailableError):
|
||||
self.c.submit("tok")
|
||||
|
||||
def test_malformed_success_body_raises(self) -> None:
|
||||
with patch(_URLOPEN, return_value=_resp({"op": "launch"})): # missing bottle_id
|
||||
with self.assertRaises(BrokerClientError):
|
||||
self.c.submit("tok")
|
||||
|
||||
def test_empty_error_body_is_tolerated(self) -> None:
|
||||
# An error with no readable JSON body still classifies by status code.
|
||||
with patch(_URLOPEN, side_effect=_http_error(401)):
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
self.c.submit("forged")
|
||||
|
||||
def test_non_json_success_body_raises(self) -> None:
|
||||
# A 200 whose body isn't JSON is tolerated into {} then fails the
|
||||
# missing-field check — a definite client error, not a crash.
|
||||
m = MagicMock()
|
||||
m.__enter__.return_value.read.return_value = b"not json at all"
|
||||
with patch(_URLOPEN, return_value=m):
|
||||
with self.assertRaises(BrokerClientError):
|
||||
self.c.submit("tok")
|
||||
|
||||
def test_unreadable_error_body_is_tolerated(self) -> None:
|
||||
# An HTTPError whose body can't be read (fp=None) still classifies by
|
||||
# status — the error detail is best-effort.
|
||||
err = urllib.error.HTTPError(
|
||||
"http://host/broker", 502, "err", {}, None) # type: ignore[arg-type]
|
||||
with patch(_URLOPEN, side_effect=err):
|
||||
with self.assertRaises(BrokerClientError):
|
||||
self.c.submit("tok")
|
||||
|
||||
|
||||
class TestListLive(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.c = BrokerClient("http://host:8091")
|
||||
|
||||
def test_returns_source_ips(self) -> None:
|
||||
with patch(_URLOPEN, return_value=_resp({"source_ips": ["10.0.0.1", "10.0.0.2"]})):
|
||||
self.assertEqual(["10.0.0.1", "10.0.0.2"], self.c.list_live("tok"))
|
||||
|
||||
def test_posts_token_to_live_endpoint(self) -> None:
|
||||
with patch(_URLOPEN, return_value=_resp({"source_ips": []})) as m:
|
||||
self.c.list_live("signed-token")
|
||||
request = m.call_args.args[0]
|
||||
self.assertEqual("POST", request.get_method())
|
||||
self.assertTrue(request.full_url.endswith("/broker/live"))
|
||||
self.assertEqual({"token": "signed-token"}, json.loads(request.data))
|
||||
|
||||
def test_non_string_entries_are_filtered(self) -> None:
|
||||
with patch(_URLOPEN, return_value=_resp({"source_ips": ["10.0.0.1", 5, None, ""]})):
|
||||
self.assertEqual(["10.0.0.1"], self.c.list_live("tok"))
|
||||
|
||||
def test_missing_source_ips_raises(self) -> None:
|
||||
# A response without the field is malformed — better to fail (reconcile
|
||||
# skips) than to reconcile against a silent empty set.
|
||||
with patch(_URLOPEN, return_value=_resp({})):
|
||||
with self.assertRaises(BrokerClientError):
|
||||
self.c.list_live("tok")
|
||||
|
||||
def test_401_raises_broker_auth_error(self) -> None:
|
||||
with patch(_URLOPEN, side_effect=_http_error(401, {"error": "bad signature"})):
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
self.c.list_live("forged")
|
||||
|
||||
def test_502_enumeration_failure_is_client_error(self) -> None:
|
||||
with patch(_URLOPEN, side_effect=_http_error(502, {"error": "docker down"})):
|
||||
with self.assertRaises(BrokerClientError):
|
||||
self.c.list_live("tok")
|
||||
|
||||
def test_unreachable_is_unavailable(self) -> None:
|
||||
with patch(_URLOPEN, side_effect=urllib.error.URLError("refused")):
|
||||
with self.assertRaises(BrokerUnavailableError):
|
||||
self.c.list_live("tok")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -159,26 +159,27 @@ class TestReconcile(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.c = OrchestratorClient("http://orch:8080")
|
||||
|
||||
def test_posts_live_ips_and_returns_reaped(self) -> None:
|
||||
def test_triggers_sweep_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"])
|
||||
got = self.c.reconcile()
|
||||
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"])
|
||||
# The live set is no longer sent — the orchestrator enumerates it.
|
||||
self.assertNotIn("live_source_ips", sent)
|
||||
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.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([]))
|
||||
self.assertEqual(["ok"], self.c.reconcile())
|
||||
with patch(_URLOPEN, return_value=_resp(200, {})):
|
||||
self.assertEqual([], self.c.reconcile([]))
|
||||
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([])
|
||||
self.c.reconcile()
|
||||
|
||||
@@ -8,6 +8,7 @@ from unittest.mock import Mock, patch
|
||||
|
||||
from bot_bottle.orchestrator.broker import (
|
||||
BrokerAuthError,
|
||||
BrokerUnavailableError,
|
||||
LaunchRequest,
|
||||
sign_request,
|
||||
)
|
||||
@@ -96,5 +97,50 @@ class TestDockerBroker(unittest.TestCase):
|
||||
m.assert_not_called()
|
||||
|
||||
|
||||
class TestDockerBrokerListLive(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.secret = secrets.token_bytes(16)
|
||||
self.broker = DockerBroker(self.secret)
|
||||
|
||||
def test_enumerates_labeled_container_ips(self) -> None:
|
||||
# `docker ps` lists two labeled containers; each inspect yields an IP.
|
||||
ps = Mock(returncode=0, stdout="c1\nc2\n", stderr="")
|
||||
i1 = Mock(returncode=0, stdout="10.0.0.1 \n", stderr="")
|
||||
i2 = Mock(returncode=0, stdout="10.0.0.2 \n", stderr="")
|
||||
with patch.object(self.broker, "_docker", side_effect=[ps, i1, i2]) as m:
|
||||
got = self.broker._list_live()
|
||||
self.assertEqual(["10.0.0.1", "10.0.0.2"], got)
|
||||
# First call is the label-filtered `docker ps`.
|
||||
ps_argv = m.call_args_list[0].args[0]
|
||||
self.assertEqual(["docker", "ps", "--filter", f"label={BOTTLE_ID_LABEL}"], ps_argv[:4])
|
||||
|
||||
def test_no_containers_is_empty(self) -> None:
|
||||
with patch.object(self.broker, "_docker",
|
||||
return_value=Mock(returncode=0, stdout="", stderr="")):
|
||||
self.assertEqual([], self.broker._list_live())
|
||||
|
||||
def test_ps_failure_raises(self) -> None:
|
||||
with patch.object(self.broker, "_docker",
|
||||
return_value=Mock(returncode=1, stdout="", stderr="daemon down")):
|
||||
with self.assertRaises(DockerBrokerError):
|
||||
self.broker._list_live()
|
||||
|
||||
def test_inspect_failure_raises(self) -> None:
|
||||
ps = Mock(returncode=0, stdout="c1\n", stderr="")
|
||||
bad = Mock(returncode=1, stdout="", stderr="no such container")
|
||||
with patch.object(self.broker, "_docker", side_effect=[ps, bad]):
|
||||
with self.assertRaises(DockerBrokerError):
|
||||
self.broker._list_live()
|
||||
|
||||
def test_enumeration_failure_surfaces_as_unavailable_via_list_live(self) -> None:
|
||||
# Through the verb, a backend failure is the single "live set unknown"
|
||||
# signal (BrokerUnavailableError) reconcile catches to skip the sweep.
|
||||
token = sign_request(LaunchRequest(op="list_live"), self.secret)
|
||||
with patch.object(self.broker, "_docker",
|
||||
return_value=Mock(returncode=1, stdout="", stderr="down")):
|
||||
with self.assertRaises(BrokerUnavailableError):
|
||||
self.broker.list_live(token)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -0,0 +1,354 @@
|
||||
"""Unit tests for the host control server (issue #468, chunk 1).
|
||||
|
||||
Mostly exercises the pure `dispatch()` (socket-free, like the orchestrator
|
||||
server tests), plus a real-socket round-trip through `BrokerClient` that proves
|
||||
the full sign -> POST -> verify -> act seam over HTTP.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import http.client
|
||||
import io
|
||||
import json
|
||||
import os
|
||||
import secrets
|
||||
import tempfile
|
||||
import threading
|
||||
import typing
|
||||
import unittest
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from bot_bottle.orchestrator.broker import (
|
||||
BrokerAuthError,
|
||||
LaunchBroker,
|
||||
LaunchRequest,
|
||||
StubBroker,
|
||||
sign_request,
|
||||
)
|
||||
from bot_bottle.orchestrator.broker_client import BrokerClient
|
||||
from bot_bottle.orchestrator.host_server import (
|
||||
MAX_BODY_BYTES,
|
||||
Handler,
|
||||
HostControlServer,
|
||||
broker_secret,
|
||||
dispatch,
|
||||
main,
|
||||
make_host_server,
|
||||
)
|
||||
from bot_bottle.paths import LAUNCH_BROKER_KEY_ENV
|
||||
|
||||
|
||||
def _body(obj: object) -> bytes:
|
||||
return json.dumps(obj).encode()
|
||||
|
||||
|
||||
class _RaisingBroker(LaunchBroker):
|
||||
"""A broker whose backend launch/enumeration always fails — exercises the
|
||||
502 path (an operational backend failure, distinct from a fail-closed
|
||||
provenance 401)."""
|
||||
|
||||
def _launch(self, req: LaunchRequest) -> None:
|
||||
raise RuntimeError("docker down")
|
||||
|
||||
def _teardown(self, req: LaunchRequest) -> None:
|
||||
raise RuntimeError("docker down")
|
||||
|
||||
def _list_live(self) -> list[str]:
|
||||
raise RuntimeError("docker down")
|
||||
|
||||
|
||||
class TestDispatch(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.secret = secrets.token_bytes(16)
|
||||
self.broker = StubBroker(self.secret)
|
||||
|
||||
def _token(self, **kwargs: object) -> str:
|
||||
return sign_request(LaunchRequest(**kwargs), self.secret) # type: ignore[arg-type]
|
||||
|
||||
def test_health(self) -> None:
|
||||
status, payload = dispatch(self.broker, "GET", "/health", b"")
|
||||
self.assertEqual(200, status)
|
||||
self.assertEqual("ok", payload["status"])
|
||||
|
||||
def test_broker_launch_verifies_and_acts(self) -> None:
|
||||
token = self._token(
|
||||
op="launch", bottle_id="b1", source_ip="10.243.0.1",
|
||||
image_ref="img", slot=2,
|
||||
)
|
||||
status, payload = dispatch(self.broker, "POST", "/broker", _body({"token": token}))
|
||||
self.assertEqual(200, status)
|
||||
self.assertEqual("launch", payload["op"])
|
||||
self.assertEqual("b1", payload["bottle_id"])
|
||||
self.assertEqual("img", payload["image_ref"])
|
||||
self.assertEqual(2, payload["slot"])
|
||||
self.assertEqual(["b1"], [r.bottle_id for r in self.broker.launched])
|
||||
|
||||
def test_broker_teardown_acts(self) -> None:
|
||||
token = self._token(op="teardown", bottle_id="b1")
|
||||
status, _ = dispatch(self.broker, "POST", "/broker", _body({"token": token}))
|
||||
self.assertEqual(200, status)
|
||||
self.assertEqual(["b1"], [r.bottle_id for r in self.broker.torn_down])
|
||||
|
||||
def test_forged_token_is_401_and_nothing_acted(self) -> None:
|
||||
forged = sign_request(
|
||||
LaunchRequest(op="launch", bottle_id="b1"), secrets.token_bytes(16))
|
||||
status, payload = dispatch(self.broker, "POST", "/broker", _body({"token": forged}))
|
||||
self.assertEqual(401, status)
|
||||
self.assertIn("broker auth failed", str(payload["error"]))
|
||||
self.assertEqual([], self.broker.launched) # fail-closed: never launched
|
||||
|
||||
def test_backend_failure_is_502(self) -> None:
|
||||
broker = _RaisingBroker(self.secret)
|
||||
token = self._token(op="launch", bottle_id="b1", image_ref="img")
|
||||
status, payload = dispatch(broker, "POST", "/broker", _body({"token": token}))
|
||||
self.assertEqual(502, status)
|
||||
self.assertIn("backend launch failed", str(payload["error"]))
|
||||
|
||||
def test_broker_live_returns_source_ips(self) -> None:
|
||||
# Two launched bottles -> the stub reports both as live.
|
||||
self.broker.submit(self._token(op="launch", bottle_id="b1", source_ip="10.0.0.1"))
|
||||
self.broker.submit(self._token(op="launch", bottle_id="b2", source_ip="10.0.0.2"))
|
||||
token = self._token(op="list_live")
|
||||
status, payload = dispatch(self.broker, "POST", "/broker/live", _body({"token": token}))
|
||||
self.assertEqual(200, status)
|
||||
self.assertEqual(
|
||||
["10.0.0.1", "10.0.0.2"],
|
||||
sorted(typing.cast("list[str]", payload["source_ips"])),
|
||||
)
|
||||
|
||||
def test_broker_live_forged_token_is_401(self) -> None:
|
||||
forged = sign_request(LaunchRequest(op="list_live"), secrets.token_bytes(16))
|
||||
status, payload = dispatch(self.broker, "POST", "/broker/live", _body({"token": forged}))
|
||||
self.assertEqual(401, status)
|
||||
self.assertIn("broker auth failed", str(payload["error"]))
|
||||
|
||||
def test_broker_live_rejects_a_mutation_token(self) -> None:
|
||||
# A launch token routed to the query endpoint is a fail-closed 401 — the
|
||||
# endpoints don't share a schema even though they share a secret.
|
||||
token = self._token(op="launch", bottle_id="b1", image_ref="img")
|
||||
status, _ = dispatch(self.broker, "POST", "/broker/live", _body({"token": token}))
|
||||
self.assertEqual(401, status)
|
||||
|
||||
def test_broker_live_backend_failure_is_502(self) -> None:
|
||||
broker = _RaisingBroker(self.secret)
|
||||
token = self._token(op="list_live")
|
||||
status, payload = dispatch(broker, "POST", "/broker/live", _body({"token": token}))
|
||||
self.assertEqual(502, status)
|
||||
self.assertIn("backend enumeration failed", str(payload["error"]))
|
||||
|
||||
def test_broker_live_missing_token_is_400(self) -> None:
|
||||
status, _ = dispatch(self.broker, "POST", "/broker/live", _body({}))
|
||||
self.assertEqual(400, status)
|
||||
|
||||
def test_missing_token_is_400(self) -> None:
|
||||
status, _ = dispatch(self.broker, "POST", "/broker", _body({}))
|
||||
self.assertEqual(400, status)
|
||||
|
||||
def test_bad_json_is_400(self) -> None:
|
||||
status, _ = dispatch(self.broker, "POST", "/broker", b"{not json")
|
||||
self.assertEqual(400, status)
|
||||
|
||||
def test_empty_body_is_missing_token_400(self) -> None:
|
||||
# Empty body parses to {} (no token) → 400, never reaching the broker.
|
||||
status, _ = dispatch(self.broker, "POST", "/broker", b"")
|
||||
self.assertEqual(400, status)
|
||||
self.assertEqual([], self.broker.launched)
|
||||
|
||||
def test_non_object_body_is_400(self) -> None:
|
||||
status, _ = dispatch(self.broker, "POST", "/broker", b"[1, 2]")
|
||||
self.assertEqual(400, status)
|
||||
|
||||
def test_unknown_route_404(self) -> None:
|
||||
status, _ = dispatch(self.broker, "GET", "/nope", b"")
|
||||
self.assertEqual(404, status)
|
||||
|
||||
def test_trailing_slash_normalized(self) -> None:
|
||||
status, _ = dispatch(self.broker, "GET", "/health/", b"")
|
||||
self.assertEqual(200, status)
|
||||
|
||||
|
||||
class TestBrokerSecret(unittest.TestCase):
|
||||
"""The durable launch-broker key (#468/#476): prefer the env-injected key,
|
||||
else the durable host key file, so signer and verifier resolve the same one."""
|
||||
|
||||
def test_reads_injected_key_from_env(self) -> None:
|
||||
# The injected key is honoured regardless of allow_host_file — both the
|
||||
# host controller and the guest orchestrator take an injected key.
|
||||
self.assertEqual(
|
||||
b"injected-key", broker_secret({LAUNCH_BROKER_KEY_ENV: "injected-key"}))
|
||||
self.assertEqual(
|
||||
b"injected-key",
|
||||
broker_secret({LAUNCH_BROKER_KEY_ENV: "injected-key"}, allow_host_file=True))
|
||||
|
||||
def test_guest_without_injection_fails_closed(self) -> None:
|
||||
# The default (guest orchestrator): no env key and NO host-file fallback,
|
||||
# so it returns None rather than mint a divergent process-local key.
|
||||
self.assertIsNone(broker_secret({}))
|
||||
|
||||
def test_host_side_falls_back_to_the_durable_key_file(self) -> None:
|
||||
# allow_host_file=True (host controller / dev-harness): mint/read the
|
||||
# durable host key file, the same key on every call (restart re-adoption).
|
||||
with tempfile.TemporaryDirectory() as root:
|
||||
with patch.dict("os.environ", {"BOT_BOTTLE_ROOT": root}, clear=False):
|
||||
os.environ.pop(LAUNCH_BROKER_KEY_ENV, None)
|
||||
first = broker_secret(allow_host_file=True)
|
||||
second = broker_secret(allow_host_file=True)
|
||||
self.assertTrue(first)
|
||||
self.assertEqual(first, second)
|
||||
|
||||
|
||||
class TestSeamRoundTrip(unittest.TestCase):
|
||||
"""The whole point of chunk 1: a request signed by the orchestrator side is
|
||||
POSTed to a real host control server, verified there, and acted on — over
|
||||
HTTP, not an in-process call."""
|
||||
|
||||
def _serve(self, broker: LaunchBroker) -> BrokerClient:
|
||||
server = make_host_server(broker, "127.0.0.1", 0)
|
||||
self.addCleanup(server.server_close)
|
||||
threading.Thread(target=server.serve_forever, daemon=True).start()
|
||||
self.addCleanup(server.shutdown)
|
||||
host, port = server.server_address[0], server.server_address[1]
|
||||
return BrokerClient(f"http://{host}:{port}")
|
||||
|
||||
def test_sign_post_verify_act_over_http(self) -> None:
|
||||
secret = secrets.token_bytes(16)
|
||||
broker = StubBroker(secret)
|
||||
client = self._serve(broker)
|
||||
req = LaunchRequest(
|
||||
op="launch", bottle_id="b1", source_ip="10.0.0.1", image_ref="img", slot=1)
|
||||
got = client.submit(sign_request(req, secret))
|
||||
self.assertEqual(req, got) # the controller echoes the verified request
|
||||
self.assertEqual(["b1"], [r.bottle_id for r in broker.launched])
|
||||
|
||||
def test_forged_token_raises_broker_auth_error_over_http(self) -> None:
|
||||
secret = secrets.token_bytes(16)
|
||||
broker = StubBroker(secret)
|
||||
client = self._serve(broker)
|
||||
forged = sign_request(
|
||||
LaunchRequest(op="launch", bottle_id="b1"), secrets.token_bytes(16))
|
||||
with self.assertRaises(BrokerAuthError):
|
||||
client.submit(forged)
|
||||
self.assertEqual([], broker.launched) # fail-closed across the wire
|
||||
|
||||
def test_list_live_enumerates_over_http(self) -> None:
|
||||
secret = secrets.token_bytes(16)
|
||||
broker = StubBroker(secret)
|
||||
broker.live_source_ips = ["10.0.0.7", "10.0.0.8"]
|
||||
client = self._serve(broker)
|
||||
got = client.list_live(sign_request(LaunchRequest(op="list_live"), secret))
|
||||
self.assertEqual(["10.0.0.7", "10.0.0.8"], got)
|
||||
|
||||
|
||||
class TestRequestLimits(unittest.TestCase):
|
||||
"""The privileged listener must not let a caller that can merely reach the
|
||||
socket (no signed token) exhaust it via an oversized declared body — and it
|
||||
rejects on the Content-Length *header*, before reading the body."""
|
||||
|
||||
def _addr(self) -> tuple[str, int]:
|
||||
self.broker = StubBroker(secrets.token_bytes(16))
|
||||
server = make_host_server(self.broker, "127.0.0.1", 0)
|
||||
self.addCleanup(server.server_close)
|
||||
threading.Thread(target=server.serve_forever, daemon=True).start()
|
||||
self.addCleanup(server.shutdown)
|
||||
host, port = server.server_address[:2]
|
||||
return typing.cast(str, host), port
|
||||
|
||||
def test_oversized_content_length_is_rejected_before_reading(self) -> None:
|
||||
host, port = self._addr()
|
||||
conn = http.client.HTTPConnection(host, port, timeout=5)
|
||||
self.addCleanup(conn.close)
|
||||
# Declare an oversized body but send only a sliver: the server must reject
|
||||
# on the header before reading, so the caller gets a clean, deterministic
|
||||
# 413 (no large unread body to race a connection reset).
|
||||
conn.putrequest("POST", "/broker", skip_accept_encoding=True)
|
||||
conn.putheader("Content-Type", "application/json")
|
||||
conn.putheader("Content-Length", str(MAX_BODY_BYTES + 1))
|
||||
conn.endheaders()
|
||||
conn.send(b"{}") # far short of the declared length; never read
|
||||
resp = conn.getresponse()
|
||||
self.assertEqual(413, resp.status)
|
||||
self.assertEqual([], self.broker.launched) # never reached the broker
|
||||
|
||||
|
||||
class TestServeUnit(unittest.TestCase):
|
||||
"""Drive `Handler._serve` directly (no socket). The real per-request handler
|
||||
runs in a daemon thread whose coverage/trace data is lost, so the
|
||||
bounded-body and error paths are exercised here in the main thread instead."""
|
||||
|
||||
def _handler(self, broker: LaunchBroker, headers: dict[str, str],
|
||||
body: bytes = b"") -> tuple[Handler, MagicMock]:
|
||||
server = HostControlServer.__new__(HostControlServer)
|
||||
server.broker = broker
|
||||
h = Handler.__new__(Handler)
|
||||
h.server = server
|
||||
h.headers = headers # type: ignore[assignment] — dict is a valid .get() stand-in
|
||||
h.path = "/broker"
|
||||
h.rfile = io.BytesIO(body)
|
||||
h.wfile = io.BytesIO()
|
||||
send_response = MagicMock()
|
||||
h.send_response = send_response # type: ignore[method-assign]
|
||||
h.send_header = MagicMock() # type: ignore[method-assign]
|
||||
h.end_headers = MagicMock() # type: ignore[method-assign]
|
||||
return h, send_response
|
||||
|
||||
def test_oversized_content_length_is_413(self) -> None:
|
||||
broker = StubBroker(secrets.token_bytes(16))
|
||||
h, send_response = self._handler(broker, {"Content-Length": str(MAX_BODY_BYTES + 1)})
|
||||
h.do_POST() # exercises do_POST -> _serve
|
||||
send_response.assert_called_once_with(413)
|
||||
self.assertEqual([], broker.launched) # rejected before the broker
|
||||
|
||||
def test_invalid_content_length_is_400(self) -> None:
|
||||
h, send_response = self._handler(StubBroker(secrets.token_bytes(16)),
|
||||
{"Content-Length": "not-a-number"})
|
||||
h._serve("POST")
|
||||
send_response.assert_called_once_with(400)
|
||||
|
||||
def test_valid_request_dispatches_200(self) -> None:
|
||||
secret = secrets.token_bytes(16)
|
||||
broker = StubBroker(secret)
|
||||
body = _body({"token": sign_request(
|
||||
LaunchRequest(op="teardown", bottle_id="b1"), secret)})
|
||||
h, send_response = self._handler(broker, {"Content-Length": str(len(body))}, body)
|
||||
h._serve("POST")
|
||||
send_response.assert_called_once_with(200)
|
||||
self.assertEqual(["b1"], [r.bottle_id for r in broker.torn_down])
|
||||
|
||||
def test_dispatch_exception_becomes_500(self) -> None:
|
||||
# dispatch is total, but the handler still guards it: a raised dispatch
|
||||
# returns 500 rather than dropping the connection.
|
||||
h, send_response = self._handler(
|
||||
StubBroker(secrets.token_bytes(16)), {"Content-Length": "0"})
|
||||
with patch("bot_bottle.orchestrator.host_server.dispatch",
|
||||
side_effect=RuntimeError("boom")):
|
||||
h._serve("POST")
|
||||
send_response.assert_called_once_with(500)
|
||||
|
||||
def test_health_over_do_get(self) -> None:
|
||||
h, send_response = self._handler(StubBroker(secrets.token_bytes(16)), {})
|
||||
h.path = "/health"
|
||||
h.do_GET()
|
||||
send_response.assert_called_once_with(200)
|
||||
|
||||
|
||||
class TestMain(unittest.TestCase):
|
||||
def test_fail_closed_without_secret(self) -> None:
|
||||
with patch("bot_bottle.orchestrator.host_server.broker_secret",
|
||||
return_value=None):
|
||||
self.assertEqual(2, main(["--port", "0"]))
|
||||
|
||||
def test_serves_then_shuts_down_cleanly(self) -> None:
|
||||
fake = MagicMock()
|
||||
fake.server_address = ("127.0.0.1", 0)
|
||||
fake.serve_forever.side_effect = KeyboardInterrupt
|
||||
with patch("bot_bottle.orchestrator.host_server.broker_secret",
|
||||
return_value=b"k"), \
|
||||
patch("bot_bottle.orchestrator.host_server.make_host_server",
|
||||
return_value=fake):
|
||||
self.assertEqual(0, main(["--port", "0"]))
|
||||
fake.serve_forever.assert_called_once()
|
||||
fake.server_close.assert_called_once()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,67 @@
|
||||
"""Unit: the orchestrator dev-harness entrypoint (`python -m bot_bottle.orchestrator`).
|
||||
|
||||
Exercises broker selection (stub / docker / http) and the fail-closed http path,
|
||||
patching `make_server` so the serve loop returns instead of blocking.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import secrets
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from bot_bottle.orchestrator.__main__ import main
|
||||
from bot_bottle.paths import LAUNCH_BROKER_KEY_ENV
|
||||
|
||||
|
||||
def _fake_server() -> MagicMock:
|
||||
fake = MagicMock()
|
||||
fake.server_address = ("127.0.0.1", 0)
|
||||
# Break out of serve_forever immediately, exercising the try/finally.
|
||||
fake.serve_forever.side_effect = KeyboardInterrupt
|
||||
return fake
|
||||
|
||||
|
||||
class TestMain(unittest.TestCase):
|
||||
def _run(self, broker: str, env: dict[str, str] | None = None) -> tuple[int, MagicMock]:
|
||||
fake = _fake_server()
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
argv = ["--db", str(Path(d) / "r.db"), "--port", "0", "--broker", broker]
|
||||
with patch("bot_bottle.orchestrator.__main__.make_server", return_value=fake), \
|
||||
patch.dict("os.environ", env or {}, clear=False):
|
||||
if env is None:
|
||||
os.environ.pop(LAUNCH_BROKER_KEY_ENV, None)
|
||||
rc = main(argv)
|
||||
return rc, fake
|
||||
|
||||
def test_stub_broker_serves_and_closes(self) -> None:
|
||||
rc, fake = self._run("stub")
|
||||
self.assertEqual(0, rc)
|
||||
fake.serve_forever.assert_called_once()
|
||||
fake.server_close.assert_called_once()
|
||||
|
||||
def test_docker_broker_serves(self) -> None:
|
||||
rc, _ = self._run("docker")
|
||||
self.assertEqual(0, rc)
|
||||
|
||||
def test_http_broker_with_injected_key_serves(self) -> None:
|
||||
# The guest orchestrator takes the launch-broker key by injection.
|
||||
rc, _ = self._run(
|
||||
"http", env={LAUNCH_BROKER_KEY_ENV: secrets.token_urlsafe(16)})
|
||||
self.assertEqual(0, rc)
|
||||
|
||||
def test_http_broker_without_injected_key_exits(self) -> None:
|
||||
# Fail-closed: no host-file fallback for the guest, so a missing injected
|
||||
# key is a usage error rather than a silently-minted divergent key.
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
with patch.dict("os.environ", {}, clear=False):
|
||||
os.environ.pop(LAUNCH_BROKER_KEY_ENV, None)
|
||||
with self.assertRaises(SystemExit):
|
||||
main(["--db", str(Path(d) / "r.db"), "--broker", "http"])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -2,10 +2,12 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import unittest
|
||||
|
||||
from bot_bottle.orchestrator.store.secret_store import (
|
||||
ENV_VAR_SECRET_NAME,
|
||||
_NONCE_BYTES,
|
||||
decrypt_value,
|
||||
encrypt_value,
|
||||
new_env_var_secret,
|
||||
@@ -65,22 +67,27 @@ class TestDecryptErrors(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.secret = new_env_var_secret()
|
||||
|
||||
def test_wrong_key_raises_value_error(self) -> None:
|
||||
def test_wrong_key_always_raises_value_error(self) -> None:
|
||||
# Deterministic: the authentication tag rejects a wrong key every time,
|
||||
# so reprovision can never inject a garbage token. Repeat across many
|
||||
# random keys (the old unauthenticated scheme let ~5% through when the
|
||||
# garbage happened to decode as valid UTF-8).
|
||||
for _ in range(200):
|
||||
ct = encrypt_value(self.secret, "secret-token")
|
||||
with self.assertRaises(ValueError):
|
||||
decrypt_value(new_env_var_secret(), ct)
|
||||
|
||||
def test_tampered_ciphertext_raises_value_error(self) -> None:
|
||||
ct = encrypt_value(self.secret, "secret-token")
|
||||
other_key = new_env_var_secret()
|
||||
# Wrong key produces garbage bytes; decrypt_value raises ValueError
|
||||
# when the result is non-UTF-8 (which is very likely for 12-char data).
|
||||
# We allow it to succeed only if garbage happens to be valid UTF-8, but
|
||||
# the plaintext must not match.
|
||||
try:
|
||||
result = decrypt_value(other_key, ct)
|
||||
self.assertNotEqual("secret-token", result)
|
||||
except ValueError:
|
||||
pass
|
||||
raw = bytearray(base64.urlsafe_b64decode(ct + "=" * (-len(ct) % 4)))
|
||||
raw[_NONCE_BYTES] ^= 0x01 # flip a bit in the ciphertext body → tag mismatch
|
||||
tampered = base64.urlsafe_b64encode(bytes(raw)).rstrip(b"=").decode()
|
||||
with self.assertRaises(ValueError):
|
||||
decrypt_value(self.secret, tampered)
|
||||
|
||||
def test_truncated_blob_raises_value_error(self) -> None:
|
||||
with self.assertRaises(ValueError):
|
||||
decrypt_value(self.secret, "dG9vc2hvcnQ") # "tooshort" — under 16 nonce bytes
|
||||
decrypt_value(self.secret, "dG9vc2hvcnQ") # "tooshort" — under nonce+tag
|
||||
|
||||
def test_invalid_base64_raises_value_error(self) -> None:
|
||||
with self.assertRaises(ValueError):
|
||||
|
||||
@@ -624,13 +624,18 @@ if __name__ == "__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."""
|
||||
"""`POST /reconcile` — a bare self-heal trigger. The orchestrator pulls its
|
||||
own live set from the broker (`list_live`); the caller no longer supplies it,
|
||||
so the request body carries at most `grace_seconds`."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.orch = _orchestrator(Path(self._tmp.name) / "r.db")
|
||||
store = RegistryStore(Path(self._tmp.name) / "r.db")
|
||||
store.migrate()
|
||||
secret = secrets.token_bytes(16)
|
||||
# Hold a typed StubBroker so tests can set its live set directly.
|
||||
self.broker = StubBroker(secret)
|
||||
self.orch = OrchestratorCore(store, self.broker, secret)
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
@@ -647,44 +652,35 @@ class TestReconcileRoute(unittest.TestCase):
|
||||
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.broker.live_source_ips = ["10.0.0.2"] # broker: only .2 is up
|
||||
status, payload = dispatch(self.orch, "POST", "/reconcile", _body({}))
|
||||
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_empty_body_triggers_the_sweep(self) -> None:
|
||||
"""No body is fine — the live set comes from the broker, not the body."""
|
||||
dead = self._old("10.0.0.1")
|
||||
self.broker.live_source_ips = [] # broker: nothing running
|
||||
status, payload = dispatch(self.orch, "POST", "/reconcile", b"")
|
||||
self.assertEqual(200, status)
|
||||
self.assertEqual([dead], payload["reaped"])
|
||||
|
||||
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")
|
||||
self.broker.live_source_ips = []
|
||||
status, payload = dispatch(
|
||||
self.orch, "POST", "/reconcile",
|
||||
_body({"live_source_ips": [], "grace_seconds": 3600}))
|
||||
self.orch, "POST", "/reconcile", _body({"grace_seconds": 3600}))
|
||||
self.assertEqual(200, status)
|
||||
self.assertEqual([], payload["reaped"])
|
||||
|
||||
def test_non_string_entries_are_rejected(self) -> None:
|
||||
status, payload = dispatch(
|
||||
self.orch, "POST", "/reconcile",
|
||||
_body({"live_source_ips": [None, 7, "10.0.0.9"]}))
|
||||
self.assertEqual(400, status)
|
||||
self.assertIn("live_source_ips", str(payload["error"]))
|
||||
|
||||
def test_empty_live_source_ip_is_rejected(self) -> None:
|
||||
status, payload = dispatch(
|
||||
self.orch, "POST", "/reconcile", _body({"live_source_ips": [""]}))
|
||||
self.assertEqual(400, status)
|
||||
self.assertIn("live_source_ips", str(payload["error"]))
|
||||
|
||||
def test_invalid_grace_seconds_is_rejected(self) -> None:
|
||||
for value in (True, "30", -1, float("inf"), float("nan")):
|
||||
with self.subTest(value=value):
|
||||
status, payload = dispatch(
|
||||
self.orch, "POST", "/reconcile",
|
||||
_body({"live_source_ips": [], "grace_seconds": value}))
|
||||
_body({"grace_seconds": value}))
|
||||
self.assertEqual(400, status)
|
||||
self.assertIn("grace_seconds", str(payload["error"]))
|
||||
|
||||
@@ -11,7 +11,12 @@ from contextlib import closing
|
||||
from pathlib import Path
|
||||
from unittest.mock import patch
|
||||
|
||||
from bot_bottle.orchestrator.broker import LaunchBroker, LaunchRequest, StubBroker
|
||||
from bot_bottle.orchestrator.broker import (
|
||||
BrokerUnavailableError,
|
||||
LaunchBroker,
|
||||
LaunchRequest,
|
||||
StubBroker,
|
||||
)
|
||||
from bot_bottle.orchestrator.store.registry_store import RegistryStore
|
||||
from bot_bottle.orchestrator.service import OrchestratorCore
|
||||
from bot_bottle.orchestrator.store.secret_store import new_env_var_secret
|
||||
@@ -25,8 +30,8 @@ from bot_bottle.orchestrator.supervisor import (
|
||||
|
||||
|
||||
class _FailingBroker(LaunchBroker):
|
||||
"""Verifies the token like any broker, then fails the launch — to
|
||||
exercise the orchestrator's registry rollback."""
|
||||
"""Verifies the token like any broker, then fails the launch *definitely* —
|
||||
to exercise the orchestrator's registry rollback."""
|
||||
|
||||
def _launch(self, req: LaunchRequest) -> None:
|
||||
raise RuntimeError("launch failed")
|
||||
@@ -34,6 +39,24 @@ class _FailingBroker(LaunchBroker):
|
||||
def _teardown(self, req: LaunchRequest) -> None:
|
||||
pass
|
||||
|
||||
def _list_live(self) -> list[str]:
|
||||
return []
|
||||
|
||||
|
||||
class _UnavailableBroker(LaunchBroker):
|
||||
"""Verifies the token, then raises the *ambiguous* BrokerUnavailableError —
|
||||
the host may already have launched — so the orchestrator must KEEP the
|
||||
registry row rather than orphan a running container."""
|
||||
|
||||
def _launch(self, req: LaunchRequest) -> None:
|
||||
raise BrokerUnavailableError("delivery dropped after send")
|
||||
|
||||
def _teardown(self, req: LaunchRequest) -> None:
|
||||
pass
|
||||
|
||||
def _list_live(self) -> list[str]:
|
||||
return []
|
||||
|
||||
|
||||
class TestOrchestrator(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
@@ -144,11 +167,20 @@ class TestOrchestrator(unittest.TestCase):
|
||||
self.assertIsNotNone(self.orch.resolve("10.243.0.3", rec.identity_token))
|
||||
self.assertIsNone(self.orch.resolve("10.243.0.3", "wrong-token"))
|
||||
|
||||
def test_launch_rolls_back_registry_on_broker_failure(self) -> None:
|
||||
def test_launch_rolls_back_registry_on_definite_broker_failure(self) -> None:
|
||||
orch = OrchestratorCore(self.store, _FailingBroker(self.secret), self.secret)
|
||||
with self.assertRaises(RuntimeError):
|
||||
orch.launch_bottle("10.243.0.9")
|
||||
self.assertEqual([], self.store.all()) # no orphan
|
||||
self.assertEqual([], self.store.all()) # no orphan row
|
||||
|
||||
def test_launch_keeps_registry_on_ambiguous_broker_failure(self) -> None:
|
||||
# The host may already have launched the bottle before the response was
|
||||
# lost, so deregistering would orphan a running container with no row.
|
||||
# The row is kept for reconcile to reap iff the bottle is not live.
|
||||
orch = OrchestratorCore(self.store, _UnavailableBroker(self.secret), self.secret)
|
||||
with self.assertRaises(BrokerUnavailableError):
|
||||
orch.launch_bottle("10.243.0.9")
|
||||
self.assertEqual(1, len(self.store.all())) # row survives — no orphan container
|
||||
|
||||
def test_gateway_status_reports_unconfigured(self) -> None:
|
||||
# The orchestrator no longer owns a standalone gateway lifecycle; the
|
||||
@@ -339,8 +371,9 @@ class TestOrchestratorSupervise(unittest.TestCase):
|
||||
pid = self.orch.supervise_queue_proposal(
|
||||
rec.bottle_id, tool=TOOL_EGRESS_ALLOW,
|
||||
proposed_file="routes:\n - host: google.com\n", justification="j")
|
||||
# No live source IPs -> the bottle is reaped (grace 0 so it's immediate).
|
||||
self.assertEqual([rec.bottle_id], self.orch.reconcile([], grace_seconds=0))
|
||||
# No live source IPs (nothing launched through the stub) -> the bottle
|
||||
# is reaped (grace 0 so it's immediate).
|
||||
self.assertEqual([rec.bottle_id], self.orch.reconcile(grace_seconds=0))
|
||||
self.assertEqual(
|
||||
{"status": "unknown"}, self.orch.supervise_poll_response(rec.bottle_id, pid))
|
||||
|
||||
@@ -385,7 +418,9 @@ class TestOrchestratorReconcile(unittest.TestCase):
|
||||
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"]))
|
||||
# The broker reports only .2 as still running -> .1 is reaped.
|
||||
self.broker.live_source_ips = ["10.243.0.2"]
|
||||
self.assertEqual([dead.bottle_id], self.orch.reconcile())
|
||||
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.
|
||||
@@ -397,14 +432,31 @@ class TestOrchestratorReconcile(unittest.TestCase):
|
||||
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.broker.live_source_ips = [] # broker reports nothing running
|
||||
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.broker.live_source_ips = ["10.243.0.1", "10.243.0.2"]
|
||||
self.assertEqual([], self.orch.reconcile())
|
||||
self.assertIsNotNone(self.store.get(a.bottle_id))
|
||||
self.assertIsNotNone(self.store.get(b.bottle_id))
|
||||
|
||||
def test_reconcile_skipped_when_broker_cannot_enumerate(self) -> None:
|
||||
"""A broker that can't return an authoritative live set must NOT be
|
||||
treated as "nothing is live" — that would reap every healthy bottle.
|
||||
The sweep is skipped instead."""
|
||||
a = self.orch.launch_bottle("10.243.0.1")
|
||||
b = self.orch.launch_bottle("10.243.0.2")
|
||||
self._age_all(600)
|
||||
|
||||
def _boom() -> list[str]:
|
||||
raise RuntimeError("docker ps failed")
|
||||
|
||||
self.broker._list_live = _boom # type: ignore[method-assign]
|
||||
self.assertEqual([], self.orch.reconcile())
|
||||
self.assertIsNotNone(self.store.get(a.bottle_id))
|
||||
self.assertIsNotNone(self.store.get(b.bottle_id))
|
||||
|
||||
@@ -6,10 +6,13 @@ import unittest
|
||||
from unittest.mock import patch
|
||||
|
||||
from bot_bottle import orchestrator_auth
|
||||
from bot_bottle.orchestrator_auth import ROLE_CLI, ROLE_GATEWAY
|
||||
from bot_bottle.orchestrator_auth import ROLE_CLI, ROLE_GATEWAY, ROLE_HOST
|
||||
from bot_bottle.trust_domain import (
|
||||
CONTROL_PLANE,
|
||||
HOST_CONTROLLER,
|
||||
LAUNCH_BROKER,
|
||||
ControlPlaneProvisioning,
|
||||
LaunchBrokerProvisioning,
|
||||
ProvisioningError,
|
||||
TrustDomain,
|
||||
)
|
||||
@@ -101,5 +104,73 @@ class TestControlPlaneProvisioning(unittest.TestCase):
|
||||
self.assertNotEqual(ROLE_CLI, CONTROL_PLANE.verify(tok, "k"))
|
||||
|
||||
|
||||
class TestLaunchBrokerAndHostControllerDomains(unittest.TestCase):
|
||||
"""The real #468 domains: the launch-broker key (shared by orchestrator +
|
||||
host controller) and the host controller's own lifecycle key."""
|
||||
|
||||
def test_launch_broker_mints_no_role_tokens(self) -> None:
|
||||
# Empty role set — it provides durable key material for the broker's own
|
||||
# launch JWT, not orchestrator_auth role tokens.
|
||||
self.assertEqual(frozenset(), LAUNCH_BROKER.roles)
|
||||
with patch("bot_bottle.trust_domain.host_signing_key", return_value="k"):
|
||||
with self.assertRaises(ValueError):
|
||||
LAUNCH_BROKER.mint(ROLE_CLI)
|
||||
|
||||
def test_host_controller_signs_host_role_only(self) -> None:
|
||||
with patch("bot_bottle.trust_domain.host_signing_key", return_value="k"):
|
||||
tok = HOST_CONTROLLER.mint(ROLE_HOST)
|
||||
self.assertEqual(ROLE_HOST, HOST_CONTROLLER.verify(tok, "k"))
|
||||
# A control-plane `cli` token (the orchestrator's key) never verifies as a
|
||||
# host-controller role — the orchestrator can't forge lifecycle creds.
|
||||
cli_tok = orchestrator_auth.mint(ROLE_CLI, "k")
|
||||
self.assertIsNone(HOST_CONTROLLER.verify(cli_tok, "k"))
|
||||
|
||||
def test_control_plane_cannot_mint_the_host_role(self) -> None:
|
||||
# `host` is outside the control-plane role set on purpose.
|
||||
with patch("bot_bottle.trust_domain.host_signing_key", return_value="k"):
|
||||
with self.assertRaises(ValueError):
|
||||
CONTROL_PLANE.mint(ROLE_HOST)
|
||||
|
||||
def test_the_three_domains_use_distinct_keys_and_env_vars(self) -> None:
|
||||
self.assertEqual(3, len({
|
||||
CONTROL_PLANE.key_filename,
|
||||
LAUNCH_BROKER.key_filename,
|
||||
HOST_CONTROLLER.key_filename,
|
||||
}))
|
||||
self.assertEqual(3, len({
|
||||
CONTROL_PLANE.key_env, LAUNCH_BROKER.key_env, HOST_CONTROLLER.key_env,
|
||||
}))
|
||||
|
||||
|
||||
class TestLaunchBrokerProvisioning(unittest.TestCase):
|
||||
def test_broker_key_returns_the_durable_key(self) -> None:
|
||||
prov = LaunchBrokerProvisioning()
|
||||
with patch("bot_bottle.trust_domain.host_signing_key", return_value="bk"):
|
||||
self.assertEqual("bk", prov.broker_key())
|
||||
|
||||
def test_broker_key_fail_closes_when_empty(self) -> None:
|
||||
# An empty key would leave the host controller unable to verify any
|
||||
# launch — fail-closed rather than hand back a useless/dangerous key.
|
||||
prov = LaunchBrokerProvisioning()
|
||||
with patch("bot_bottle.trust_domain.host_signing_key", return_value=""):
|
||||
with self.assertRaises(ProvisioningError):
|
||||
prov.broker_key()
|
||||
|
||||
def test_controller_key_is_distinct_from_the_broker_key(self) -> None:
|
||||
# The orchestrator holds the broker key but NEVER the controller key.
|
||||
prov = LaunchBrokerProvisioning()
|
||||
keys = {"launch-broker-key": "bk", "host-controller-key": "ck"}
|
||||
with patch("bot_bottle.trust_domain.host_signing_key",
|
||||
side_effect=keys.__getitem__):
|
||||
self.assertEqual("bk", prov.broker_key())
|
||||
self.assertEqual("ck", prov.controller_key())
|
||||
|
||||
def test_controller_key_fail_closes_when_empty(self) -> None:
|
||||
prov = LaunchBrokerProvisioning()
|
||||
with patch("bot_bottle.trust_domain.host_signing_key", return_value=""):
|
||||
with self.assertRaises(ProvisioningError):
|
||||
prov.controller_key()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
Reference in New Issue
Block a user