Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 016e59e029 | |||
| f3664dea9f | |||
| 88b82a169e | |||
| 5a9428cc86 | |||
| 60039f2eb3 | |||
| bc42836327 |
@@ -17,9 +17,7 @@ from pathlib import Path
|
||||
|
||||
from .. import log
|
||||
from .store.store_manager import StoreManager
|
||||
from .broker import StubBroker, SubmitBroker
|
||||
from .broker_client import BrokerClient
|
||||
from .host_server import BROKER_SECRET_ENV, DEFAULT_PORT, broker_secret_from_env
|
||||
from .broker import LaunchBroker, StubBroker
|
||||
from .server import make_server
|
||||
from .docker_broker import DockerBroker
|
||||
from .store.registry_store import RegistryStore, default_db_path
|
||||
@@ -36,13 +34,8 @@ def main(argv: list[str] | None = None) -> int:
|
||||
help=f"registry DB path (default: {default_db_path()})",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--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)",
|
||||
"--broker", choices=("stub", "docker"), default="stub",
|
||||
help="launch broker: 'stub' records requests; 'docker' runs containers",
|
||||
)
|
||||
args = parser.parse_args(argv)
|
||||
|
||||
@@ -54,25 +47,11 @@ def main(argv: list[str] | None = None) -> int:
|
||||
# operator reaches it over HTTP (never a second, disconnected DB).
|
||||
StoreManager(registry.db_path).migrate()
|
||||
|
||||
# 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 secret the host controller holds, so it is read from
|
||||
# the shared env var (the chunk-1 stand-in for out-of-band provisioning).
|
||||
broker: SubmitBroker
|
||||
if args.broker == "http":
|
||||
secret = broker_secret_from_env()
|
||||
if secret is None:
|
||||
parser.error(
|
||||
f"--broker http requires a shared signing secret in "
|
||||
f"${BROKER_SECRET_ENV} (hex), matching the host control server"
|
||||
)
|
||||
broker = BrokerClient(args.host_controller_url)
|
||||
else:
|
||||
secret = secrets.token_bytes(32)
|
||||
broker = DockerBroker(secret) if args.broker == "docker" else StubBroker(secret)
|
||||
# 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)
|
||||
orchestrator = OrchestratorCore(registry, broker, secret)
|
||||
|
||||
server = make_server(orchestrator, host=args.host, port=args.port)
|
||||
|
||||
@@ -29,7 +29,6 @@ import json
|
||||
import secrets
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from typing import Protocol
|
||||
|
||||
_JWT_HEADER = {"alg": "HS256", "typ": "JWT"}
|
||||
_ALLOWED_OPS = ("launch", "teardown")
|
||||
@@ -38,21 +37,7 @@ _ALLOWED_OPS = ("launch", "teardown")
|
||||
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.
|
||||
|
||||
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."""
|
||||
the fixed launch-request shape. Fail-closed: the broker must not act."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -138,16 +123,6 @@ def verify_request(token: str, secret: bytes) -> LaunchRequest:
|
||||
|
||||
# --- the broker itself ------------------------------------------------------
|
||||
|
||||
class SubmitBroker(Protocol):
|
||||
"""The single method `OrchestratorCore` depends on: verify a signed token and
|
||||
perform its op, returning the verified request. 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: ...
|
||||
|
||||
|
||||
class LaunchBroker(abc.ABC):
|
||||
"""Verifies a signed request came from the orchestrator, then performs
|
||||
the backend-native launch/teardown. Subclasses implement `_launch` /
|
||||
@@ -193,9 +168,7 @@ class StubBroker(LaunchBroker):
|
||||
|
||||
__all__ = [
|
||||
"BrokerAuthError",
|
||||
"BrokerUnavailableError",
|
||||
"LaunchRequest",
|
||||
"SubmitBroker",
|
||||
"LaunchBroker",
|
||||
"StubBroker",
|
||||
"sign_request",
|
||||
|
||||
@@ -1,126 +0,0 @@
|
||||
"""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 _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 _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",
|
||||
]
|
||||
@@ -1,270 +0,0 @@
|
||||
"""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>"}
|
||||
|
||||
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 durable
|
||||
`TrustDomain` key in a later chunk.
|
||||
|
||||
The shared signing secret is read from `$BOT_BOTTLE_BROKER_SECRET` (hex). That is
|
||||
a **chunk-1 stopgap**: it must be provisioned to signer and verifier out of band,
|
||||
which is exactly what the durable `TrustDomain` key in chunk 2 (#476) replaces.
|
||||
"""
|
||||
|
||||
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 .broker import BrokerAuthError, LaunchBroker
|
||||
from .docker_broker import DockerBroker
|
||||
|
||||
# JSON body payload type (parsed request / rendered response).
|
||||
Json = dict[str, object]
|
||||
|
||||
# The hex-encoded HS256 secret shared with the request signer (the orchestrator).
|
||||
# Chunk-1 stopgap for the durable, out-of-band `TrustDomain` key of chunk 2.
|
||||
BROKER_SECRET_ENV = "BOT_BOTTLE_BROKER_SECRET"
|
||||
|
||||
# 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_from_env(environ: typing.Mapping[str, str] | None = None) -> bytes | None:
|
||||
"""The shared HS256 secret from `$BOT_BOTTLE_BROKER_SECRET` (hex), or None
|
||||
when unset or not valid hex. The signer (orchestrator, `--broker http`) and
|
||||
the verifier (this server) read the same env var so both hold the same key —
|
||||
the chunk-1 stand-in for out-of-band provisioning."""
|
||||
env = os.environ if environ is None else environ
|
||||
raw = env.get(BROKER_SECRET_ENV, "").strip()
|
||||
if not raw:
|
||||
return None
|
||||
try:
|
||||
return bytes.fromhex(raw)
|
||||
except ValueError:
|
||||
return 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,
|
||||
}
|
||||
|
||||
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 a shared `$BOT_BOTTLE_BROKER_SECRET` the server can
|
||||
verify no request's provenance, so it refuses to start rather than run a
|
||||
launcher that accepts unsigned input."""
|
||||
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_from_env()
|
||||
if secret is None:
|
||||
sys.stderr.write(
|
||||
f"host controller: refusing to start without a shared signing secret "
|
||||
f"(${BROKER_SECRET_ENV}, hex) — 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_from_env",
|
||||
"main",
|
||||
"Json",
|
||||
"BROKER_SECRET_ENV",
|
||||
"DEFAULT_PORT",
|
||||
]
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -25,7 +25,7 @@ import json
|
||||
from collections.abc import Iterable
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from .broker import BrokerUnavailableError, LaunchRequest, SubmitBroker, sign_request
|
||||
from .broker import LaunchBroker, LaunchRequest, sign_request
|
||||
from .store.registry_store import DEFAULT_REAP_GRACE_SECONDS, BottleRecord, RegistryStore
|
||||
from .supervisor import (
|
||||
AuditEntry,
|
||||
@@ -62,7 +62,7 @@ class OrchestratorCore:
|
||||
def __init__(
|
||||
self,
|
||||
registry: RegistryStore,
|
||||
broker: SubmitBroker,
|
||||
broker: LaunchBroker,
|
||||
sign_secret: bytes,
|
||||
supervisor: Supervisor | None = None,
|
||||
) -> None:
|
||||
@@ -111,23 +111,14 @@ class OrchestratorCore:
|
||||
image_ref=image_ref,
|
||||
slot=slot,
|
||||
)
|
||||
launched = False
|
||||
try:
|
||||
self._broker.submit(sign_request(req, self._secret))
|
||||
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
|
||||
launched = True
|
||||
finally:
|
||||
if not launched:
|
||||
self.registry.deregister(rec.bottle_id)
|
||||
self._tokens.pop(rec.bottle_id, None)
|
||||
return rec
|
||||
|
||||
def teardown_bottle(self, bottle_id: str) -> bool:
|
||||
|
||||
@@ -0,0 +1,732 @@
|
||||
# PRD prd-new: Canonical tamper-evident audit-event schema and local query contract
|
||||
|
||||
- **Status:** Draft
|
||||
- **Author:** didericis-claude
|
||||
- **Created:** 2026-07-26
|
||||
- **Issue:** #487
|
||||
|
||||
## Summary
|
||||
|
||||
bot-bottle already emits security- and provenance-relevant events from
|
||||
several producers — supervise operator decisions (PRD 0013's
|
||||
`AuditStore`), egress allow/block enforcement, git-gate push decisions,
|
||||
control-plane token minting, and (next) host-controller lifecycle
|
||||
transitions — but each writes its own shape to its own sink. There is no
|
||||
shared envelope, no tamper-evidence, and no single place to search. Local
|
||||
incident reconstruction means grepping several stores that don't agree on
|
||||
field names, timestamps, or how a bottled agent is identified.
|
||||
|
||||
This PRD defines **one canonical audit-event contract** every producer
|
||||
emits into:
|
||||
|
||||
1. A **versioned envelope** — schema version, event id, event/observed
|
||||
timestamps, host-attributed identity (`bottle`/`bottled_agent`/
|
||||
`activation`), provenance (`manifest_digest` — the manifest is the
|
||||
policy — and `engine` = bot-bottle version/SHA),
|
||||
host-observed `actor`/`action`/`resource`/`outcome`, correlation/
|
||||
causation ids, a sensitivity class, a typed payload, and an explicit
|
||||
**trust boundary** between host-supplied and agent-claimed fields.
|
||||
2. **Canonical JSON serialization + a per-writer hash chain** (with
|
||||
normative test vectors), so any deletion, edit, or reorder of a past
|
||||
record breaks the chain and is detectable offline.
|
||||
3. An **append-only JSONL journal as the source of truth**, with a
|
||||
**rebuildable SQLite index** and a local `audit query`/`verify` surface
|
||||
— no paid platform, no network dependency.
|
||||
4. An **initial event registry** covering lifecycle, host-controller,
|
||||
supervise decision, egress (request/decision/cutoff/anomaly), git-gate
|
||||
and signed-commit (#480), auth/authz, and audit self-events — each with
|
||||
its trusted-vs-claimed fields and redaction rules.
|
||||
5. A **stable export projection** (CloudEvents / OpenTelemetry Logs) and
|
||||
the **#324 delivery contract** (payload, per-chain
|
||||
`(host, epoch, seq)` cursor, dedup,
|
||||
backpressure, retention ordering).
|
||||
|
||||
It is explicitly scheduled to land **immediately after the host
|
||||
controller (#468)** so the host controller's lifecycle transitions are the
|
||||
first producer wired onto the new contract (per the directive on #487).
|
||||
|
||||
## Problem
|
||||
|
||||
Audit infrastructure is fragmented across #468, #324, and #480 with no
|
||||
shared schema. Concretely:
|
||||
|
||||
- **No shared envelope.** `supervise_audit_entries` (PRD 0013) has
|
||||
`timestamp, bottle_slug, component, operator_action, ...`. The egress
|
||||
proxy and git-gate log their own ad-hoc lines. There is no common
|
||||
`event_id`, `event_type`, or version, so cross-producer correlation
|
||||
("what did bottled agent X do between its start and this rejected push?") is
|
||||
manual and lossy.
|
||||
- **No tamper-evidence.** The audit store is a plain SQLite table. Anyone
|
||||
who can write the DB can delete or rewrite a row and leave no trace.
|
||||
Audit that an attacker (or a buggy agent) can silently rewrite is not
|
||||
audit.
|
||||
- **Trusted and untrusted data are mixed.** A bottled agent is attributed by
|
||||
**source IP → slug** at the gateway (host-supplied, trustworthy). An
|
||||
agent can also *claim* things about itself in a tool call
|
||||
(agent-claimed, adversarial). Today nothing in the record marks which is
|
||||
which, so a reader can be misled by an agent-supplied field that looks
|
||||
authoritative.
|
||||
- **No local search.** Reconstructing an incident means reading multiple
|
||||
sinks with different schemas. There is no query contract and no promise
|
||||
that the index can be rebuilt from the journal if it drifts or is lost.
|
||||
- **No redaction rule.** Nothing prohibits a producer from writing a raw
|
||||
token or secret into an audit record, which would turn the audit log
|
||||
itself into a credential store.
|
||||
|
||||
## Goals / Success Criteria
|
||||
|
||||
- A single `AuditEvent` envelope type, versioned, that every producer
|
||||
emits. The trust boundary is **one `untrusted` region**: everything
|
||||
outside it is host-established and trusted (source-IP → bottled-agent
|
||||
attribution, host wall-clock, producer identity, chain metadata);
|
||||
`untrusted` is the sole place anything an agent or a remote claimed may
|
||||
go. The boundary is structural, not a convention.
|
||||
- **Canonical serialization** (`sort_keys`, `(",", ":")` separators,
|
||||
UTF-8, `ensure_ascii=False`) is defined once and reused, so the same
|
||||
logical event always hashes identically across producers and hosts.
|
||||
- Each writer maintains a **hash chain**: `hash = sha256(prev_hash ||
|
||||
canonical(event))`. Deleting or editing any past record breaks every
|
||||
subsequent link; a standalone verifier detects the break offline with no
|
||||
secret material.
|
||||
- The **JSONL journal is the source of truth**; the **SQLite index is
|
||||
fully rebuildable** from it (`audit rebuild` reconstructs the DB and
|
||||
re-verifies the chain).
|
||||
- **Local query** works with no paid platform and no egress: filter by
|
||||
bottled agent, event type, time range, and producer, and follow a bottled
|
||||
agent's events in order.
|
||||
- A **redaction rule** is enforced at the envelope boundary: known
|
||||
credential-shaped fields are rejected/redacted before a record is
|
||||
written; the writer refuses raw secrets rather than storing them.
|
||||
- The envelope **projects onto the OpenTelemetry Logs data model and a
|
||||
CloudEvents JSON envelope** by field re-mapping alone (no reformat),
|
||||
preserving the trust boundary and carrying the integrity fields — per
|
||||
#487's export/interop requirement. (The export adapters are follow-up;
|
||||
the *schema* must make them a re-map.)
|
||||
- The **host controller (#468)** emits `lifecycle.*` events through this
|
||||
contract as the first consumer; existing supervise/egress producers are
|
||||
migrated behind the same envelope without changing operator-facing
|
||||
behavior.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- **Cross-host aggregation / shipping.** This PRD makes each host's journal
|
||||
canonical and correlatable *by construction* (stable ids, hash chain),
|
||||
but the transport that merges multiple hosts into one timeline is a
|
||||
follow-up (#324). The schema is designed so that merge is a later append,
|
||||
not a reformat.
|
||||
- **Cryptographic signing / external anchoring.** Hash-chaining gives
|
||||
tamper-**evidence** (you can detect edits), not tamper-**resistance**
|
||||
against an attacker who can rewrite the whole chain. Per-writer signing
|
||||
keys and periodic external anchoring are a follow-up; the chain-head hash
|
||||
is the seam they attach to.
|
||||
- **Real-time alerting / SIEM rules.** Query is local and pull-based here.
|
||||
- **Retention / rotation policy.** Journal rotation and TTL are operator
|
||||
policy, tracked separately; the format must survive rotation (chain head
|
||||
carried across segments) but this PRD does not set the schedule.
|
||||
- **Replacing PRD 0013's operator queue.** The supervise proposal/response
|
||||
queue is unchanged; only its terminal *audit* record is re-emitted onto
|
||||
the new envelope.
|
||||
|
||||
## Design
|
||||
|
||||
### The envelope
|
||||
|
||||
One dataclass, `AuditEvent`, serialized to a JSON object with a small,
|
||||
stable top level:
|
||||
|
||||
```
|
||||
{
|
||||
// ---- schema + integrity (host-owned) ----
|
||||
"v": 1, // schema version — bumped only on a breaking change
|
||||
"id": "<uuid4>", // globally unique event id; stable across export/replay (dedup key)
|
||||
"type": "egress.decision", // dotted event type from the registry
|
||||
"epoch": 7, // writer-boot counter, bumped once per host-controller (writer) start
|
||||
"seq": 1287, // monotonic sequence within this epoch (gap-detectable)
|
||||
"segment": "20260726T000000Z", // journal segment id (rotation boundary); chain continues across segments
|
||||
"prev": "<hex>", // prior record hash ("" only for the first record of a new chain)
|
||||
"hash": "<hex>", // sha256(prev + canonical(this event with hash=""))
|
||||
|
||||
// ---- timestamps (host-owned) ----
|
||||
"ts_event": "2026-07-26T18:22:04.061Z", // when the underlying event occurred at the boundary
|
||||
"ts_recorded": "2026-07-26T18:22:04.113Z", // when the single writer appended it (authoritative)
|
||||
"ts_mono": 90142.55, // monotonic secs since this epoch's boot (intra-epoch ordering only)
|
||||
|
||||
// ---- attribution + provenance (host-established) ----
|
||||
"producer": "egress", // host component that emitted the event
|
||||
"host": "mac-studio-1",
|
||||
"engine": "bot-bottle/0.1.0+abc1234", // bot-bottle version + git SHA of the enforcing host code
|
||||
"bottle": "amber-fox", // bottle (container/VM) identity
|
||||
"bottled_agent": "amber-fox-12", // bottled-agent slug from source-IP attribution (null for host-level events)
|
||||
"activation": "01J8Z...", // activation id: one run/session of the bottled agent (null if n/a)
|
||||
"manifest_digest": "sha256:9f2…", // digest of the manifest — which IS the policy (egress routes etc.); null if n/a
|
||||
|
||||
// ---- semantics: host-observed facts of what happened ----
|
||||
"actor": "bottled-agent:amber-fox-12", // who acted, as a host-attributed identity
|
||||
"action": "egress.connect", // what was attempted / done
|
||||
"resource": "registry.npmjs.org:443", // what it acted on, as observed at the boundary
|
||||
"outcome": "blocked", // host-decided result: allowed|blocked|deferred|success|failure
|
||||
"sensitivity": "security", // classification: normal|security|restricted (drives redaction + export)
|
||||
|
||||
// ---- correlation (host-assigned) ----
|
||||
"correlation_id": "flow-9c2a…", // groups a related flow (request → decision → cutoff)
|
||||
"causation_id": "<event id>", // the event that directly caused this one ("" if root)
|
||||
|
||||
// ---- typed, trusted, event-specific payload (shape fixed per type in the registry) ----
|
||||
"payload": {
|
||||
"route_id": 4,
|
||||
"detector": "token_patterns"
|
||||
},
|
||||
|
||||
// ---- the ONLY untrusted region: agent- or remote-claimed data ----
|
||||
"untrusted": {
|
||||
"reason": "npm install needs registry.npmjs.org" // the agent's stated justification
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
The **trust boundary is a single region, not a split.** Everything outside
|
||||
`untrusted` is trusted by construction — the host established it: schema and
|
||||
chain metadata, both timestamps, the attribution/provenance fields
|
||||
(source-IP → `bottled_agent`/`bottle`/`activation`, `manifest_digest`,
|
||||
`engine`), the host-observed semantics
|
||||
(`actor`/`action`/`resource`/`outcome`), the host-assigned correlation ids,
|
||||
and the typed `payload`. `untrusted` is the **one** place anything an agent
|
||||
or a remote claimed may go (e.g. the agent's free-text `reason`). A reader
|
||||
(or a future policy engine) trusts every field outside `untrusted` for
|
||||
attribution and treats `untrusted.*` — and only `untrusted.*` — as
|
||||
adversarial claims.
|
||||
|
||||
Framing it as "one untrusted region, everything else trusted" removes the
|
||||
mistake where a producer forgets to mark a claimed field: a field is
|
||||
trusted unless it is deliberately placed inside `untrusted`. The
|
||||
construction API enforces this — producers pass trusted fields explicitly
|
||||
and hand all agent/remote-claimed data as the single `untrusted` mapping,
|
||||
so there is no way to emit a top-level field that *looks* authoritative but
|
||||
isn't.
|
||||
|
||||
**Two timestamps** because they answer different questions and can diverge
|
||||
under backpressure: `ts_event` is when the thing happened at the boundary
|
||||
(the proxy saw the connect, the gate saw the push); `ts_recorded` is when
|
||||
the single writer durably appended it. Ordering and the chain use
|
||||
`(epoch, seq)`, never either wall clock. Both are host-set — a bottled
|
||||
agent never supplies a timestamp.
|
||||
|
||||
**The manifest *is* the policy.** bot-bottle has no separate policy
|
||||
artifact — a bottled agent's egress routes and other constraints are
|
||||
declared in its manifest (`bot_bottle/manifest/egress.py`), so
|
||||
`manifest_digest` already pins the ruleset in force; there is no distinct
|
||||
`policy_version`. Given a fixed manifest, the only other thing that can
|
||||
change a decision's outcome is the enforcing code — captured by `engine`
|
||||
(bot-bottle version + git SHA). So two `egress.decision` records with the
|
||||
same `resource` but different `outcome` are explained by exactly one of:
|
||||
different `manifest_digest` (the rules changed) or different `engine` (the
|
||||
enforcer changed). Runtime operator overrides (a supervise `egress-allow`)
|
||||
are themselves audit events, so the effective ruleset at any instant is
|
||||
`manifest_digest` plus the logged, approved deltas — reconstructable from
|
||||
the chain, not from a version stamp.
|
||||
|
||||
**Optionality.** `bottle`/`bottled_agent`/`activation`, `manifest_digest`,
|
||||
and `payload`/`untrusted` are `null`/absent for events that have no such
|
||||
subject (a host-level `hostctl.*` or `audit.*` event has no bottled agent).
|
||||
Absent ≠ empty: a reader distinguishes "no subject" from "unknown". `id`,
|
||||
`type`, the chain fields, both timestamps, `producer`, `host`, `engine`,
|
||||
`actor`, `action`, `outcome`, and `sensitivity` are always present.
|
||||
|
||||
#### Trust provenance of every common field
|
||||
|
||||
| Field | Trust | Set by |
|
||||
|---|---|---|
|
||||
| `v` `id` `type` `epoch` `seq` `segment` `prev` `hash` | trusted | the single writer |
|
||||
| `ts_event` | trusted | emitting host component (boundary) |
|
||||
| `ts_recorded` `ts_mono` | trusted | the single writer |
|
||||
| `producer` `host` `engine` | trusted | the single writer |
|
||||
| `bottle` `bottled_agent` `activation` | trusted | gateway source-IP → slug attribution |
|
||||
| `manifest_digest` | trusted | control plane (the manifest = the policy in force) |
|
||||
| `actor` `action` `resource` `outcome` | trusted | host component that observed/decided it |
|
||||
| `sensitivity` | trusted | registry default for `type`, overridable up (never down) by the producer |
|
||||
| `correlation_id` `causation_id` | trusted | the single writer (assigned as it threads the flow) |
|
||||
| `payload.*` | trusted | emitting host component (shape fixed per `type`) |
|
||||
| `untrusted.*` | **claimed** | copied verbatim from a bottle / gateway / forge / remote |
|
||||
|
||||
Every registry entry (below) restates, per event type, which `payload`
|
||||
keys are required and names any `untrusted` keys it carries — so "trusted
|
||||
vs claimed" is explicit for every event-specific attribute, not just the
|
||||
common ones.
|
||||
|
||||
### Canonical serialization + hash chain
|
||||
|
||||
Serialization is defined once (extends the existing `sha256_hex` /
|
||||
`util.py` helpers):
|
||||
|
||||
```
|
||||
def canonical(event: dict) -> str:
|
||||
return json.dumps(event, sort_keys=True, separators=(",", ":"),
|
||||
ensure_ascii=False)
|
||||
```
|
||||
|
||||
The `hash` field is computed over the canonical form of the event **with
|
||||
`hash` set to `""`**, prefixed by the previous record's hash:
|
||||
|
||||
```
|
||||
digest = sha256_hex(prev_hash + canonical({**event, "hash": ""}))
|
||||
```
|
||||
|
||||
`prev` is the prior record's `hash`; only the first record of a brand-new
|
||||
chain uses `prev = ""`. The first record of a rotated segment carries the
|
||||
preceding segment's head, as specified under *Rotation* below.
|
||||
Two exact rules pin the bytes so the chain is reproducible anywhere:
|
||||
|
||||
1. **Serialize the record with its own `hash` field set to `""`** (present,
|
||||
empty), never omitted — the key set is identical before and after
|
||||
hashing.
|
||||
2. **Digest = `sha256_hex(prev + canonical(record_with_empty_hash))`**,
|
||||
where `prev` is the previous record's `hash` string (`""` at genesis),
|
||||
`+` is string concatenation, and `canonical` is the function above.
|
||||
`ts_mono`, being a float, is serialized by Python's shortest-round-trip
|
||||
`repr` via `json.dumps`; producers therefore emit it as a JSON number
|
||||
they do not post-process. (All other fields are strings/ints/objects,
|
||||
which serialize unambiguously.)
|
||||
|
||||
Editing or deleting record *n* changes its hash, so record *n+1*'s `prev`
|
||||
no longer matches — the break is local and names the tampered record.
|
||||
Verification needs only the journal itself (no keys), so it runs offline
|
||||
and in CI.
|
||||
|
||||
#### Test vectors (normative)
|
||||
|
||||
Two records, reduced to the chain-relevant fields, demonstrate the exact
|
||||
serialization and linkage. An implementation is conformant iff it
|
||||
reproduces these bytes and hashes.
|
||||
|
||||
```
|
||||
# Record 0 — brand-new chain genesis (prev = "")
|
||||
canonical(record0, hash=""):
|
||||
{"hash":"","id":"11111111-1111-4111-8111-111111111111","prev":"","seq":0,"type":"audit.segment_open"}
|
||||
hash0 = sha256("" + canonical) =
|
||||
942ea5729bcac6efdbdea942396bfa574ab0d6ebf5615402595359422f2aeb83
|
||||
|
||||
# Record 1 — chains onto record 0 (prev = hash0)
|
||||
canonical(record1, hash=""):
|
||||
{"hash":"","id":"22222222-2222-4222-8222-222222222222","prev":"942ea5729bcac6efdbdea942396bfa574ab0d6ebf5615402595359422f2aeb83","seq":1,"type":"lifecycle.bottled_agent_start"}
|
||||
hash1 = sha256(hash0 + canonical) =
|
||||
bc082347680405fee50b60a9c304611aa026950b15d869b7e3ae56e1c451b856
|
||||
|
||||
# Tamper check: flip record0.type → recompute →
|
||||
# 4553eda647f33f0c608cfea44be28efbbaca45ed30b873fcbd4405fa5ce737ed
|
||||
# which no longer equals record1.prev (942ea5…) — the break is detected at record1.
|
||||
```
|
||||
|
||||
The implementation PR ships these plus full-envelope vectors (every field
|
||||
populated, and a redaction case) as committed fixtures, so a schema-version
|
||||
bump that changes the bytes fails a golden test loudly.
|
||||
|
||||
### Ordering, idempotency, and duplicate handling
|
||||
|
||||
- **Ordering.** `(epoch, seq)` is a strict total order within one host's
|
||||
native chain and, because
|
||||
the writer is single, a strict order per bottle/activation within that
|
||||
host — satisfying "at least strict causal order per activation/bottle".
|
||||
`causation_id` records the explicit cause edges (a DAG) on top of the
|
||||
total order, so a consumer can reconstruct request → decision → cutoff
|
||||
even if unrelated events interleave between them. There is deliberately
|
||||
no invented total order across imported host chains; a cross-host key is
|
||||
`(host, epoch, seq)`, and consumers use correlation/causation edges where
|
||||
causal ordering across hosts is known.
|
||||
- **Idempotency.** `id` is the idempotency key. A producer that retries an
|
||||
emit (e.g. after a writer restart mid-handoff) **reuses the same `id`**;
|
||||
before appending, the writer checks a durable, host-wide id ledger that
|
||||
spans every segment and survives restart/retention, and drops an id
|
||||
already present. The ledger is operational metadata, not an audit source
|
||||
of truth: after a crash it is reconciled from the journal before appends
|
||||
resume, and retention preserves id tombstones after journal segments are
|
||||
pruned. An in-memory set may cache the ledger but is never the authority.
|
||||
The index additionally has a unique key on `id`; its `UPSERT` is
|
||||
defensive and does not substitute for the pre-append check. Thus a
|
||||
duplicate never enters the source-of-truth journal, double-counts, or
|
||||
forks the chain.
|
||||
- **Deduplication downstream.** Because `id` is stable across export and
|
||||
replay, #324's cursor replay and any cross-host merge dedup on `id` — no
|
||||
consumer needs to invent a second identity.
|
||||
|
||||
### Behavior across rotation, restart, import, truncation
|
||||
|
||||
- **Rotation.** At a segment boundary the writer opens a new segment file,
|
||||
sets its `segment` id, and carries the rotated-out segment's head as the
|
||||
new segment's first `prev` — so the chain is continuous *across* segments
|
||||
(`prev = ""` is reserved for the first record of a brand-new chain) while
|
||||
each file stays independently openable. `verify` walks segments in order
|
||||
and checks the preceding-head-to-first-record link at each seam.
|
||||
- **Restart.** Covered above: read last line → adopt its `hash` as `prev`,
|
||||
bump `epoch`, reset `seq`. The chain never restarts even though the
|
||||
counters do.
|
||||
- **Import.** `audit import <segment>` appends an externally supplied
|
||||
segment (e.g. recovered from another host or a backup). Import verifies
|
||||
the incoming chain in isolation first. A continuation whose first `prev`
|
||||
matches a known head extends that chain. A foreign chain is registered as
|
||||
a separate immutable chain namespace rather than rewriting or grafting
|
||||
its records (which would invalidate their hashes). Imported records keep
|
||||
their original `host`, `id`, `epoch`, and `seq`; the index keys their
|
||||
native order by `(host, epoch, seq)` so attribution is not laundered and
|
||||
tuples from different hosts cannot collide.
|
||||
- **Truncation.** A crash can leave a partial final line; `verify` reports
|
||||
it as `truncated-tail` (recoverable — replay resumes from the last intact
|
||||
record). A chain that ends before a persisted head, or a missing interior
|
||||
`seq`, is reported as `gap`/`missing-suffix` (evidence of deletion, not a
|
||||
clean crash). The two are distinguished so an operator can tell "power
|
||||
loss" from "someone trimmed the log".
|
||||
|
||||
### Single writer; ordering across restarts
|
||||
|
||||
**Decided: one writer per host** (reviewed — the host controller owns it).
|
||||
Producers hand events to the host controller, which is the sole appender,
|
||||
so the chain has one well-defined total order and one `seq`/`epoch`
|
||||
counter. This ties audit availability to the host controller being up,
|
||||
which is acceptable because the host controller already gates every
|
||||
lifecycle transition; per-producer chains are noted only as a future
|
||||
scaling path, not built now.
|
||||
|
||||
**Restarts** are handled by the chain, not the clock. `ts_mono` resets to
|
||||
~0 on every writer start, so it orders events only *within* one boot. On
|
||||
start the writer:
|
||||
|
||||
1. reads the last line of the journal, adopts its `hash` as the next
|
||||
record's `prev` (the chain is continuous across the restart), and
|
||||
2. bumps `epoch` (persisted alongside the chain head) and resets `seq` to
|
||||
0 for the new boot.
|
||||
|
||||
Total order within this host's native chain is therefore `(epoch, seq)` —
|
||||
monotonic across restarts by construction — with
|
||||
`ts_event`/`ts_recorded` for human reading and `ts_mono` for sub-second
|
||||
ordering inside an epoch. Imported chains retain their own
|
||||
`(host, epoch, seq)` order and do not acquire a fictional order relative to
|
||||
the local chain. A crash mid-append truncates at most the last (partial)
|
||||
line; the verifier flags it and replay resumes from the last intact record.
|
||||
|
||||
### Journal (source of truth) + SQLite index (rebuildable)
|
||||
|
||||
- **Journal:** one append-only JSONL file per host (path from `paths.py`,
|
||||
alongside `host_db_path()`), one canonical event per line, opened
|
||||
`O_APPEND`. This is authoritative.
|
||||
- **Index:** a new `audit_events` table via the existing `DbStore` /
|
||||
`TableMigrations` machinery. It is a **derived cache, not a second source
|
||||
of truth**: `audit rebuild` truncates and replays the journal,
|
||||
re-verifying the chain as it goes, so a deleted or drifted DB is
|
||||
regenerated from the journal with no data loss. On a `verify` failure
|
||||
during rebuild it stops and reports rather than indexing past a break.
|
||||
(This supersedes the free-standing `supervise_audit_entries` table, which
|
||||
becomes a producer onto the new index.)
|
||||
|
||||
**Indexable fields** (columns + indices): `ts_event`, `ts_recorded`,
|
||||
`type`, `host`, `bottle`, `bottled_agent`, `activation`, `actor`,
|
||||
`outcome`, `sensitivity`, `correlation_id`, `causation_id`, plus two
|
||||
event-specific projections promoted out of `payload` for query —
|
||||
`repository` and `commit_sha` (populated for `forge.*`/`commit.*`, null
|
||||
otherwise). The unique event id and native-order index are respectively
|
||||
`id` and `(host, epoch, seq)`; a local-only `ingest_seq` provides stable
|
||||
display order when a query intentionally mixes chains without pretending
|
||||
that it is causal order. The full canonical record is stored verbatim in a
|
||||
`raw` column so the index never loses fidelity to the journal.
|
||||
|
||||
**Local query surface** — `audit query`, no egress, no paid platform:
|
||||
|
||||
```
|
||||
audit query \
|
||||
[--since T] [--until T] [--type egress.*] [--host H] [--bottle B] \
|
||||
[--activation A] [--agent SLUG] [--actor ID] [--outcome blocked] \
|
||||
[--repository R] [--correlation-id C] [--commit SHA] \
|
||||
[--follow BOTTLE] # one host chain's events in (epoch, seq) order
|
||||
[--json | --table]
|
||||
audit verify [--segment S] # offline chain check; exit non-zero on any break
|
||||
audit rebuild # drop + replay journal → index
|
||||
audit import <segment> # graft an external segment (see above)
|
||||
```
|
||||
|
||||
Type filters accept a `group.*` glob. A read-only local HTTP endpoint
|
||||
mirrors the same filters for the future review console; both are pure reads
|
||||
over the index and can never mutate the journal.
|
||||
|
||||
### Event registry (initial)
|
||||
|
||||
Dotted `type` names, grouped. The registry is a table mapping each type to
|
||||
its required `payload` keys, its `untrusted` keys (if any), a default
|
||||
`sensitivity`, and its correlation behavior — so producers and the verifier
|
||||
agree on shape and "trusted vs claimed" is pinned per type. Initial
|
||||
coverage (the issue's mandated set):
|
||||
|
||||
| Group / type | Producer | Required `payload` (trusted) | `untrusted` | Default sensitivity |
|
||||
|---|---|---|---|---|
|
||||
| **lifecycle.*** — `bottled_agent_start` / `_stop` / `_crash` | host-controller (#468) | `manifest_digest`, `exit` (for stop/crash) | — | normal |
|
||||
| **hostctl.*** — `broker_launch`, `broker_teardown`, `broker_reject` | host-controller (#468) | `op`, `request_digest` | — | security |
|
||||
| **decision.*** — `proposed`, `resolved` | supervise | `tool`, `operator_action`, `justification`, `diff_digest` | `agent_rationale` | security |
|
||||
| **egress.*** — `request`, `decision`, `cutoff`, `anomaly` | egress proxy | `route_id`, `detector` (on match), `bytes` (cutoff) | `reason`, `target_claimed` | security |
|
||||
| **forge.*** — `push_accepted`, `push_rejected`, `pr_opened` | git-gate | `repository`, `ref`, `gitleaks_result` | `title`, `description` | security |
|
||||
| **commit.signed** (#480) | git-gate | `repository`, `commit_sha`, `activation_key_id`, `signature_ref` | `commit_message` | security |
|
||||
| **auth.*** — `token_minted`, `token_rejected`, `authz_denied` | control plane | `role`, `token_id`, `reason_code` | — | security |
|
||||
| **audit.*** — `segment_open`, `verify_failed`, `truncation_detected`, `export_failed` | audit writer/verifier | `segment`, `detail` | — | security |
|
||||
|
||||
Notes:
|
||||
- **`egress.request` vs `egress.decision`** share a `correlation_id`; the
|
||||
`decision`'s `causation_id` points at the `request`, and a later `cutoff`
|
||||
chains onto the `decision` — so a flow is reconstructable.
|
||||
- **`audit.*` self-events** make the audit subsystem audit itself: a failed
|
||||
verification, a detected truncation, or a dropped export is itself a
|
||||
chained, tamper-evident record — you cannot silence the alarm without
|
||||
breaking the chain that carries it.
|
||||
- **Free-text and remote-echoed fields are always `untrusted`** (`reason`,
|
||||
`agent_rationale`, PR `title`/`description`, `target_claimed`), because
|
||||
they originate in the bottle or a remote response; the host-observed
|
||||
counterpart (`resource`, `outcome`, `gitleaks_result`) is the trusted
|
||||
fact.
|
||||
|
||||
**Sensitivity + redaction per type.** Every type's default `sensitivity`
|
||||
is listed above; a producer may raise it (never lower it). `restricted`
|
||||
events keep their `payload` in the journal but the export projection ships
|
||||
only the envelope + a payload digest unless the consumer is authorized —
|
||||
so a `security`/`restricted` record is still counted and correlated
|
||||
downstream without leaking its body. The credential-shape redaction rules
|
||||
(next) apply to **every** type regardless of sensitivity.
|
||||
|
||||
**Schema evolution & backward-compatible readers.** The registry is
|
||||
append-only: **adding** a type, an optional `payload` key, or an
|
||||
`untrusted` key does **not** bump `v`; readers ignore unknown fields
|
||||
(forward-compatible) and treat absent optional fields as `null`.
|
||||
**Removing** or **re-typing** a field, or making an optional field
|
||||
required, bumps `v`. A reader declares the max `v` it understands and
|
||||
refuses to *interpret* a higher-`v` record, but the **verifier is
|
||||
version-agnostic** — the hash covers whatever fields exist, so chain
|
||||
integrity is checkable across versions without understanding semantics.
|
||||
Every `v` bump ships a migration note and updated golden vectors.
|
||||
|
||||
### Redaction rule
|
||||
|
||||
Redaction runs at the envelope boundary, before a record is written, in two
|
||||
layers:
|
||||
|
||||
1. **Key deny-list (structural).** A field whose *key* matches a known
|
||||
credential shape (`token`, `secret`, `password`, `authorization`,
|
||||
`*_key`) is refused — the producer must pass a reference (a token *id*
|
||||
or `sha256` fingerprint), never the raw value. `auth.token_minted`
|
||||
therefore records the token id and role, not the JWT. This is the
|
||||
primary guard: it is cheap, deterministic, and catches the intended
|
||||
mistake (a producer stuffing a credential into a named field).
|
||||
|
||||
2. **Value scan — reuse the egress DLP detectors.** Per review, the value
|
||||
layer reuses the *same* deterministic credential-shape detectors the
|
||||
egress proxy already ships:
|
||||
`bot_bottle/gateway/egress/dlp_detectors.py` —
|
||||
`scan_token_patterns` / `redact_tokens` (and `scan_known_secrets` for
|
||||
host-known secret material). They are pure-Python, mitmproxy-free, and
|
||||
already the project's source of truth for "what a leaked credential
|
||||
looks like," so a single detector set governs both what may leave over
|
||||
the wire and what may land in the journal — they can't drift apart.
|
||||
|
||||
**Scoped deliberately:** only the pattern/known-secret detectors are
|
||||
reused, **not** `scan_entropy`. Entropy scoring is tuned for large
|
||||
streamed request bodies; on the short, high-entropy structured values an
|
||||
audit event legitimately carries (hashes, uuids, base64 ids) it would
|
||||
false-positive and start redacting the very fingerprints the log needs.
|
||||
So the shared layer is the deterministic detectors; entropy stays an
|
||||
egress-only concern. (This is the "evaluate how reasonable that is" from
|
||||
review: reuse the deterministic detectors — yes; share the entropy
|
||||
heuristic — no.)
|
||||
|
||||
On a value-layer match the default is **redact** (scrub to a placeholder
|
||||
and keep the event) rather than drop, so a producer bug can never make an
|
||||
audit event vanish; the key deny-list stays a hard refusal because a
|
||||
credential in a named field is always a producer bug worth surfacing.
|
||||
|
||||
**No raw-payload capture by default.** The envelope carries *decisions and
|
||||
metadata*, not traffic. Prompts, model responses, request/response bodies,
|
||||
and file contents are **not** recorded unless a producer opts a specific,
|
||||
reviewed field in — and such a field is `untrusted` and subject to both
|
||||
redaction layers. This keeps the audit log from becoming a covert copy of
|
||||
the very data the sandbox exists to contain (the issue's "unsafe payload
|
||||
capture" non-goal).
|
||||
|
||||
### Export / interoperability (CloudEvents, OpenTelemetry Logs)
|
||||
|
||||
#487 requires the envelope to map onto the **OpenTelemetry Logs data
|
||||
model** and/or a **CloudEvents JSON** envelope *without losing integrity or
|
||||
attribution semantics*. The flattened shape (one `untrusted` region,
|
||||
everything else trusted at top level) does **not** conflict with either — it
|
||||
maps *more* cleanly than a nested `trusted`/`untrusted` pair would, because
|
||||
both target models expect a flat set of top-level fields plus one payload
|
||||
subtree.
|
||||
|
||||
**CloudEvents.** Context attributes MUST be scalar simple types — a map
|
||||
cannot be a context attribute — so a nested `trusted` block would have had
|
||||
to be flattened for CloudEvents anyway. Our flat top level maps directly:
|
||||
`id`→`id`, `type`→`type`, `producer`+`host`→`source`,
|
||||
`bottled_agent`→`subject`, `ts_event`→`time` (`ts_recorded` as an extension); the integrity/chain fields
|
||||
(`epoch`, `seq`, `prev`, `hash`, `v`) ride as **extension attributes**
|
||||
(scalars — legal). The `untrusted` map goes in `data`. Only mechanical
|
||||
transform needed: extension attribute names must be lowercase-alphanumeric,
|
||||
so `bottled_agent`/`ts_mono`/etc. are renamed at export (e.g. a
|
||||
`botbottle`-prefixed form) — a naming rule, not a schema conflict.
|
||||
|
||||
**OpenTelemetry Logs.** `ts_event`→`Timestamp`, `ts_recorded`→`ObservedTimestamp`; `type`→the `event.name`
|
||||
attribute; the flat trusted fields → `Attributes` under a `botbottle.*`
|
||||
namespace (`botbottle.bottled_agent`, `botbottle.producer`,
|
||||
`botbottle.chain.hash`, …); `untrusted.*` → `Attributes` under
|
||||
`botbottle.untrusted.*` (or `Body`). OTel attributes are a dotted map that
|
||||
happily carries the nested subtree.
|
||||
|
||||
**Attribution is preserved** precisely because the boundary is now
|
||||
structural: on export, top-level fields become trusted context/attributes
|
||||
and the `untrusted` subtree stays a single, clearly-named region — so a
|
||||
downstream consumer still sees exactly which fields an agent claimed.
|
||||
Nothing agent-claimed is promoted to a trusted-looking position.
|
||||
|
||||
**Integrity has one deliberate caveat.** CloudEvents/OTel are
|
||||
representation envelopes with their own (or no) canonicalization; `hash`
|
||||
and `prev` are computed over **our** canonical JSON, not over the exported
|
||||
form. So the chain fields travel *as data* for reference, but
|
||||
tamper-evidence is always verified against the **native journal** (the
|
||||
source of truth) — never re-derived from an exported CloudEvents/OTel
|
||||
record, whose key ordering / number formatting the exporter may change.
|
||||
Export is thus a lossless-for-attribution **projection** that carries the
|
||||
integrity fields along; verification stays on the canonical journal. This
|
||||
satisfies "without losing integrity or attribution semantics": both are
|
||||
carried, neither is *relied upon* in the foreign format.
|
||||
|
||||
The export adapters themselves are follow-up implementation — this PRD
|
||||
fixes the *schema* so that projection is a field re-map, never a reformat.
|
||||
|
||||
#### The #324 delivery contract (payload, cursor, backpressure)
|
||||
|
||||
#324 transports events off-box; it must not invent a second envelope. This
|
||||
PRD fixes the contract it depends on:
|
||||
|
||||
- **Payload.** #324 ships the **native canonical record verbatim** (the
|
||||
exact bytes the hash covers), optionally wrapped in the CloudEvents
|
||||
projection whose `data` *is* that record. Either way the integrity fields
|
||||
travel intact and the receiver can verify against the same bytes.
|
||||
- **Cursor.** Export maintains one cursor per native host chain:
|
||||
`(host, epoch, seq)` (equivalently that chain's last exported `hash`).
|
||||
It advances **only on acknowledgement**, so delivery is at-least-once and
|
||||
gap-free; a crash re-sends from the last acked cursor. Imported foreign
|
||||
chains use independent cursors and never share a bare `(epoch, seq)`
|
||||
namespace with the local chain.
|
||||
- **Idempotency / replay.** Dedup is on `id` (stable across replay), so
|
||||
at-least-once delivery is safe — the receiver collapses re-sends.
|
||||
- **Backpressure.** The outbox is the journal itself plus a cursor; when
|
||||
the endpoint is slow the cursor simply lags — the writer never blocks on
|
||||
export, and audit never applies backpressure to the data plane it
|
||||
records.
|
||||
- **Retention interaction.** Retention/rotation **must not** prune a
|
||||
segment whose records are still behind the export cursor; the reaper
|
||||
honors `min(cursor)` across all configured consumers. (The schedule
|
||||
itself stays the retention follow-up; this is the *ordering* constraint
|
||||
that follow-up must respect.)
|
||||
|
||||
#### #480 signed-commit attribution maps in without weakening it
|
||||
|
||||
#480 binds a commit's bytes to a per-activation signing key. It maps to the
|
||||
`commit.signed` event: `payload` carries `repository`, `commit_sha`,
|
||||
`activation_key_id`, and a `signature_ref` (the detached-signature
|
||||
location or its digest) — **not** the private key and not a re-derived
|
||||
signature. The audit event therefore *references and timestamps* #480's
|
||||
existing byte-to-activation-key proof inside the tamper-evident chain; it
|
||||
does not re-implement or replace it, so #480's guarantee is unweakened —
|
||||
the signature still verifies against the commit bytes independently, and
|
||||
the audit record adds only "this binding was observed at this point in the
|
||||
chain". The trusted `actor`/`activation` fields and the `commit_sha`
|
||||
payload are host-observed at the gate, so attribution cannot be forged by
|
||||
the committing agent.
|
||||
|
||||
## Implementation chunks
|
||||
|
||||
1. **(this PR — PRD only.)** The contract above. No code; scheduled to land
|
||||
right after #468.
|
||||
2. **Envelope + canonical + chain core.** `AuditEvent` dataclass,
|
||||
`canonical()`, chain hashing, and the single-writer journal appender in
|
||||
`bot_bottle/store/` (reusing `sha256_hex`); redaction wired to the
|
||||
existing `gateway/egress/dlp_detectors` (`scan_token_patterns` /
|
||||
`redact_tokens`); embed the git SHA at build so `engine` is populated
|
||||
(only `version = "0.1.0"` exists in `pyproject.toml` today — the build
|
||||
must stamp the SHA); unit tests for determinism, chain-break detection,
|
||||
`epoch`/`seq` continuity across a simulated restart, and redaction of
|
||||
both a deny-listed key and a token-shaped value.
|
||||
3. **SQLite index + `audit` CLI.** New `audit_events` migration (indexable
|
||||
fields above); replay-from-journal; offline chain verifier
|
||||
(`truncated-tail` vs `gap`); `query` / `verify` / `rebuild` / `import`;
|
||||
idempotent `UPSERT` by `id`.
|
||||
4. **Host controller as first producer (#468).** Wire
|
||||
`lifecycle.bottled_agent_*` and `hostctl.*` emission into the host
|
||||
controller; establish the `epoch` bump + chain-head carry + segment
|
||||
rotation on writer restart here (it owns the single writer).
|
||||
5. **Migrate existing producers.** Re-emit supervise `decision.*` (retiring
|
||||
the standalone `supervise_audit_entries` shape behind the index), egress
|
||||
`egress.*`, git-gate `forge.*` + `commit.signed` (#480), control-plane
|
||||
`auth.*`; add the `audit.*` self-events (verify/truncation/export
|
||||
failure).
|
||||
6. **CloudEvents / OTel export adapters + #324 delivery.** Projection layer
|
||||
(field re-map per *Export / interoperability*) plus the outbox cursor,
|
||||
ack-driven advance, and retention-ordering guard the #324 contract
|
||||
specifies.
|
||||
7. **(follow-up.)** Cross-host merge transport; per-writer signing +
|
||||
external anchoring on the chain head; retention/rotation *schedule*.
|
||||
|
||||
## Acceptance-criteria coverage (#487)
|
||||
|
||||
The issue defines the contract; implementation is explicitly split into
|
||||
follow-up PRs. This PRD is the durable decision record; each acceptance box
|
||||
maps to a section:
|
||||
|
||||
| #487 acceptance criterion | Where |
|
||||
|---|---|
|
||||
| Durable PRD defines versioned envelope + initial registry | *The envelope*, *Event registry* |
|
||||
| Canonical JSON + hash-chain rules, unambiguous, with test vectors | *Canonical serialization + hash chain* → *Test vectors* |
|
||||
| Trust provenance explicit for every common + event-specific field | *Trust provenance of every common field*; per-type `untrusted` in *Event registry* |
|
||||
| Redaction prohibits credentials / raw secrets / unsafe capture by default | *Redaction rule*; `untrusted`-only claims; sensitivity classes |
|
||||
| JSONL journal canonical; SQLite index fully rebuildable | *Journal + SQLite index* (`audit rebuild`) |
|
||||
| Minimum local search/query contract | *Journal + SQLite index* → *Local query surface* |
|
||||
| #324 can transport/replay without a second envelope | *The #324 delivery contract* |
|
||||
| #480 maps in without weakening its byte-to-activation-key guarantee | *#480 signed-commit attribution maps in…* |
|
||||
| Schema evolution + backward-compatible readers | *Schema evolution & backward-compatible readers* |
|
||||
| Integrity detects modification / deletion / reorder / bad continuation | *Test vectors* (tamper), *Behavior across rotation…truncation*, `audit verify` |
|
||||
|
||||
Two acceptance items are **specified here, implemented later** by design
|
||||
(the issue permits this): the concrete test-vector *fixtures* and the
|
||||
`audit` CLI land in impl chunks 2–3; the #324 outbox lands in chunk 6.
|
||||
Nothing in the contract is left undefined — only its code is deferred.
|
||||
|
||||
## Resolved in review (#495)
|
||||
|
||||
- **Single writer per host — decided.** The host controller owns the sole
|
||||
appender; per-producer chains are a future scaling path only. (Design →
|
||||
*Single writer; ordering across restarts*.)
|
||||
- **Restarts — decided.** An `epoch` counter (bumped per writer boot) plus
|
||||
carrying the last chain head as the next `prev` gives a total order of
|
||||
`(epoch, seq)` that survives restarts; `ts_mono` orders only within an
|
||||
epoch. (Design → *ordering across restarts*.)
|
||||
- **Flatten to one `untrusted` region — decided.** Everything outside
|
||||
`untrusted` (chain metadata, `producer`/`host`, `bottled_agent`, `ts_*`)
|
||||
is trusted by construction, so the separate `trusted` sub-block is
|
||||
removed; a field is trusted unless deliberately placed under `untrusted`.
|
||||
(Design → *The envelope*.)
|
||||
- **Subject term is `bottled_agent` everywhere** — the top-level field and
|
||||
the `lifecycle.bottled_agent_*` leaf names. (Design → *The envelope* /
|
||||
*Event registry*.)
|
||||
- **Retention head-carry — yes.** When a journal segment is rotated out,
|
||||
the new segment's first record carries the rotated-out head as `prev`, so
|
||||
the verifier still trusts the current head across a rotation. (Folds into the
|
||||
retention follow-up.)
|
||||
- **Redaction reuses the egress detectors — yes, scoped.** Reuse the
|
||||
deterministic `dlp_detectors` (`scan_token_patterns` / `redact_tokens` /
|
||||
`scan_known_secrets`); exclude `scan_entropy` as brittle on the short,
|
||||
high-entropy structured values audit records carry. (Design → *Redaction
|
||||
rule*.)
|
||||
|
||||
## Open questions
|
||||
|
||||
- **Value-scan cost on the hot path.** The single writer runs the reused
|
||||
detectors on every event's `untrusted` block inline. Is that cheap enough
|
||||
at lifecycle-event volume, or should the value scan move to index-build
|
||||
time (journal stays raw, index stores the redacted view)? Leaning inline
|
||||
so the raw journal never contains a leaked value in the first place.
|
||||
- **`epoch` persistence location.** Store the per-writer `epoch` + chain
|
||||
head in the SQLite index (rebuildable, but then the writer needs the DB
|
||||
at boot) or in a tiny sidecar file next to the journal (independent of
|
||||
the index)? Leaning sidecar, so the writer can start and append without
|
||||
the index present.
|
||||
@@ -1,273 +0,0 @@
|
||||
# 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.
|
||||
@@ -1,117 +0,0 @@
|
||||
"""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")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -1,279 +0,0 @@
|
||||
"""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 io
|
||||
import json
|
||||
import secrets
|
||||
import threading
|
||||
import unittest
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
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_from_env,
|
||||
dispatch,
|
||||
main,
|
||||
make_host_server,
|
||||
)
|
||||
|
||||
|
||||
def _body(obj: object) -> bytes:
|
||||
return json.dumps(obj).encode()
|
||||
|
||||
|
||||
class _RaisingBroker(LaunchBroker):
|
||||
"""A broker whose backend launch 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")
|
||||
|
||||
|
||||
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_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 TestBrokerSecretFromEnv(unittest.TestCase):
|
||||
def test_reads_hex_secret(self) -> None:
|
||||
s = secrets.token_bytes(16)
|
||||
self.assertEqual(s, broker_secret_from_env({"BOT_BOTTLE_BROKER_SECRET": s.hex()}))
|
||||
|
||||
def test_unset_is_none(self) -> None:
|
||||
self.assertIsNone(broker_secret_from_env({}))
|
||||
|
||||
def test_invalid_hex_is_none(self) -> None:
|
||||
self.assertIsNone(broker_secret_from_env({"BOT_BOTTLE_BROKER_SECRET": "not-hex"}))
|
||||
|
||||
|
||||
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
|
||||
|
||||
|
||||
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."""
|
||||
|
||||
def _base(self) -> str:
|
||||
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[0], server.server_address[1]
|
||||
return f"http://{host}:{port}"
|
||||
|
||||
def test_oversized_body_is_rejected_before_acting(self) -> None:
|
||||
base = self._base()
|
||||
big = b"x" * (MAX_BODY_BYTES + 1)
|
||||
req = urllib.request.Request(
|
||||
f"{base}/broker", data=big, method="POST",
|
||||
headers={"Content-Type": "application/json"})
|
||||
with self.assertRaises(urllib.error.HTTPError) as cm:
|
||||
urllib.request.urlopen(req, timeout=5)
|
||||
self.assertEqual(413, cm.exception.code)
|
||||
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_from_env",
|
||||
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_from_env",
|
||||
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()
|
||||
@@ -1,64 +0,0 @@
|
||||
"""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
|
||||
|
||||
|
||||
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("BOT_BOTTLE_BROKER_SECRET", 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_secret_serves(self) -> None:
|
||||
rc, _ = self._run(
|
||||
"http", env={"BOT_BOTTLE_BROKER_SECRET": secrets.token_bytes(16).hex()})
|
||||
self.assertEqual(0, rc)
|
||||
|
||||
def test_http_broker_without_secret_exits(self) -> None:
|
||||
# Fail-closed: --broker http with no shared secret is a usage error.
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
with patch.dict("os.environ", {}, clear=False):
|
||||
os.environ.pop("BOT_BOTTLE_BROKER_SECRET", None)
|
||||
with self.assertRaises(SystemExit):
|
||||
main(["--db", str(Path(d) / "r.db"), "--broker", "http"])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -11,12 +11,7 @@ from contextlib import closing
|
||||
from pathlib import Path
|
||||
from unittest.mock import patch
|
||||
|
||||
from bot_bottle.orchestrator.broker import (
|
||||
BrokerUnavailableError,
|
||||
LaunchBroker,
|
||||
LaunchRequest,
|
||||
StubBroker,
|
||||
)
|
||||
from bot_bottle.orchestrator.broker import 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
|
||||
@@ -30,8 +25,8 @@ from bot_bottle.orchestrator.supervisor import (
|
||||
|
||||
|
||||
class _FailingBroker(LaunchBroker):
|
||||
"""Verifies the token like any broker, then fails the launch *definitely* —
|
||||
to exercise the orchestrator's registry rollback."""
|
||||
"""Verifies the token like any broker, then fails the launch — to
|
||||
exercise the orchestrator's registry rollback."""
|
||||
|
||||
def _launch(self, req: LaunchRequest) -> None:
|
||||
raise RuntimeError("launch failed")
|
||||
@@ -40,18 +35,6 @@ class _FailingBroker(LaunchBroker):
|
||||
pass
|
||||
|
||||
|
||||
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
|
||||
|
||||
|
||||
class TestOrchestrator(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
@@ -153,20 +136,11 @@ 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_definite_broker_failure(self) -> None:
|
||||
def test_launch_rolls_back_registry_on_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 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
|
||||
self.assertEqual([], self.store.all()) # no orphan
|
||||
|
||||
def test_gateway_status_reports_unconfigured(self) -> None:
|
||||
# The orchestrator no longer owns a standalone gateway lifecycle; the
|
||||
|
||||
Reference in New Issue
Block a user