Compare commits

...

15 Commits

Author SHA1 Message Date
didericis-codex 3dbf1780b4 docs(prd): require shared storage permissions
prd-number-check / require-numbered-prds (pull_request) Successful in 11s
test / image-input-builds (pull_request) Successful in 45s
test / integration-docker (pull_request) Successful in 1m8s
test / unit (pull_request) Successful in 2m30s
test / coverage (pull_request) Successful in 25s
tracker-policy-pr / check-pr (pull_request) Failing after 12m49s
2026-07-27 04:45:45 +00:00
didericis-codex ff4da6f41e docs(prd): bound heavy gateway operations
prd-number-check / require-numbered-prds (pull_request) Successful in 10s
test / unit (pull_request) Successful in 54s
test / coverage (pull_request) Has been skipped
test / integration-docker (pull_request) Has been cancelled
test / image-input-builds (pull_request) Successful in 53s
tracker-policy-pr / check-pr (pull_request) Failing after 14m19s
2026-07-27 04:16:56 +00:00
didericis-codex 3bb90da11c docs(prd): require cleanup execution integrity
test / image-input-builds (pull_request) Failing after 14m11s
test / unit (pull_request) Failing after 14m19s
prd-number-check / require-numbered-prds (pull_request) Failing after 14m25s
test / integration-docker (pull_request) Has been cancelled
test / coverage (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Failing after 12m36s
2026-07-27 03:55:18 +00:00
didericis-codex 0146450951 docs(prd): extend authoritative cleanup boundaries
prd-number-check / require-numbered-prds (pull_request) Successful in 13s
test / image-input-builds (pull_request) Successful in 3m18s
test / unit (pull_request) Successful in 57s
test / integration-docker (pull_request) Successful in 1m9s
test / coverage (pull_request) Successful in 19s
tracker-policy-pr / check-pr (pull_request) Failing after 10m30s
2026-07-27 03:31:46 +00:00
didericis-codex ecaf23cdb5 docs(prd): define authoritative failure boundaries
prd-number-check / require-numbered-prds (pull_request) Successful in 9s
test / unit (pull_request) Successful in 52s
test / integration-docker (pull_request) Waiting to run
test / coverage (pull_request) Blocked by required conditions
test / image-input-builds (pull_request) Successful in 2m42s
tracker-policy-pr / check-pr (pull_request) Failing after 12s
2026-07-27 02:58:07 +00:00
didericis-codex 6e46a9b191 fix(orchestrator): satisfy adapter type contracts
prd-number-check / require-numbered-prds (pull_request) Successful in 12s
lint / lint (push) Successful in 53s
test / image-input-builds (pull_request) Successful in 1m2s
test / integration-docker (pull_request) Successful in 58s
test / unit (pull_request) Successful in 46s
test / coverage (pull_request) Successful in 21s
tracker-policy-pr / check-pr (pull_request) Failing after 7s
2026-07-27 02:19:11 +00:00
didericis-codex a59e495faa refactor(egress): split request policy pipeline stages 2026-07-27 02:19:11 +00:00
didericis-codex b8818948a0 feat(firecracker): enumerate running bottles 2026-07-27 02:19:11 +00:00
didericis-codex b09952045a fix(orchestrator): bound unauthenticated HTTP requests 2026-07-27 02:19:11 +00:00
didericis-codex de192359ee fix(security): authenticate persisted egress secrets 2026-07-27 02:19:11 +00:00
didericis-codex 7d9933edc0 fix(docker): fail closed on network address scan errors 2026-07-27 02:19:11 +00:00
didericis-codex 2bc9ef8ec0 fix(firecracker): abort cleanup on scan errors 2026-07-27 02:19:11 +00:00
didericis-codex 47b6bead69 refactor(backend): type enumeration failures 2026-07-27 02:19:11 +00:00
didericis-codex 15ecada022 fix(docker): surface enumeration failures 2026-07-27 02:19:11 +00:00
didericis-codex e2222bd96b fix(security): prohibit unauthenticated orchestrator 2026-07-27 02:19:11 +00:00
25 changed files with 774 additions and 164 deletions
+3
View File
@@ -30,6 +30,7 @@ if TYPE_CHECKING:
BottleImages,
BottlePlan,
BottleSpec,
EnumerationError,
ExecResult,
)
from .selection import (
@@ -59,6 +60,7 @@ _LAZY_MODULES: dict[str, str] = {
"BottleImages": "base",
"BottleBackend": "base",
"BackendStatus": "base",
"EnumerationError": "base",
"get_bottle_backend": "selection",
"known_backend_names": "selection",
"has_backend": "selection",
@@ -100,6 +102,7 @@ __all__ = [
"BottlePlan",
"BottleSpec",
"ExecResult",
"EnumerationError",
"CommitCancelled",
"Freezer",
"get_freezer",
+4
View File
@@ -42,6 +42,10 @@ class BackendStatus(enum.IntEnum):
READY = 0
class EnumerationError(RuntimeError):
"""A backend could not produce an authoritative live-resource snapshot."""
@dataclass(frozen=True)
class BottleSpec:
"""CLI-supplied intent. Backend-agnostic — each backend's prepare
+27 -13
View File
@@ -16,6 +16,7 @@ from pathlib import Path
from typing import Any
from ...log import die, warn
from ..base import EnumerationError
# --- Lifecycle helpers (PRD 0018 chunk 3) ----------------------------------
@@ -52,19 +53,20 @@ def slug_from_compose_project(project: str) -> str:
def list_compose_projects(
*, include_stopped: bool = True, warn_on_error: bool = True,
*,
include_stopped: bool = True,
warn_on_error: bool = True,
raise_on_error: bool = False,
) -> list[str]:
"""All compose project names starting with `bot-bottle-`.
`include_stopped=True` (default) runs `docker compose ls --all`
so exited projects appear too; pass False to get only projects
with at least one running container.
Returns [] on docker daemon errors or malformed output rather
than raising — callers should treat the empty list as "no
projects discoverable", not "no projects exist". `warn_on_error`
stays true for explicit operator commands like cleanup, but active
discovery paths set it false so dashboard refreshes don't spam
stderr while Docker Desktop is stopped."""
Best-effort callers get ``[]`` on Docker errors or malformed output.
Enumeration callers pass ``raise_on_error=True`` so a failed query is not
reported as an authoritative empty result.
"""
argv = ["docker", "compose", "ls", "--format", "json"]
if include_stopped:
argv.insert(3, "--all")
@@ -72,19 +74,27 @@ def list_compose_projects(
result = subprocess.run(
argv, capture_output=True, text=True, check=False,
)
except FileNotFoundError:
# docker binary not on PATH — same shape as a daemon-down
# error from the caller's POV: no projects discoverable.
except FileNotFoundError as exc:
if raise_on_error:
raise EnumerationError(
"docker compose ls failed: docker not found"
) from exc
return []
if result.returncode != 0:
message = f"docker compose ls failed: {result.stderr.strip()}"
if raise_on_error:
raise EnumerationError(message)
if warn_on_error:
warn(f"docker compose ls failed: {result.stderr.strip()}")
warn(message)
return []
try:
projects = json.loads(result.stdout or "[]")
except json.JSONDecodeError as e:
message = f"docker compose ls returned malformed JSON: {e}"
if raise_on_error:
raise EnumerationError(message) from e
if warn_on_error:
warn(f"docker compose ls returned malformed JSON: {e}")
warn(message)
return []
names: list[str] = []
for p in projects:
@@ -97,7 +107,10 @@ def list_compose_projects(
def list_active_slugs(
*, include_stopped: bool = False, warn_on_error: bool = True,
*,
include_stopped: bool = False,
warn_on_error: bool = True,
raise_on_error: bool = False,
) -> list[str]:
"""Slugs (project name minus prefix) of currently-running
bottles. Used by the dashboard's operator-edit verbs to choose
@@ -108,6 +121,7 @@ def list_active_slugs(
for p in list_compose_projects(
include_stopped=include_stopped,
warn_on_error=warn_on_error,
raise_on_error=raise_on_error,
)
) if slug
)
@@ -68,6 +68,11 @@ def _network_container_ips(network: str) -> list[str]:
"docker", "network", "inspect", "--format",
"{{range .Containers}}{{.IPv4Address}} {{end}}", network,
])
if proc.returncode != 0:
detail = proc.stderr.strip() or f"exit {proc.returncode}"
raise ConsolidatedLaunchError(
f"could not inspect addresses on gateway network {network}: {detail}"
)
ips: list[str] = []
for entry in proc.stdout.split():
ips.append(entry.split("/", 1)[0])
+12 -12
View File
@@ -1,9 +1,8 @@
"""Active-agent enumeration for the docker backend.
Returns `ActiveAgent` records the CLI `active` command and the
dashboard agents pane consume. Empty when docker isn't reachable
— gated by `has_backend('docker')` at the cross-backend caller
so this module trusts that docker is available when called.
dashboard agents pane consume. Docker query failures raise rather
than masquerading as an authoritative empty result.
The parser (`_parse_services_by_project`) is exposed for direct
unit testing; the docker `docker ps` invocation is in
@@ -13,17 +12,18 @@ from __future__ import annotations
import subprocess
from .. import ActiveAgent
from .. import ActiveAgent, EnumerationError
from ...bottle_state import read_metadata
from .compose import compose_project_name, list_active_slugs
def enumerate_active() -> list[ActiveAgent]:
"""All currently-running docker-backed agents. Caller is
responsible for gating on `has_backend('docker')` if it
matters; if docker is missing the `docker ps` call below
returns an empty list silently."""
slugs = list_active_slugs(include_stopped=False, warn_on_error=False)
"""All currently-running docker-backed agents."""
slugs = list_active_slugs(
include_stopped=False,
warn_on_error=False,
raise_on_error=True,
)
if not slugs:
return []
services_by_project = _query_services_by_project()
@@ -74,8 +74,8 @@ def _query_services_by_project() -> dict[str, set[str]]:
],
capture_output=True, text=True, check=False,
)
except FileNotFoundError:
return {}
except FileNotFoundError as exc:
raise EnumerationError("docker ps failed: docker not found") from exc
if r.returncode != 0:
return {}
raise EnumerationError(f"docker ps failed: {r.stderr.strip()}")
return _parse_services_by_project(r.stdout or "")
+17 -5
View File
@@ -29,6 +29,7 @@ import subprocess
from pathlib import Path
from ...log import info
from .. import EnumerationError
from . import util
from .bottle_cleanup_plan import FirecrackerBottleCleanupPlan
@@ -62,12 +63,23 @@ def _scan_processes(run_root: Path) -> tuple[set[str], list[int]]:
* ``orphan_pids`` — firecracker pids whose run dir no longer exists
(a lingering VMM to kill).
"""
result = subprocess.run(
["pgrep", "-a", "firecracker"],
capture_output=True, text=True, check=False,
)
if result.returncode != 0:
try:
result = subprocess.run(
["pgrep", "-a", "firecracker"],
capture_output=True, text=True, check=False,
)
except OSError as exc:
raise EnumerationError(
f"could not enumerate Firecracker processes: {exc}"
) from exc
if result.returncode == 1:
# pgrep's documented "no processes matched" result.
return set(), []
if result.returncode != 0:
detail = (result.stderr or "").strip() or f"exit {result.returncode}"
raise EnumerationError(
f"could not enumerate Firecracker processes: {detail}"
)
live: set[str] = set()
orphan_pids: list[int] = []
for line in result.stdout.splitlines():
+22 -4
View File
@@ -1,14 +1,32 @@
"""Active-agent enumeration for the Firecracker backend.
The backend is disabled during the companion-container removal (#385) — it can't
launch bottles, so there are none to enumerate. Real enumeration returns
with the backend's consolidated relaunch (#354).
Running bottles are the Firecracker processes whose ``--config-file`` points
at an existing per-bottle run directory. The same authoritative process scan
protects cleanup from deleting live VMs; operational scan failures propagate
as ``EnumerationError`` instead of masquerading as an empty host.
"""
from __future__ import annotations
from ...bottle_state import read_metadata
from .. import ActiveAgent
from .cleanup import live_run_dirs
def enumerate_active() -> list[ActiveAgent]:
return []
out: list[ActiveAgent] = []
for run_dir in live_run_dirs():
slug = run_dir.name
metadata = read_metadata(slug)
out.append(ActiveAgent(
backend_name="firecracker",
slug=slug,
agent_name=metadata.agent_name if metadata else "?",
started_at=metadata.started_at if metadata else "",
# Firecracker uses the shared gateway, so there are no
# per-bottle gateway service containers to report.
services=(),
label=metadata.label if metadata else "",
color=metadata.color if metadata else "",
))
return out
+12 -11
View File
@@ -5,7 +5,7 @@ from __future__ import annotations
import subprocess
from ...bottle_state import read_metadata
from .. import ActiveAgent
from .. import ActiveAgent, EnumerationError
from .infra import INFRA_NAME, ORCHESTRATOR_NAME
# The name every agent container carries: `bot-bottle-<slug>`. Exported
@@ -20,17 +20,18 @@ CONTAINER_NAME_PREFIX = "bot-bottle-"
_INFRA_NAMES = frozenset({INFRA_NAME, ORCHESTRATOR_NAME})
class EnumerationError(RuntimeError):
"""container list failed; the resulting live set is not authoritative."""
def enumerate_active() -> list[ActiveAgent]:
result = subprocess.run(
["container", "list", "--quiet"],
capture_output=True,
text=True,
check=False,
)
try:
result = subprocess.run(
["container", "list", "--quiet"],
capture_output=True,
text=True,
check=False,
)
except FileNotFoundError as exc:
raise EnumerationError(
"container list failed: container CLI not found"
) from exc
if result.returncode != 0:
raise EnumerationError(
f"container list failed: "
+37 -24
View File
@@ -389,19 +389,9 @@ class EgressAddon:
self._passthrough_conns.discard(conn_id)
async def request(self, flow: http.HTTPFlow) -> None:
config, slug, env = self._request_context(flow)
request_path, _, query = flow.request.path.partition("?")
# Reuse the context stashed by http_connect for HTTPS flows (one
# orchestrator round-trip per connection). Plain-HTTP flows have no
# prior CONNECT stash, so resolve now and stash for response/websocket.
meta = getattr(flow, "metadata", None)
if isinstance(meta, dict) and _FLOW_CTX_KEY in meta:
config, slug, env = meta[_FLOW_CTX_KEY]
self._request_token(flow) # strip identity headers; token already resolved
else:
config, slug, env = self._resolve_flow(flow)
self._stash_flow_ctx(flow, config, slug, env)
# Introspection ("_egress.local/allowlist") reports the calling bottle's
# own resolved routes — served after resolution so it reflects this
# bottle's policy, not a stale global.
@@ -422,6 +412,29 @@ class EgressAddon:
# the path/query the git checks below rely on.
request_path, _, query = flow.request.path.partition("?")
if not self._allow_git_request(flow, config, request_path, query):
return
self._apply_route_policy(flow, config, route, request_path, env)
def _request_context(
self, flow: http.HTTPFlow,
) -> tuple[Config, str, "typing.Mapping[str, str]"]:
"""Resolve one bottle context, reusing the HTTPS CONNECT snapshot."""
meta = getattr(flow, "metadata", None)
if isinstance(meta, dict) and _FLOW_CTX_KEY in meta:
config, slug, env = meta[_FLOW_CTX_KEY]
self._request_token(flow)
return config, slug, env
config, slug, env = self._resolve_flow(flow)
self._stash_flow_ctx(flow, config, slug, env)
return config, slug, env
def _allow_git_request(
self, flow: http.HTTPFlow, config: Config,
request_path: str, query: str,
) -> bool:
"""Apply the HTTPS Git push/fetch boundary before general routing."""
if is_git_push_request(request_path, query):
self._block(
flow,
@@ -430,20 +443,20 @@ class EgressAddon:
"git-gate's pre-receive hook).",
ctx=self._req_ctx(flow),
)
return
if is_git_fetch_request(request_path, query):
git_decision = decide_git_fetch(
config.routes, flow.request.pretty_host,
)
if git_decision.action == "block":
self._block(
flow,
git_decision.reason,
ctx=self._req_ctx(flow),
)
return
return False
if not is_git_fetch_request(request_path, query):
return True
git_decision = decide_git_fetch(config.routes, flow.request.pretty_host)
if git_decision.action != "block":
return True
self._block(flow, git_decision.reason, ctx=self._req_ctx(flow))
return False
def _apply_route_policy(
self, flow: http.HTTPFlow, config: Config, route: Route | None,
request_path: str, env: "typing.Mapping[str, str]",
) -> None:
"""Strip agent auth, evaluate the route, then inject gateway auth."""
# Strip agent-set Authorization after DLP scan so smuggled tokens
# are caught above; the route may inject gateway-owned auth below.
# Routes with preserve_auth=True pass the header through as-is so the
+11 -4
View File
@@ -1,6 +1,7 @@
"""Run the orchestrator control plane as a plain process (PRD 0070 dev-harness).
python -m bot_bottle.orchestrator [--host H] [--port P] [--db PATH]
BOT_BOTTLE_ORCHESTRATOR_TOKEN=<signing-key> \
python -m bot_bottle.orchestrator [--host H] [--port P] [--db PATH]
The PRD sequences the orchestrator as a plain-process dev-harness first, so
the consolidation core (registry + attribution + HTTP control plane + live
@@ -16,12 +17,13 @@ import secrets
from pathlib import Path
from .. import log
from .store.store_manager import StoreManager
from ..trust_domain import CONTROL_PLANE
from .broker import LaunchBroker, StubBroker
from .server import make_server
from .docker_broker import DockerBroker
from .store.registry_store import RegistryStore, default_db_path
from .server import make_server
from .service import OrchestratorCore
from .store.store_manager import StoreManager
from .store.registry_store import RegistryStore, default_db_path
def main(argv: list[str] | None = None) -> int:
@@ -38,6 +40,11 @@ def main(argv: list[str] | None = None) -> int:
help="launch broker: 'stub' records requests; 'docker' runs containers",
)
args = parser.parse_args(argv)
if not CONTROL_PLANE.key_from_env():
log.die(
f"{CONTROL_PLANE.key_env} is required; refusing to start the "
"orchestrator without caller authentication"
)
registry = RegistryStore(args.db)
registry.migrate()
+92 -33
View File
@@ -59,8 +59,10 @@ import http.server
import json
import math
import os
import socket
import socketserver
import sys
import threading
import typing
from urllib.parse import urlsplit
@@ -80,6 +82,9 @@ Json = dict[str, object]
# token at all, and a compromised gateway holds only `gateway` — neither can
# drive the operator routes (approve proposals, rewrite policy, read tokens).
ORCHESTRATOR_AUTH_HEADER = "x-bot-bottle-orchestrator-auth"
MAX_BODY_BYTES = 1 * 1024 * 1024
REQUEST_TIMEOUT_SECONDS = 10.0
MAX_REQUEST_THREADS = 32
# The routes the data plane (role `gateway`) is allowed to reach — exactly the
# per-request lookups PolicyResolver makes. Every other authenticated route is
@@ -116,9 +121,8 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
no I/O beyond the orchestrator so it is fully testable without a socket.
`role` is the caller's verified control-plane role (`gateway` or `cli`), or
None for an unauthenticated request; an open-mode server (no signing key
configured see `OrchestratorServer`) passes `cli`. Every route except
`GET /health` requires a role: a missing role is 401, and a role that
None for an unauthenticated request. Every route except `GET /health`
requires a role: a missing role is 401, and a role that
doesn't cover the route is 403 — so a `gateway` data-plane token can reach
`/resolve` + `/supervise/{propose,poll}` but not the operator routes
(rewrite policy, read injected tokens, approve its own supervise proposals).
@@ -372,10 +376,33 @@ class Handler(http.server.BaseHTTPRequestHandler):
plane down for the caller."""
server = self.server
assert isinstance(server, OrchestratorServer)
length = int(self.headers.get("Content-Length") or 0)
body = self.rfile.read(length) if length > 0 else b""
role = server.role_for(self.headers.get(ORCHESTRATOR_AUTH_HEADER, ""))
route = urlsplit(self.path).path.rstrip("/") or "/"
if not (method == "GET" and route == "/health") and role is None:
self._write_json(
401, {"error": "control-plane authentication required"},
)
return
length_header = self.headers.get("Content-Length")
try:
length = int(length_header) if length_header is not None else 0
except ValueError:
self._write_json(400, {"error": "invalid Content-Length"})
return
if length < 0:
self._write_json(400, {"error": "invalid Content-Length"})
return
if length > MAX_BODY_BYTES:
self._write_json(413, {"error": "request body too large"})
return
try:
body = self.rfile.read(length) if length else b""
except (TimeoutError, socket.timeout):
self._write_json(408, {"error": "request body read timed out"})
return
try:
status: int
payload: Json
status, payload = dispatch(
server.orchestrator, method, self.path, body, role=role)
except Exception as e: # noqa: BLE001 — the control plane must stay up
@@ -388,6 +415,9 @@ class Handler(http.server.BaseHTTPRequestHandler):
)
sys.stderr.flush()
status, payload = 500, {"error": "internal error"}
self._write_json(status, payload)
def _write_json(self, status: int, payload: Json) -> None:
data = json.dumps(payload).encode()
self.send_response(status)
self.send_header("Content-Type", "application/json")
@@ -414,51 +444,80 @@ class OrchestratorServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
Holds the per-host control-plane *signing key* (from
`$BOT_BOTTLE_ORCHESTRATOR_TOKEN`, injected by the launcher into the
orchestrator process only) and verifies each request's role-scoped token
against it. When a key is set, every route but `/health` requires a valid
token whose role covers the route; when it is unset the server runs **open**
(full `cli` access) and says so loudly at startup a fail-visible fallback
for tests and any backend that hasn't wired the key yet (e.g. Firecracker,
whose nft boundary already blocks agents from the control-plane port)."""
against it. Every route but `/health` requires a valid token whose role
covers the route. Construction fails when the key is absent so a new or
misconfigured launcher cannot accidentally expose an open control plane."""
daemon_threads = True
allow_reuse_address = True
def __init__(self, address: tuple[str, int], orchestrator: OrchestratorCore) -> None:
def __init__(
self,
address: tuple[str, int],
orchestrator: OrchestratorCore,
*,
signing_key: str,
) -> None:
self.orchestrator = orchestrator
# The control-plane trust domain's signing key, as injected into THIS
# (the owning) process by the launcher (#476). Unset → open mode below.
self._signing_key = CONTROL_PLANE.key_from_env()
self._signing_key = signing_key.strip()
if not self._signing_key:
sys.stderr.write(
"orchestrator: WARNING — no control-plane signing key "
f"(${CONTROL_PLANE.key_env}); running WITHOUT caller "
"authentication. Any client that can reach this port can drive "
"it. Backends that put the control plane on an agent-reachable "
"network MUST set this.\n"
raise ValueError(
"orchestrator control-plane signing key is required; "
"refusing to start without caller authentication"
)
sys.stderr.flush()
self._request_slots = threading.BoundedSemaphore(MAX_REQUEST_THREADS)
super().__init__(address, Handler)
def get_request(self) -> tuple[socket.socket, typing.Any]:
request, client_address = super().get_request()
request.settimeout(REQUEST_TIMEOUT_SECONDS)
return request, client_address
def process_request(
self, request: typing.Any, client_address: typing.Any,
) -> None:
# Bound concurrency before ThreadingMixIn creates a worker. Backpressure
# stays in the accept loop instead of allocating an unbounded thread per
# slow or malicious connection.
self._request_slots.acquire()
try:
super().process_request(request, client_address)
except BaseException:
self._request_slots.release()
raise
def process_request_thread(
self, request: typing.Any, client_address: typing.Any,
) -> None:
try:
super().process_request_thread(request, client_address)
finally:
self._request_slots.release()
def role_for(self, presented: str) -> str | None:
"""The role the request is authorized as, or None if unauthenticated.
Open mode (no signing key) grants full `cli` access the fail-visible
fallback. Otherwise verify the presented signed token; a missing/invalid
token yields None ( 401), a valid one yields its `gateway`/`cli`
role ( per-route 401/403 in `dispatch`)."""
if not self._signing_key:
return ROLE_CLI
"""The verified caller role, or None for a missing/invalid token."""
return CONTROL_PLANE.verify(presented, self._signing_key)
def make_server(
orchestrator: OrchestratorCore, host: str = "127.0.0.1", port: int = 0
orchestrator: OrchestratorCore,
host: str = "127.0.0.1",
port: int = 0,
*,
signing_key: str | None = None,
) -> OrchestratorServer:
"""Build (but do not start) a control-plane server. `port=0` binds an
ephemeral port read `server.server_address` for the actual one."""
return OrchestratorServer((host, port), orchestrator)
"""Build an authenticated control-plane server.
``signing_key=None`` reads the owning process's injected environment.
Empty or missing keys are rejected by :class:`OrchestratorServer`.
"""
key = CONTROL_PLANE.key_from_env() if signing_key is None else signing_key
return OrchestratorServer(
(host, port), orchestrator, signing_key=key,
)
__all__ = [
"dispatch", "Handler", "OrchestratorServer", "make_server", "Json",
"ORCHESTRATOR_AUTH_HEADER",
"ORCHESTRATOR_AUTH_HEADER", "MAX_BODY_BYTES",
]
+11 -3
View File
@@ -366,15 +366,23 @@ class OrchestratorCore:
value with *env_var_secret*, and restores ``_tokens[bottle_id]``.
Returns True on success, False when no stored secrets exist for this
bottle or decryption fails (wrong key / corrupt data)."""
from .store.secret_store import decrypt_value
from .store.secret_store import decrypt_value, encrypt_value, is_legacy_blob
encrypted = self.registry.get_agent_secrets(bottle_id)
if not encrypted:
return False
try:
self._tokens[bottle_id] = {k: decrypt_value(env_var_secret, v)
for k, v in encrypted.items()}
decrypted = {
k: decrypt_value(env_var_secret, v) for k, v in encrypted.items()
}
except ValueError:
return False
self._tokens[bottle_id] = decrypted
if any(is_legacy_blob(value) for value in encrypted.values()):
migrated = {
key: encrypt_value(env_var_secret, value)
for key, value in decrypted.items()
}
self.registry.store_agent_secrets(bottle_id, migrated)
return True
# --- consolidated gateway ----------------------------------------------
+80 -18
View File
@@ -12,9 +12,15 @@ reattachment path reads ENV_VAR_SECRET from the running agent container via
``POST /bottles/<id>/reprovision_gateway``; the orchestrator decrypts the
stored rows and re-populates ``_tokens``.
Encryption scheme: HMAC-SHA256 used as a PRF in CTR mode (stdlib-only,
no external deps). Each value is encrypted independently. The output blob is
``nonce (16 bytes) || ciphertext`` encoded as URL-safe base64 (no padding).
Encryption scheme: encrypt-then-MAC using independent HMAC-SHA256-derived
encryption and authentication subkeys (stdlib-only, no external deps). Each
value is encrypted independently. New output blobs are:
``version || nonce (16 bytes) || ciphertext || tag (32 bytes)``
encoded as URL-safe base64 (no padding). The version marker lets the reader
accept legacy ``nonce || ciphertext`` rows long enough to rewrite them in the
authenticated format after a successful reprovision.
keystream_block_i = HMAC-SHA256(key, nonce || i.to_bytes(4, "big"))
ciphertext_i = plaintext_i XOR keystream_block_i[:len(plaintext_i)]
@@ -30,6 +36,8 @@ import secrets
_KEY_BYTES = 32 # 256-bit key from ENV_VAR_SECRET
_NONCE_BYTES = 16 # 128-bit random nonce per encrypt call
_BLOCK = 32 # HMAC-SHA256 output width == one keystream block
_TAG_BYTES = 32
_VERSION = b"BBSE1"
# Env-var name the agent container receives at startup.
ENV_VAR_SECRET_NAME = "ENV_VAR_SECRET"
@@ -41,7 +49,13 @@ def new_env_var_secret() -> str:
def _b64dec(s: str) -> bytes:
return base64.urlsafe_b64decode(s + "=" * (-len(s) % 4))
return base64.b64decode(
s + "=" * (-len(s) % 4), altchars=b"-_", validate=True,
)
def _subkey(key: bytes, purpose: bytes) -> bytes:
return hmac.new(key, b"bot-bottle-secret-store:" + purpose, hashlib.sha256).digest()
def _keystream(key: bytes, nonce: bytes, block_index: int) -> bytes:
@@ -53,42 +67,90 @@ def _keystream(key: bytes, nonce: bytes, block_index: int) -> bytes:
def encrypt_value(secret_b64: str, plaintext: str) -> str:
"""Encrypt a single string value with *secret_b64* (the ENV_VAR_SECRET).
Returns a URL-safe base64 blob ``nonce || ciphertext`` suitable for
Returns a URL-safe base64 authenticated blob suitable for
the ``bottled_agent_secrets.value`` column."""
key = _b64dec(secret_b64)
encryption_key = _subkey(key, b"encryption")
authentication_key = _subkey(key, b"authentication")
pt = plaintext.encode()
nonce = secrets.token_bytes(_NONCE_BYTES)
ct = bytearray()
for i in range(0, len(pt), _BLOCK):
chunk = pt[i : i + _BLOCK]
ks = _keystream(key, nonce, i)[: len(chunk)]
ks = _keystream(encryption_key, nonce, i // _BLOCK)[: len(chunk)]
ct.extend(p ^ k for p, k in zip(chunk, ks))
return base64.urlsafe_b64encode(nonce + bytes(ct)).rstrip(b"=").decode()
authenticated = _VERSION + nonce + bytes(ct)
tag = hmac.new(authentication_key, authenticated, hashlib.sha256).digest()
return base64.urlsafe_b64encode(authenticated + tag).rstrip(b"=").decode()
def decrypt_value(secret_b64: str, blob_b64: str) -> str:
"""Decrypt a blob produced by :func:`encrypt_value`.
Returns the original plaintext string. Raises ``ValueError`` for malformed
input or a key mismatch (wrong key produces garbage, not an error, unless
the plaintext is non-UTF-8 treat all such failures as wrong key)."""
key = _b64dec(secret_b64)
def is_legacy_blob(blob_b64: str) -> bool:
"""Whether *blob_b64* uses the pre-authentication storage format."""
try:
blob = _b64dec(blob_b64)
except Exception as exc:
raise ValueError(f"invalid ciphertext blob: {exc}") from exc
return not _b64dec(blob_b64).startswith(_VERSION)
except (ValueError, TypeError):
return False
def _decrypt_legacy(key: bytes, blob: bytes) -> str:
"""Read the original ``nonce || ciphertext`` format for migration only."""
if len(blob) < _NONCE_BYTES:
raise ValueError("ciphertext blob too short")
nonce, ciphertext = blob[:_NONCE_BYTES], blob[_NONCE_BYTES:]
pt = bytearray()
# The legacy format used the byte offset as the PRF counter.
for i in range(0, len(ciphertext), _BLOCK):
chunk = ciphertext[i : i + _BLOCK]
ks = _keystream(key, nonce, i)[: len(chunk)]
pt.extend(c ^ k for c, k in zip(chunk, ks))
try:
return bytes(pt).decode()
except UnicodeDecodeError as exc:
raise ValueError(f"decryption produced non-UTF-8 output: {exc}") from exc
def decrypt_value(secret_b64: str, blob_b64: str) -> str:
"""Decrypt a blob produced by :func:`encrypt_value`.
Returns the original plaintext string. Raises ``ValueError`` for malformed
input, authentication failure, or a key mismatch. Legacy unauthenticated
rows remain readable so callers can migrate them immediately."""
key = _b64dec(secret_b64)
try:
blob = _b64dec(blob_b64)
except (ValueError, TypeError) as exc:
raise ValueError(f"invalid ciphertext blob: {exc}") from exc
if not blob.startswith(_VERSION):
return _decrypt_legacy(key, blob)
minimum = len(_VERSION) + _NONCE_BYTES + _TAG_BYTES
if len(blob) < minimum:
raise ValueError("ciphertext blob too short")
authenticated, supplied_tag = blob[:-_TAG_BYTES], blob[-_TAG_BYTES:]
authentication_key = _subkey(key, b"authentication")
expected_tag = hmac.new(
authentication_key, authenticated, hashlib.sha256,
).digest()
if not hmac.compare_digest(supplied_tag, expected_tag):
raise ValueError("ciphertext authentication failed")
nonce_start = len(_VERSION)
nonce = blob[nonce_start : nonce_start + _NONCE_BYTES]
ciphertext = blob[nonce_start + _NONCE_BYTES : -_TAG_BYTES]
encryption_key = _subkey(key, b"encryption")
pt = bytearray()
for i in range(0, len(ciphertext), _BLOCK):
chunk = ciphertext[i : i + _BLOCK]
ks = _keystream(encryption_key, nonce, i // _BLOCK)[: len(chunk)]
pt.extend(c ^ k for c, k in zip(chunk, ks))
try:
return bytes(pt).decode()
except UnicodeDecodeError as exc:
raise ValueError(f"decryption produced non-UTF-8 output (wrong key?): {exc}") from exc
__all__ = ["ENV_VAR_SECRET_NAME", "new_env_var_secret", "encrypt_value", "decrypt_value"]
__all__ = [
"ENV_VAR_SECRET_NAME",
"new_env_var_secret",
"encrypt_value",
"decrypt_value",
"is_legacy_blob",
]
+3 -4
View File
@@ -40,7 +40,7 @@ from .paths import (
class ProvisioningError(RuntimeError):
"""A control-plane auth invariant would be violated (e.g. starting the
orchestrator without its signing key which would run OPEN)."""
orchestrator without its signing key)."""
@dataclass(frozen=True)
@@ -67,9 +67,8 @@ class TrustDomain:
def key_from_env(self, environ: Mapping[str, str] | None = None) -> str:
"""The signing key as the owning process sees it — read from `key_env`
(default `os.environ`). "" when unset; the caller decides whether that is
fatal (`ControlPlaneProvisioning`) or the open-mode fallback
(`OrchestratorServer`)."""
(default `os.environ`). ``""`` when unset; owning services reject that
value rather than start without authentication."""
env = os.environ if environ is None else environ
return env.get(self.key_env, "").strip()
@@ -56,8 +56,7 @@ key.
- Rewriting the HMAC primitive: `orchestrator_auth.mint/verify` gain an optional
`roles=` arg (default unchanged) so a key can carry a different role set;
nothing else changes.
- Network topology, the plane split (#469), or the server's open-mode fallback
for tests.
- Network topology or the plane split (#469).
## Design
@@ -0,0 +1,197 @@
# PRD 0082: Authoritative failure boundaries
- **Status:** Draft
- **Author:** codex
- **Created:** 2026-07-27
- **Issue:** #444
## Summary
Make every security- or lifecycle-sensitive snapshot distinguish authoritative
empty state from unavailable state, and make every destructive or
resource-consuming boundary revalidate the assumptions it acts on. This
finishes the focused quality work begun under #444 without broad rewrites:
cleanup cannot act on stale identities, policy introspection cannot publish a
fabricated empty policy, gateway servers bound untrusted work, and daemon
shutdown does not emit uncaught background-thread failures. Shared
control-plane storage and gateway credential provisioning also enforce their
filesystem security contract before sensitive data is written.
## Problem
Several paths are individually fail-closed but compose into unsafe or
misleading behavior:
1. `cleanup` prepares a plan, waits indefinitely for operator confirmation,
then kills stored PIDs and removes stored paths without checking that those
identities still describe the same orphan. A PID may be reused or a run
directory may become active during the prompt.
2. The supervisor reuses egress's deny-all fallback for
`list-egress-routes`. Deny-all is correct for enforcement, but presenting it
as a successful empty route table can cause a later replace-all proposal to
discard live routes.
3. The supervisor and Git HTTP services accept bounded declared body sizes but
use blocking reads and unbounded request threads. An untrusted bottle can
exhaust the shared gateway with slow or parallel requests.
4. macOS cleanup enumerates containers and networks independently and treats a
failed query as an empty class, so a partial snapshot can still become a
destructive plan.
5. Gateway log-pump threads race stream closure during shutdown and emit
uncaught exceptions even when shutdown otherwise succeeds.
6. Firecracker discovers VMs through whitespace-split `pgrep -a` output.
A configured cache path containing spaces can hide a live VM from the
snapshot and make its run directory appear orphaned.
7. Docker cleanup asks compose for its project snapshot in best-effort mode.
A transient query failure can therefore become an empty stopped-project
set and authorize deletion of associated state directories.
8. Firecracker artifact downloads and registry publication have no network
deadline, so an unresponsive registry can hold setup or release work
indefinitely.
9. Authenticated secret blobs select the unauthenticated legacy decoder when
their in-band version prefix is changed, allowing storage tampering to
bypass tag verification.
10. Cleanup executes the entire post-confirmation snapshot rather than the
intersection with what the operator saw, and mutation failures are not
reflected in the command result.
11. Git smart-HTTP can retain sixteen 100 MiB request bodies concurrently,
cleanup mutations have no subprocess deadline, and Firecracker signalling
failures bypass shared mutation accounting.
12. SQLite creates the shared control-plane database before its mode is
restricted, then suppresses permission-repair failures. Gateway transports
also differ in whether copied deploy-key modes are preserved.
These are one design problem: state used to authorize deletion, replacement,
or resource allocation must be authoritative at the point of use.
## Goals / Success Criteria
- Cleanup never signals a PID or recursively deletes a path solely because it
appeared in a pre-confirmation snapshot.
- Firecracker cleanup proves immediately before action that a PID is still the
same Firecracker process and that a run directory is still orphaned.
- Firecracker process discovery reads NUL-delimited argv from `/proc`; paths
are never reconstructed from whitespace-delimited process listings.
- All backend cleanup discovery primitives raise a typed enumeration error on
operational failure. No backend may independently continue from a partial
snapshot.
- Shared cleanup control flow lives in the backend layer; concrete backends
override resource-specific discovery and validation primitives rather than
each implementing a bespoke failure policy.
- `list-egress-routes` returns an MCP error when attribution or policy
resolution is unavailable. A genuine, authoritatively resolved empty policy
remains a successful empty list.
- Supervisor and Git HTTP request bodies have total read deadlines, and each
service bounds concurrent request work. Limits apply to authenticated
callers because bottles themselves are untrusted.
- Gateway child-output pumping treats expected stream closure during shutdown
as completion while preserving diagnostics for unexpected failures.
- Artifact pull, existence-check, and publication requests use explicit
network deadlines.
- Persisted secrets accept only the authenticated format. The schema migration
intentionally clears legacy rows; local agents are reprovisioned rather
than retaining a ciphertext-controlled downgrade path.
- Cleanup executes only resources present in both the displayed and current
authoritative plans, attempts every approved mutation, and returns failure
when any mutation does not complete.
- Git request bodies spool to disk behind a separate heavy-work semaphore;
cleanup commands have configurable deadlines; Firecracker signalling
failures aggregate while identity-verification uncertainty still aborts.
- The shared database directory and file are private before SQLite writes any
control-plane state; an inability to enforce those modes aborts startup.
- Gateway credential directories and files receive explicit private modes
inside the gateway, independent of Docker, Apple Container, or SSH copy
semantics.
- Unit tests cover PID/path reuse, partial backend enumeration, transient
policy resolution failure, slow bodies, concurrency saturation, and stream
closure races.
## Non-goals
- Further decomposition solely to reduce module line counts.
- Replacing gateway stdlib HTTP services with a web framework.
- Changing egress matching, DLP decisions, proposal semantics, or backend
launch behavior beyond the synchronization required for safe cleanup.
- Making cleanup silently skip uncertain resources. Uncertainty is an
operator-visible failure.
## Design
### Shared backend control flow
Follow the backend architecture rule used by gateway attachment: shared
behavior lives above concrete backends; subclasses provide primitives, not
control flow.
Cleanup remains previewable, but confirmation authorizes a *new authoritative
evaluation*, not blind execution of the displayed object. The shared flow:
1. asks each available backend for a preview;
2. displays the union and asks for confirmation;
3. refreshes each non-empty backend plan;
4. validates destructive identities immediately before action;
5. aborts loudly if the refreshed plan or any identity cannot be proven safe.
Backend-specific primitives define how to identify a resource. Firecracker
uses process start identity plus canonical config/run paths; container
backends use authoritative CLI queries and stable resource names/labels.
Container engines expose destructive name-based commands without a portable
compare-and-delete operation. Cleanup therefore refreshes after confirmation
and requires every discovery query to succeed, minimizing but not claiming to
eliminate the final name-reuse race. A future engine-specific stable-ID
primitive may close that residual window without moving control flow back
into each backend.
### Enforcement state versus introspection state
Egress enforcement retains its deny-all fallback because uncertainty must not
grant network access. Supervisor introspection uses a strict resolver path:
unattributed callers and resolver failures become typed MCP errors, while a
successfully resolved policy containing zero routes returns `routes: []`.
### Gateway resource boundaries
Both stdlib servers set a per-connection body deadline before reading and use a
bounded request executor or semaphore. Saturated capacity fails quickly with a
service-unavailable response. Existing size caps remain independent:
supervisor proposals retain the 1 MiB cap and Git pack requests retain their
larger protocol-appropriate cap.
### Shutdown diagnostics
The gateway output pump catches only stream-closure exceptions expected after
the supervisor closes child pipes. Other I/O failures remain visible and are
reported through the supervisor's normal diagnostic channel.
### Shared filesystem security
The common SQLite store owns database creation for every backend. It creates
the parent directory and an empty database with private modes before opening
SQLite, repairs existing modes, verifies the resulting state, and propagates
every enforcement failure. Backend launchers do not duplicate this policy.
The backend-neutral gateway provisioner likewise applies directory and file
modes after transport copies complete. This avoids relying on copy behavior
that differs among Docker, Apple Container, and Firecracker's SSH transport.
## Implementation chunks
1. Existing fail-closed security and backend enumeration fixes.
2. FastAPI orchestrator transport and bounded control-plane bodies.
3. Egress request-policy and outbound-DLP pipeline extraction.
4. Supervisor MCP dispatch extraction.
5. Shared cleanup refresh/revalidation plus authoritative macOS discovery.
6. Strict supervisor introspection and bounded supervisor/Git HTTP work.
7. Gateway shutdown log-pump closure handling.
8. Lossless Firecracker process identities, authoritative Docker cleanup
queries, and bounded Firecracker artifact transfers.
9. Mandatory authenticated secret storage, shared cleanup-plan intersection
and mutation accounting, and contained Git backend process failures.
10. Disk-spooled and separately bounded Git bodies, cleanup command deadlines,
and classified Firecracker signalling failures.
11. Fail-closed shared database creation and backend-neutral gateway credential
permissions.
## Open questions
None.
+14
View File
@@ -18,6 +18,7 @@ from bot_bottle.backend.docker.compose import (
list_compose_projects,
slug_from_compose_project,
)
from bot_bottle.backend import EnumerationError
class TestProjectNaming(unittest.TestCase):
@@ -69,6 +70,19 @@ class TestComposeProjectListing(unittest.TestCase):
self.assertEqual([], list_active_slugs(warn_on_error=False))
warn.assert_not_called()
def test_compose_ls_error_can_be_raised_for_enumeration(self):
with mock.patch(
"bot_bottle.backend.docker.compose.subprocess.run",
return_value=subprocess.CompletedProcess(
args=["docker"], returncode=1, stdout="", stderr="no daemon",
),
):
with self.assertRaisesRegex(EnumerationError, "no daemon"):
list_active_slugs(
warn_on_error=False,
raise_on_error=True,
)
if __name__ == "__main__":
unittest.main()
+12
View File
@@ -7,6 +7,8 @@ from pathlib import Path
from unittest.mock import MagicMock, Mock, patch
from bot_bottle.backend.docker.consolidated_launch import (
ConsolidatedLaunchError,
_network_container_ips,
launch_consolidated,
deprovision_consolidated,
)
@@ -86,6 +88,16 @@ class TestLaunchConsolidated(unittest.TestCase):
client.teardown_bottle.assert_called_once_with("b1") # no orphan left
class TestNetworkContainerIps(unittest.TestCase):
def test_fails_closed_when_network_inspection_fails(self) -> None:
result = Mock(returncode=1, stdout="", stderr="daemon unavailable")
with (
patch(f"{_MOD}.run_docker", return_value=result),
self.assertRaisesRegex(ConsolidatedLaunchError, "daemon unavailable"),
):
_network_container_ips("bot-bottle-gateway")
class TestTeardownConsolidated(unittest.TestCase):
def test_deregisters_and_deprovisions(self) -> None:
client = Mock()
@@ -19,13 +19,16 @@ of issue #77 — the dashboard now delegates to this layer.
from __future__ import annotations
import subprocess
import tempfile
import unittest
from pathlib import Path
from unittest.mock import patch
from tests.unit import use_bottle_root
from bot_bottle import bottle_state
from bot_bottle.backend.docker import enumerate as _enumerate
from bot_bottle.backend import EnumerationError
class TestParseServicesByProject(unittest.TestCase):
@@ -72,6 +75,18 @@ class TestParseServicesByProject(unittest.TestCase):
self.assertEqual({"bot-bottle-dev-abc": {"egress"}}, out)
class TestQueryServicesByProject(unittest.TestCase):
def test_docker_ps_failure_is_not_an_empty_result(self):
with patch(
"bot_bottle.backend.docker.enumerate.subprocess.run",
return_value=subprocess.CompletedProcess(
args=["docker"], returncode=1, stdout="", stderr="daemon down",
),
):
with self.assertRaisesRegex(EnumerationError, "daemon down"):
_enumerate._query_services_by_project()
class _FakeHomeMixin:
def _setup_fake_home(self) -> None:
self._tmp = tempfile.TemporaryDirectory(prefix="enum-active.")
+29 -1
View File
@@ -15,6 +15,7 @@ import unittest
from pathlib import Path
from unittest.mock import patch
from bot_bottle.backend import EnumerationError
from bot_bottle.backend.firecracker import cleanup as fc_cleanup
from bot_bottle.backend.firecracker.bottle_cleanup_plan import (
FirecrackerBottleCleanupPlan,
@@ -58,10 +59,37 @@ class TestProcessScan(unittest.TestCase):
self.assertEqual({str(run_root / "live-a")}, live)
self.assertEqual([222], orphan_pids)
def test_scan_empty_when_pgrep_fails(self):
def test_scan_empty_when_pgrep_finds_no_processes(self):
with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(returncode=1)):
self.assertEqual((set(), []), fc_cleanup._scan_processes(Path("/x")))
def test_scan_raises_when_pgrep_errors(self):
proc = subprocess.CompletedProcess(
[], 2, stdout="", stderr="invalid process expression",
)
with (
patch.object(fc_cleanup.subprocess, "run", return_value=proc),
self.assertRaisesRegex(EnumerationError, "invalid process expression"),
):
fc_cleanup._scan_processes(Path("/x"))
def test_prepare_cleanup_does_not_plan_deletions_when_scan_errors(self):
with tempfile.TemporaryDirectory() as tmp:
run_root = Path(tmp)
live = run_root / "live-a"
live.mkdir()
with (
patch.object(fc_cleanup, "_run_root", return_value=run_root),
patch.object(
fc_cleanup.subprocess,
"run",
return_value=_proc(returncode=2),
),
self.assertRaises(EnumerationError),
):
fc_cleanup.prepare_cleanup()
self.assertTrue(live.is_dir())
def test_live_run_dirs_returns_paths_in_stable_order(self):
with patch.object(fc_cleanup, "_run_root", return_value=Path("/run")), \
patch.object(fc_cleanup, "_scan_processes",
+64
View File
@@ -0,0 +1,64 @@
"""Unit tests for Firecracker active-agent enumeration."""
from __future__ import annotations
import unittest
from pathlib import Path
from unittest.mock import patch
from bot_bottle.backend import EnumerationError
from bot_bottle.backend.firecracker import enumerate as fc_enumerate
from bot_bottle.bottle_state import BottleMetadata
class TestEnumerateActive(unittest.TestCase):
def test_maps_live_run_dirs_to_active_agents(self) -> None:
metadata = BottleMetadata(
identity="dev-a",
agent_name="claude",
cwd="",
copy_cwd=False,
started_at="2026-07-26T12:00:00Z",
label="review",
color="blue",
)
with (
patch.object(
fc_enumerate, "live_run_dirs",
return_value=(Path("/cache/run/dev-a"),),
),
patch.object(fc_enumerate, "read_metadata", return_value=metadata),
):
agents = fc_enumerate.enumerate_active()
self.assertEqual(1, len(agents))
self.assertEqual("firecracker", agents[0].backend_name)
self.assertEqual("dev-a", agents[0].slug)
self.assertEqual("claude", agents[0].agent_name)
self.assertEqual("review", agents[0].label)
self.assertEqual((), agents[0].services)
def test_missing_metadata_uses_safe_defaults(self) -> None:
with (
patch.object(
fc_enumerate, "live_run_dirs",
return_value=(Path("/cache/run/dev-a"),),
),
patch.object(fc_enumerate, "read_metadata", return_value=None),
):
agent = fc_enumerate.enumerate_active()[0]
self.assertEqual("?", agent.agent_name)
self.assertEqual("", agent.started_at)
def test_process_scan_failure_propagates(self) -> None:
with (
patch.object(
fc_enumerate, "live_run_dirs",
side_effect=EnumerationError("pgrep failed"),
),
self.assertRaisesRegex(EnumerationError, "pgrep failed"),
):
fc_enumerate.enumerate_active()
if __name__ == "__main__":
unittest.main()
+10 -1
View File
@@ -5,6 +5,7 @@ from __future__ import annotations
import unittest
from unittest.mock import patch
from bot_bottle.backend import EnumerationError
from bot_bottle.backend.macos_container import cleanup, enumerate as enum_mod
from bot_bottle.backend.macos_container.bottle_cleanup_plan import (
MacosContainerBottleCleanupPlan,
@@ -69,10 +70,18 @@ class TestMacosContainerEnumerate(unittest.TestCase):
self.assertEqual(["dev-abc"], [a.slug for a in agents])
def test_raises_when_the_cli_fails(self):
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with self.assertRaises(EnumerationError):
self._enumerate("", returncode=1)
def test_raises_typed_error_when_cli_is_missing(self):
with (
patch.object(
enum_mod.subprocess, "run", side_effect=FileNotFoundError,
),
self.assertRaisesRegex(EnumerationError, "CLI not found"),
):
enum_mod.enumerate_active()
if __name__ == "__main__":
unittest.main()
+25 -9
View File
@@ -2,6 +2,9 @@
from __future__ import annotations
import base64
import hashlib
import hmac
import unittest
from bot_bottle.orchestrator.store.secret_store import (
@@ -68,15 +71,28 @@ class TestDecryptErrors(unittest.TestCase):
def test_wrong_key_raises_value_error(self) -> None:
ct = encrypt_value(self.secret, "secret-token")
other_key = new_env_var_secret()
# Wrong key produces garbage bytes; decrypt_value raises ValueError
# when the result is non-UTF-8 (which is very likely for 12-char data).
# We allow it to succeed only if garbage happens to be valid UTF-8, but
# the plaintext must not match.
try:
result = decrypt_value(other_key, ct)
self.assertNotEqual("secret-token", result)
except ValueError:
pass
with self.assertRaisesRegex(ValueError, "authentication failed"):
decrypt_value(other_key, ct)
def test_tampered_ciphertext_raises_value_error(self) -> None:
raw = bytearray(base64.urlsafe_b64decode(
encrypt_value(self.secret, "secret-token") + "=="
))
raw[22] ^= 1
tampered = base64.urlsafe_b64encode(raw).rstrip(b"=").decode()
with self.assertRaisesRegex(ValueError, "authentication failed"):
decrypt_value(self.secret, tampered)
def test_reads_legacy_ciphertext_for_migration(self) -> None:
key = base64.urlsafe_b64decode(self.secret + "==")
nonce = b"0123456789abcdef"
plaintext = b"legacy-token"
stream = hmac.new(
key, nonce + (0).to_bytes(4, "big"), hashlib.sha256,
).digest()
ciphertext = bytes(p ^ k for p, k in zip(plaintext, stream))
legacy = base64.urlsafe_b64encode(nonce + ciphertext).rstrip(b"=").decode()
self.assertEqual("legacy-token", decrypt_value(self.secret, legacy))
def test_truncated_blob_raises_value_error(self) -> None:
with self.assertRaises(ValueError):
+70 -18
View File
@@ -7,6 +7,7 @@ server tests), plus one real-socket round-trip to prove the handler wiring.
from __future__ import annotations
import base64
import http.client
import io
import json
import secrets
@@ -22,7 +23,7 @@ from unittest.mock import MagicMock, patch
from bot_bottle.orchestrator_auth import ROLE_CLI, ROLE_GATEWAY, mint
from bot_bottle.orchestrator.broker import StubBroker
from bot_bottle.orchestrator.server import dispatch, make_server
from bot_bottle.orchestrator.server import MAX_BODY_BYTES, dispatch, make_server
from bot_bottle.orchestrator.store.registry_store import BottleRecord, RegistryStore
from bot_bottle.orchestrator.service import OrchestratorCore
from bot_bottle.orchestrator.store.store_manager import StoreManager
@@ -251,11 +252,51 @@ class TestDispatch(unittest.TestCase):
class TestServerRoundTrip(unittest.TestCase):
def _raw_status(self, content_length: str, *, authenticated: bool = True) -> int:
tmp = tempfile.TemporaryDirectory()
self.addCleanup(tmp.cleanup)
key = "request-limits-key"
server = make_server(
_orchestrator(Path(tmp.name) / "r.db"),
"127.0.0.1", 0, signing_key=key,
)
self.addCleanup(server.server_close)
threading.Thread(target=server.serve_forever, daemon=True).start()
self.addCleanup(server.shutdown)
conn = http.client.HTTPConnection(
str(server.server_address[0]), server.server_address[1], timeout=5,
)
self.addCleanup(conn.close)
conn.putrequest("POST", "/bottles")
conn.putheader("Content-Length", content_length)
if authenticated:
conn.putheader(
"x-bot-bottle-orchestrator-auth", mint(ROLE_CLI, key),
)
conn.endheaders()
return conn.getresponse().status
def test_rejects_malformed_content_length(self) -> None:
self.assertEqual(400, self._raw_status("not-a-number"))
def test_rejects_oversized_body_without_reading_it(self) -> None:
self.assertEqual(413, self._raw_status(str(MAX_BODY_BYTES + 1)))
def test_rejects_unauthenticated_request_before_reading_body(self) -> None:
self.assertEqual(
401,
self._raw_status(str(MAX_BODY_BYTES), authenticated=False),
)
def test_http_register_health_attribute(self) -> None:
tmp = tempfile.TemporaryDirectory()
self.addCleanup(tmp.cleanup)
orch = _orchestrator(Path(tmp.name) / "r.db")
server = make_server(orch, "127.0.0.1", 0)
signing_key = "round-trip-key"
auth = mint(ROLE_CLI, signing_key)
server = make_server(
orch, "127.0.0.1", 0, signing_key=signing_key,
)
self.addCleanup(server.server_close)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
@@ -267,7 +308,10 @@ class TestServerRoundTrip(unittest.TestCase):
reg = json.load(urllib.request.urlopen(
urllib.request.Request(
f"{base}/bottles", data=_body({"source_ip": "10.243.0.7"}),
method="POST", headers={"Content-Type": "application/json"},
method="POST", headers={
"Content-Type": "application/json",
"x-bot-bottle-orchestrator-auth": auth,
},
), timeout=5,
))
self.assertTrue(reg["bottle_id"])
@@ -279,7 +323,10 @@ class TestServerRoundTrip(unittest.TestCase):
urllib.request.Request(
f"{base}/attribute",
data=_body({"source_ip": "10.243.0.7", "identity_token": reg["identity_token"]}),
method="POST", headers={"Content-Type": "application/json"},
method="POST", headers={
"Content-Type": "application/json",
"x-bot-bottle-orchestrator-auth": auth,
},
), timeout=5,
))
self.assertEqual(reg["bottle_id"], attr["bottle_id"])
@@ -287,15 +334,25 @@ class TestServerRoundTrip(unittest.TestCase):
def test_internal_failure_is_contextual_but_redacted(self) -> None:
orch = MagicMock()
orch.registry.all.side_effect = RuntimeError("SENSITIVE request value")
signing_key = "failure-path-key"
auth = mint(ROLE_CLI, signing_key)
with patch("sys.stderr", io.StringIO()) as stderr:
server = make_server(orch, "127.0.0.1", 0)
server = make_server(
orch, "127.0.0.1", 0, signing_key=signing_key,
)
self.addCleanup(server.server_close)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
self.addCleanup(server.shutdown)
host, port = server.server_address[0], server.server_address[1]
with self.assertRaises(urllib.error.HTTPError) as raised:
urllib.request.urlopen(f"http://{host}:{port}/bottles", timeout=5)
urllib.request.urlopen(
urllib.request.Request(
f"http://{host}:{port}/bottles",
headers={"x-bot-bottle-orchestrator-auth": auth},
),
timeout=5,
)
payload = json.loads(raised.exception.read())
output = stderr.getvalue()
self.assertEqual({"error": "internal error"}, payload)
@@ -369,8 +426,9 @@ class TestOrchestratorAuth(unittest.TestCase):
self.assertIsNotNone(self.orch.registry.get(rec.bottle_id))
def _server_with_key(self, signing_key: str):
with patch.dict("os.environ", {"BOT_BOTTLE_ORCHESTRATOR_TOKEN": signing_key}):
server = make_server(self.orch, "127.0.0.1", 0)
server = make_server(
self.orch, "127.0.0.1", 0, signing_key=signing_key,
)
self.addCleanup(server.server_close)
threading.Thread(target=server.serve_forever, daemon=True).start()
self.addCleanup(server.shutdown)
@@ -399,16 +457,10 @@ class TestOrchestratorAuth(unittest.TestCase):
self.assertEqual(403, self._status(f"{base}/bottles", header=gateway_tok))
self.assertEqual(200, self._status(f"{base}/bottles", header=cli_tok))
def test_unconfigured_server_runs_open(self) -> None:
"""No signing key set (tests / nft-protected Firecracker): open mode
grants full cli access, so existing round-trip behavior is unchanged."""
with patch.dict("os.environ", {}, clear=False):
import os
os.environ.pop("BOT_BOTTLE_ORCHESTRATOR_TOKEN", None)
server = make_server(self.orch, "127.0.0.1", 0)
self.addCleanup(server.server_close)
self.assertEqual(ROLE_CLI, server.role_for(""))
self.assertEqual(ROLE_CLI, server.role_for("anything"))
def test_unconfigured_server_refuses_to_start(self) -> None:
with patch.dict("os.environ", {}, clear=True):
with self.assertRaisesRegex(ValueError, "signing key is required"):
make_server(self.orch, "127.0.0.1", 0)
class TestDispatchSupervise(unittest.TestCase):
+1 -2
View File
@@ -78,8 +78,7 @@ class TestControlPlaneProvisioning(unittest.TestCase):
self.assertEqual("key", prov.orchestrator_key())
def test_orchestrator_key_fail_closes_when_empty(self) -> None:
# Invariant 4: the orchestrator must never start without a key — it would
# run OPEN and grant every caller that reaches it full `cli`. There is no
# Invariant 4: the orchestrator must never start without a key. There is no
# topology opt-out: a separate host does not stop the gateway (or any
# other caller) from reaching the control-plane listener.
prov = ControlPlaneProvisioning()