Compare commits

..

7 Commits

Author SHA1 Message Date
didericis-claude 755a11a608 firecracker status(): add binary and KVM readiness checks
tracker-policy-pr / check-pr (pull_request) Successful in 11s
test / integration-docker (pull_request) Successful in 19s
lint / lint (push) Successful in 1m0s
test / unit (pull_request) Successful in 1m43s
test / integration-firecracker (pull_request) Successful in 3m22s
test / coverage (pull_request) Successful in 16s
test / publish-infra (pull_request) Has been skipped
Adds `_firecracker_binary_ok()` (runs `firecracker --version`) and
`_kvm_accessible()` (opens /dev/kvm, issues KVM_GET_API_VERSION ioctl)
and gates `status()` on both, so the launch preflight reports missing
binary or inaccessible KVM before the caller ever attempts a boot.

Regression tests in TestFirecrackerBinaryCheck, TestFirecrackerKvmCheck,
and TestFirecrackerStatusRuntime cover all branches. Updated the
pre-existing TestFirecrackerStatus fixture to also stub the two new
helpers so it still exercises only the tap-pool/nft logic it was
designed for.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-24 06:51:58 +00:00
didericis-claude e074c6959f fix: add missing type annotations and remove unused variables in new tests
test / integration-docker (pull_request) Successful in 23s
test / unit (pull_request) Successful in 47s
lint / lint (push) Successful in 3m3s
test / integration-firecracker (pull_request) Successful in 3m24s
test / coverage (pull_request) Successful in 20s
test / publish-infra (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Successful in 17s
Pyright flagged three inner class `status` methods missing a type
annotation on the `quiet` parameter, and two unused variables
(`written`, `real_status`) left over from an earlier draft of
test_docker_quiet_true_suppresses_stderr.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-24 06:18:51 +00:00
didericis-claude 05b62ce805 test: cover is_backend_available, is_backend_ready, and status(quiet=True)
test / integration-docker (pull_request) Successful in 21s
lint / lint (push) Failing after 59s
test / unit (pull_request) Successful in 1m47s
test / integration-firecracker (pull_request) Successful in 3m22s
test / coverage (pull_request) Successful in 18s
test / publish-infra (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Failing after 13m0s
Add unit tests for the three new public symbols introduced by this PR:

- TestIsBackendAvailable — delegates to has_backend; returns True/False
  correctly for available and unavailable backends
- TestIsBackendReady — unknown names return False; known names return
  True/False based on status() return code; quiet flag is forwarded
- TestBackendStatusQuiet — verifies the quiet=True branch in
  DockerBottleBackend, FirecrackerBottleBackend, and
  MacosContainerBottleBackend: return code is propagated and
  diagnostic stderr is suppressed

Together these bring diff-coverage for the PR's changed lines to 100%.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-24 06:12:02 +00:00
didericis-claude 1f9a15dede fix: align skip-guard tests with is_backend_ready and drop unused import
tracker-policy-pr / check-pr (pull_request) Successful in 14s
test / integration-docker (pull_request) Successful in 16s
test / unit (pull_request) Successful in 37s
lint / lint (push) Successful in 53s
test / integration-firecracker (pull_request) Successful in 3m18s
test / coverage (pull_request) Failing after 35s
test / publish-infra (pull_request) Has been skipped
tests/_backend.py imported is_backend_available but never called it,
causing a pyright reportUnusedImport error. Remove the unused import.

test_backend_skip_guards.py was patching tests._backend.has_backend but
_backend.py calls is_backend_ready — the mock never intercepted the real
call, producing AttributeError at test runtime. Update all patch targets
to is_backend_ready and add quiet=False to the assert_called_once_with
assertion to match the actual call signature.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-24 04:25:04 +00:00
didericis-claude 0ca39cf4d5 feat(backend): BackendStatus enum, quiet status(), is_backend_ready()
tracker-policy-pr / check-pr (pull_request) Successful in 11s
test / integration-docker (pull_request) Successful in 17s
lint / lint (push) Failing after 54s
test / unit (pull_request) Failing after 1m33s
test / integration-firecracker (pull_request) Successful in 3m17s
test / coverage (pull_request) Has been skipped
test / publish-infra (pull_request) Has been skipped
Add a module-level BackendStatus(IntEnum) with READY=0 to replace magic
0 comparisons. Extend BottleBackend.status() with a quiet parameter:
quiet=False (default) prints diagnostics; quiet=True returns the code
silently for programmatic checks.

Add is_backend_available() (cheap PATH check, alias for has_backend) and
is_backend_ready(name, *, quiet=False) (full status() check) to the
backend package for callers that need to distinguish availability from
readiness.

Update tests/_backend.py guards to use is_backend_ready(quiet=False) so
diagnostic output is printed during test discovery, giving the operator a
concrete reason for each skip rather than a bare "prerequisites
unavailable" message.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-24 03:03:00 +00:00
didericis-claude d4d45f835e ci: reuse backend.has_backend / backend status for readiness
test / integration-docker (pull_request) Successful in 25s
test / unit (pull_request) Successful in 52s
lint / lint (push) Successful in 1m0s
test / integration-firecracker (pull_request) Successful in 2m5s
test / coverage (pull_request) Successful in 24s
test / publish-infra (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Successful in 6s
Review feedback: the hand-rolled `/dev/kvm` + docker-daemon capability
probes duplicated logic the CLI already owns. Delegate instead.

- tests/_backend.py skip guards now gate on `bot_bottle.backend.has_backend`
  (each backend's `is_available()` classmethod — the same probe behind
  `./cli.py backend status`), dropping the bespoke `Capability` probes.
- Remove tests/backend_preflight.py; the docker integration job runs
  `./cli.py backend status --backend=docker` as its preflight (clear
  per-check summary, non-zero exit when unready), matching the firecracker
  job. The firecracker preflight reverts to its original binary/KVM checks
  (backend status doesn't cover those).
- Unit test + docs updated to match.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-24 01:42:04 +00:00
didericis-claude 9fdaba4bd4 ci: backend-agnostic integration guards + per-backend preflight
tracker-policy-pr / check-pr (pull_request) Successful in 15s
test / integration-docker (pull_request) Successful in 20s
lint / lint (push) Successful in 1m5s
test / unit (pull_request) Successful in 1m49s
test / integration-firecracker (pull_request) Successful in 2m0s
test / coverage (pull_request) Successful in 16s
test / publish-infra (pull_request) Has been skipped
Integration tests now select their backend from BOT_BOTTLE_BACKEND and
skip on the capability that backend actually needs, instead of gating
every backend on unrelated Docker availability.

Task 1 — backend-agnostic guards (tests/_backend.py):
- Capability probes: docker_capability() (reachable daemon) and
  firecracker_capability() (accessible /dev/kvm + firecracker on PATH,
  Docker-independent). backend_capability()/selected_backend() resolve
  the target from BOT_BOTTLE_BACKEND (default docker).
- skip_unless_selected_backend_available() for backend-agnostic tests
  (test_sandbox_escape) — runs through whichever backend is selected and
  checks that backend's real capability.
- skip_unless_backend("docker") for Docker-implementation tests
  (DockerBroker, DockerGateway, backend.docker.*) — they no-op under a
  non-Docker run rather than testing internals that run doesn't target.
- Retires tests/_docker.py; the KVM job no longer needs SKIP_DOCKER_TESTS
  to steer Docker-only classes.

Task 2 — explicit per-backend skip visibility:
- tests/backend_preflight.py prints a clear PASS/FAIL capability line and
  exits non-zero when the selected backend is missing.
- Both integration jobs run it as a preflight, so absent infrastructure
  is surfaced at the job level instead of hidden among unittest.skip
  lines. The docker job replaces its soft "Show environment" step; the
  firecracker job keeps its richer backend-status check.

Docs (tests/README.md, docs/ci.md) updated; unit coverage for the probes,
guards, and preflight in test_backend_skip_guards.py.

Closes #414

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-23 21:33:14 -04:00
55 changed files with 1161 additions and 1836 deletions
+8 -6
View File
@@ -103,14 +103,16 @@ jobs:
- name: Install coverage
run: python3 -m pip install --break-system-packages coverage
- name: Show environment
# Fail loudly if the backend this job promises isn't actually usable,
# rather than letting every test silently `unittest.skip` and the job
# go green on zero coverage. `backend status` prints a clear per-check
# summary (docker on PATH, daemon reachable) and exits non-zero when a
# prerequisite is missing — the same readiness check the skip guards
# gate on via `has_backend`.
- name: Preflight — Docker backend is ready
run: |
python3 --version
if command -v docker >/dev/null 2>&1; then
docker version || true
else
echo "docker not on PATH — integration tests will skip"
fi
python3 cli.py backend status --backend=docker
- name: Run integration tests (docker) with coverage
env:
+1 -2
View File
@@ -26,9 +26,8 @@
# /git-gate-entrypoint.sh docker-cp'd at start time
# /git-gate/creds/* docker-cp'd at start time
# /git/* bare repos, populated at runtime
# /run/supervise/bot-bottle.db bind-mounted at run time
# /home/mitmproxy/.mitmproxy/ mitmproxy CA dir
# (No bot-bottle.db mount: the data plane reaches the supervise queue over the
# control-plane RPC and never opens the DB — PRD 0070 / issue #469.)
#
# Exposed ports inside the container:
# 9099 egress (mitmproxy, agent-facing HTTPS proxy)
+43 -4
View File
@@ -33,6 +33,7 @@ backend field; the host picks.
from __future__ import annotations
import enum
import os
import shlex
import sys
@@ -59,6 +60,12 @@ if TYPE_CHECKING:
from .freeze import CommitCancelled, Freezer, get_freezer
class BackendStatus(enum.IntEnum):
"""Return codes for BottleBackend.status(). READY == 0 so callsites
can compare against 0 or the named constant interchangeably."""
READY = 0
@dataclass(frozen=True)
class BottleSpec:
"""CLI-supplied intent. Backend-agnostic — each backend's prepare
@@ -611,12 +618,18 @@ class BottleBackend(ABC, Generic[PlanT, CleanupT]):
@classmethod
@abstractmethod
def status(cls) -> int:
def status(cls, *, quiet: bool = False) -> int:
"""Report whether this backend's prerequisites are satisfied on
the host — binaries, daemon reachability, network pool, range
conflicts, etc. Prints a human-readable summary; returns 0 when
the backend is ready to launch and non-zero when something is
missing. Invoked by `./cli.py backend status [--backend=…]`."""
conflicts, etc. Returns BackendStatus.READY (0) when the backend
is ready to launch and non-zero when something is missing.
When quiet=False (default) prints a human-readable summary to
stderr. When quiet=True returns the status code silently —
useful for cheap programmatic checks.
Invoked by `./cli.py backend status [--backend=…]` (quiet=False)
and by is_backend_ready() (caller-controlled)."""
@classmethod
@abstractmethod
@@ -811,6 +824,29 @@ def has_backend(name: str) -> bool:
return backends[name].is_available()
def is_backend_available(name: str) -> bool:
"""Cheap availability check: is the backend's binary on PATH?
Suitable for cleanup enumeration and auto-selection — does NOT probe
the daemon or network pool. Use is_backend_ready() for a full
readiness check before launching tests."""
return has_backend(name)
def is_backend_ready(name: str, *, quiet: bool = False) -> bool:
"""Full readiness check: passes all of the backend's status() checks.
When quiet=False the backend prints diagnostic output explaining what
is missing — intended for test-suite guards that run at discovery time
so the operator sees a concrete failure reason for each skip.
Returns False for unknown backend names."""
backends = _get_backends()
if name not in backends:
return False
return backends[name].status(quiet=quiet) == BackendStatus.READY
def enumerate_active_agents() -> list[ActiveAgent]:
"""All currently-running agents, across every available
backend. Used by CLI `list active` and the dashboard's agents
@@ -835,6 +871,7 @@ def enumerate_active_agents() -> list[ActiveAgent]:
__all__ = [
"ActiveAgent",
"BackendStatus",
"Bottle",
"BottleBackend",
"BottleCleanupPlan",
@@ -847,5 +884,7 @@ __all__ = [
"get_bottle_backend",
"get_freezer",
"has_backend",
"is_backend_available",
"is_backend_ready",
"known_backend_names",
]
+6 -2
View File
@@ -20,7 +20,8 @@ infrastructure: CA install and git copy-in.
from __future__ import annotations
import shutil
from contextlib import contextmanager
import io
from contextlib import contextmanager, redirect_stderr
from pathlib import Path
from typing import Generator, Sequence
@@ -60,8 +61,11 @@ class DockerBottleBackend(BottleBackend["DockerBottlePlan", "DockerBottleCleanup
return _setup.setup()
@classmethod
def status(cls) -> int:
def status(cls, *, quiet: bool = False) -> int:
from . import setup as _setup
if quiet:
with redirect_stderr(io.StringIO()):
return _setup.status()
return _setup.status()
@classmethod
+6 -2
View File
@@ -8,7 +8,8 @@ fail-closed nftables egress boundary. Selected by
from __future__ import annotations
from contextlib import contextmanager
import io
from contextlib import contextmanager, redirect_stderr
from pathlib import Path
from typing import Generator, Sequence
@@ -52,8 +53,11 @@ class FirecrackerBottleBackend(
return _setup.setup()
@classmethod
def status(cls) -> int:
def status(cls, *, quiet: bool = False) -> int:
from . import setup as _setup
if quiet:
with redirect_stderr(io.StringIO()):
return _setup.status()
return _setup.status()
@classmethod
+5 -50
View File
@@ -33,17 +33,11 @@ from pathlib import Path
from typing import Generator
from ...log import die, info
from ...paths import CONTROL_PLANE_TOKEN_FILENAME, bot_bottle_root
from .. import util as backend_util
from ..docker import util as docker_mod
from ..docker.gateway_provision import GatewayProvisionError
from . import firecracker_vm, infra_artifact, netpool, util
# Where the infra VM keeps its control-plane signing key (generated on the
# persistent /dev/vdb volume mounted at BOT_BOTTLE_ROOT). The host mirrors it
# back so the CLI signs `cli` tokens the VM verifies (issue #469 review).
_GUEST_SIGNING_KEY_PATH = "/var/lib/bot-bottle/control-plane-token"
# The single infra-VM image: gateway data plane + baked control-plane source
# (Dockerfile.infra FROM the gateway image). Built from source by default;
# a pull-from-registry mode lands later.
@@ -173,45 +167,21 @@ def ensure_running() -> InfraVm:
want = _expected_version()
if _adoptable(key, url, want):
info(f"adopting running infra VM at {url}")
return _with_signing_key(InfraVm(guest_ip=slot.guest_ip, private_key=key))
return InfraVm(guest_ip=slot.guest_ip, private_key=key)
with _singleton_lock():
# Re-check under the lock: another launcher may have booted it while
# we waited for the lock (double-checked, so we adopt not re-boot).
if _adoptable(key, url, want):
info(f"adopting running infra VM at {url}")
return _with_signing_key(InfraVm(guest_ip=slot.guest_ip, private_key=key))
return InfraVm(guest_ip=slot.guest_ip, private_key=key)
# Clear a stale/hung/OUTDATED VM holding the link before booting fresh.
stop()
ensure_built()
infra = boot()
wait_for_health(infra)
_record_booted_version(want)
return _with_signing_key(infra)
def _with_signing_key(infra: InfraVm) -> InfraVm:
"""Mirror the infra VM's control-plane signing key (generated on its
persistent volume) into the host's control-plane-token file, so the host CLI
signs `cli` tokens the VM verifies (issue #469 review). Best-effort: an
unreadable key is logged, not fatal — the VM still enforces auth, but the CLI
may then be rejected until the key is readable. Returns `infra` for chaining."""
proc = subprocess.run(
util.ssh_base_argv(infra.private_key, infra.guest_ip)
+ [f"cat {_GUEST_SIGNING_KEY_PATH}"],
capture_output=True, text=True, check=False,
)
signing_key = proc.stdout.strip()
if proc.returncode != 0 or not signing_key:
info("infra signing key not yet readable; control-plane auth may fail")
return infra
path = bot_bottle_root() / CONTROL_PLANE_TOKEN_FILENAME
path.parent.mkdir(parents=True, exist_ok=True)
fd = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w") as f:
f.write(signing_key)
os.chmod(path, stat.S_IRUSR | stat.S_IWUSR)
return infra
@contextmanager
@@ -517,31 +487,16 @@ mount -t ext4 /dev/vdb /var/lib/bot-bottle 2>/dev/null || true
# Control plane. Source is baked at /app; the package is stdlib-only.
cd /app
# Control-plane signing key + a pre-minted `gateway` JWT, scoped per-process
# (issue #469 review). The key is generated once on the persistent volume and
# handed ONLY to the orchestrator (to verify tokens); the data-plane daemons get
# the `gateway` JWT they present, never the key. Without this the control plane
# would run OPEN and a compromised egress / supervise / git-http daemon in this
# same VM — reaching the orchestrator over 127.0.0.1, past the nft boundary that
# only fences off the separate agent VM — could drive the operator routes
# (approve its own supervise proposals, rewrite policy, read injected tokens).
CP_KEY=$(BOT_BOTTLE_ROOT=/var/lib/bot-bottle python3 -c 'from bot_bottle.paths import host_control_plane_token as t; print(t())')
GW_JWT=$(BB_SIGNING_KEY="$CP_KEY" python3 -c 'import os; from bot_bottle.control_auth import mint, ROLE_GATEWAY; print(mint(ROLE_GATEWAY, os.environ["BB_SIGNING_KEY"]))')
BOT_BOTTLE_ROOT=/var/lib/bot-bottle BOT_BOTTLE_CONTROL_PLANE_TOKEN="$CP_KEY" python3 -m bot_bottle.orchestrator \\
BOT_BOTTLE_ROOT=/var/lib/bot-bottle python3 -m bot_bottle.orchestrator \\
--host 0.0.0.0 --port {CONTROL_PLANE_PORT} --broker stub &
# Gateway data plane, multi-tenant: each request resolves source-IP ->
# policy against the local control plane. The VM backend reaches git over
# git-http (9420), so the git:// daemon (git-gate, needs a per-bottle
# entrypoint the consolidated model doesn't use) is left out. No
# SUPERVISE_DB_PATH: the data plane reaches the supervise queue over the
# control-plane RPC and never opens bot-bottle.db (PRD 0070 / #469). It presents
# the pre-minted `gateway` JWT; gateway_init keeps the signing key out of the
# data-plane daemons' env.
# entrypoint the consolidated model doesn't use) is left out.
BOT_BOTTLE_GATEWAY_DAEMONS=egress,git-http,supervise \\
BOT_BOTTLE_ORCHESTRATOR_URL=http://127.0.0.1:{CONTROL_PLANE_PORT} \\
BOT_BOTTLE_CONTROL_AUTH_JWT="$GW_JWT" \\
SUPERVISE_DB_PATH=/var/lib/bot-bottle/db/bot-bottle.db \\
python3 -m bot_bottle.gateway_init &
# Reap as PID 1; children are backgrounded, so `wait` blocks.
+60 -6
View File
@@ -13,6 +13,7 @@ generic `./cli.py backend {setup,status}` command dispatches to.
from __future__ import annotations
import fcntl
import os
import shutil
import subprocess
@@ -22,6 +23,9 @@ from pathlib import Path
from . import netpool
from . import util
# KVM_GET_API_VERSION = _IO(KVMIO=0xAE, 0x00): cheapest proof of KVM access.
_KVM_GET_API_VERSION = 0xAE00
_FC_RELEASES = "https://github.com/firecracker-microvm/firecracker/releases"
_UNIT_PATH = Path("/etc/systemd/system") / netpool.SYSTEMD_UNIT
@@ -219,14 +223,64 @@ def teardown() -> int:
return 0
def _firecracker_binary_ok() -> bool:
"""True iff the firecracker binary is on PATH and `--version` exits 0."""
if shutil.which("firecracker") is None:
return False
try:
return subprocess.run(
["firecracker", "--version"],
capture_output=True, check=False, timeout=5,
).returncode == 0
except (OSError, subprocess.TimeoutExpired):
return False
def _kvm_accessible() -> bool:
"""True iff /dev/kvm can be opened and responds to KVM_GET_API_VERSION.
Uses an ioctl rather than os.access() so that a device with correct
permission bits but a non-functional KVM subsystem is caught here
rather than at VM boot time."""
if not os.path.exists(util._KVM_DEVICE):
return False
try:
with open(util._KVM_DEVICE, "rb") as kvm:
fcntl.ioctl(kvm, _KVM_GET_API_VERSION)
return True
except OSError:
return False
def status() -> int:
# Readiness == what the launch preflight hard-requires: the TAP pool
# present (unprivileged, authoritative) and no range overlap. Listing
# the nft table usually needs root, so — like the preflight — an
# unconfirmable table is reported but NOT treated as not-ready; the
# post-boot isolation probe is the authoritative check. This keeps an
# unprivileged `backend status` usable as a launch gate.
# Readiness == what the launch preflight hard-requires: the binary
# executable, /dev/kvm accessible, the TAP pool present, and no range
# overlap. Listing the nft table usually needs root, so — like the
# preflight — an unconfirmable table is reported but NOT treated as
# not-ready; the post-boot isolation probe is the authoritative check.
# This keeps an unprivileged `backend status` usable as a launch gate.
ok = True
if _firecracker_binary_ok():
sys.stderr.write(f"firecracker binary: ok ({shutil.which('firecracker')})\n")
else:
fc_path = shutil.which("firecracker")
if fc_path is None:
sys.stderr.write("firecracker binary: NOT found on PATH\n")
else:
sys.stderr.write(
f"firecracker binary: found ({fc_path}) but `--version` failed\n"
)
ok = False
if _kvm_accessible():
sys.stderr.write(f"KVM: {util._KVM_DEVICE} accessible\n")
else:
if not os.path.exists(util._KVM_DEVICE):
sys.stderr.write(f"KVM: {util._KVM_DEVICE} not present\n")
else:
sys.stderr.write(
f"KVM: {util._KVM_DEVICE} not accessible (open/ioctl failed)\n"
)
ok = False
missing = netpool.missing_taps()
total = netpool.pool_size()
if missing:
@@ -2,7 +2,8 @@
from __future__ import annotations
from contextlib import contextmanager
import io
from contextlib import contextmanager, redirect_stderr
from pathlib import Path
from typing import Generator, Sequence
@@ -43,8 +44,11 @@ class MacosContainerBottleBackend(
return _setup.setup()
@classmethod
def status(cls) -> int:
def status(cls, *, quiet: bool = False) -> int:
from . import setup as _setup
if quiet:
with redirect_stderr(io.StringIO()):
return _setup.status()
return _setup.status()
@classmethod
+9 -21
View File
@@ -48,9 +48,7 @@ from ...orchestrator.lifecycle import (
OrchestratorStartError,
source_hash,
)
from ...control_auth import ROLE_GATEWAY, mint
from ...paths import (
CONTROL_AUTH_JWT_ENV,
CONTROL_PLANE_TOKEN_ENV,
HOST_DB_FILENAME,
host_control_plane_token,
@@ -102,12 +100,10 @@ def _init_script(port: int) -> str:
f"( cd {_SRC_IN_CONTAINER} && BOT_BOTTLE_ROOT={_DB_ROOT_IN_CONTAINER} "
f"python3 -m bot_bottle.orchestrator --host 0.0.0.0 --port {port} "
"--broker stub ) &\n"
# Gateway data plane, multi-tenant against the local control plane. No
# SUPERVISE_DB_PATH: the data plane reaches the supervise queue over the
# control-plane RPC and never opens bot-bottle.db (PRD 0070 / #469).
# Gateway data plane, multi-tenant against the local control plane.
f"( cd /app && BOT_BOTTLE_GATEWAY_DAEMONS={_GATEWAY_DAEMONS} "
f"BOT_BOTTLE_ORCHESTRATOR_URL=http://127.0.0.1:{port} "
f"python3 -m bot_bottle.gateway_init ) &\n"
f"SUPERVISE_DB_PATH={_DB_PATH_IN_CONTAINER} python3 -m bot_bottle.gateway_init ) &\n"
"while : ; do wait ; done\n"
)
@@ -239,26 +235,18 @@ class MacosInfraService:
# Baked onto the container so `_source_current` can detect a real
# control-plane code change and recreate.
"--env", f"BOT_BOTTLE_SOURCE_HASH={current_hash}",
# The control-plane signing key (control plane: verifies tokens) and
# the pre-minted `gateway` JWT (the gateway's PolicyResolver: presents
# it) — they share this one container, and gateway_init scopes each to
# its process so a compromised data-plane daemon never sees the key
# (issue #469 review). Bare `--env NAME` inherits the value from the
# run process below, so neither lands on argv or in `container
# inspect`'s command line. The agent runs in a SEPARATE container that
# is never given these vars, which is the whole point.
# The control-plane secret, for BOTH the control plane (to require
# it) and the gateway's PolicyResolver (to present it) — they share
# this one container. Bare `--env NAME` inherits the value from the
# run process below, so the secret never lands on argv or in
# `container inspect`'s command line. The agent runs in a SEPARATE
# container that is never given this var, which is the whole point.
"--env", CONTROL_PLANE_TOKEN_ENV,
"--env", CONTROL_AUTH_JWT_ENV,
"--entrypoint", "sh",
self.image,
"-c", _init_script(self.port),
]
_signing_key = host_control_plane_token()
run_env = {
**os.environ,
CONTROL_PLANE_TOKEN_ENV: _signing_key,
CONTROL_AUTH_JWT_ENV: mint(ROLE_GATEWAY, _signing_key),
}
run_env = {**os.environ, CONTROL_PLANE_TOKEN_ENV: host_control_plane_token()}
result = container_mod.run_container_argv(argv, env=run_env)
if result.returncode != 0:
raise OrchestratorStartError(
-100
View File
@@ -1,100 +0,0 @@
"""Role-scoped control-plane credentials (issue #469 review follow-up).
The control plane no longer trusts a single shared bearer secret for every
route. Instead the orchestrator holds a *signing key* and issues short,
HMAC-signed tokens (compact HS256 JWTs) that embed a **role** naming the kind
of caller:
* ``gateway`` the data plane (egress / git-gate / supervise). Restricted to
the agent-facing routes it actually needs (``/resolve``,
``/supervise/propose``, ``/supervise/poll``).
* ``cli`` the host operator / launcher. Full access to the mutating and
operator routes (launch/teardown, policy, ``/supervise/respond``, ).
Only the orchestrator (and the host CLI, which shares the host trust domain)
holds the signing key; the gateway is handed a pre-minted ``gateway`` token it
cannot rewrite into a ``cli`` token. So a compromised data-plane process can no
longer approve its own supervise proposals or drive operator routes it can
only present the ``gateway`` role it was issued (control_plane rejects it on
operator routes with 403).
Stdlib-only (HMAC-SHA256 over a JSON payload); no JWT dependency the project
carries no runtime pip deps. Tokens are **signed, not encrypted** (the role is
not a secret) and **long-lived** (parity with the static token they replace;
the security win is the unforgeable role claim, not rotation).
"""
from __future__ import annotations
import base64
import binascii
import hashlib
import hmac
import json
ROLE_GATEWAY = "gateway"
ROLE_CLI = "cli"
ROLES: frozenset[str] = frozenset({ROLE_GATEWAY, ROLE_CLI})
_ALG = "HS256"
def _b64url_encode(raw: bytes) -> str:
return base64.urlsafe_b64encode(raw).rstrip(b"=").decode("ascii")
def _b64url_decode(text: str) -> bytes:
padded = text + "=" * (-len(text) % 4)
return base64.urlsafe_b64decode(padded.encode("ascii"))
def _sign(secret: str, signing_input: str) -> str:
mac = hmac.new(secret.encode("utf-8"), signing_input.encode("ascii"), hashlib.sha256)
return _b64url_encode(mac.digest())
# The fixed, canonical JOSE header — the same for every token we mint.
_HEADER_SEGMENT = _b64url_encode(
json.dumps({"alg": _ALG, "typ": "JWT"}, separators=(",", ":")).encode("utf-8")
)
def mint(role: str, secret: str) -> str:
"""A compact HS256 token asserting `role`, signed with `secret`.
Raises ValueError for an unknown role (mint only what the control plane will
accept) or an empty signing key (an unsigned credential is never valid)."""
if role not in ROLES:
raise ValueError(f"unknown control-plane role {role!r}")
if not secret:
raise ValueError("cannot mint a control-plane token without a signing key")
payload = _b64url_encode(json.dumps({"role": role}, separators=(",", ":")).encode("utf-8"))
signing_input = f"{_HEADER_SEGMENT}.{payload}"
return f"{signing_input}.{_sign(secret, signing_input)}"
def verify(token: str, secret: str) -> str | None:
"""The role a valid `token` carries, or None if it is malformed, wrongly
signed, or names an unknown role. Constant-time signature check; rejects any
header whose alg isn't HS256 (no alg-confusion / `none`)."""
if not token or not secret:
return None
parts = token.split(".")
if len(parts) != 3:
return None
header_b64, payload_b64, sig = parts
expected = _sign(secret, f"{header_b64}.{payload_b64}")
if not hmac.compare_digest(sig, expected):
return None
try:
header = json.loads(_b64url_decode(header_b64))
payload = json.loads(_b64url_decode(payload_b64))
except (ValueError, binascii.Error):
return None
if not isinstance(header, dict) or header.get("alg") != _ALG:
return None
role = payload.get("role") if isinstance(payload, dict) else None
return role if isinstance(role, str) and role in ROLES else None
__all__ = ["ROLE_GATEWAY", "ROLE_CLI", "ROLES", "mint", "verify"]
+28 -50
View File
@@ -40,13 +40,8 @@ from bot_bottle.egress_addon_core import (
scan_inbound,
scan_outbound,
)
from bot_bottle.policy_resolver import PolicyResolveError, PolicyResolver
from bot_bottle.supervise_types import (
STATUS_APPROVED,
STATUS_MODIFIED,
STATUSES,
TOOL_EGRESS_TOKEN_ALLOW,
)
from bot_bottle import supervise as _sv
from bot_bottle.policy_resolver import PolicyResolver
INTROSPECT_HOST = "_egress.local"
@@ -319,15 +314,8 @@ class EgressAddon:
flow.request.headers.pop("Proxy-Authorization", None)
flow.request.headers.pop(IDENTITY_HEADER, None)
conn = flow.client_conn
conn_id = getattr(conn, "id", "") if conn is not None else ""
if not token and conn_id:
token = self._conn_tokens.get(conn_id, "")
# Remember the token per connection so a later token-block on this flow
# can attribute its supervise proposal by (source_ip, identity_token)
# over RPC — plain-HTTP requests carry it here, HTTPS tunnels captured
# it at CONNECT (http_connect); both land in `_conn_tokens`.
if token and conn_id:
self._conn_tokens[conn_id] = token
if not token and conn is not None:
token = self._conn_tokens.get(getattr(conn, "id", ""), "")
return token
def http_connect(self, flow: http.HTTPFlow) -> None:
@@ -604,48 +592,42 @@ class EgressAddon:
redact_tokens(request_path, env=env),
result,
)
# Attribute the proposal by (source_ip, identity_token) over the control
# plane — the data plane no longer opens bot-bottle.db (PRD 0070 / #469).
# source_ip + token come from this bottle's connection; `slug` is kept
# only for its own DLP safelist.
conn = flow.client_conn
source_ip = conn.peername[0] if conn is not None and conn.peername else ""
token = self._conn_tokens.get(getattr(conn, "id", ""), "") if conn is not None else ""
proposal = _sv.Proposal.new(
bottle_slug=slug,
tool=_sv.TOOL_EGRESS_TOKEN_ALLOW,
proposed_file=payload,
justification=_TOKEN_ALLOW_JUSTIFICATION,
current_file_hash=_sv.sha256_hex(payload),
)
try:
proposal_id = self._resolver.propose_supervise(
source_ip, token,
tool=TOOL_EGRESS_TOKEN_ALLOW,
proposed_file=payload,
justification=_TOKEN_ALLOW_JUSTIFICATION,
)
except PolicyResolveError as e:
_sv.write_proposal(proposal)
except OSError as e:
sys.stderr.write(
f"egress: could not queue token-allow proposal: {e}; "
"blocking request\n"
)
proposal_id = None
if proposal_id is None:
self._block(flow, f"egress DLP: {result.reason}", ctx=self._req_ctx(flow))
return False
sys.stderr.write(json.dumps({
"event": "egress_token_supervise",
"reason": f"egress DLP: {result.reason}",
"proposal": proposal_id,
"proposal": proposal.id,
**self._req_ctx(flow),
}) + "\n")
response = await self._await_token_response(proposal_id, source_ip, token)
response = await self._await_token_response(proposal.id, slug)
_sv.archive_proposal(slug, proposal.id)
if response is not None and response.get("status") in (
STATUS_APPROVED, STATUS_MODIFIED,
if response is not None and response.status in (
_sv.STATUS_APPROVED, _sv.STATUS_MODIFIED,
):
self._safe_tokens_for(slug).add(result.matched)
if self._flow_log(flow) >= LOG_BLOCKS:
sys.stderr.write(json.dumps({
"event": "egress_token_allowed",
"reason": f"egress DLP: {result.reason}",
"proposal": proposal_id,
"proposal": proposal.id,
**self._req_ctx(flow),
}) + "\n")
return True
@@ -663,23 +645,19 @@ class EgressAddon:
async def _await_token_response(
self,
proposal_id: str,
source_ip: str,
token: str,
) -> "dict[str, object] | None":
"""Poll the control plane for the operator's decision without blocking
the proxy event loop. Returns the terminal `{status, ...}` payload once
decided, or None on timeout. A transient orchestrator error (or a
`pending`/`unknown` status) is retried until the deadline, then fails
closed the caller blocks the request."""
slug: str,
) -> "_sv.Response | None":
"""Poll the DB for the operator's response without blocking the
proxy event loop. Returns the Response, or None on timeout."""
loop = asyncio.get_running_loop()
deadline = loop.time() + self._token_allow_timeout
while True:
try:
result = self._resolver.poll_supervise(source_ip, token, proposal_id)
except PolicyResolveError:
result = None
if result is not None and result.get("status") in STATUSES:
return result
return _sv.read_response(slug, proposal_id)
except (OSError, ValueError, KeyError):
# Not written yet, or a partial/malformed write — retry until
# the deadline, then fail closed.
pass
if loop.time() >= deadline:
return None
await asyncio.sleep(TOKEN_ALLOW_POLL_INTERVAL_SECONDS)
+10 -27
View File
@@ -54,15 +54,6 @@ class _DaemonSpec:
_EGRESS_ONLY_ENV_PREFIXES: tuple[str, ...] = ("EGRESS_TOKEN_",)
_READY_GATED_DAEMONS: tuple[str, ...] = ("git-gate", "git-http")
# The control-plane signing key is the orchestrator's alone — verifying tokens.
# The data-plane daemons instead hold the pre-minted `gateway` JWT they present.
# Scoping each to its process (even in the combined infra container) keeps a
# compromised data-plane daemon from reading the key and minting a `cli` token
# (issue #469 review). Values match paths.CONTROL_PLANE_TOKEN_ENV /
# CONTROL_AUTH_JWT_ENV; hardcoded here so this supervisor stays import-light.
_SIGNING_KEY_ENV = "BOT_BOTTLE_CONTROL_PLANE_TOKEN"
_GATEWAY_JWT_ENV = "BOT_BOTTLE_CONTROL_AUTH_JWT"
# Daemons that must be requested explicitly via BOT_BOTTLE_GATEWAY_DAEMONS
# and are NOT started in the default (env-var-unset) case. The orchestrator
# only runs in the combined infra container, never in a standalone gateway.
@@ -70,24 +61,16 @@ _OPT_IN_DAEMONS: frozenset[str] = frozenset({"orchestrator"})
def _env_for_daemon(name: str, base_env: dict[str, str]) -> dict[str, str]:
"""Per-daemon env, scoped to what each process legitimately needs.
Returns a fresh dict callers can mutate without affecting `base_env`.
* `EGRESS_TOKEN_*` upstream-auth slots go to egress only.
* the control-plane signing key goes to the orchestrator only.
* the pre-minted `gateway` JWT goes to the data-plane daemons only (the
orchestrator verifies tokens, it never presents one)."""
env = dict(base_env)
if name != "egress":
env = {
k: v for k, v in env.items()
if not any(k.startswith(p) for p in _EGRESS_ONLY_ENV_PREFIXES)
}
if name == "orchestrator":
env.pop(_GATEWAY_JWT_ENV, None)
else:
env.pop(_SIGNING_KEY_ENV, None)
return env
"""Egress sees the full bundle env. Everyone else gets a copy
with `EGRESS_TOKEN_*` (and any other future egress-only
credential slots) stripped. Returns a fresh dict callers
can mutate without affecting `base_env`."""
if name == "egress":
return dict(base_env)
return {
k: v for k, v in base_env.items()
if not any(k.startswith(p) for p in _EGRESS_ONLY_ENV_PREFIXES)
}
# The orchestrator is listed first so it starts before the gateway daemons,
+37 -44
View File
@@ -272,27 +272,21 @@ supervise_gitleaks_allow() {
proposal_id=$(
PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" GITLEAKS_ALLOW_REF="$ref" python3 - "$report_file" <<'PY'
import datetime
import hashlib
import json
import os
import sys
from pathlib import Path
# Queue over the control-plane RPC — the git-gate data plane no longer opens
# bot-bottle.db (PRD 0070 / issue #469). Attribution is by (source_ip,
# identity_token), resolved server-side, so the proposal lands under the
# calling bottle exactly as a direct write once did.
try:
from bot_bottle.policy_resolver import PolicyResolver, PolicyResolveError
from bot_bottle.supervise_types import TOOL_GITLEAKS_ALLOW
import supervise as _sv
except ImportError:
from policy_resolver import PolicyResolver, PolicyResolveError
from supervise_types import TOOL_GITLEAKS_ALLOW
from bot_bottle import supervise as _sv
report_path = Path(sys.argv[1])
source_ip = os.environ.get("SUPERVISE_SOURCE_IP", "")
token = os.environ.get("SUPERVISE_IDENTITY_TOKEN", "")
orch_url = os.environ.get("BOT_BOTTLE_ORCHESTRATOR_URL", "")
if not source_ip or not orch_url:
slug = os.environ.get("SUPERVISE_BOTTLE_SLUG", "")
if not slug:
sys.exit(2)
try:
@@ -329,21 +323,19 @@ for i, finding in enumerate(raw, 1):
])
payload = "\n".join(lines).rstrip() + "\n"
try:
proposal_id = PolicyResolver(orch_url).propose_supervise(
source_ip, token,
tool=TOOL_GITLEAKS_ALLOW,
proposed_file=payload,
justification=(
"git-gate found gitleaks findings hidden by # gitleaks:allow; "
"approve only for dummy test fixtures or confirmed false positives"
),
)
except PolicyResolveError:
sys.exit(4)
if not proposal_id:
sys.exit(2)
print(proposal_id)
proposal = _sv.Proposal.new(
bottle_slug=slug,
tool=_sv.TOOL_GITLEAKS_ALLOW,
proposed_file=payload,
justification=(
"git-gate found gitleaks findings hidden by # gitleaks:allow; "
"approve only for dummy test fixtures or confirmed false positives"
),
current_file_hash=hashlib.sha256(payload.encode("utf-8")).hexdigest(),
now=datetime.datetime.now(datetime.timezone.utc),
)
_sv.write_proposal(proposal)
print(proposal.id)
PY
)
rc=$?
@@ -356,6 +348,7 @@ PY
return 1
fi
slug=${SUPERVISE_BOTTLE_SLUG:-}
timeout=${SUPERVISE_GITLEAKS_ALLOW_TIMEOUT_SECONDS:-300}
case "$timeout" in
''|*[!0-9]*)
@@ -367,30 +360,20 @@ PY
echo "git-gate: approve with './cli.py supervise' to continue this push" >&2
waited=0
while [ "$waited" -lt "$timeout" ]; do
status=$(PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" python3 - "$proposal_id" <<'PY'
import os
status=$(PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" python3 - "$slug" "$proposal_id" <<'PY'
import sys
# Non-blocking poll over the control plane. A decided proposal is archived
# server-side on read, so no separate archive step is needed here.
try:
from bot_bottle.policy_resolver import PolicyResolver, PolicyResolveError
import supervise as _sv
except ImportError:
from policy_resolver import PolicyResolver, PolicyResolveError
from bot_bottle import supervise as _sv
source_ip = os.environ.get("SUPERVISE_SOURCE_IP", "")
token = os.environ.get("SUPERVISE_IDENTITY_TOKEN", "")
orch_url = os.environ.get("BOT_BOTTLE_ORCHESTRATOR_URL", "")
slug = sys.argv[1]
try:
result = PolicyResolver(orch_url).poll_supervise(source_ip, token, sys.argv[1])
except PolicyResolveError:
response = _sv.read_response(slug, sys.argv[2])
except FileNotFoundError:
sys.exit(2)
if result is None:
sys.exit(2)
status = result.get("status", "")
if status in ("pending", "unknown"):
sys.exit(2)
print(status)
print(response.status)
PY
)
rc=$?
@@ -402,6 +385,16 @@ PY
if [ -n "$status" ]; then
case "$status" in
approved|modified)
PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}" python3 - "$slug" "$proposal_id" <<'PY' || true
import sys
try:
import supervise as _sv
except ImportError:
from bot_bottle import supervise as _sv
_sv.archive_proposal(sys.argv[1], sys.argv[2])
PY
echo "git-gate: supervisor approved # gitleaks:allow for $ref" >&2
return 0
;;
+4 -7
View File
@@ -167,14 +167,11 @@ class GitHttpHandler(BaseHTTPRequestHandler):
"SERVER_PORT": str(self.server.server_port), # type: ignore
"SERVER_PROTOCOL": self.request_version,
})
# Attribute the gitleaks-allow supervise proposal (queued by
# Attribute the gitleaks-allow supervise proposal (written by
# receive-pack's pre-receive hook, a child of the CGI we spawn below) to
# the calling bottle. The hook queues over the control-plane RPC (PRD
# 0070 / issue #469), which re-resolves the bottle from these — the same
# (source_ip, identity_token) pair that selected the namespaced root
# above — so it can only ever queue its own proposals.
env["SUPERVISE_SOURCE_IP"] = self.client_address[0]
env["SUPERVISE_IDENTITY_TOKEN"] = self.headers.get(IDENTITY_HEADER, "")
# the calling bottle. The namespaced root is `<base>/<bottle_id>`, so its
# final component is the bottle id — the same per-bottle key egress uses.
env["SUPERVISE_BOTTLE_SLUG"] = sandbox_root.name
for header, variable in (
("accept", "HTTP_ACCEPT"),
("content-encoding", "HTTP_CONTENT_ENCODING"),
+5 -8
View File
@@ -18,7 +18,6 @@ import urllib.request
from collections.abc import Iterable
from dataclasses import dataclass
from ..control_auth import ROLE_CLI, mint
from ..paths import host_control_plane_token
from .control_plane import CONTROL_AUTH_HEADER
@@ -26,14 +25,12 @@ DEFAULT_TIMEOUT_SECONDS = 5.0
def _host_auth_token() -> str:
"""A freshly minted `cli`-role token, signed with the host control-plane
key, or "" if the key can't be read. The host CLI shares the orchestrator's
trust domain, so it holds the signing key and mints its own operator token;
"" means 'send no auth header' correct against an open (unconfigured)
control plane, and harmlessly rejected by a secured one."""
"""The per-host control-plane secret, or "" if it can't be read. "" means
'send no auth header' correct against an open (unconfigured) control
plane, and harmlessly rejected by a secured one."""
try:
return mint(ROLE_CLI, host_control_plane_token())
except (OSError, ValueError):
return host_control_plane_token()
except OSError:
return ""
+37 -129
View File
@@ -27,20 +27,6 @@ vsock / unix-socket portability caveats):
POST /supervise/respond -> 200 {"responded": true} | 409 (operator)
body: {"proposal_id","bottle_slug",
"decision", ["notes"],["final_file"]}
POST /supervise/propose -> 201 {"proposal_id"} | 403 (agent)
body: {"source_ip","identity_token",
"tool","proposed_file","justification"}
POST /supervise/poll -> 200 {"status", ["notes"],["final_file"]} | 403
body: {"source_ip","identity_token",
"proposal_id"}
The `/supervise/propose` + `/supervise/poll` pair is the **agent** half of the
supervise flow: the data plane (supervise / egress / git-gate) queues a proposal
and polls for its response over RPC instead of opening `bot-bottle.db` directly.
`poll` is idempotent it never archives, so a dropped connection can't lose an
operator decision (the row is reaped when the bottle is torn down / reconciled).
Both attribute the caller by `(source_ip, identity_token)` exactly like
`/resolve`, so a bottle can only ever queue or read its own proposals.
`POST /bottles` / `DELETE` drive the full launch lifecycle: they mint (or
tear down) the bottle in the registry AND broker the backend-native launch
@@ -55,6 +41,7 @@ returned only once, to the caller that launches the bottle.
from __future__ import annotations
import hmac
import http.server
import json
import os
@@ -63,40 +50,20 @@ import sys
import typing
from urllib.parse import urlsplit
from ..control_auth import ROLE_CLI, ROLES, verify
from ..paths import CONTROL_PLANE_TOKEN_ENV
from ..supervise_types import TOOLS
from .service import Orchestrator
# JSON body payload type (parsed request / rendered response).
Json = dict[str, object]
# The request header carrying the caller's role-scoped control-plane token (a
# signed JWT naming the caller's role — see control_auth). The role gates which
# routes the caller may reach: the data plane holds a `gateway` token good only
# for the agent-facing lookups; the host CLI holds a `cli` token for the
# operator/mutating routes. An agent that can merely *reach* the port holds no
# token at all, and a compromised gateway holds only `gateway` — neither can
# drive the operator routes (approve proposals, rewrite policy, read tokens).
# The request header carrying the per-host control-plane secret. Every route
# except `GET /health` requires it (see `dispatch`). The trusted callers hold
# the secret (the gateway's PolicyResolver, the host CLI's OrchestratorClient);
# an agent that can merely *reach* the port cannot present it, so it can't
# enumerate bottles, rewrite policy, read injected upstream tokens, or approve
# its own supervise proposals.
CONTROL_AUTH_HEADER = "x-bot-bottle-control-auth"
# The routes the data plane (role `gateway`) is allowed to reach — exactly the
# per-request lookups PolicyResolver makes. Every other authenticated route is
# operator-only. `cli` is a superset role: it may reach any route.
_GATEWAY_ROUTES: frozenset[tuple[str, str]] = frozenset({
("POST", "/resolve"),
("POST", "/supervise/propose"),
("POST", "/supervise/poll"),
})
def _allowed_roles(method: str, route: str) -> frozenset[str]:
"""The roles permitted on `(method, route)`: `gateway` or `cli` on the
data-plane routes, `cli`-only everywhere else."""
if (method, route) in _GATEWAY_ROUTES:
return ROLES
return frozenset({ROLE_CLI})
def _parse_json_object(body: bytes) -> Json:
"""Parse a JSON object body. Raises ValueError for non-objects / bad JSON."""
@@ -109,33 +76,28 @@ def _parse_json_object(body: bytes) -> Json:
def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
orch: Orchestrator, method: str, path: str, body: bytes, *, role: str | None = ROLE_CLI,
orch: Orchestrator, method: str, path: str, body: bytes, *, authorized: bool = True,
) -> tuple[int, Json]:
"""Route one control-plane request to a (status, payload) pair. Pure —
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 `ControlPlaneServer`) passes `cli`. 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).
The source-IP + identity-token checks inside `/resolve` and `/attribute`
authenticate the *bottle* a request is about, not the *caller*, so this role
gate is what protects the caller-privileged routes. Defaults `cli` so unit
tests of the routing logic don't have to thread it through."""
`authorized` is whether the request presented the control-plane secret (or
no secret is configured see `ControlPlaneServer`). Every route except
`GET /health` requires it: the source-IP + identity-token checks inside
`/resolve` and `/attribute` authenticate the *bottle* a request is about,
not the *caller*, so without this gate any agent that can reach the port
could rewrite another bottle's policy, read the injected upstream tokens,
or approve its own supervise proposals. Defaults True so unit tests of the
routing logic don't have to thread it through."""
route = urlsplit(path).path.rstrip("/") or "/"
if method == "GET" and route == "/health":
return 200, {"status": "ok"}
# Role gate — every route below is a trusted-caller operation. Deny before
# touching the registry / broker / supervise store.
if role is None:
if not authorized:
# Everything below is a trusted-caller operation. Deny before touching
# the registry / broker / supervise store.
return 401, {"error": "control-plane authentication required"}
if role not in _allowed_roles(method, route):
return 403, {"error": "insufficient role for this route"}
if method == "GET" and route == "/gateway":
return 200, orch.gateway_status()
@@ -273,56 +235,6 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
return 200, {"responded": True}
return 409, {"error": err}
if method == "POST" and route == "/supervise/propose":
# Agent half: queue a proposal, attributed to the caller resolved from
# (source_ip, identity_token) — never a caller-supplied slug — so the
# data plane can't forge attribution. Fail-closed 403 when unattributed.
try:
data = _parse_json_object(body)
except ValueError as e:
return 400, {"error": f"invalid JSON: {e}"}
source_ip = data.get("source_ip")
token = data.get("identity_token")
tool = data.get("tool")
proposed_file = data.get("proposed_file")
justification = data.get("justification")
if not isinstance(source_ip, str) or not source_ip:
return 400, {"error": "source_ip (string) is required"}
if not isinstance(tool, str) or tool not in TOOLS:
return 400, {"error": f"tool (string) must be one of {TOOLS}"}
if not isinstance(proposed_file, str) or not proposed_file:
return 400, {"error": "proposed_file (string) is required"}
if not isinstance(justification, str) or not justification:
return 400, {"error": "justification (string) is required"}
rec = orch.resolve(source_ip, token if isinstance(token, str) else "")
if rec is None:
return 403, {"error": "unattributed"}
proposal_id = orch.supervise_queue_proposal(
rec.bottle_id, tool=tool, proposed_file=proposed_file,
justification=justification,
)
return 201, {"proposal_id": proposal_id}
if method == "POST" and route == "/supervise/poll":
# Agent half: non-blocking read of the caller's own proposal decision.
# Attributed like /propose, and scoped to the resolved bottle id, so a
# guessed proposal_id can never read another bottle's response.
try:
data = _parse_json_object(body)
except ValueError as e:
return 400, {"error": f"invalid JSON: {e}"}
source_ip = data.get("source_ip")
token = data.get("identity_token")
proposal_id = data.get("proposal_id")
if not isinstance(source_ip, str) or not source_ip:
return 400, {"error": "source_ip (string) is required"}
if not isinstance(proposal_id, str) or not proposal_id:
return 400, {"error": "proposal_id (string) is required"}
rec = orch.resolve(source_ip, token if isinstance(token, str) else "")
if rec is None:
return 403, {"error": "unattributed"}
return 200, orch.supervise_poll_response(rec.bottle_id, proposal_id)
if method == "POST" and route == "/resolve":
# The per-request lookup the multi-tenant gateway makes: returns the
# bottle's policy. Requires a matching (source_ip, identity_token)
@@ -368,10 +280,10 @@ class Handler(http.server.BaseHTTPRequestHandler):
assert isinstance(server, ControlPlaneServer)
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(CONTROL_AUTH_HEADER, ""))
authorized = server.is_authorized(self.headers.get(CONTROL_AUTH_HEADER, ""))
try:
status, payload = dispatch(
server.orchestrator, method, self.path, body, role=role)
server.orchestrator, method, self.path, body, authorized=authorized)
except Exception as e: # noqa: BLE001 — the control plane must stay up
sys.stderr.write(f"orchestrator: {method} {self.path} failed: {e!r}\n")
sys.stderr.flush()
@@ -399,24 +311,22 @@ class Handler(http.server.BaseHTTPRequestHandler):
class ControlPlaneServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
"""Threading HTTP server that carries the orchestrator for its handlers.
Holds the per-host control-plane *signing key* (from
`$BOT_BOTTLE_CONTROL_PLANE_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)."""
Holds the per-host control-plane secret (from `$BOT_BOTTLE_CONTROL_PLANE_TOKEN`,
injected by the launcher into this container only). When a secret is set,
every route but `/health` requires it; when it is unset the server runs
**open** and says so loudly at startup a fail-visible fallback for tests
and any backend that hasn't wired the secret yet (e.g. Firecracker, whose
nft boundary already blocks agents from the control-plane port)."""
daemon_threads = True
allow_reuse_address = True
def __init__(self, address: tuple[str, int], orchestrator: Orchestrator) -> None:
self.orchestrator = orchestrator
self._signing_key = os.environ.get(CONTROL_PLANE_TOKEN_ENV, "").strip()
if not self._signing_key:
self._auth_token = os.environ.get(CONTROL_PLANE_TOKEN_ENV, "").strip()
if not self._auth_token:
sys.stderr.write(
"orchestrator: WARNING — no control-plane signing key "
"orchestrator: WARNING — no control-plane secret "
f"(${CONTROL_PLANE_TOKEN_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 "
@@ -425,15 +335,13 @@ class ControlPlaneServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
sys.stderr.flush()
super().__init__(address, Handler)
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
return verify(presented, self._signing_key)
def is_authorized(self, presented: str) -> bool:
"""True iff the request may proceed past `/health`: either no secret is
configured (open mode) or the presented header matches it. Constant-time
compare so a wrong token leaks nothing timing-wise."""
if not self._auth_token:
return True
return hmac.compare_digest(presented, self._auth_token)
def make_server(
+26 -14
View File
@@ -22,13 +22,19 @@ import os
import time
from pathlib import Path
from ..control_auth import ROLE_GATEWAY, mint
from ..docker_cmd import run_docker
from ..paths import (
CONTROL_AUTH_JWT_ENV,
CONTROL_PLANE_TOKEN_ENV,
host_control_plane_token,
host_db_path,
host_gateway_ca_dir,
)
from ..supervise import DB_PATH_IN_CONTAINER
# The host DB dir is bind-mounted here so the gateway's supervise daemon
# writes its queued proposals into the ONE host DB (the same file the
# orchestrator container opens and the operator reaches over HTTP).
_SUPERVISE_DB_DIR_IN_CONTAINER = os.path.dirname(DB_PATH_IN_CONTAINER)
# The gateway's mitmproxy writes its CA a beat after the container starts, so
# reads poll for it rather than assuming it's there on a fresh launch.
@@ -69,6 +75,14 @@ GATEWAY_DOCKERFILE = "Dockerfile.gateway"
_REPO_ROOT = Path(__file__).resolve().parents[2]
def _host_db_dir() -> str:
"""The host DB directory (created if missing), for the gateway's
supervise-DB bind-mount."""
db_dir = host_db_path().parent
db_dir.mkdir(parents=True, exist_ok=True)
return str(db_dir)
def rotate_gateway_ca(ca_dir: Path | None = None) -> list[Path]:
"""Delete the persisted mitmproxy CA so the next gateway start mints a
fresh one the explicit, deliberate CA-rollover path (issue #450).
@@ -241,10 +255,11 @@ class DockerGateway(Gateway):
# container recreation AND docker volume pruning (agents trust it)
# — see host_gateway_ca_dir / issue #450.
"--volume", f"{host_gateway_ca_dir()}:{MITMPROXY_HOME}",
# No DB mount: the data plane (egress / supervise / git-gate) reaches
# the supervise queue over the control-plane RPC and never opens
# bot-bottle.db, so the gateway container gets no file handle on it
# (PRD 0070 / issue #469).
# Share the one host DB: the supervise daemon queues proposals
# into the same file the orchestrator (and the operator, over
# HTTP) reads — no second, disconnected DB in the container.
"--volume", f"{_host_db_dir()}:{_SUPERVISE_DB_DIR_IN_CONTAINER}",
"--env", f"SUPERVISE_DB_PATH={DB_PATH_IN_CONTAINER}",
]
for port in self._host_port_bindings:
argv += ["--publish", f"0.0.0.0:{port}:{port}"]
@@ -253,14 +268,11 @@ class DockerGateway(Gateway):
# policy against the control plane per request (guaranteed non-empty by
# the check above).
argv += ["--env", f"BOT_BOTTLE_ORCHESTRATOR_URL={self._orchestrator_url}"]
# ...presenting a role-scoped `gateway` token (a signed JWT minted from
# the host signing key) on those calls. The gateway never receives the
# signing key — only this pre-minted token, which it cannot rewrite into
# a `cli` token, so a compromised data-plane process can't drive the
# operator routes (issue #469 review). Bare `--env NAME` keeps the value
# off argv / `docker inspect`; only the gateway (not the agent) is given it.
argv += ["--env", CONTROL_AUTH_JWT_ENV]
run_env[CONTROL_AUTH_JWT_ENV] = mint(ROLE_GATEWAY, host_control_plane_token())
# ...and present the control-plane secret on those /resolve calls (the
# control plane requires it). Bare `--env NAME` keeps the value off argv
# / `docker inspect`; only the gateway (not the agent) is given it.
argv += ["--env", CONTROL_PLANE_TOKEN_ENV]
run_env[CONTROL_PLANE_TOKEN_ENV] = host_control_plane_token()
argv.append(self.image_ref)
proc = run_docker(argv, env=run_env)
if proc.returncode != 0:
+12 -17
View File
@@ -22,21 +22,21 @@ import urllib.request
from pathlib import Path
from .. import log
from ..control_auth import ROLE_GATEWAY, mint
from ..docker_cmd import run_docker
from ..paths import (
CONTROL_AUTH_JWT_ENV,
CONTROL_PLANE_TOKEN_ENV,
bot_bottle_root,
host_control_plane_token,
host_gateway_ca_dir,
)
from ..supervise import DB_PATH_IN_CONTAINER
from .gateway import (
GATEWAY_DOCKERFILE,
GATEWAY_IMAGE,
GATEWAY_NETWORK,
GatewayError,
MITMPROXY_HOME,
_host_db_dir,
)
DEFAULT_PORT = 8099
@@ -69,12 +69,12 @@ _INFRA_DAEMONS = "egress,git-http,supervise,orchestrator"
# container. Separate from /app so the gateway's baked scripts
# (egress_addon.py, egress-entrypoint.sh) are not overlaid.
_SRC_IN_CONTAINER = "/bot-bottle-src"
# Bot-bottle host-root bind-mount inside the container (DB + state). The
# orchestrator control plane opens bot-bottle.db under here (via BOT_BOTTLE_ROOT
# -> host_db_path()); it is the ONLY process in the container with a file
# handle on it (PRD 0070 / issue #469).
# Bot-bottle host-root bind-mount inside the container (DB + state).
_ROOT_IN_CONTAINER = "/bot-bottle-root"
# The supervise daemon writes proposals into the host DB directory.
_SUPERVISE_DB_DIR_IN_CONTAINER = os.path.dirname(DB_PATH_IN_CONTAINER)
_HEALTH_POLL_SECONDS = 0.25
_HEALTH_REQUEST_TIMEOUT_SECONDS = 1.0
@@ -188,7 +188,6 @@ class OrchestratorService:
so a later `ensure_running` can detect a real code change."""
self._ensure_network()
run_docker(["docker", "rm", "--force", self._infra_name])
_signing_key = host_control_plane_token()
proc = run_docker([
"docker", "run", "--detach",
"--name", self._infra_name,
@@ -203,6 +202,9 @@ class OrchestratorService:
# recreation AND docker volume pruning (issue #450): every agent
# trusts this one CA, so a fresh one would break all running bottles.
"--volume", f"{host_gateway_ca_dir()}:{MITMPROXY_HOME}",
# Shared supervise DB (same file the operator reads over HTTP).
"--volume", f"{_host_db_dir()}:{_SUPERVISE_DB_DIR_IN_CONTAINER}",
"--env", f"SUPERVISE_DB_PATH={DB_PATH_IN_CONTAINER}",
# Live control-plane source, mounted to a path that does not
# overlay the gateway's baked /app scripts.
"--volume", f"{self._repo_root}:{_SRC_IN_CONTAINER}:ro",
@@ -212,23 +214,16 @@ class OrchestratorService:
# Orchestrator registry DB on the host (sole writer: control plane).
"--volume", f"{self._host_root}:{_ROOT_IN_CONTAINER}",
"--env", f"BOT_BOTTLE_ROOT={_ROOT_IN_CONTAINER}",
# Control-plane signing key (orchestrator: verifies tokens) + the
# pre-minted `gateway` JWT (data-plane daemons: present it). gateway_init
# scopes each to its process, so a compromised data-plane daemon never
# sees the key and can't mint a `cli` token (issue #469 review).
# Control-plane secret: required by the orchestrator (to enforce)
# and by the gateway daemons (to present on /resolve calls).
"--env", CONTROL_PLANE_TOKEN_ENV,
"--env", CONTROL_AUTH_JWT_ENV,
# Gateway daemons reach the orchestrator over loopback at its
# fixed internal port (DEFAULT_PORT), independent of self.port.
"--env", f"BOT_BOTTLE_ORCHESTRATOR_URL=http://127.0.0.1:{DEFAULT_PORT}",
# Opt the orchestrator into gateway_init's supervise tree.
"--env", f"BOT_BOTTLE_GATEWAY_DAEMONS={_INFRA_DAEMONS}",
self.image,
], env={
**os.environ,
CONTROL_PLANE_TOKEN_ENV: _signing_key,
CONTROL_AUTH_JWT_ENV: mint(ROLE_GATEWAY, _signing_key),
})
], env={**os.environ, CONTROL_PLANE_TOKEN_ENV: host_control_plane_token()})
if proc.returncode != 0:
raise OrchestratorStartError(
f"infra container failed to start: {proc.stderr.strip()}"
-68
View File
@@ -31,23 +31,16 @@ from .gateway import Gateway
from ..supervise import (
AuditEntry,
COMPONENT_FOR_TOOL,
POLL_STATUS_PENDING,
POLL_STATUS_UNKNOWN,
Proposal,
Response,
STATUS_APPROVED,
STATUS_MODIFIED,
STATUS_REJECTED,
TOOL_EGRESS_ALLOW,
TOOL_EGRESS_BLOCK,
archive_all_proposals,
list_all_pending_proposals,
read_proposal,
read_response,
render_diff,
sha256_hex,
write_audit_entry,
write_proposal,
write_response,
)
@@ -136,8 +129,6 @@ class Orchestrator:
self._broker.submit(sign_request(req, self._secret))
self.registry.deregister(bottle_id)
self._tokens.pop(bottle_id, None)
# Reap any supervise proposals the gone bottle never acknowledged.
archive_all_proposals(bottle_id)
return True
def reconcile(
@@ -162,8 +153,6 @@ class Orchestrator:
live_source_ips, grace_seconds=grace_seconds)
for rec in reaped:
self._tokens.pop(rec.bottle_id, None)
# Reap any supervise proposals the gone bottle never acknowledged.
archive_all_proposals(rec.bottle_id)
return [rec.bottle_id for rec in reaped]
def tokens_for(self, bottle_id: str) -> dict[str, str]:
@@ -199,63 +188,6 @@ class Orchestrator:
# The orchestrator owns the single DB *and* the live policy, so operator
# decisions are applied here, server-side, and reached over HTTP by the
# host TUI (no direct-DB access, one path for every backend).
#
# The *agent* half of the queue (queue a proposal, poll for its response)
# is served here too (PRD 0070): the data plane no longer opens the DB, so
# supervise_server / egress / git-gate reach the queue over the control
# plane. Both are keyed on the caller's orchestrator-resolved bottle id —
# never a caller-supplied slug — so a bottle can only queue or read its own
# proposals even if the data plane is compromised.
def supervise_queue_proposal(
self, bottle_id: str, *, tool: str, proposed_file: str, justification: str,
) -> str:
"""Queue an agent-side proposal under `bottle_id` and return its id.
`bottle_id` is the caller resolved from `(source_ip, identity_token)` by
the control plane, so the proposal is attributed to the calling bottle,
not to anything the (possibly hostile) data plane asserts. The
proposed-file self-hash is computed here so the agent never supplies it."""
proposal = Proposal.new(
bottle_slug=bottle_id,
tool=tool,
proposed_file=proposed_file,
justification=justification,
current_file_hash=sha256_hex(proposed_file),
)
write_proposal(proposal)
return proposal.id
def supervise_poll_response(self, bottle_id: str, proposal_id: str) -> dict[str, object]:
"""Non-blocking, **idempotent** poll of one of `bottle_id`'s proposals.
Returns `{"status": ...}`: a terminal decision (`approved`/`modified`/
`rejected`, with `notes` + `final_file`) once the operator has responded,
`pending` while it's still queued undecided, or `unknown` when no such
queued proposal exists for this bottle (wrong id, or already archived).
Poll does **not** archive a decided proposal keeps returning the same
decision so a client whose connection drops after the decision but before
it consumes the response still gets it on retry (issue #469 review). A
decided proposal is already excluded from the operator's pending list (a
response row exists), so it doesn't linger there; the row itself is reaped
when the bottle is torn down / reconciled.
Reads are scoped to `bottle_id` (the queue key), so a caller can never
read another bottle's proposal even with a guessed id."""
try:
response = read_response(bottle_id, proposal_id)
except FileNotFoundError:
try:
read_proposal(bottle_id, proposal_id)
except FileNotFoundError:
return {"status": POLL_STATUS_UNKNOWN}
return {"status": POLL_STATUS_PENDING}
return {
"status": response.status,
"notes": response.notes,
"final_file": response.final_file,
}
def supervise_pending(self) -> list[dict[str, object]]:
"""All pending proposals across bottles, FIFO, as JSON dicts
-10
View File
@@ -37,16 +37,7 @@ HOST_DB_FILENAME = "bot-bottle.db"
# trusted callers (control plane, gateway, host CLI) and never handed to an
# agent, so an agent that can reach the control-plane port still can't drive it.
CONTROL_PLANE_TOKEN_FILENAME = "control-plane-token"
# The env var carrying the control-plane *signing key* — held only by the
# orchestrator (to verify tokens) and the host CLI (to mint its own), never by
# the data plane. Same value as the host token file; the name is unchanged for
# backward compatibility with existing launchers.
CONTROL_PLANE_TOKEN_ENV = "BOT_BOTTLE_CONTROL_PLANE_TOKEN"
# The env var carrying the data plane's pre-minted `gateway`-role token (a
# signed JWT the launcher mints from the signing key). The gateway presents this
# on /resolve + /supervise/{propose,poll}; it never holds the signing key, so it
# cannot forge a higher-privilege `cli` token (issue #469 review).
CONTROL_AUTH_JWT_ENV = "BOT_BOTTLE_CONTROL_AUTH_JWT"
# The host directory holding the gateway's persistent mitmproxy CA. Bind-mounted
# into the infra/gateway container at mitmproxy's confdir so the self-generated
@@ -133,7 +124,6 @@ __all__ = [
"HOST_DB_FILENAME",
"CONTROL_PLANE_TOKEN_FILENAME",
"CONTROL_PLANE_TOKEN_ENV",
"CONTROL_AUTH_JWT_ENV",
"GATEWAY_CA_DIRNAME",
"bot_bottle_root",
"host_db_path",
+25 -78
View File
@@ -34,22 +34,21 @@ import urllib.request
DEFAULT_TIMEOUT_SECONDS = 2.0
# The role-scoped `gateway` token this data plane presents on every control-plane
# call, read from the env the launcher injects into the gateway (a pre-minted
# signed JWT — the gateway never holds the signing key, so it can't forge a
# higher-privilege `cli` token). Constant + env-var name are duplicated here
# rather than imported because this module is COPYed flat into the gateway image,
# free of bot-bottle imports — same rationale as IDENTITY_HEADER in egress_addon
# / git_http_backend.
# The control-plane secret this gateway presents on every /resolve call, read
# from the env the launcher injects into the gateway container. The control
# plane requires it (orchestrator/control_plane.py). Constant + env-var name are
# duplicated here rather than imported because this module is COPYed flat into
# the gateway image, free of bot-bottle imports — same rationale as
# IDENTITY_HEADER in egress_addon / git_http_backend.
CONTROL_AUTH_HEADER = "x-bot-bottle-control-auth"
CONTROL_AUTH_JWT_ENV = "BOT_BOTTLE_CONTROL_AUTH_JWT"
CONTROL_PLANE_TOKEN_ENV = "BOT_BOTTLE_CONTROL_PLANE_TOKEN"
def _control_auth_headers() -> dict[str, str]:
"""The auth header to send, or {} when no token is configured (an open
"""The auth header to send, or {} when no secret is configured (an open
control plane, e.g. Firecracker behind its nft boundary sending nothing
is correct there and harmlessly ignored)."""
token = os.environ.get(CONTROL_AUTH_JWT_ENV, "").strip()
token = os.environ.get(CONTROL_PLANE_TOKEN_ENV, "").strip()
return {CONTROL_AUTH_HEADER: token} if token else {}
@@ -65,80 +64,28 @@ class PolicyResolver:
self._base = base_url.rstrip("/")
self._timeout = timeout
def _post_json(self, path: str, payload: dict[str, object]) -> dict[str, object] | None:
"""POST `payload` to control-plane `path`, returning the JSON object body
or None when the orchestrator answers `403` (unattributed / fail
closed). Raises `PolicyResolveError` on an unreachable / unexpected-status
/ malformed response so every caller can fail closed. Shared by
`_post_resolve` and the supervise propose/poll RPCs."""
body = json.dumps(payload).encode()
req = urllib.request.Request(
f"{self._base}{path}", data=body, method="POST",
headers={"Content-Type": "application/json", **_control_auth_headers()},
)
try:
with urllib.request.urlopen(req, timeout=self._timeout) as resp:
data = json.loads(resp.read())
except urllib.error.HTTPError as e:
if e.code == 403:
return None # unattributed → fail closed (caller denies)
raise PolicyResolveError(f"{path} returned HTTP {e.code}") from e
except (urllib.error.URLError, TimeoutError, OSError, ValueError) as e:
raise PolicyResolveError(f"{path} unreachable or malformed: {e}") from e
return data if isinstance(data, dict) else {}
def _post_resolve(self, source_ip: str, identity_token: str) -> dict[str, object] | None:
"""The orchestrator's `/resolve` payload for this client, or None if
unattributed (a clean `403`). Raises `PolicyResolveError` on an
unreachable / unexpected-status / malformed response so every caller
can fail closed. Shared by `resolve` and `resolve_bottle_id`."""
return self._post_json(
"/resolve", {"source_ip": source_ip, "identity_token": identity_token},
body = json.dumps(
{"source_ip": source_ip, "identity_token": identity_token}
).encode()
req = urllib.request.Request(
f"{self._base}/resolve", data=body, method="POST",
headers={"Content-Type": "application/json", **_control_auth_headers()},
)
def propose_supervise(
self,
source_ip: str,
identity_token: str,
*,
tool: str,
proposed_file: str,
justification: str,
) -> str | None:
"""Queue a supervise proposal on the control plane and return its
`proposal_id`. The orchestrator attributes the proposal to the calling
bottle by `(source_ip, identity_token)` exactly like `/resolve` so a
bottle can only ever queue its *own* proposals. Returns None when the
pair is unattributed (a clean `403`); raises `PolicyResolveError` if the
orchestrator can't be reached, so the data-plane caller fails closed
(blocks / refuses) rather than silently dropping the proposal."""
payload = self._post_json("/supervise/propose", {
"source_ip": source_ip,
"identity_token": identity_token,
"tool": tool,
"proposed_file": proposed_file,
"justification": justification,
})
if payload is None:
return None
proposal_id = payload.get("proposal_id")
return proposal_id if isinstance(proposal_id, str) and proposal_id else None
def poll_supervise(
self, source_ip: str, identity_token: str, proposal_id: str,
) -> dict[str, object] | None:
"""Poll a queued proposal for the operator's decision, **non-blocking**.
Attributed by `(source_ip, identity_token)` so a bottle can only read its
*own* proposal's response. Returns `{"status": ...}` where status is one
of the terminal decisions (`approved`/`modified`/`rejected`, carrying
`notes` + `final_file`), `pending` (queued, no decision yet), or
`unknown` (no such queued proposal for this bottle). Returns None when
unattributed; raises `PolicyResolveError` if unreachable."""
return self._post_json("/supervise/poll", {
"source_ip": source_ip,
"identity_token": identity_token,
"proposal_id": proposal_id,
})
try:
with urllib.request.urlopen(req, timeout=self._timeout) as resp:
payload = json.loads(resp.read())
except urllib.error.HTTPError as e:
if e.code == 403:
return None # unattributed → fail closed (caller denies)
raise PolicyResolveError(f"/resolve returned HTTP {e.code}") from e
except (urllib.error.URLError, TimeoutError, OSError, ValueError) as e:
raise PolicyResolveError(f"/resolve unreachable or malformed: {e}") from e
return payload if isinstance(payload, dict) else {}
def resolve(self, source_ip: str, identity_token: str = "") -> str | None:
"""The calling bottle's policy blob, or None if unattributed. Always
-17
View File
@@ -192,23 +192,6 @@ class QueueStore(DbStore):
(self.queue_key, proposal_id),
)
def archive_all(self) -> None:
"""Archive every proposal + response for this queue key. Used to reap a
bottle's supervise records when it is torn down or reconciled away —
the retention path for a client that received a decision but died before
acknowledging it. Idempotent."""
if not self.db_path.is_file():
return
with self._connection() as conn:
conn.execute(
"UPDATE supervise_proposals SET archived = 1 WHERE queue_key = ?",
(self.queue_key,),
)
conn.execute(
"UPDATE supervise_responses SET archived = 1 WHERE queue_key = ?",
(self.queue_key,),
)
@staticmethod
def _row_to_proposal(row: sqlite3.Row) -> Proposal:
return Proposal(
-11
View File
@@ -40,8 +40,6 @@ from pathlib import Path
from .supervise_types import (
ACTION_OPERATOR_EDIT,
AuditEntry,
POLL_STATUS_PENDING,
POLL_STATUS_UNKNOWN,
Proposal,
Response,
STATUSES,
@@ -171,12 +169,6 @@ def archive_proposal(bottle_slug: str, proposal_id: str) -> None:
QueueStore(bottle_slug).archive_proposal(proposal_id)
def archive_all_proposals(bottle_slug: str) -> None:
"""Archive every proposal + response for `bottle_slug` (bottle teardown /
reconcile cleanup). Idempotent."""
QueueStore(bottle_slug).archive_all()
# --- Audit log -------------------------------------------------------------
@@ -265,8 +257,6 @@ __all__ = [
"STATUS_APPROVED",
"STATUS_MODIFIED",
"STATUS_REJECTED",
"POLL_STATUS_PENDING",
"POLL_STATUS_UNKNOWN",
"SUPERVISE_HOSTNAME",
"SUPERVISE_PORT",
"Supervise",
@@ -281,7 +271,6 @@ __all__ = [
"TOOL_EGRESS_TOKEN_ALLOW",
"TOOL_LIST_EGRESS_ROUTES",
"archive_proposal",
"archive_all_proposals",
"audit_dir",
"audit_log_path",
"list_pending_proposals",
+107 -158
View File
@@ -7,13 +7,9 @@ config changes when stuck. The tools are `egress-allow`,
Each queued proposal tool call:
1. Validates the proposed file syntactically.
2. Queues a Proposal over the control-plane RPC (`POST /supervise/propose`),
which attributes it to the calling bottle by (source_ip, identity_token)
and writes it to the one orchestrator-owned SQLite database. This daemon
never opens `bot-bottle.db` itself (PRD 0070 / issue #469).
3. Polls the control plane (`POST /supervise/poll`) for a matching Response,
up to a short grace window (`SUPERVISE_RESPONSE_TIMEOUT_SECONDS`,
default 30s).
2. Writes a Proposal to the host SQLite database.
3. Blocks polling for a matching Response row, up to a short grace
window (`SUPERVISE_RESPONSE_TIMEOUT_SECONDS`, default 30s).
4. On a decision within the window, returns the operator's
`{status, notes}`. On timeout, returns `status: pending` **with the
proposal id** and leaves the proposal queued the flow is
@@ -28,8 +24,8 @@ without holding an HTTP request open.
One shared server fronts every bottle (PRD 0070) and attributes each
proposal to the calling bottle by source IP, resolved from the orchestrator
an unattributed or unreachable source fails closed. BOT_BOTTLE_ORCHESTRATOR_URL
is mandatory: there is no fixed-slug single-tenant fallback, and the queue is
reached only over that control-plane RPC (no direct database handle).
is mandatory: there is no fixed-slug single-tenant fallback. SUPERVISE_DB_PATH
points at the bind-mounted host database.
Speaks MCP over HTTP+JSON-RPC. Methods handled:
@@ -55,7 +51,7 @@ import socketserver
import sys
import time
import typing
from dataclasses import dataclass
from dataclasses import dataclass, replace
from bot_bottle.constants import IDENTITY_HEADER
from bot_bottle.egress_addon_core import (
@@ -82,11 +78,7 @@ ERR_INVALID_PARAMS = -32602
ERR_INTERNAL = -32603
DEFAULT_RESPONSE_TIMEOUT_SECONDS = 30.0
# Grace-window poll cadence. 250 ms keeps the operator-visible latency low while
# staying gentle on the shared control plane — each poll is now an HTTP round
# trip, not a local file read, so the old 50 ms would be ~20 req/s per pending
# proposal across every waiting agent (issue #469 review).
MIN_RESPONSE_POLL_INTERVAL_SECONDS = 0.25
MIN_RESPONSE_POLL_INTERVAL_SECONDS = 0.05
# The per-host orchestrator control plane the shared supervise server attributes
# each proposal to, by source IP. Mandatory — there is no single-tenant
@@ -320,9 +312,7 @@ def validate_proposed_file(tool: str, content: str) -> None:
@dataclass(frozen=True)
class ServerConfig:
# No bottle_slug: the calling bottle is attributed per request by the
# control plane from (source_ip, identity_token), so nothing on this daemon
# is keyed by a single slug (PRD 0070 / issue #469).
bottle_slug: str
response_timeout_seconds: float = DEFAULT_RESPONSE_TIMEOUT_SECONDS
@@ -341,22 +331,14 @@ def handle_tools_list(_params: dict[str, object]) -> dict[str, object]:
def handle_tools_call(
params: dict[str, object],
config: ServerConfig,
*,
resolver: "PolicyResolver",
source_ip: str,
identity_token: str,
) -> dict[str, object]:
"""Validates the proposal, queues it over the control-plane RPC, polls for
a Response through the grace window, and returns the result wrapped in MCP
`content`.
"""Validates the proposal, writes it to the queue, blocks waiting
for a Response, returns the result wrapped in MCP `content`.
The queue lives behind the orchestrator now `propose_supervise` attributes
the proposal to the calling bottle by `(source_ip, identity_token)` and
`poll_supervise` reads only that bottle's own responses, so this daemon
never opens `bot-bottle.db`. `list-egress-routes` never reaches here the
handler answers it from the calling bottle's resolved policy before
dispatching (see `MCPHandler._dispatch`); this path is the queued,
operator-approved `egress-allow` / `egress-block` tools."""
`list-egress-routes` never reaches here the handler answers it from
the calling bottle's resolved policy before dispatching (see
`MCPHandler._dispatch`); this path is the queued, operator-approved
`egress-allow` / `egress-block` tools."""
name = params.get("name")
if not isinstance(name, str):
raise _RpcClientError(ERR_INVALID_PARAMS, "tools/call missing 'name'")
@@ -384,100 +366,59 @@ def handle_tools_call(
else:
raise _RpcClientError(ERR_INVALID_PARAMS, f"unknown tool {name!r}")
proposal_id = _queue_proposal(
resolver, source_ip, identity_token,
tool=name, proposed_file=proposed_file, justification=justification,
proposal = _sv.Proposal.new(
bottle_slug=config.bottle_slug,
tool=name,
proposed_file=proposed_file,
justification=justification,
current_file_hash=_sv.sha256_hex(proposed_file),
)
try:
_sv.write_proposal(proposal)
except OSError as e:
raise _RpcInternalError(f"failed to write proposal to queue: {e}") from e
sys.stderr.write(
f"supervise: queued proposal {proposal_id} ({name}); waiting for operator...\n"
f"supervise: queued proposal {proposal.id} ({name}) "
f"for bottle {config.bottle_slug}; waiting for operator...\n"
)
sys.stderr.flush()
deadline = time.monotonic() + config.response_timeout_seconds
while time.monotonic() < deadline:
response = _poll_response(resolver, source_ip, identity_token, proposal_id)
if response is not None:
return {
"content": [{"type": "text", "text": format_response_text(response)}],
"isError": response.status == _sv.STATUS_REJECTED,
}
time.sleep(MIN_RESPONSE_POLL_INTERVAL_SECONDS)
try:
response = _sv.wait_for_response(
config.bottle_slug,
proposal.id,
poll_interval=MIN_RESPONSE_POLL_INTERVAL_SECONDS,
deadline=deadline,
)
except TimeoutError:
text = format_pending_response_text(proposal.id, config.response_timeout_seconds)
return {
"content": [{"type": "text", "text": text}],
"isError": False,
}
try:
_sv.archive_proposal(config.bottle_slug, proposal.id)
except OSError as e:
raise _RpcInternalError(f"failed to archive proposal: {e}") from e
text = format_pending_response_text(proposal_id, config.response_timeout_seconds)
text = format_response_text(response)
return {
"content": [{"type": "text", "text": text}],
"isError": False,
"isError": response.status == _sv.STATUS_REJECTED,
}
def _queue_proposal(
resolver: "PolicyResolver",
source_ip: str,
identity_token: str,
*,
tool: str,
proposed_file: str,
justification: str,
) -> str:
"""Queue a proposal over the control plane and return its id, mapping the
resolver's fail-closed outcomes to internal errors so `do_POST` returns a
clean -32603 rather than leaking the cause to the agent."""
try:
proposal_id = resolver.propose_supervise(
source_ip, identity_token,
tool=tool, proposed_file=proposed_file, justification=justification,
)
except PolicyResolveError as e:
raise _RpcInternalError(f"orchestrator unreachable, cannot queue proposal: {e}") from e
if proposal_id is None:
raise _RpcInternalError("request source is not attributed to a bottle")
return proposal_id
def _poll_response(
resolver: "PolicyResolver", source_ip: str, identity_token: str, proposal_id: str,
) -> "_sv.Response | None":
"""One non-blocking control-plane poll for `proposal_id`'s decision. Returns
the operator's Response once decided, or None while it is still pending
(or `unknown` treated as still-waiting so a transient blip degrades to the
pending-timeout path rather than a spurious error mid grace window)."""
try:
result = resolver.poll_supervise(source_ip, identity_token, proposal_id)
except PolicyResolveError as e:
raise _RpcInternalError(f"orchestrator unreachable, cannot poll proposal: {e}") from e
if result is None:
raise _RpcInternalError("request source is not attributed to a bottle")
if result.get("status") in _sv.STATUSES:
return _response_from_poll(proposal_id, result)
return None
def _response_from_poll(proposal_id: str, result: "dict[str, object]") -> "_sv.Response":
"""Build a Response from a decided `/supervise/poll` payload."""
final_file = result.get("final_file")
return _sv.Response(
proposal_id=proposal_id,
status=str(result.get("status")),
notes=str(result.get("notes") or ""),
final_file=final_file if isinstance(final_file, str) else None,
)
def handle_check_proposal(
params: dict[str, object],
*,
resolver: "PolicyResolver",
source_ip: str,
identity_token: str,
config: ServerConfig,
) -> dict[str, object]:
"""Non-blocking poll of a queued proposal's decision, by id.
Never creates a Proposal (so `check-proposal` isn't in `TOOLS`); it only
reads the queue over the control-plane RPC, attributed by
`(source_ip, identity_token)` so it can only read *this* bottle's own
proposals. The orchestrator archives a decided proposal on read the same
terminal step the synchronous path triggers so `pending` proposals stay
visible to the operator until they're both decided *and* polled."""
reads the queue. Resolution order mirrors the synchronous path's terminal
step a decided proposal is archived here exactly as `handle_tools_call`
archives it after `wait_for_response`, so `pending` proposals stay visible
to the operator until they're both decided *and* polled."""
args_raw = params.get("arguments", {})
if not isinstance(args_raw, dict):
raise _RpcClientError(ERR_INVALID_PARAMS, "tools/call 'arguments' must be an object")
@@ -490,24 +431,25 @@ def handle_check_proposal(
proposal_id = proposal_id.strip()
try:
result = resolver.poll_supervise(source_ip, identity_token, proposal_id)
except PolicyResolveError as e:
raise _RpcInternalError(f"orchestrator unreachable, cannot poll proposal: {e}") from e
if result is None:
raise _RpcInternalError("request source is not attributed to a bottle")
status = result.get("status")
if status == _sv.POLL_STATUS_UNKNOWN:
return {
"content": [{"type": "text", "text": format_unknown_proposal_text(proposal_id)}],
"isError": True,
}
if status == _sv.POLL_STATUS_PENDING:
response = _sv.read_response(config.bottle_slug, proposal_id)
except FileNotFoundError:
# No decision yet — distinguish "still queued" from "unknown id".
try:
_sv.read_proposal(config.bottle_slug, proposal_id)
except FileNotFoundError:
return {
"content": [{"type": "text", "text": format_unknown_proposal_text(proposal_id)}],
"isError": True,
}
return {
"content": [{"type": "text", "text": format_still_pending_text(proposal_id)}],
"isError": False,
}
response = _response_from_poll(proposal_id, result)
try:
_sv.archive_proposal(config.bottle_slug, proposal_id)
except OSError as e:
raise _RpcInternalError(f"failed to archive proposal: {e}") from e
return {
"content": [{"type": "text", "text": format_response_text(response)}],
"isError": response.status == _sv.STATUS_REJECTED,
@@ -649,34 +591,16 @@ class MCPHandler(http.server.BaseHTTPRequestHandler):
# — silently dropping base routes like api.anthropic.com on approval.
if req.params.get("name") == _sv.TOOL_LIST_EGRESS_ROUTES:
return self._resolved_routes_payload()
resolver = self._resolver_or_fail()
source_ip = self.client_address[0]
token = self._identity_token()
# `check-proposal` is a non-blocking read of the calling bottle's
# own queue — attributed by (source_ip, identity_token) like a
# proposal, but it never queues or blocks.
# own queue — attributed by source IP like a proposal, but it
# never queues or blocks.
if req.params.get("name") == _sv.TOOL_CHECK_PROPOSAL:
return handle_check_proposal(
req.params, resolver=resolver,
source_ip=source_ip, identity_token=token,
)
# The control plane attributes the proposal to the source-IP + token
# resolved bottle, so the one shared queue holds each bottle's
# proposal under its own id — no slug is asserted by this daemon.
return handle_tools_call(
req.params, config, resolver=resolver,
source_ip=source_ip, identity_token=token,
)
return handle_check_proposal(req.params, self._attributed_config(config))
# Attribute the proposal to the source-IP-resolved bottle, so the one
# shared server queues each bottle's proposal under its own slug.
return handle_tools_call(req.params, self._attributed_config(config))
raise _RpcClientError(ERR_METHOD_NOT_FOUND, f"method not found: {method}")
def _identity_token(self) -> str:
"""The agent's per-bottle identity token from the request header (the
MCP client sends it via `mcp add --header`), or empty. The control plane
requires the (source_ip, token) pair, so a missing/wrong token
fail-closes there."""
headers = getattr(self, "headers", None)
return headers.get(IDENTITY_HEADER, "") if headers is not None else ""
def _resolver_or_fail(self) -> "PolicyResolver":
"""This server's policy resolver. A server started without one is a
misconfiguration, not a tenancy mode fail closed rather than
@@ -688,18 +612,40 @@ class MCPHandler(http.server.BaseHTTPRequestHandler):
def _resolved_routes_payload(self) -> dict[str, object]:
"""The calling bottle's live egress routes as the `list-egress-routes`
JSON payload, resolved by (source_ip, identity token). Fail-closed: an
unattributed source or an unreachable orchestrator yields an empty route
list (never another bottle's), courtesy of `resolve_client_context`."""
JSON payload, resolved by (source_ip, identity token). Fail-closed like
`_attributed_config`: an unattributed source or an unreachable
orchestrator yields an empty route list (never another bottle's),
courtesy of `resolve_client_context`."""
resolver = self._resolver_or_fail()
headers = getattr(self, "headers", None)
token = headers.get(IDENTITY_HEADER, "") if headers is not None else ""
conf, _slug, _tokens = resolve_client_context(
resolver, self.client_address[0], self._identity_token(),
resolver, self.client_address[0], token,
)
body = json.dumps(
{"routes": [route_to_yaml_dict(r) for r in conf.routes]}, indent=2,
)
return {"content": [{"type": "text", "text": body}], "isError": False}
def _attributed_config(self, config: ServerConfig) -> ServerConfig:
"""The ServerConfig with `bottle_slug` bound to *this request's* bottle:
the bottle id attributed from the source IP **fail-closed**, an
unattributed or unreachable source raises so no proposal is queued under
the wrong (or empty) slug."""
resolver = self._resolver_or_fail()
# The agent's MCP client sends the identity token as a request header
# (provisioned via `mcp add --header`); the orchestrator requires the
# (source_ip, token) pair, so a missing/wrong token fail-closes below.
headers = getattr(self, "headers", None)
token = headers.get(IDENTITY_HEADER, "") if headers is not None else ""
try:
bottle_id = resolver.resolve_bottle_id(self.client_address[0], token)
except PolicyResolveError as e:
raise _RpcInternalError(f"orchestrator unreachable, cannot attribute: {e}") from e
if not bottle_id:
raise _RpcInternalError("request source is not attributed to a bottle")
return replace(config, bottle_slug=bottle_id)
def _write_jsonrpc(self, body: bytes) -> None:
self.send_response(200)
self.send_header("Content-Type", "application/json")
@@ -722,10 +668,10 @@ class MCPHandler(http.server.BaseHTTPRequestHandler):
class MCPServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
allow_reuse_address = True
daemon_threads = True
config: ServerConfig = ServerConfig()
# Set by `serve`; every proposal is attributed to the source-IP + token
# resolved bottle by the control plane. A server without a resolver fails
# closed per request (see `_resolver_or_fail`).
config: ServerConfig = ServerConfig(bottle_slug="")
# Set by `serve`; every proposal is attributed to the source-IP-resolved
# bottle. The class default is a placeholder — a server without a resolver
# fails closed per request (see `_resolver_or_fail`).
policy_resolver: "PolicyResolver | None" = None
@@ -740,9 +686,12 @@ def serve(
response_timeout_seconds: float = DEFAULT_RESPONSE_TIMEOUT_SECONDS,
) -> typing.NoReturn:
server = MCPServer((bind, port), MCPHandler)
# Every request's proposal is attributed to the source-IP + token resolved
# bottle by the control plane (see MCPHandler._dispatch).
server.config = ServerConfig(response_timeout_seconds=response_timeout_seconds)
# bottle_slug is a placeholder: every request's proposal is attributed to
# the source-IP-resolved bottle (see MCPHandler._attributed_config).
server.config = ServerConfig(
bottle_slug="",
response_timeout_seconds=response_timeout_seconds,
)
server.policy_resolver = resolver
sys.stderr.write(
f"supervise listening on {bind}:{port}; multi-tenant; "
-11
View File
@@ -37,15 +37,6 @@ STATUS_MODIFIED = "modified"
STATUS_REJECTED = "rejected"
STATUSES: tuple[str, ...] = (STATUS_APPROVED, STATUS_MODIFIED, STATUS_REJECTED)
# Non-terminal markers the control-plane supervise-poll RPC returns when a
# proposal has no operator decision yet (`PENDING`) or no queued proposal with
# that id exists for the calling bottle (`UNKNOWN`). They are never a
# `Response.status` — only the poll wire contract (see
# `Orchestrator.supervise_poll_response`) — so they are deliberately outside
# `STATUSES`.
POLL_STATUS_PENDING = "pending"
POLL_STATUS_UNKNOWN = "unknown"
ACTION_OPERATOR_EDIT = "operator-edit"
@@ -166,8 +157,6 @@ __all__ = [
"STATUS_APPROVED",
"STATUS_MODIFIED",
"STATUS_REJECTED",
"POLL_STATUS_PENDING",
"POLL_STATUS_UNKNOWN",
"TOOLS",
"TOOL_EGRESS_ALLOW",
"TOOL_EGRESS_BLOCK",
+13 -3
View File
@@ -1,13 +1,23 @@
# CI
The test workflow lives at [`.gitea/workflows/test.yml`](../.gitea/workflows/test.yml).
It runs `tests/run_tests.py` (full suite — unit + integration) on:
It runs the unit suite plus one integration job per backend
(`integration-docker`, `integration-firecracker`) on:
- every push to a branch with an open pull request, and
- every push to `main`.
Integration tests need Docker on the runner; they skip cleanly via
`tests/_docker.skip_unless_docker` when no daemon is reachable.
Each integration job selects its backend via `BOT_BOTTLE_BACKEND` and
runs a **preflight** (`./cli.py backend status --backend=<name>`) that
prints a clear per-check readiness summary and fails the job when the
backend is missing — so absent infrastructure is visible at the job level
rather than hidden among per-test `unittest.skip` lines. The skip guards in
[`tests/_backend.py`](../tests/_backend.py) gate on the same readiness
check (`bot_bottle.backend.has_backend`): backend-agnostic tests use
`skip_unless_selected_backend_available()` and run through whichever
backend is selected (checking, e.g., Linux + `/dev/kvm` for Firecracker
rather than unrelated Docker availability); Docker-implementation tests use
`skip_unless_backend("docker")` and no-op under a non-Docker run.
A small subset of integration tests skip when running specifically
under Gitea Actions (`GITEA_ACTIONS=true`), because `act_runner` runs
+5 -8
View File
@@ -299,14 +299,11 @@ and integrate a console against.
the data plane and the console reach state through the control-plane RPC,
never a direct file handle.** No agent-facing component gets the file, so
none can forge attribution. (This supersedes the earlier `ro`-mount idea.)
This rule is now **in force** across the data plane (issue #469): the
supervise daemon, the egress addon, and the git-gate pre-receive hook queue
proposals and poll for their responses over the agent-side supervise RPC
(`POST /supervise/propose` + `POST /supervise/poll`, attributed by
`(source_ip, identity_token)` exactly like `/resolve`), so a bottle can only
ever queue or read its own proposals. The gateway/data-plane containers no
longer bind-mount the DB directory or carry `SUPERVISE_DB_PATH`; the
orchestrator is the sole opener of the one file.
- *Transitional caveat:* today the per-bottle **supervise gateway
rw-bind-mounts `bot-bottle.db`** to write proposals — exactly the
pattern the orchestrator removes (supervise consolidates into the
orchestrator; gateway writes become RPC calls). Until that lands, don't
put the attribution registry behind a data-plane-writable mount.
Implementation note for the VM slices: SQLite **WAL** over a guest share
(virtiofs/9p) is finicky (the `-shm`/`-wal` files need real mmap/locking),
+20 -5
View File
@@ -2,14 +2,15 @@
Plain-Python test suite using stdlib `unittest`. No external
dependencies. Unit tests run anywhere Python 3 is present; integration
tests need Docker and skip cleanly otherwise.
tests run through the backend named by `BOT_BOTTLE_BACKEND` (default
`docker`) and skip cleanly when that backend isn't available on the host.
## Layout
```
tests/
fixtures.py # JSON manifest builders (shared)
_docker.py # docker-availability skip helper (shared)
_backend.py # backend selection + skip guards (shared)
unit/
test_egress.py
test_egress_addon_core.py
@@ -73,7 +74,7 @@ BOT_BOTTLE_RUN_CANARIES=1 python -m unittest discover -t . -s tests/canaries -v
## Adding a test
1. Pick the directory: `tests/unit/` for a pure unit test,
`tests/integration/` for one that needs Docker.
`tests/integration/` for one that needs a backend.
2. Filename: `test_<topic>.py`.
3. Boilerplate:
```python
@@ -88,5 +89,19 @@ BOT_BOTTLE_RUN_CANARIES=1 python -m unittest discover -t . -s tests/canaries -v
if __name__ == "__main__":
unittest.main()
```
4. For Docker-dependent tests, decorate the class with
`@skip_unless_docker()` from `tests._docker`.
4. Skip guards live in `tests._backend` and gate on the backend's own
readiness check, `bot_bottle.backend.has_backend` — the same probe
behind `./cli.py backend status`:
- Backend-agnostic tests (go through `get_bottle_backend()`) decorate
the class with `@skip_unless_selected_backend_available()` — the test
runs against whichever backend `BOT_BOTTLE_BACKEND` selects and skips
unless that backend is available (checking, e.g., Linux + `/dev/kvm`
for Firecracker rather than unrelated Docker availability).
- Backend-specific tests (exercise `DockerBroker`, `DockerGateway`,
`backend.docker.*`, …) decorate with `@skip_unless_backend("docker")`
so they no-op under a run targeting a different backend.
Each CI integration job runs `./cli.py backend status --backend=<name>`
as a preflight, which prints a clear per-check summary and exits non-zero
when the backend is missing — so absent infrastructure fails the job
instead of hiding among per-test `unittest.skip` lines.
+71
View File
@@ -0,0 +1,71 @@
"""Backend selection + readiness-aware skip guards for the integration suite.
Each integration test targets the backend named by ``BOT_BOTTLE_BACKEND``
(default ``docker``) and gates on that backend's full readiness check —
``is_backend_ready()`` (equivalent to ``./cli.py backend status``), not just
a binary-on-PATH probe. When the backend is not ready, diagnostic output is
printed during test discovery so the operator sees a concrete reason for each
skip.
"""
from __future__ import annotations
import os
import unittest
from bot_bottle.backend import is_backend_ready
# Default when ``BOT_BOTTLE_BACKEND`` is unset. Docker preserves the historical
# Docker-backed CI path (and mirrors the pin in ``test_sandbox_escape``).
DEFAULT_BACKEND = "docker"
def selected_backend() -> str:
"""The backend this test run targets, from ``BOT_BOTTLE_BACKEND``.
Mirrors the CLI's env selector; unset means ``docker`` so an
unconfigured run behaves exactly as the suite did before backends were
pluggable.
"""
return os.environ.get("BOT_BOTTLE_BACKEND") or DEFAULT_BACKEND
def skip_unless_backend(backend: str):
"""Skip a backend-specific test unless the selected backend matches AND
that backend is fully ready.
Docker-implementation tests (``DockerBroker``, ``DockerGateway``,
``backend.docker.*``) use ``skip_unless_backend("docker")`` so they no-op
under a run targeting a different backend instead of testing Docker
internals that run doesn't exercise — the guard reads
``BOT_BOTTLE_BACKEND`` rather than "is Docker installed".
When the backend is not ready, ``status()`` output is printed so the
operator sees a concrete diagnostic for each skipped test module.
"""
sel = selected_backend()
if sel != backend:
return unittest.skip(
f"backend {backend!r} not selected (BOT_BOTTLE_BACKEND={sel})"
)
return unittest.skipUnless(
is_backend_ready(backend, quiet=False),
f"{backend} backend not ready",
)
def skip_unless_selected_backend_available():
"""Skip a backend-agnostic test unless the *selected* backend is fully ready.
The test then runs through whichever backend ``BOT_BOTTLE_BACKEND`` names,
gated on that backend's full status() check (e.g. daemon reachable, TAP
pool present for Firecracker) rather than just a binary-on-PATH probe.
When the backend is not ready, ``status()`` output is printed so the
operator sees a concrete diagnostic for each skipped test module.
"""
backend = selected_backend()
return unittest.skipUnless(
is_backend_ready(backend, quiet=False),
f"selected backend {backend!r} not ready",
)
-45
View File
@@ -1,45 +0,0 @@
"""Docker availability check used by integration tests."""
from __future__ import annotations
import os
import shutil
import subprocess
import unittest
def docker_available() -> bool:
if os.environ.get("SKIP_DOCKER_TESTS"):
return False
if shutil.which("docker") is None:
return False
try:
return (
subprocess.run(
["docker", "info"],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
check=False,
timeout=5,
).returncode
== 0
)
except subprocess.TimeoutExpired:
return False
def skip_unless_docker(reason: str = "docker unreachable"):
return unittest.skipUnless(docker_available(), reason)
def skip_unless_docker_or_firecracker(
reason: str = "neither Docker nor Firecracker selected",
):
"""Skip a backend-agnostic test unless one supported backend can run.
Firecracker does not require the host Docker daemon. The KVM coverage job
deliberately sets ``SKIP_DOCKER_TESTS`` to exclude Docker-only integration
classes while still exercising this path.
"""
firecracker_selected = os.environ.get("BOT_BOTTLE_BACKEND") == "firecracker"
return unittest.skipUnless(firecracker_selected or docker_available(), reason)
+2 -2
View File
@@ -25,14 +25,14 @@ import os
import subprocess
import unittest
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
_IMAGE = "bot-bottle-gateway-test:chunk1"
_DOCKERFILE = "Dockerfile.gateway"
@skip_unless_docker()
@skip_unless_backend("docker")
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: multi-stage build pulls a 200+MB "
@@ -31,7 +31,7 @@ from bot_bottle.backend.docker.gateway_net import next_free_ip
from bot_bottle.orchestrator.client import OrchestratorClient
from bot_bottle.orchestrator.gateway import GATEWAY_IMAGE, GATEWAY_NAME, GATEWAY_NETWORK
from bot_bottle.orchestrator.lifecycle import OrchestratorService
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
# One upstream reachable under two names; reflects the Authorization header so
# the probe can see exactly what the gateway injected (or didn't).
@@ -72,7 +72,7 @@ _PROBE_SRC = (
)
@skip_unless_docker()
@skip_unless_backend("docker")
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: the orchestrator container bind-mounts the repo "
@@ -13,12 +13,12 @@ import unittest
from bot_bottle.orchestrator.broker import LaunchRequest, sign_request
from bot_bottle.orchestrator.docker_broker import DockerBroker, container_name
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
IMAGE = "busybox"
@skip_unless_docker()
@skip_unless_backend("docker")
class TestDockerBrokerIntegration(unittest.TestCase):
def setUp(self) -> None:
self.secret = secrets.token_bytes(16)
@@ -25,11 +25,10 @@ import tempfile
import unittest
from pathlib import Path
from bot_bottle.control_auth import ROLE_CLI, ROLE_GATEWAY, mint
from bot_bottle.orchestrator.client import OrchestratorClient
from bot_bottle.orchestrator.lifecycle import OrchestratorService
from bot_bottle.paths import host_control_plane_token
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
# Fixed (not per-run-suffixed) so repeated runs reuse the same layer-cached
# image instead of leaking a new dangling tag on every invocation.
@@ -38,7 +37,7 @@ _TEST_GATEWAY_IMAGE = "bot-bottle-gateway:itest"
_TEST_INFRA_IMAGE = "bot-bottle-infra:itest"
@skip_unless_docker()
@skip_unless_backend("docker")
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: the orchestrator container bind-mounts the repo "
@@ -86,11 +85,7 @@ class TestDockerControlPlaneAuthIntegration(unittest.TestCase):
host_root=host_root,
)
cls.svc.ensure_running()
# The control plane now verifies role-scoped signed tokens, not the raw
# key. Mint one of each role from the host signing key (issue #469 review).
signing_key = host_control_plane_token()
cls.cli_token = mint(ROLE_CLI, signing_key)
cls.gateway_token = mint(ROLE_GATEWAY, signing_key)
cls.token = host_control_plane_token()
@staticmethod
def _teardown_docker(
@@ -141,17 +136,11 @@ class TestDockerControlPlaneAuthIntegration(unittest.TestCase):
status, _ = self._request("GET", "/bottles", token="not-the-real-secret")
self.assertEqual(401, status)
def test_bottles_accepts_the_cli_token(self) -> None:
status, payload = self._request("GET", "/bottles", token=self.cli_token)
def test_bottles_accepts_the_real_token(self) -> None:
status, payload = self._request("GET", "/bottles", token=self.token)
self.assertEqual(200, status)
self.assertEqual([], payload["bottles"])
def test_bottles_rejects_the_gateway_token(self) -> None:
"""A compromised gateway holds only a `gateway` token — it must not be
able to enumerate bottles (an operator route) with it (issue #469)."""
status, _ = self._request("GET", "/bottles", token=self.gateway_token)
self.assertEqual(403, status)
def test_resolve_rejects_a_caller_with_no_token(self) -> None:
"""The credential-lift attack from issue #400: an agent could POST
/resolve directly and read back the upstream tokens it's never meant
@@ -159,13 +148,6 @@ class TestDockerControlPlaneAuthIntegration(unittest.TestCase):
status, _ = self._request("POST", "/resolve")
self.assertEqual(401, status)
def test_resolve_accepts_the_gateway_token(self) -> None:
"""The gateway's own token DOES reach /resolve: with an empty body it
gets 400 (missing source_ip) rather than the 401 a rejected caller sees
proving the role gate let the gateway token through."""
status, _ = self._request("POST", "/resolve", token=self.gateway_token)
self.assertEqual(400, status) # past the role gate; body validation fails
if __name__ == "__main__":
unittest.main()
@@ -11,12 +11,12 @@ import subprocess
import unittest
from bot_bottle.orchestrator.gateway import DockerGateway
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
IMAGE = "busybox"
@skip_unless_docker()
@skip_unless_backend("docker")
class TestDockerGatewayIntegration(unittest.TestCase):
def setUp(self) -> None:
self.name = "bot-bottle-orch-gateway-itest-" + secrets.token_hex(4)
@@ -12,12 +12,12 @@ import subprocess
import unittest
from bot_bottle.orchestrator.gateway import DockerGateway
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
IMAGE = "busybox"
@skip_unless_docker()
@skip_unless_backend("docker")
class TestDockerGatewayImageExists(unittest.TestCase):
def test_image_exists_true_for_present_false_for_absent(self) -> None:
# Ensure the tiny image is present (build_if_missing is disabled here
+2 -2
View File
@@ -18,10 +18,10 @@ from bot_bottle.backend.docker.network import (
network_create_internal,
network_remove,
)
from tests._docker import skip_unless_docker
from tests._backend import skip_unless_backend
@skip_unless_docker()
@skip_unless_backend("docker")
class TestOrphanCleanup(unittest.TestCase):
def setUp(self):
self.slug = f"cb-test-orphan-{os.getpid()}"
+2 -2
View File
@@ -31,7 +31,7 @@ from pathlib import Path
from bot_bottle.backend import BottleSpec, get_bottle_backend
from bot_bottle.bottle_state import cleanup_state
from bot_bottle.manifest import ManifestIndex
from tests._docker import skip_unless_docker_or_firecracker
from tests._backend import skip_unless_selected_backend_available
# Secrets planted in the bottle env as literals (agents substitute via
@@ -67,7 +67,7 @@ _DUMMY_HOST_KEY = (
)
@skip_unless_docker_or_firecracker()
@skip_unless_selected_backend_available()
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true"
and os.environ.get("BOT_BOTTLE_BACKEND") != "firecracker",
+54
View File
@@ -385,6 +385,60 @@ class TestHasBackend(unittest.TestCase):
self.assertFalse(has_backend("nonexistent"))
class TestIsBackendAvailable(unittest.TestCase):
def test_delegates_to_has_backend(self):
from bot_bottle.backend import is_backend_available
with patch.object(backend_mod, "_backends", {}):
self.assertFalse(is_backend_available("docker"))
def test_known_and_available(self):
class _Ready:
def is_available(self):
return True
from bot_bottle.backend import is_backend_available
with patch.object(backend_mod, "_backends", {"docker": _Ready()}):
self.assertTrue(is_backend_available("docker"))
class TestIsBackendReady(unittest.TestCase):
def test_unknown_backend_returns_false(self):
from bot_bottle.backend import is_backend_ready
with patch.object(backend_mod, "_backends", {}):
self.assertFalse(is_backend_ready("docker"))
def test_ready_when_status_returns_zero(self):
class _ReadyBackend:
def status(self, *, quiet: bool = False) -> int:
return 0
from bot_bottle.backend import is_backend_ready
with patch.object(backend_mod, "_backends", {"docker": _ReadyBackend()}):
self.assertTrue(is_backend_ready("docker"))
def test_not_ready_when_status_nonzero(self):
class _BrokenBackend:
def status(self, *, quiet: bool = False) -> int:
return 1
from bot_bottle.backend import is_backend_ready
with patch.object(backend_mod, "_backends", {"docker": _BrokenBackend()}):
self.assertFalse(is_backend_ready("docker"))
def test_quiet_flag_forwarded(self):
calls = []
class _SpyBackend:
def status(self, *, quiet: bool = False) -> int:
calls.append(quiet)
return 0
from bot_bottle.backend import is_backend_ready
with patch.object(backend_mod, "_backends", {"docker": _SpyBackend()}):
is_backend_ready("docker", quiet=True)
self.assertEqual([True], calls)
class TestEnsureOrchestrator(unittest.TestCase):
"""The backend-agnostic orchestrator bring-up entry point. Docker starts
the orchestrator + gateway containers; firecracker boots the infra VM;
+138
View File
@@ -328,5 +328,143 @@ class TestNetpoolShellRenderers(unittest.TestCase):
self.assertIn("BOT_BOTTLE_FC_POOL_SIZE=4", out)
class TestFirecrackerBinaryCheck(unittest.TestCase):
def test_binary_missing_returns_false(self):
with patch.object(fc.shutil, "which", return_value=None):
self.assertFalse(fc._firecracker_binary_ok())
def test_binary_present_and_runs_ok(self):
with patch.object(fc.shutil, "which", return_value="/usr/bin/firecracker"), \
patch.object(fc.subprocess, "run",
return_value=subprocess.CompletedProcess([], 0)):
self.assertTrue(fc._firecracker_binary_ok())
def test_binary_found_but_exits_nonzero(self):
with patch.object(fc.shutil, "which", return_value="/usr/bin/firecracker"), \
patch.object(fc.subprocess, "run",
return_value=subprocess.CompletedProcess([], 1)):
self.assertFalse(fc._firecracker_binary_ok())
def test_binary_found_but_oserror(self):
with patch.object(fc.shutil, "which", return_value="/usr/bin/firecracker"), \
patch.object(fc.subprocess, "run", side_effect=OSError("exec failed")):
self.assertFalse(fc._firecracker_binary_ok())
class TestFirecrackerKvmCheck(unittest.TestCase):
def test_kvm_device_absent_returns_false(self):
with patch.object(fc.os.path, "exists", return_value=False):
self.assertFalse(fc._kvm_accessible())
def test_kvm_accessible_when_ioctl_succeeds(self):
m = MagicMock()
m.__enter__ = MagicMock(return_value=m)
m.__exit__ = MagicMock(return_value=False)
with patch.object(fc.os.path, "exists", return_value=True), \
patch("builtins.open", return_value=m), \
patch.object(fc.fcntl, "ioctl", return_value=12):
self.assertTrue(fc._kvm_accessible())
def test_kvm_present_but_ioctl_fails(self):
m = MagicMock()
m.__enter__ = MagicMock(return_value=m)
m.__exit__ = MagicMock(return_value=False)
with patch.object(fc.os.path, "exists", return_value=True), \
patch("builtins.open", return_value=m), \
patch.object(fc.fcntl, "ioctl", side_effect=OSError("permission denied")):
self.assertFalse(fc._kvm_accessible())
class TestFirecrackerStatusRuntime(unittest.TestCase):
"""status() reports binary and KVM problems and returns non-zero."""
def _apply_pool_ok(self, stack: contextlib.ExitStack) -> None:
stack.enter_context(patch.object(netpool, "missing_taps", return_value=[]))
stack.enter_context(patch.object(netpool, "pool_size", return_value=8))
stack.enter_context(patch.object(netpool, "overlapping_routes", return_value=[]))
stack.enter_context(patch.object(fc, "_report_persistence", lambda: None))
def test_status_fails_when_binary_missing(self):
with contextlib.ExitStack() as stack:
stack.enter_context(
patch.object(fc, "_firecracker_binary_ok", return_value=False))
stack.enter_context(
patch.object(fc, "_kvm_accessible", return_value=True))
stack.enter_context(
patch.object(fc.shutil, "which", return_value=None))
self._apply_pool_ok(stack)
rc, out = _cap(fc.status)
self.assertEqual(1, rc)
self.assertIn("NOT found on PATH", out)
def test_status_fails_when_kvm_not_accessible(self):
with contextlib.ExitStack() as stack:
stack.enter_context(
patch.object(fc, "_firecracker_binary_ok", return_value=True))
stack.enter_context(
patch.object(fc.shutil, "which", return_value="/usr/bin/firecracker"))
stack.enter_context(
patch.object(fc, "_kvm_accessible", return_value=False))
stack.enter_context(
patch.object(fc.os.path, "exists", return_value=True))
self._apply_pool_ok(stack)
rc, out = _cap(fc.status)
self.assertEqual(1, rc)
self.assertIn("not accessible", out)
def test_status_ok_when_binary_and_kvm_ready(self):
with contextlib.ExitStack() as stack:
stack.enter_context(
patch.object(fc, "_firecracker_binary_ok", return_value=True))
stack.enter_context(
patch.object(fc.shutil, "which", return_value="/usr/bin/firecracker"))
stack.enter_context(
patch.object(fc, "_kvm_accessible", return_value=True))
stack.enter_context(
patch.object(netpool, "nft_table_present", return_value=True))
self._apply_pool_ok(stack)
rc, out = _cap(fc.status)
self.assertEqual(0, rc)
self.assertIn("firecracker binary: ok", out)
self.assertIn("KVM:", out)
class TestBackendStatusQuiet(unittest.TestCase):
"""status(quiet=True) routes the underlying status() call through
redirect_stderr so diagnostic output is suppressed. Verify the return
code is propagated and that calling with quiet=False leaves the non-quiet
path active (covered by the other TestDockerSetupStatus tests)."""
def test_docker_quiet_true_propagates_return_code(self):
from bot_bottle.backend.docker.backend import DockerBottleBackend
with patch.object(dk, "status", return_value=0):
self.assertEqual(0, DockerBottleBackend.status(quiet=True))
def test_docker_quiet_true_suppresses_stderr(self):
import io as _io
from bot_bottle.backend.docker.backend import DockerBottleBackend
def _loud_status():
import sys
print("should be suppressed", file=sys.stderr)
return 0
with patch.object(dk, "status", side_effect=_loud_status):
buf = _io.StringIO()
with contextlib.redirect_stderr(buf):
DockerBottleBackend.status(quiet=True)
self.assertEqual("", buf.getvalue())
def test_firecracker_quiet_true_propagates_return_code(self):
from bot_bottle.backend.firecracker.backend import FirecrackerBottleBackend
with patch.object(fc, "status", return_value=1):
self.assertEqual(1, FirecrackerBottleBackend.status(quiet=True))
def test_macos_quiet_true_propagates_return_code(self):
from bot_bottle.backend.macos_container.backend import MacosContainerBottleBackend
with patch.object(mc, "status", return_value=0):
self.assertEqual(0, MacosContainerBottleBackend.status(quiet=True))
if __name__ == "__main__":
unittest.main()
+78
View File
@@ -0,0 +1,78 @@
"""Tests for the backend-aware skip guards in ``tests/_backend.py``.
The guards delegate their readiness check to
``bot_bottle.backend.is_backend_ready`` (the probe behind ``./cli.py backend
status``); here that probe is mocked so the unit job asserts the
selection/skip logic without either backend present on the runner.
"""
from __future__ import annotations
import os
import unittest
from unittest.mock import patch
from tests._backend import (
selected_backend,
skip_unless_backend,
skip_unless_selected_backend_available,
)
def _skipped(decorated: type) -> bool:
return getattr(decorated, "__unittest_skip__", False)
def _new_case() -> type:
return type("Case", (unittest.TestCase,), {})
class TestSelectedBackend(unittest.TestCase):
def test_defaults_to_docker_when_unset(self):
with patch.dict(os.environ, {}, clear=True):
self.assertEqual("docker", selected_backend())
def test_reads_env(self):
with patch.dict(os.environ, {"BOT_BOTTLE_BACKEND": "firecracker"}, clear=True):
self.assertEqual("firecracker", selected_backend())
class TestSkipUnlessBackend(unittest.TestCase):
def test_skips_when_other_backend_selected(self):
# A different backend is selected — no host probe needed, skip.
with patch.dict(os.environ, {"BOT_BOTTLE_BACKEND": "firecracker"}, clear=True), \
patch("tests._backend.is_backend_ready") as has:
decorated = skip_unless_backend("docker")(_new_case())
self.assertTrue(_skipped(decorated))
has.assert_not_called()
def test_runs_when_selected_and_available(self):
with patch.dict(os.environ, {"BOT_BOTTLE_BACKEND": "docker"}, clear=True), \
patch("tests._backend.is_backend_ready", return_value=True):
decorated = skip_unless_backend("docker")(_new_case())
self.assertFalse(_skipped(decorated))
def test_skips_when_selected_but_unavailable(self):
with patch.dict(os.environ, {"BOT_BOTTLE_BACKEND": "docker"}, clear=True), \
patch("tests._backend.is_backend_ready", return_value=False):
decorated = skip_unless_backend("docker")(_new_case())
self.assertTrue(_skipped(decorated))
class TestSkipUnlessSelectedBackendAvailable(unittest.TestCase):
def test_runs_when_selected_backend_available(self):
with patch.dict(os.environ, {"BOT_BOTTLE_BACKEND": "firecracker"}, clear=True), \
patch("tests._backend.is_backend_ready", return_value=True) as has:
decorated = skip_unless_selected_backend_available()(_new_case())
self.assertFalse(_skipped(decorated))
has.assert_called_once_with("firecracker", quiet=False)
def test_skips_when_selected_backend_unavailable(self):
with patch.dict(os.environ, {"BOT_BOTTLE_BACKEND": "firecracker"}, clear=True), \
patch("tests._backend.is_backend_ready", return_value=False):
decorated = skip_unless_selected_backend_available()(_new_case())
self.assertTrue(_skipped(decorated))
if __name__ == "__main__":
unittest.main()
-86
View File
@@ -1,86 +0,0 @@
"""Unit: role-scoped control-plane tokens (issue #469 review)."""
from __future__ import annotations
import base64
import json
import unittest
from bot_bottle.control_auth import ROLE_CLI, ROLE_GATEWAY, ROLES, mint, verify
_KEY = "test-key"
def _b64(obj: object) -> str:
return base64.urlsafe_b64encode(
json.dumps(obj).encode()).rstrip(b"=").decode("ascii")
class TestMintVerify(unittest.TestCase):
def test_round_trip_each_role(self) -> None:
for role in ROLES:
self.assertEqual(role, verify(mint(role, _KEY), _KEY))
def test_gateway_and_cli_are_distinct(self) -> None:
self.assertEqual(ROLE_GATEWAY, verify(mint(ROLE_GATEWAY, _KEY), _KEY))
self.assertEqual(ROLE_CLI, verify(mint(ROLE_CLI, _KEY), _KEY))
def test_wrong_key_rejected(self) -> None:
self.assertIsNone(verify(mint(ROLE_GATEWAY, _KEY), "other-key"))
def test_tampered_signature_rejected(self) -> None:
tok = mint(ROLE_CLI, _KEY)
self.assertIsNone(verify(tok[:-1] + ("A" if tok[-1] != "A" else "B"), _KEY))
def test_tampered_payload_rejected(self) -> None:
header, _payload, sig = mint(ROLE_GATEWAY, _KEY).split(".")
forged = f"{header}.{_b64({'role': 'cli'})}.{sig}"
self.assertIsNone(verify(forged, _KEY))
def test_alg_none_rejected(self) -> None:
# An unsigned "alg: none" token must never verify.
tok = f"{_b64({'alg': 'none', 'typ': 'JWT'})}.{_b64({'role': 'cli'})}."
self.assertIsNone(verify(tok, _KEY))
def test_validly_signed_but_wrong_alg_rejected(self) -> None:
# Alg-confusion: even a *correctly signed* token whose header claims a
# non-HS256 alg must be rejected.
from bot_bottle.control_auth import _sign # noqa: PLC0415
signing_input = f"{_b64({'alg': 'none', 'typ': 'JWT'})}.{_b64({'role': 'cli'})}"
self.assertIsNone(verify(f"{signing_input}.{_sign(_KEY, signing_input)}", _KEY))
def test_validly_signed_but_undecodable_rejected(self) -> None:
# A correct signature over a header that isn't valid base64/JSON still
# fails closed rather than raising.
from bot_bottle.control_auth import _sign # noqa: PLC0415
signing_input = f"!!!not-base64!!!.{_b64({'role': 'cli'})}"
self.assertIsNone(verify(f"{signing_input}.{_sign(_KEY, signing_input)}", _KEY))
def test_malformed_tokens_rejected(self) -> None:
for bad in ("", "a", "a.b", "a.b.c.d", "not-base64.$$.$$"):
self.assertIsNone(verify(bad, _KEY))
def test_empty_key_never_verifies(self) -> None:
self.assertIsNone(verify(mint(ROLE_CLI, _KEY), ""))
def test_unknown_role_in_payload_rejected(self) -> None:
header, _p, _s = mint(ROLE_CLI, _KEY).split(".")
# Re-sign a token carrying an unknown role — a valid signature but a
# role the control plane doesn't recognise must still be rejected.
from bot_bottle.control_auth import _sign # noqa: PLC0415
payload = _b64({"role": "root"})
signing_input = f"{header}.{payload}"
forged = f"{signing_input}.{_sign(_KEY, signing_input)}"
self.assertIsNone(verify(forged, _KEY))
def test_mint_rejects_unknown_role(self) -> None:
with self.assertRaises(ValueError):
mint("root", _KEY)
def test_mint_rejects_empty_key(self) -> None:
with self.assertRaises(ValueError):
mint(ROLE_GATEWAY, "")
if __name__ == "__main__":
unittest.main()
-35
View File
@@ -1,35 +0,0 @@
"""Tests for integration-test backend selection helpers."""
from __future__ import annotations
import os
import unittest
from unittest.mock import patch
from tests._docker import skip_unless_docker_or_firecracker
class TestSkipUnlessDockerOrFirecracker(unittest.TestCase):
def test_firecracker_runs_when_docker_tests_are_disabled(self):
with patch.dict(
os.environ,
{"BOT_BOTTLE_BACKEND": "firecracker", "SKIP_DOCKER_TESTS": "1"},
clear=True,
):
decorated = skip_unless_docker_or_firecracker()(type("Case", (), {}))
self.assertFalse(getattr(decorated, "__unittest_skip__", False))
def test_non_firecracker_still_skips_when_docker_tests_are_disabled(self):
with patch.dict(
os.environ,
{"BOT_BOTTLE_BACKEND": "docker", "SKIP_DOCKER_TESTS": "1"},
clear=True,
):
decorated = skip_unless_docker_or_firecracker()(type("Case", (), {}))
self.assertTrue(getattr(decorated, "__unittest_skip__", False))
if __name__ == "__main__":
unittest.main()
+73 -86
View File
@@ -209,7 +209,6 @@ from bot_bottle.egress_addon_core import ( # noqa: E402
Route,
route_to_yaml_dict,
)
from bot_bottle.policy_resolver import PolicyResolveError # noqa: E402
# ---------------------------------------------------------------------------
@@ -267,48 +266,7 @@ def _config_to_policy(config: Config) -> str:
}) + "\n"
class _SuperviseRpcFake:
"""Mixin adding the control-plane supervise RPCs to a fake resolver, now
that the egress data plane queues/polls proposals over RPC instead of
opening the DB (issue #469). `supervise_status` is the operator's eventual
decision (None models a timeout poll stays `pending`); `propose_error`
models an unreachable orchestrator. `propose_calls` records what was queued
so a test can assert the proposal was attributed by the caller's source IP."""
supervise_status: "str | None" = None
propose_error: bool = False
poll_error: bool = False
@property
def propose_calls(self) -> list[dict[str, str]]:
if not hasattr(self, "_propose_calls"):
self._propose_calls: list[dict[str, str]] = []
return self._propose_calls
def propose_supervise(
self, source_ip: str, identity_token: str, *,
tool: str, proposed_file: str, justification: str,
) -> str:
del proposed_file, justification
self.propose_calls.append(
{"source_ip": source_ip, "identity_token": identity_token, "tool": tool}
)
if self.propose_error:
raise PolicyResolveError("orchestrator down")
return "prop-1"
def poll_supervise(
self, source_ip: str, identity_token: str, proposal_id: str,
) -> dict[str, object]:
del source_ip, identity_token, proposal_id
if self.poll_error:
raise PolicyResolveError("orchestrator down")
if self.supervise_status is None:
return {"status": "pending"}
return {"status": self.supervise_status, "notes": "", "final_file": None}
class _StaticResolver(_SuperviseRpcFake):
class _StaticResolver:
"""Fake orchestrator resolver that serves one Config (+ optional bottle id
and per-bottle tokens) for every client the host-test stand-in for a
bottle's policy now that egress is resolver-only."""
@@ -365,7 +323,7 @@ def _with_client_ip(flow: _Flow, ip: str) -> _Flow:
return flow
class _CtxResolver(_SuperviseRpcFake):
class _CtxResolver:
"""Fake orchestrator resolver: maps source IP -> bottle id, and grants the
same allow-list to any attributed bottle (unattributed -> deny)."""
@@ -558,47 +516,70 @@ class TestOutboundDlpPolicy(unittest.TestCase):
# ---------------------------------------------------------------------------
def _fake_sv(response_status: str | None) -> types.SimpleNamespace:
"""Stand-in for the `supervise` module the adapter queues proposals to.
`response_status` of None models a timeout (read_response never returns a
decision); a status string models the operator's eventual answer."""
def _new_proposal(**_kw: Any) -> Any:
return types.SimpleNamespace(id="prop-1")
def _sha256_hex(_payload: Any) -> str:
return "hash"
def _noop(*_args: Any) -> None:
return None
def _read_response(_slug: Any, _pid: Any) -> Any:
if response_status is None:
raise OSError("not written yet") # forces poll -> timeout
return types.SimpleNamespace(status=response_status)
ns = types.SimpleNamespace()
ns.STATUS_APPROVED = "approved"
ns.STATUS_MODIFIED = "modified"
ns.TOOL_EGRESS_TOKEN_ALLOW = "egress_token_allow"
ns.Proposal = types.SimpleNamespace(new=_new_proposal)
ns.sha256_hex = _sha256_hex
ns.write_proposal = _noop
ns.archive_proposal = _noop
ns.read_response = _read_response
return ns
class TestSuperviseBranch(unittest.TestCase):
def _supervised_addon(self, status: str | None) -> EgressAddon:
def _supervised_addon(self) -> EgressAddon:
addon = _addon(Config(routes=(Route(host="api.example.com"),)), slug="test-bottle")
addon._token_allow_timeout = 0.05
cast(Any, addon._resolver).supervise_status = status
return addon
def test_operator_approval_allows_token_and_forwards(self) -> None:
addon = self._supervised_addon("approved")
addon = self._supervised_addon()
flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}"))
_run_request(addon, flow)
with patch.object(_ea_mod, "_sv", _fake_sv("approved")):
_run_request(addon, flow)
self.assertIsNone(flow.response) # forwarded after approval
# Approval lands in the calling bottle's safelist (keyed by slug).
self.assertIn(_OPENAI_KEY, addon._safe_tokens_for("test-bottle"))
def test_operator_rejection_blocks(self) -> None:
addon = self._supervised_addon("rejected")
addon = self._supervised_addon()
flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}"))
_run_request(addon, flow)
with patch.object(_ea_mod, "_sv", _fake_sv("rejected")):
_run_request(addon, flow)
assert flow.response is not None
self.assertEqual(403, flow.response.status_code)
self.assertIn("rejected", flow.response.get_text())
def test_supervise_timeout_blocks(self) -> None:
addon = self._supervised_addon(None) # poll stays pending -> timeout
addon = self._supervised_addon()
flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}"))
_run_request(addon, flow)
with patch.object(_ea_mod, "_sv", _fake_sv(None)):
_run_request(addon, flow)
assert flow.response is not None
self.assertEqual(403, flow.response.status_code)
self.assertIn("timed out", flow.response.get_text())
def test_poll_error_during_wait_times_out_and_blocks(self) -> None:
# A transient orchestrator error on each poll is retried until the
# deadline, then fails closed (blocked) — never forwarded unsupervised.
addon = self._supervised_addon("approved")
cast(Any, addon._resolver).poll_error = True
flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}"))
_run_request(addon, flow)
assert flow.response is not None
self.assertEqual(403, flow.response.status_code)
# ---------------------------------------------------------------------------
# Inbound DLP on responses
@@ -791,14 +772,19 @@ class TestRedactSurfaces(unittest.TestCase):
class TestSuperviseWriteFailure(unittest.TestCase):
def test_propose_rpc_error_blocks(self) -> None:
# An unreachable orchestrator (propose RPC raises) fails closed: the
# request is blocked rather than forwarded unsupervised.
def test_write_proposal_oserror_blocks(self) -> None:
addon = _addon(Config(routes=(Route(host="api.example.com"),)), slug="test-bottle")
addon._token_allow_timeout = 0.05
cast(Any, addon._resolver).propose_error = True
flow = _Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}"))
_run_request(addon, flow)
fake = _fake_sv("approved")
def _raise(_p: Any) -> None:
raise OSError("disk full")
fake.write_proposal = _raise
with patch.object(_ea_mod, "_sv", fake):
_run_request(addon, flow)
assert flow.response is not None
self.assertEqual(403, flow.response.status_code)
@@ -858,23 +844,22 @@ class TestSuperviseMultiTenant(unittest.TestCase):
"""Consolidated gateway: supervise proposals + the DLP safelist are keyed
per bottle, resolved by source IP (PRD 0070)."""
def _consolidated_addon(self, status: str | None = "approved") -> EgressAddon:
def _consolidated_addon(self) -> EgressAddon:
# Static config is empty; the resolver supplies each bottle's config.
addon = _addon(Config(routes=()))
resolver = _CtxResolver({"10.0.0.1": "bottle-a", "10.0.0.2": "bottle-b"})
resolver.supervise_status = status
addon._resolver = cast(Any, resolver)
addon._resolver = cast(Any, _CtxResolver({"10.0.0.1": "bottle-a", "10.0.0.2": "bottle-b"}))
addon._token_allow_timeout = 0.05
return addon
def test_approval_is_scoped_to_the_calling_bottle(self) -> None:
addon = self._consolidated_addon("approved")
addon = self._consolidated_addon()
# bottle-a (10.0.0.1) sends the token; the operator approves.
flow = _with_client_ip(
_Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")),
"10.0.0.1",
)
_run_request(addon, flow)
with patch.object(_ea_mod, "_sv", _fake_sv("approved")):
_run_request(addon, flow)
self.assertIsNone(flow.response) # forwarded after approval
# The approval lands ONLY in bottle-a's safelist — never bottle-b's.
# A global set here would be the cross-tenant leak this slice closes.
@@ -882,19 +867,22 @@ class TestSuperviseMultiTenant(unittest.TestCase):
self.assertNotIn(_OPENAI_KEY, addon._safe_tokens_for("bottle-b"))
def test_proposal_is_attributed_to_the_source_ip_bottle(self) -> None:
addon = self._consolidated_addon("approved")
addon = self._consolidated_addon()
seen: list[str] = []
fake = _fake_sv("approved")
def _capture(**kw: Any) -> Any:
seen.append(kw["bottle_slug"])
return types.SimpleNamespace(id="p")
fake.Proposal = types.SimpleNamespace(new=_capture)
flow = _with_client_ip(
_Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")),
"10.0.0.2",
)
_run_request(addon, flow)
# The addon forwards the caller's source IP to the control plane, which
# attributes the proposal server-side by (source_ip, identity_token) —
# the addon never asserts a slug. The resolved bottle keys the safelist.
calls = cast(Any, addon._resolver).propose_calls
self.assertEqual(["10.0.0.2"], [c["source_ip"] for c in calls])
self.assertIn(_OPENAI_KEY, addon._safe_tokens_for("bottle-b"))
self.assertNotIn(_OPENAI_KEY, addon._safe_tokens_for("bottle-a"))
with patch.object(_ea_mod, "_sv", fake):
_run_request(addon, flow)
self.assertEqual(["bottle-b"], seen) # proposal keyed by the resolved bottle
def test_auth_token_injected_from_resolved_tokens(self) -> None:
# The bottle's upstream token comes from /resolve (in-memory on the
@@ -937,17 +925,16 @@ class TestSuperviseMultiTenant(unittest.TestCase):
self.assertIsNotNone(flow.response) # blocked — token unset
def test_unattributed_source_ip_cannot_supervise(self) -> None:
addon = self._consolidated_addon("approved")
addon = self._consolidated_addon()
# 10.9.9.9 is not in the resolver map -> deny-all config, empty slug.
flow = _with_client_ip(
_Flow(_Request(host="api.example.com", method="POST", body=f"k={_OPENAI_KEY}")),
"10.9.9.9",
)
_run_request(addon, flow)
with patch.object(_ea_mod, "_sv", _fake_sv("approved")):
_run_request(addon, flow)
self.assertIsNotNone(flow.response) # blocked (no route, no supervise)
self.assertNotIn(_OPENAI_KEY, addon._safe_tokens_for(""))
# Never even reached the queue — no route means no supervise proposal.
self.assertEqual([], cast(Any, addon._resolver).propose_calls)
class TestMultiTenantInboundDlp(unittest.TestCase):
+3 -1
View File
@@ -136,7 +136,9 @@ class TestFirecrackerStatus(unittest.TestCase):
from bot_bottle.backend.firecracker import setup as fc_setup
with patch.object(fc_setup.netpool, "missing_taps", return_value=[]), \
patch.object(fc_setup.netpool, "overlapping_routes", return_value=[]), \
patch.object(fc_setup.shutil, "which", return_value=None):
patch.object(fc_setup.shutil, "which", return_value=None), \
patch.object(fc_setup, "_firecracker_binary_ok", return_value=True), \
patch.object(fc_setup, "_kvm_accessible", return_value=True):
rc, out = self._run()
self.assertEqual(0, rc)
self.assertIn("unverified", out)
-40
View File
@@ -8,7 +8,6 @@ decisions that must hold without a VM.
from __future__ import annotations
import os
import tempfile
import unittest
from pathlib import Path
from unittest.mock import MagicMock, patch
@@ -16,37 +15,6 @@ from unittest.mock import MagicMock, patch
from bot_bottle.backend.firecracker import infra_vm
class TestSigningKeySync(unittest.TestCase):
"""`_with_signing_key` mirrors the VM's control-plane signing key into the
host token file so the CLI signs `cli` tokens the VM verifies (issue #469)."""
def _infra(self) -> infra_vm.InfraVm:
return infra_vm.InfraVm(guest_ip="10.0.0.1", private_key=Path("/k"))
def test_mirrors_guest_key_to_host_file(self):
with tempfile.TemporaryDirectory() as d:
root = Path(d)
proc = MagicMock(returncode=0, stdout="the-signing-key\n")
with patch.object(infra_vm.subprocess, "run", return_value=proc) as run, \
patch.object(infra_vm, "bot_bottle_root", return_value=root):
out = infra_vm._with_signing_key(self._infra())
self.assertIsInstance(out, infra_vm.InfraVm) # returned for chaining
token = root / infra_vm.CONTROL_PLANE_TOKEN_FILENAME
self.assertEqual("the-signing-key", token.read_text())
self.assertEqual(0o600, token.stat().st_mode & 0o777)
# It cat'd the guest volume path over SSH.
self.assertIn(infra_vm._GUEST_SIGNING_KEY_PATH, run.call_args.args[0][-1])
def test_unreadable_key_is_not_fatal(self):
with tempfile.TemporaryDirectory() as d:
root = Path(d)
proc = MagicMock(returncode=255, stdout="")
with patch.object(infra_vm.subprocess, "run", return_value=proc), \
patch.object(infra_vm, "bot_bottle_root", return_value=root):
infra_vm._with_signing_key(self._infra()) # no raise
self.assertFalse((root / infra_vm.CONTROL_PLANE_TOKEN_FILENAME).exists())
class TestControlPlaneUrl(unittest.TestCase):
def test_url_uses_guest_ip_and_port(self):
infra = infra_vm.InfraVm(
@@ -78,14 +46,6 @@ class TestBuildInfraRootfs(unittest.TestCase):
self.assertIn("/dev/vdb", init)
# VM backend uses git-http (9420); the git:// daemon is left out.
self.assertIn("BOT_BOTTLE_GATEWAY_DAEMONS=egress,git-http,supervise", init)
# Role-scoped control-plane auth (issue #469 review): the orchestrator
# gets the signing key, the gateway daemons get a pre-minted `gateway`
# JWT — never open mode in the infra VM.
self.assertIn("host_control_plane_token", init) # key generated on the volume
self.assertIn("mint, ROLE_GATEWAY", init) # gateway JWT minted from it
self.assertIn('BOT_BOTTLE_CONTROL_PLANE_TOKEN="$CP_KEY" python3 -m bot_bottle.orchestrator',
init) # key -> orchestrator only
self.assertIn('BOT_BOTTLE_CONTROL_AUTH_JWT="$GW_JWT"', init) # JWT -> gateway daemons
class TestSshGatewayTransport(unittest.TestCase):
-24
View File
@@ -64,30 +64,6 @@ class TestEnvForDaemon(unittest.TestCase):
self.assertNotIn("X", self._BASE)
class TestControlPlaneEnvScoping(unittest.TestCase):
"""The control-plane signing key stays with the orchestrator; the pre-minted
`gateway` JWT goes to the data-plane daemons (issue #469 review). Scoping
them per-process keeps a compromised data-plane daemon from reading the key
and minting a higher-privilege token, even in the combined infra container."""
_BASE = {
"PATH": "/usr/bin",
"BOT_BOTTLE_CONTROL_PLANE_TOKEN": "sk-x",
"BOT_BOTTLE_CONTROL_AUTH_JWT": "gw-jwt",
}
def test_orchestrator_gets_key_not_jwt(self):
env = _env_for_daemon("orchestrator", self._BASE)
self.assertEqual("sk-x", env["BOT_BOTTLE_CONTROL_PLANE_TOKEN"])
self.assertNotIn("BOT_BOTTLE_CONTROL_AUTH_JWT", env)
def test_data_plane_daemons_get_jwt_not_key(self):
for name in ("egress", "git-gate", "git-http", "supervise"):
env = _env_for_daemon(name, self._BASE)
self.assertNotIn("BOT_BOTTLE_CONTROL_PLANE_TOKEN", env, name)
self.assertEqual("gw-jwt", env["BOT_BOTTLE_CONTROL_AUTH_JWT"], name)
class TestSelectedDaemons(unittest.TestCase):
"""Env-var subset filtering. The compose renderer is the source
of truth for which daemons are wired; the supervisor just
+10 -13
View File
@@ -213,25 +213,22 @@ class TestHookRender(unittest.TestCase):
# the suppressed findings for human approval.
self.assertIn("--ignore-gitleaks-allow", hook)
self.assertIn("--report-format=json", hook)
# The hook queues + polls over the control-plane RPC — it no longer
# opens the DB directly (PRD 0070 / issue #469).
self.assertIn("tool=TOOL_GITLEAKS_ALLOW", hook)
self.assertIn("propose_supervise", hook)
self.assertIn("poll_supervise", hook)
self.assertIn("SUPERVISE_SOURCE_IP", hook)
self.assertIn("SUPERVISE_IDENTITY_TOKEN", hook)
self.assertIn("tool=_sv.TOOL_GITLEAKS_ALLOW", hook)
self.assertIn("_sv.write_proposal", hook)
self.assertIn("_sv.read_response", hook)
self.assertIn("SUPERVISE_BOTTLE_SLUG", hook)
self.assertIn("supervisor approved # gitleaks:allow", hook)
self.assertIn("supervisor rejected # gitleaks:allow", hook)
def test_inline_gitleaks_allow_python_imports_work_in_gateway_layout(self):
hook = git_gate_render_hook()
# The gateway image copies the package modules flat under /app, while
# host-side tests import them through the bot_bottle package. Hooks
# execute from the bare repo directory, so the embedded Python must
# include /app and support both import layouts.
# The gateway image copies supervise.py flat under /app, while
# host-side tests import it through the bot_bottle package.
# Hooks execute from the bare repo directory, so the embedded
# Python must include /app and support both import layouts.
self.assertIn('PYTHONPATH="/app${PYTHONPATH:+:$PYTHONPATH}"', hook)
self.assertIn("from bot_bottle.policy_resolver import PolicyResolver", hook)
self.assertIn("from policy_resolver import PolicyResolver", hook)
self.assertIn("import supervise as _sv", hook)
self.assertIn("from bot_bottle import supervise as _sv", hook)
def test_inline_gitleaks_allow_fails_closed_without_supervisor(self):
hook = git_gate_render_hook()
+8 -10
View File
@@ -112,13 +112,11 @@ class TestGitHttpBackend(unittest.TestCase):
).strip()
self.assertEqual(head, cloned)
def test_consolidated_push_stamps_supervise_attribution_for_the_hook(self):
# In consolidated mode the backend attributes the push by (source IP,
# identity token) and stamps SUPERVISE_SOURCE_IP + SUPERVISE_IDENTITY_TOKEN
# into the CGI env, so the gitleaks-allow pre-receive hook can queue its
# proposal over the control-plane RPC (which re-resolves the bottle from
# exactly that pair — PRD 0070 / issue #469). The hook here just records
# the source IP it received.
def test_consolidated_push_stamps_bottle_slug_for_the_hook(self):
# In consolidated mode the backend attributes the push by source IP and
# stamps SUPERVISE_BOTTLE_SLUG=<bottle_id> into the CGI env, so the
# gitleaks-allow pre-receive hook queues its proposal under the right
# bottle. The hook here just records what it received.
from http.server import ThreadingHTTPServer
bottle_id = "bottleab12"
@@ -132,10 +130,10 @@ class TestGitHttpBackend(unittest.TestCase):
["git", "-C", str(bare), "config", "http.receivepack", "true"],
check=True,
)
capture = root / "source-ip-capture"
capture = root / "slug-capture"
hook = bare / "hooks" / "pre-receive"
hook.write_text(
f"#!/bin/sh\nprintf '%s' \"${{SUPERVISE_SOURCE_IP:-UNSET}}\" > "
f"#!/bin/sh\nprintf '%s' \"${{SUPERVISE_BOTTLE_SLUG:-UNSET}}\" > "
f"{capture}\ncat >/dev/null\nexit 0\n"
)
hook.chmod(0o755)
@@ -168,7 +166,7 @@ class TestGitHttpBackend(unittest.TestCase):
["git", "push", url, "HEAD:refs/heads/main"],
cwd=work, check=True, capture_output=True, text=True, timeout=5,
)
self.assertEqual("127.0.0.1", capture.read_text())
self.assertEqual(bottle_id, capture.read_text())
def test_post_forwards_git_cgi_headers(self):
from http.server import ThreadingHTTPServer
-15
View File
@@ -7,30 +7,15 @@ import unittest
import urllib.error
from unittest.mock import MagicMock, patch
from bot_bottle.control_auth import ROLE_CLI, verify
from bot_bottle.orchestrator.client import (
OrchestratorClient,
OrchestratorClientError,
RegisteredBottle,
_host_auth_token,
)
_URLOPEN = "bot_bottle.orchestrator.client.urllib.request.urlopen"
class TestHostAuthToken(unittest.TestCase):
def test_mints_a_cli_token_from_the_host_key(self) -> None:
with patch("bot_bottle.orchestrator.client.host_control_plane_token",
return_value="signing-key"):
tok = _host_auth_token()
self.assertEqual(ROLE_CLI, verify(tok, "signing-key"))
def test_returns_empty_when_key_unreadable(self) -> None:
with patch("bot_bottle.orchestrator.client.host_control_plane_token",
side_effect=OSError("no host root")):
self.assertEqual("", _host_auth_token())
def _resp(status: int, payload: object) -> MagicMock:
m = MagicMock()
inner = m.__enter__.return_value
+31 -205
View File
@@ -19,10 +19,9 @@ from contextlib import closing
from pathlib import Path
from unittest.mock import patch
from bot_bottle.control_auth import ROLE_CLI, ROLE_GATEWAY, mint
from bot_bottle.orchestrator.broker import StubBroker
from bot_bottle.orchestrator.control_plane import dispatch, make_server
from bot_bottle.orchestrator.registry import BottleRecord, RegistryStore
from bot_bottle.orchestrator.registry import RegistryStore
from bot_bottle.orchestrator.service import Orchestrator
from bot_bottle.store_manager import StoreManager
from bot_bottle.supervise import (
@@ -286,71 +285,45 @@ class TestServerRoundTrip(unittest.TestCase):
class TestControlPlaneAuth(unittest.TestCase):
"""Role-scoped control-plane tokens (issue #400 / #469 review): every route
but /health needs a valid token, and the token's role gates which routes it
reaches a `gateway` data-plane token can't drive the operator routes."""
_OPERATOR_ROUTES = [
("GET", "/bottles", b""),
("POST", "/bottles", _body({"source_ip": "10.0.0.1"})),
("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})),
("DELETE", "/bottles/x", b""),
("POST", "/reconcile", _body({"live_source_ips": []})),
("POST", "/attribute", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("GET", "/supervise/proposals", b""),
("POST", "/supervise/respond",
_body({"proposal_id": "p", "bottle_slug": "s", "decision": "a"})),
]
"""The per-host control-plane secret (issue #400): every route but /health
is a trusted-caller op an agent must not be able to drive just because it
can reach the port."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.orch = _orchestrator(Path(self._tmp.name) / "r.db")
def test_health_is_public_even_unauthenticated(self) -> None:
status, _ = dispatch(self.orch, "GET", "/health", b"", role=None)
def test_health_is_public_even_unauthorized(self) -> None:
status, _ = dispatch(self.orch, "GET", "/health", b"", authorized=False)
self.assertEqual(200, status)
def test_unauthenticated_denies_every_other_route(self) -> None:
routes = self._OPERATOR_ROUTES + [
def test_unauthorized_denies_every_other_route(self) -> None:
for method, path, body in [
("GET", "/bottles", b""),
("POST", "/bottles", _body({"source_ip": "10.0.0.1"})),
("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})),
("DELETE", "/bottles/x", b""),
("POST", "/reconcile", _body({"live_source_ips": []})),
("POST", "/resolve", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("POST", "/supervise/propose",
_body({"source_ip": "1", "tool": "egress-allow", "proposed_file": "x", "justification": "j"})),
("POST", "/supervise/poll", _body({"source_ip": "1", "proposal_id": "p"})),
]
for method, path, body in routes:
status, _ = dispatch(self.orch, method, path, body, role=None)
self.assertEqual(401, status, f"{method} {path} should be 401 unauthenticated")
def test_gateway_role_denied_on_operator_routes(self) -> None:
# A compromised data-plane process holds only a `gateway` token — it must
# not reach the operator routes (approve proposals, rewrite policy, …).
for method, path, body in self._OPERATOR_ROUTES:
status, payload = dispatch(self.orch, method, path, body, role=ROLE_GATEWAY)
self.assertEqual(403, status, f"{method} {path} should be 403 for gateway")
self.assertIn("insufficient role", str(payload.get("error")))
def test_gateway_role_allowed_on_data_routes(self) -> None:
# The gateway CAN reach its own lookups. Register a bottle so /resolve
# returns 200 rather than a 403 — proving the role gate let it through.
rec = self.orch.registry.register("10.0.0.9", policy="routes: []", metadata="")
status, _ = dispatch(
self.orch, "POST", "/resolve",
_body({"source_ip": "10.0.0.9", "identity_token": rec.identity_token}),
role=ROLE_GATEWAY)
self.assertEqual(200, status)
("POST", "/attribute", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("GET", "/supervise/proposals", b""),
("POST", "/supervise/respond", _body({"proposal_id": "p", "bottle_slug": "s", "decision": "approve"})),
]:
status, _ = dispatch(self.orch, method, path, body, authorized=False)
self.assertEqual(401, status, f"{method} {path} should be 401 unauthorized")
def test_deny_happens_before_the_registry_is_touched(self) -> None:
"""An unauthenticated DELETE must not tear a bottle down. 401, and the
"""An unauthorized DELETE must not tear a bottle down. 401, and the
bottle is still there."""
rec = self.orch.registry.register("10.0.0.9", policy="", metadata="")
status, _ = dispatch(
self.orch, "DELETE", f"/bottles/{rec.bottle_id}", b"", role=None)
self.orch, "DELETE", f"/bottles/{rec.bottle_id}", b"", authorized=False)
self.assertEqual(401, status)
self.assertIsNotNone(self.orch.registry.get(rec.bottle_id))
def _server_with_key(self, signing_key: str):
with patch.dict("os.environ", {"BOT_BOTTLE_CONTROL_PLANE_TOKEN": signing_key}):
def _server_with_secret(self, secret: str):
with patch.dict("os.environ", {"BOT_BOTTLE_CONTROL_PLANE_TOKEN": secret}):
server = make_server(self.orch, "127.0.0.1", 0)
self.addCleanup(server.server_close)
threading.Thread(target=server.serve_forever, daemon=True).start()
@@ -367,29 +340,25 @@ class TestControlPlaneAuth(unittest.TestCase):
except urllib.error.HTTPError as e:
return e.code
def test_configured_server_enforces_roles_over_http(self) -> None:
key = "test-key"
base = self._server_with_key(key)
gateway_tok = mint(ROLE_GATEWAY, key)
cli_tok = mint(ROLE_CLI, key)
def test_configured_server_enforces_the_header_over_http(self) -> None:
base = self._server_with_secret("s3cret-admin")
# /health is public — no header needed.
self.assertEqual(200, self._status(f"{base}/health"))
# /bottles (operator) needs a valid cli token.
# /bottles requires the secret.
self.assertEqual(401, self._status(f"{base}/bottles"))
self.assertEqual(401, self._status(f"{base}/bottles", header="wrong"))
self.assertEqual(403, self._status(f"{base}/bottles", header=gateway_tok))
self.assertEqual(200, self._status(f"{base}/bottles", header=cli_tok))
self.assertEqual(200, self._status(f"{base}/bottles", header="s3cret-admin"))
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."""
"""No secret set (tests / nft-protected Firecracker): open mode, so the
existing round-trip and unit behavior are unchanged."""
with patch.dict("os.environ", {}, clear=False):
import os
os.environ.pop("BOT_BOTTLE_CONTROL_PLANE_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"))
self.assertTrue(server.is_authorized(""))
self.assertTrue(server.is_authorized("anything"))
class TestDispatchSupervise(unittest.TestCase):
@@ -457,149 +426,6 @@ class TestDispatchSupervise(unittest.TestCase):
self.assertIn("no such proposal", str(payload["error"]))
class TestDispatchSuperviseAgentRpc(unittest.TestCase):
"""The agent half — `/supervise/propose` + `/supervise/poll` — attributed by
(source_ip, identity_token) like /resolve, so a bottle can only ever queue
or read its own proposals (PRD 0070 / issue #469)."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
root = Path(self._tmp.name)
db = root / "db" / "bot-bottle.db"
db.parent.mkdir(parents=True)
self._env = patch.dict("os.environ", {
"BOT_BOTTLE_ROOT": str(root),
"SUPERVISE_DB_PATH": str(db),
})
self._env.start()
self.store = RegistryStore(db)
self.store.migrate()
StoreManager(db).migrate()
secret = secrets.token_bytes(16)
self.orch = Orchestrator(self.store, StubBroker(secret), secret)
def tearDown(self) -> None:
self._env.stop()
self._tmp.cleanup()
def _register(self, source_ip: str = "10.243.0.9", slug: str = "demo"):
return self.store.register(
source_ip, metadata=json.dumps({"slug": slug}), policy="routes: []\n")
def _propose(self, rec: BottleRecord, proposed: str = "routes:\n - host: g.com\n"):
return dispatch(self.orch, "POST", "/supervise/propose", _body({
"source_ip": rec.source_ip, "identity_token": rec.identity_token,
"tool": TOOL_EGRESS_ALLOW, "proposed_file": proposed, "justification": "need it",
}))
def _poll(self, rec: BottleRecord, proposal_id: str):
return dispatch(self.orch, "POST", "/supervise/poll", _body({
"source_ip": rec.source_ip, "identity_token": rec.identity_token,
"proposal_id": proposal_id,
}))
def test_propose_queues_under_the_resolved_bottle(self) -> None:
rec = self._register()
status, payload = self._propose(rec)
self.assertEqual(201, status)
pid = payload["proposal_id"]
assert isinstance(pid, str) and pid
_, listing = dispatch(self.orch, "GET", "/supervise/proposals", b"")
proposals = listing["proposals"]
assert isinstance(proposals, list)
self.assertEqual(pid, proposals[0]["id"])
# Queued under the orchestrator-resolved bottle id, never a caller slug.
self.assertEqual(rec.bottle_id, proposals[0]["bottle_slug"])
def test_propose_unattributed_is_403(self) -> None:
status, payload = dispatch(self.orch, "POST", "/supervise/propose", _body({
"source_ip": "10.9.9.9", "identity_token": "wrong",
"tool": TOOL_EGRESS_ALLOW, "proposed_file": "x\n", "justification": "j",
}))
self.assertEqual(403, status)
self.assertIn("unattributed", str(payload["error"]))
def test_propose_rejects_unknown_tool(self) -> None:
rec = self._register()
status, _ = dispatch(self.orch, "POST", "/supervise/propose", _body({
"source_ip": rec.source_ip, "identity_token": rec.identity_token,
"tool": "not-a-tool", "proposed_file": "x\n", "justification": "j",
}))
self.assertEqual(400, status)
def test_poll_pending_then_decided_then_idempotent(self) -> None:
rec = self._register()
_, proposed = self._propose(rec)
pid = proposed["proposal_id"]
assert isinstance(pid, str)
status, poll = self._poll(rec, pid)
self.assertEqual(200, status)
self.assertEqual("pending", poll["status"])
# Operator decides server-side.
dispatch(self.orch, "POST", "/supervise/respond", _body({
"proposal_id": pid, "bottle_slug": rec.bottle_id,
"decision": "approve", "notes": "ok",
}))
_, decided = self._poll(rec, pid)
self.assertEqual("approved", decided["status"])
self.assertEqual("ok", decided["notes"])
# Poll doesn't archive: it's gone from the operator's pending list, but a
# re-poll returns the same decision so a dropped connection can't lose it.
_, listing = dispatch(self.orch, "GET", "/supervise/proposals", b"")
self.assertEqual([], listing["proposals"])
_, again = self._poll(rec, pid)
self.assertEqual(decided, again)
def test_poll_cannot_read_another_bottles_proposal(self) -> None:
rec_a = self._register("10.0.0.1", "a")
rec_b = self._register("10.0.0.2", "b")
_, proposed = self._propose(rec_a)
pid = proposed["proposal_id"]
assert isinstance(pid, str)
# b polls a's proposal id: scoped to b's own queue → never a's response.
status, poll = self._poll(rec_b, pid)
self.assertEqual(200, status)
self.assertEqual("unknown", poll["status"])
# --- request validation (400) + fail-closed (403) ----------------------
def test_propose_invalid_json_is_400(self) -> None:
status, _ = dispatch(self.orch, "POST", "/supervise/propose", b"{not json")
self.assertEqual(400, status)
def test_propose_missing_fields_are_400(self) -> None:
rec = self._register()
base = {
"source_ip": rec.source_ip, "identity_token": rec.identity_token,
"tool": TOOL_EGRESS_ALLOW, "proposed_file": "routes:\n", "justification": "j",
}
for drop in ("source_ip", "proposed_file", "justification"):
body = {k: v for k, v in base.items() if k != drop}
status, _ = dispatch(self.orch, "POST", "/supervise/propose", _body(body))
self.assertEqual(400, status, drop)
def test_poll_invalid_json_is_400(self) -> None:
status, _ = dispatch(self.orch, "POST", "/supervise/poll", b"{not json")
self.assertEqual(400, status)
def test_poll_missing_fields_are_400(self) -> None:
rec = self._register()
for body in ({"identity_token": rec.identity_token, "proposal_id": "p"},
{"source_ip": rec.source_ip, "identity_token": rec.identity_token}):
status, _ = dispatch(self.orch, "POST", "/supervise/poll", _body(body))
self.assertEqual(400, status)
def test_poll_unattributed_is_403(self) -> None:
status, payload = dispatch(self.orch, "POST", "/supervise/poll", _body({
"source_ip": "10.9.9.9", "identity_token": "wrong", "proposal_id": "p",
}))
self.assertEqual(403, status)
self.assertIn("unattributed", str(payload["error"]))
if __name__ == "__main__":
unittest.main()
+7 -5
View File
@@ -124,11 +124,13 @@ class TestDockerGateway(unittest.TestCase):
src = ca_mounts[0].rsplit(":", 1)[0]
self.assertTrue(src.endswith("/" + GATEWAY_CA_DIRNAME), src)
self.assertTrue(Path(src).is_absolute(), src)
# No DB handle on the data plane: the supervise queue is reached over
# the control-plane RPC, so the gateway container carries neither the
# DB bind-mount nor SUPERVISE_DB_PATH (PRD 0070 / issue #469).
self.assertFalse(any(a.startswith("SUPERVISE_DB_PATH=") for a in runs[0]))
self.assertFalse(any(a.endswith(":/run/supervise") for a in runs[0]))
# Shares the ONE host DB: the supervise daemon queues into the same
# file the orchestrator + operator (over HTTP) use.
self.assertTrue(any(
a.startswith("SUPERVISE_DB_PATH=") and a.endswith("/run/supervise/bot-bottle.db")
for a in runs[0]))
self.assertTrue(any(
a.endswith(":/run/supervise") for a in runs[0]))
# Data plane resolves policy against the orchestrator control plane.
self.assertIn(f"BOT_BOTTLE_ORCHESTRATOR_URL={_ORCH_URL}", runs[0])
-61
View File
@@ -323,67 +323,6 @@ class TestOrchestratorSupervise(unittest.TestCase):
self.assertFalse(ok)
self.assertIn("no longer registered", err)
# --- agent half: queue + poll (issue #469) -----------------------------
def test_queue_proposal_then_poll_pending(self) -> None:
bottle_id = self._register("demo", "routes: []\n")
pid = self.orch.supervise_queue_proposal(
bottle_id, tool=TOOL_EGRESS_ALLOW,
proposed_file="routes:\n - host: google.com\n", justification="need it")
# Visible to the operator, keyed by the bottle id.
pending = self.orch.supervise_pending()
self.assertEqual([pid], [p["id"] for p in pending])
self.assertEqual(
{"status": "pending"}, self.orch.supervise_poll_response(bottle_id, pid))
def test_poll_is_idempotent_and_leaves_no_pending(self) -> None:
bottle_id = self._register("demo", "routes: []\n")
pid = self.orch.supervise_queue_proposal(
bottle_id, tool=TOOL_EGRESS_ALLOW,
proposed_file="routes:\n - host: google.com\n", justification="need it")
self.orch.supervise_respond(
pid, bottle_slug=bottle_id, decision="approve", notes="ok")
decided = self.orch.supervise_poll_response(bottle_id, pid)
self.assertEqual("approved", decided["status"])
self.assertEqual("ok", decided["notes"])
# Decided proposal drops off the operator's pending list (a response row
# exists), but poll does NOT archive — a re-poll returns the same
# decision so a dropped connection can't lose it (issue #469 review).
self.assertEqual([], self.orch.supervise_pending())
self.assertEqual(decided, self.orch.supervise_poll_response(bottle_id, pid))
def test_teardown_reaps_the_bottles_proposals(self) -> None:
rec = self.store.register("10.243.0.20", policy="routes: []\n")
pid = self.orch.supervise_queue_proposal(
rec.bottle_id, tool=TOOL_EGRESS_ALLOW,
proposed_file="routes:\n - host: google.com\n", justification="j")
self.orch.supervise_respond(
pid, bottle_slug=rec.bottle_id, decision="approve", notes="ok")
self.assertTrue(self.orch.teardown_bottle(rec.bottle_id))
# The gone bottle's decided-but-unconsumed proposal is archived, so a
# late poll returns 'unknown' rather than lingering forever.
self.assertEqual(
{"status": "unknown"}, self.orch.supervise_poll_response(rec.bottle_id, pid))
def test_reconcile_reaps_the_bottles_proposals(self) -> None:
rec = self.store.register("10.243.0.21", policy="routes: []\n")
pid = self.orch.supervise_queue_proposal(
rec.bottle_id, tool=TOOL_EGRESS_ALLOW,
proposed_file="routes:\n - host: google.com\n", justification="j")
# No live source IPs -> the bottle is reaped (grace 0 so it's immediate).
self.assertEqual([rec.bottle_id], self.orch.reconcile([], grace_seconds=0))
self.assertEqual(
{"status": "unknown"}, self.orch.supervise_poll_response(rec.bottle_id, pid))
def test_poll_unknown_for_other_bottle(self) -> None:
bottle_id = self._register("demo", "routes: []\n")
pid = self.orch.supervise_queue_proposal(
bottle_id, tool=TOOL_EGRESS_ALLOW,
proposed_file="routes:\n - host: google.com\n", justification="j")
# A different bottle id can't read demo's proposal (scoped by queue key).
self.assertEqual(
{"status": "unknown"}, self.orch.supervise_poll_response("other-bottle", pid))
if __name__ == "__main__":
unittest.main()
+1 -72
View File
@@ -7,13 +7,7 @@ import unittest
import urllib.error
from unittest.mock import MagicMock, patch
from bot_bottle.policy_resolver import (
CONTROL_AUTH_HEADER,
CONTROL_AUTH_JWT_ENV,
PolicyResolveError,
PolicyResolver,
_control_auth_headers,
)
from bot_bottle.policy_resolver import PolicyResolveError, PolicyResolver
_URLOPEN = "bot_bottle.policy_resolver.urllib.request.urlopen"
@@ -29,18 +23,6 @@ def _http_error(code: int) -> urllib.error.HTTPError:
return urllib.error.HTTPError("http://x/resolve", code, "err", {}, None) # type: ignore[arg-type]
class TestControlAuthHeaders(unittest.TestCase):
def test_sends_the_gateway_jwt_when_configured(self) -> None:
with patch.dict("os.environ", {CONTROL_AUTH_JWT_ENV: "gateway.jwt.tok"}):
self.assertEqual({CONTROL_AUTH_HEADER: "gateway.jwt.tok"}, _control_auth_headers())
def test_sends_nothing_when_unset(self) -> None:
import os
with patch.dict("os.environ", {}, clear=False):
os.environ.pop(CONTROL_AUTH_JWT_ENV, None)
self.assertEqual({}, _control_auth_headers())
class TestPolicyResolver(unittest.TestCase):
def setUp(self) -> None:
self.r = PolicyResolver("http://orch:8080")
@@ -124,59 +106,6 @@ class TestPolicyResolver(unittest.TestCase):
with self.assertRaises(PolicyResolveError):
self.r.resolve_policy_and_bottle_id("10.243.0.1")
# --- supervise agent RPCs (issue #469) ---------------------------------
def test_propose_supervise_returns_id_and_posts_payload(self) -> None:
with patch(_URLOPEN, return_value=_resp({"proposal_id": "p-7"})) as m:
pid = self.r.propose_supervise(
"10.243.0.7", "the-token",
tool="egress-allow", proposed_file="routes:\n", justification="j",
)
self.assertEqual("p-7", pid)
req = m.call_args.args[0]
self.assertTrue(req.full_url.endswith("/supervise/propose"))
sent = json.loads(req.data)
self.assertEqual("10.243.0.7", sent["source_ip"])
self.assertEqual("the-token", sent["identity_token"])
self.assertEqual("egress-allow", sent["tool"])
self.assertEqual("routes:\n", sent["proposed_file"])
def test_propose_supervise_unattributed_is_none(self) -> None:
with patch(_URLOPEN, side_effect=_http_error(403)):
self.assertIsNone(self.r.propose_supervise(
"10.9.9.9", "t", tool="egress-allow", proposed_file="x", justification="j"))
def test_propose_supervise_missing_id_is_none(self) -> None:
with patch(_URLOPEN, return_value=_resp({})):
self.assertIsNone(self.r.propose_supervise(
"10.243.0.1", "t", tool="egress-allow", proposed_file="x", justification="j"))
def test_propose_supervise_unreachable_raises(self) -> None:
with patch(_URLOPEN, side_effect=urllib.error.URLError("refused")):
with self.assertRaises(PolicyResolveError):
self.r.propose_supervise(
"10.243.0.1", "t", tool="egress-allow", proposed_file="x", justification="j")
def test_poll_supervise_returns_status(self) -> None:
with patch(_URLOPEN, return_value=_resp(
{"status": "approved", "notes": "ok", "final_file": None})
) as m:
result = self.r.poll_supervise("10.243.0.7", "tok", "p-7")
assert result is not None
self.assertEqual("approved", result["status"])
req = m.call_args.args[0]
self.assertTrue(req.full_url.endswith("/supervise/poll"))
self.assertEqual("p-7", json.loads(req.data)["proposal_id"])
def test_poll_supervise_unattributed_is_none(self) -> None:
with patch(_URLOPEN, side_effect=_http_error(403)):
self.assertIsNone(self.r.poll_supervise("10.9.9.9", "t", "p-7"))
def test_poll_supervise_unreachable_raises(self) -> None:
with patch(_URLOPEN, side_effect=urllib.error.URLError("refused")):
with self.assertRaises(PolicyResolveError):
self.r.poll_supervise("10.243.0.1", "t", "p-7")
if __name__ == "__main__":
unittest.main()
+195 -241
View File
@@ -1,11 +1,4 @@
"""Unit: supervise daemon MCP server (PRD 0013, PRD 0070).
The daemon no longer opens bot-bottle.db: it queues proposals and polls for
their responses over the control-plane RPC (issue #469). These tests drive the
handlers with `_FakeSuperviseResolver`, an in-process stand-in for
`PolicyResolver.propose_supervise` / `poll_supervise` backed by the real queue
store so the operator-response and archive contracts are still exercised
end-to-end, just through the RPC seam instead of a direct file handle."""
"""Unit: supervise daemon MCP server (PRD 0013)."""
import http.client
import json
@@ -15,6 +8,7 @@ import time
import types
import unittest
from pathlib import Path
from unittest.mock import patch
from tests.unit import use_bottle_root
@@ -50,95 +44,6 @@ from bot_bottle.supervise_server import (
validate_proposed_file,
)
# Fixed caller identity for the handler tests. The control plane attributes by
# (source_ip, identity_token); the fake resolver ignores them and answers for a
# fixed bottle, since attribution itself is covered by the orchestrator tests.
_SRC = "10.0.0.7"
_TOK = "tok"
class _FakeSuperviseResolver:
"""Stand-in for `PolicyResolver`'s supervise RPCs, backed by the real queue
store (as the orchestrator is). `bottle_id=None` models an unattributed
caller (a clean 403 None); `raises=True` models an unreachable
orchestrator (`PolicyResolveError`)."""
def __init__(
self, bottle_id: str | None = "dev", raises: bool = False, policy: str = "",
poll_raises: bool = False, poll_none: bool = False,
) -> None:
self.bottle_id = bottle_id
self.raises = raises
self._policy = policy
# propose succeeds, but the later poll fails — models an orchestrator
# that becomes unreachable (poll_raises) or unattributes mid-flight
# (poll_none) after the proposal was queued.
self.poll_raises = poll_raises
self.poll_none = poll_none
def propose_supervise(
self, source_ip: str, identity_token: str, *,
tool: str, proposed_file: str, justification: str,
) -> str | None:
del source_ip, identity_token
if self.raises:
raise supervise_server.PolicyResolveError("orchestrator down")
if self.bottle_id is None:
return None
proposal = _sv.Proposal.new(
bottle_slug=self.bottle_id, tool=tool, proposed_file=proposed_file,
justification=justification, current_file_hash=_sv.sha256_hex(proposed_file),
)
_sv.write_proposal(proposal)
return proposal.id
def poll_supervise(
self, source_ip: str, identity_token: str, proposal_id: str,
) -> dict[str, object] | None:
del source_ip, identity_token
if self.raises or self.poll_raises:
raise supervise_server.PolicyResolveError("orchestrator down")
if self.bottle_id is None or self.poll_none:
return None
try:
response = _sv.read_response(self.bottle_id, proposal_id)
except FileNotFoundError:
try:
_sv.read_proposal(self.bottle_id, proposal_id)
except FileNotFoundError:
return {"status": _sv.POLL_STATUS_UNKNOWN}
return {"status": _sv.POLL_STATUS_PENDING}
# Idempotent: poll never archives (issue #469 review).
return {
"status": response.status, "notes": response.notes,
"final_file": response.final_file,
}
# Used by list-egress-routes (`_resolved_routes_payload`), unchanged path.
def resolve_policy_and_bottle_id(
self, source_ip: str, identity_token: str = "",
) -> tuple[str | None, str | None, dict[str, str]]:
del source_ip, identity_token
if self.raises:
raise supervise_server.PolicyResolveError("orchestrator down")
return self._policy, self.bottle_id, {}
def _tools_call(
resolver: object, params: dict[str, object],
config: "ServerConfig | None" = None,
) -> dict[str, object]:
return handle_tools_call(
params, config or ServerConfig(),
resolver=resolver, source_ip=_SRC, identity_token=_TOK, # type: ignore[arg-type]
)
def _check(resolver: object, params: dict[str, object]) -> dict[str, object]:
return handle_check_proposal(
params, resolver=resolver, source_ip=_SRC, identity_token=_TOK, # type: ignore[arg-type]
)
# --- Validation ------------------------------------------------------------
@@ -206,35 +111,32 @@ class TestRpcErrorTaxonomy(unittest.TestCase):
validate_proposed_file(_sv.TOOL_EGRESS_ALLOW, "routes: nope\n")
def test_unknown_tool_in_tools_call_is_client_error(self):
config = ServerConfig(bottle_slug="dev")
with self.assertRaises(_RpcClientError) as cm:
_tools_call(_FakeSuperviseResolver(), {"name": "no-such-tool", "arguments": {}})
handle_tools_call({"name": "no-such-tool", "arguments": {}}, config)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
class TestRpcInternalErrorOnRpcFailure(unittest.TestCase):
"""A queue RPC that can't reach the orchestrator (or returns unattributed)
surfaces as ERR_INTERNAL the daemon fails closed rather than leaking the
cause to the agent."""
_ARGS: dict[str, object] = {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "x",
},
}
def test_unreachable_orchestrator_raises_internal(self):
with self.assertRaises(_RpcInternalError) as cm:
_tools_call(_FakeSuperviseResolver(raises=True), self._ARGS)
class TestRpcInternalErrorOnIoFailure(unittest.TestCase):
def test_write_proposal_os_error_raises_internal(self):
config = ServerConfig(
bottle_slug="dev",
)
with patch.object(_sv, "write_proposal", side_effect=OSError("disk full")), \
self.assertRaises(_RpcInternalError) as cm:
handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "x",
},
},
config,
)
self.assertEqual(ERR_INTERNAL, cm.exception.code)
self.assertIsNotNone(cm.exception.__cause__)
def test_unattributed_source_raises_internal(self):
with self.assertRaises(_RpcInternalError) as cm:
_tools_call(_FakeSuperviseResolver(bottle_id=None), self._ARGS)
self.assertEqual(ERR_INTERNAL, cm.exception.code)
# --- JSON-RPC parsing ------------------------------------------------------
@@ -363,8 +265,8 @@ class TestHandleToolsList(unittest.TestCase):
class TestHandleToolsCall(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory(prefix="supervise-server-test.")
self._home_patch = use_bottle_root(Path(self._tmp.name) / ".bot-bottle")
self.resolver = _FakeSuperviseResolver("dev")
self._home_patch = self._patch_home(Path(self._tmp.name))
self.config = ServerConfig(bottle_slug="dev")
_qs.QueueStore("dev").migrate()
_as.AuditStore().migrate()
@@ -372,9 +274,12 @@ class TestHandleToolsCall(unittest.TestCase):
self._home_patch()
self._tmp.cleanup()
def _patch_home(self, fake_home: Path):
return use_bottle_root(fake_home / ".bot-bottle")
def _respond_when_proposal_appears(self, status: str, notes: str = "") -> threading.Thread:
"""Background thread: poll the queue for a fresh proposal, write a
matching response the operator half, out of band."""
matching response. Returns the thread so the test can join it."""
def runner():
for _ in range(200):
pending = _sv.list_pending_proposals("dev")
@@ -393,13 +298,16 @@ class TestHandleToolsCall(unittest.TestCase):
def test_call_round_trips_through_queue(self):
responder = self._respond_when_proposal_appears(_sv.STATUS_APPROVED, notes="lgtm")
try:
result = _tools_call(self.resolver, {
"name": _sv.TOOL_EGRESS_BLOCK,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "need example.com",
result = handle_tools_call(
{
"name": _sv.TOOL_EGRESS_BLOCK,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "need example.com",
},
},
})
self.config,
)
finally:
responder.join()
self.assertFalse(result["isError"]) # type: ignore[index]
@@ -410,13 +318,16 @@ class TestHandleToolsCall(unittest.TestCase):
def test_allow_round_trips_through_queue(self):
responder = self._respond_when_proposal_appears(_sv.STATUS_APPROVED, notes="ok")
try:
result = _tools_call(self.resolver, {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "need example.com",
result = handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "need example.com",
},
},
})
self.config,
)
finally:
responder.join()
self.assertFalse(result["isError"]) # type: ignore[index]
@@ -427,71 +338,94 @@ class TestHandleToolsCall(unittest.TestCase):
def test_rejected_response_sets_isError(self):
responder = self._respond_when_proposal_appears(_sv.STATUS_REJECTED, notes="nope")
try:
result = _tools_call(self.resolver, {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "needed for tests",
result = handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "needed for tests",
},
},
})
self.config,
)
finally:
responder.join()
self.assertTrue(result["isError"]) # type: ignore[index]
def test_invalid_tool_name_raises(self):
with self.assertRaises(_RpcError) as cm:
_tools_call(self.resolver, {"name": "not-a-tool", "arguments": {}})
handle_tools_call(
{"name": "not-a-tool", "arguments": {}},
self.config,
)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
def test_missing_justification_raises(self):
with self.assertRaises(_RpcError):
_tools_call(self.resolver, {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {"routes_yaml": "routes:\n - host: example.com\n"},
})
handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {"routes_yaml": "routes:\n - host: example.com\n"},
},
self.config,
)
def test_missing_name_raises(self):
with self.assertRaises(_RpcError) as cm:
_tools_call(self.resolver, {"arguments": {}})
handle_tools_call({"arguments": {}}, self.config)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
def test_arguments_must_be_object(self):
with self.assertRaises(_RpcError) as cm:
_tools_call(self.resolver, {"name": _sv.TOOL_EGRESS_ALLOW, "arguments": []})
handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": [],
},
self.config,
)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
self.assertIn("must be an object", cm.exception.message)
def test_capability_block_call_raises_unknown_tool(self):
with self.assertRaises(_RpcError) as cm:
_tools_call(self.resolver, {
"name": "capability-block",
"arguments": {
"dockerfile": "FROM python:3.13\n",
"justification": "need git",
handle_tools_call(
{
"name": "capability-block",
"arguments": {
"dockerfile": "FROM python:3.13\n",
"justification": "need git",
},
},
})
self.config,
)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
self.assertIn("unknown tool", cm.exception.message)
def test_decided_proposal_drops_off_pending(self):
def test_archives_proposal_after_response(self):
responder = self._respond_when_proposal_appears(_sv.STATUS_APPROVED)
try:
_tools_call(self.resolver, {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "x",
handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "x",
},
},
})
self.config,
)
finally:
responder.join()
# A decided proposal drops off the operator's pending list (a response
# row exists) — poll itself no longer archives (issue #469 review).
# No pending proposals left after archive.
self.assertEqual([], _sv.list_pending_proposals("dev"))
def test_pending_response_times_out_without_archive(self):
result = _tools_call(
self.resolver,
config = ServerConfig(
bottle_slug="dev",
response_timeout_seconds=0.05,
)
result = handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
@@ -499,33 +433,15 @@ class TestHandleToolsCall(unittest.TestCase):
"justification": "need egress",
},
},
ServerConfig(response_timeout_seconds=0.05),
config,
)
self.assertFalse(result["isError"]) # type: ignore[index]
text = result["content"][0]["text"] # type: ignore[index]
self.assertIn("status: pending", text)
self.assertIn("proposal remains queued", text)
self.assertEqual(1, len(_sv.list_pending_proposals("dev")))
_ALLOW: dict[str, object] = {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {"routes_yaml": "routes:\n - host: example.com\n", "justification": "x"},
}
def test_poll_unreachable_after_queue_raises_internal(self):
# Proposal queues, then the orchestrator becomes unreachable on poll.
self.resolver.poll_raises = True
with self.assertRaises(_RpcInternalError) as cm:
_tools_call(self.resolver, self._ALLOW, ServerConfig(response_timeout_seconds=5))
self.assertEqual(ERR_INTERNAL, cm.exception.code)
def test_poll_unattributed_after_queue_raises_internal(self):
# Proposal queues, then poll comes back unattributed (fail-closed).
self.resolver.poll_none = True
with self.assertRaises(_RpcInternalError) as cm:
_tools_call(self.resolver, self._ALLOW, ServerConfig(response_timeout_seconds=5))
self.assertEqual(ERR_INTERNAL, cm.exception.code)
class TestResponseTimeoutEnv(unittest.TestCase):
def test_unset_uses_default(self):
@@ -594,7 +510,7 @@ class TestHttpEndToEnd(unittest.TestCase):
self.port = s.getsockname()[1]
s.close()
self.server = MCPServer(("127.0.0.1", self.port), MCPHandler)
self.server.config = ServerConfig()
self.server.config = ServerConfig(bottle_slug="dev")
self.thread = threading.Thread(
target=self.server.serve_forever, daemon=True,
)
@@ -636,21 +552,23 @@ class TestHttpEndToEnd(unittest.TestCase):
)
self.assertEqual(ERR_METHOD_NOT_FOUND, result["error"]["code"]) # type: ignore[index]
def test_no_resolver_fails_closed_over_http(self):
# The test server has no policy_resolver wired, so a proposal tools/call
# fails closed with ERR_INTERNAL rather than queuing anything.
result = self._post_jsonrpc({
"jsonrpc": "2.0",
"id": 99,
"method": "tools/call",
"params": {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "x",
def test_internal_error_returns_err_internal_over_http(self):
with patch.object(
supervise_server._sv, "write_proposal",
side_effect=OSError("disk full"),
):
result = self._post_jsonrpc({
"jsonrpc": "2.0",
"id": 99,
"method": "tools/call",
"params": {
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {
"routes_yaml": "routes:\n - host: example.com\n",
"justification": "x",
},
},
},
})
})
self.assertIn("error", result)
self.assertEqual(ERR_INTERNAL, result["error"]["code"]) # type: ignore[index]
@@ -665,23 +583,72 @@ class TestHttpEndToEnd(unittest.TestCase):
conn.close()
class _FakeResolver:
def __init__(
self,
bottle_id: str | None = None,
raises: bool = False,
policy: str = "",
) -> None:
self._bottle_id = bottle_id
self._raises = raises
self._policy = policy
self.calls: list[str] = []
def resolve_bottle_id(self, source_ip: str, identity_token: str = "") -> str | None:
del identity_token
self.calls.append(source_ip)
if self._raises:
# Raise the exact class supervise_server catches (it imports
# policy_resolver flat inside the bundle, package-side in tests).
raise supervise_server.PolicyResolveError("orchestrator down")
return self._bottle_id
def resolve_policy_and_bottle_id(
self, source_ip: str, identity_token: str = "",
) -> "tuple[str, str | None, dict[str, str]]":
del identity_token
self.calls.append(source_ip)
if self._raises:
raise supervise_server.PolicyResolveError("orchestrator down")
return self._policy, self._bottle_id, {}
def _handler(resolver: object) -> MCPHandler:
"""A bare MCPHandler wired with a server (carrying the resolver) and a
client address, enough to exercise the resolver-backed paths off-socket."""
client address, enough to exercise `_attributed_config` off-socket."""
h: MCPHandler = MCPHandler.__new__(MCPHandler)
h.server = types.SimpleNamespace(policy_resolver=resolver) # type: ignore[assignment]
h.client_address = ("10.0.0.7", 4321)
h.headers = {} # type: ignore[assignment]
return h
class TestResolverFailClosed(unittest.TestCase):
"""A dispatch without a resolver is a misconfiguration, not a tenancy mode
fail closed rather than queue (or list) anything (PRD 0070)."""
class TestAttributedConfig(unittest.TestCase):
"""Each proposal is attributed to the calling bottle by source IP (PRD
0070); a server without a resolver fails closed rather than queuing under an
unattributed slug."""
def test_missing_resolver_raises(self) -> None:
def test_missing_resolver_fails_closed(self) -> None:
with self.assertRaises(_RpcInternalError):
_handler(None)._resolver_or_fail()
_handler(None)._attributed_config(ServerConfig(bottle_slug="dev"))
def test_consolidated_binds_source_ip_bottle(self) -> None:
r = _FakeResolver(bottle_id="bottle-x")
cfg = _handler(r)._attributed_config(ServerConfig(bottle_slug="ignored"))
self.assertEqual("bottle-x", cfg.bottle_slug) # resolved slug wins
self.assertEqual(["10.0.0.7"], r.calls)
def test_unattributed_source_fails_closed(self) -> None:
with self.assertRaises(_RpcInternalError):
_handler(_FakeResolver(bottle_id=None))._attributed_config(
ServerConfig(bottle_slug="x")
)
def test_resolver_error_fails_closed(self) -> None:
with self.assertRaises(_RpcInternalError):
_handler(_FakeResolver(raises=True))._attributed_config(
ServerConfig(bottle_slug="x")
)
class TestResolvedRoutesPayload(unittest.TestCase):
@@ -697,7 +664,7 @@ class TestResolvedRoutesPayload(unittest.TestCase):
" - host: www.google.com\n"
)
payload = _handler(
_FakeSuperviseResolver(bottle_id="b1", policy=policy)
_FakeResolver(bottle_id="b1", policy=policy)
)._resolved_routes_payload()
assert payload is not None
self.assertFalse(payload["isError"]) # type: ignore[index]
@@ -709,7 +676,7 @@ class TestResolvedRoutesPayload(unittest.TestCase):
# resolve_client_context swallows resolver errors → deny-all (empty),
# never another bottle's routes.
payload = _handler(
_FakeSuperviseResolver(raises=True)
_FakeResolver(raises=True)
)._resolved_routes_payload()
assert payload is not None
data = json.loads(payload["content"][0]["text"]) # type: ignore[index]
@@ -731,7 +698,7 @@ class TestNonBlockingSupervise(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory(prefix="supervise-nonblock-test.")
self._home_patch = use_bottle_root(Path(self._tmp.name) / ".bot-bottle")
self.resolver = _FakeSuperviseResolver("dev")
self.config = ServerConfig(bottle_slug="dev")
_qs.QueueStore("dev").migrate()
_as.AuditStore().migrate()
@@ -750,6 +717,9 @@ class TestNonBlockingSupervise(unittest.TestCase):
_sv.write_proposal(p)
return p
def _check(self, proposal_id: str) -> dict[str, object]:
return handle_check_proposal({"arguments": {"proposal_id": proposal_id}}, self.config)
# --- pending response carries the id ---
def test_pending_text_includes_id_and_pointer(self):
@@ -760,13 +730,12 @@ class TestNonBlockingSupervise(unittest.TestCase):
def test_tools_call_timeout_returns_pending_with_id_and_stays_queued(self):
# No responder → the grace window expires → pending, not blocked forever.
result = _tools_call(
self.resolver,
result = handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {"routes_yaml": self._ROUTES, "justification": "x"},
},
ServerConfig(response_timeout_seconds=0.05),
ServerConfig(bottle_slug="dev", response_timeout_seconds=0.05),
)
self.assertFalse(result["isError"]) # type: ignore[index]
text = result["content"][0]["text"] # type: ignore[index]
@@ -777,29 +746,27 @@ class TestNonBlockingSupervise(unittest.TestCase):
# --- check-proposal poll ---
def test_check_returns_approved_idempotently(self):
def test_check_returns_approved_and_archives(self):
p = self._seed_proposal()
_sv.write_response("dev", _sv.Response(proposal_id=p.id, status=_sv.STATUS_APPROVED, notes="ok"))
result = _check(self.resolver, {"arguments": {"proposal_id": p.id}})
result = self._check(p.id)
self.assertFalse(result["isError"])
text = result["content"][0]["text"] # type: ignore[index]
self.assertIn("status: approved", text)
self.assertIn("notes: ok", text)
# Poll doesn't archive, so a re-check returns the same decision
# (issue #469 review) rather than 'unknown'.
again = _check(self.resolver, {"arguments": {"proposal_id": p.id}})
self.assertEqual(result, again)
with self.assertRaises(FileNotFoundError): # archived on read
_sv.read_proposal("dev", p.id)
def test_check_rejected_sets_isError(self):
p = self._seed_proposal()
_sv.write_response("dev", _sv.Response(proposal_id=p.id, status=_sv.STATUS_REJECTED, notes="no"))
result = _check(self.resolver, {"arguments": {"proposal_id": p.id}})
result = self._check(p.id)
self.assertTrue(result["isError"])
self.assertIn("status: rejected", result["content"][0]["text"]) # type: ignore[index]
def test_check_pending_when_no_decision_yet(self):
p = self._seed_proposal()
result = _check(self.resolver, {"arguments": {"proposal_id": p.id}})
result = self._check(p.id)
self.assertFalse(result["isError"])
text = result["content"][0]["text"] # type: ignore[index]
self.assertIn("status: pending", text)
@@ -807,56 +774,43 @@ class TestNonBlockingSupervise(unittest.TestCase):
self.assertEqual(1, len(_sv.list_pending_proposals("dev"))) # not archived
def test_check_unknown_id_is_error(self):
result = _check(self.resolver, {"arguments": {"proposal_id": "no-such-proposal"}})
result = self._check("no-such-proposal")
self.assertTrue(result["isError"])
self.assertIn("status: unknown", result["content"][0]["text"]) # type: ignore[index]
def test_check_poll_unreachable_raises_internal(self):
self.resolver.poll_raises = True
with self.assertRaises(_RpcInternalError) as cm:
_check(self.resolver, {"arguments": {"proposal_id": "p"}})
self.assertEqual(ERR_INTERNAL, cm.exception.code)
def test_check_poll_unattributed_raises_internal(self):
self.resolver.poll_none = True
with self.assertRaises(_RpcInternalError) as cm:
_check(self.resolver, {"arguments": {"proposal_id": "p"}})
self.assertEqual(ERR_INTERNAL, cm.exception.code)
def test_check_missing_id_raises(self):
with self.assertRaises(_RpcClientError) as cm:
_check(self.resolver, {"arguments": {}})
handle_check_proposal({"arguments": {}}, self.config)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
def test_check_empty_id_raises(self):
with self.assertRaises(_RpcClientError) as cm:
_check(self.resolver, {"arguments": {"proposal_id": " "}})
handle_check_proposal({"arguments": {"proposal_id": " "}}, self.config)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
def test_check_arguments_must_be_object(self):
with self.assertRaises(_RpcClientError) as cm:
_check(self.resolver, {"arguments": []})
handle_check_proposal({"arguments": []}, self.config)
self.assertEqual(ERR_INVALID_PARAMS, cm.exception.code)
def test_full_nonblocking_round_trip(self):
# 1. tools/call times out → pending with id
result = _tools_call(
self.resolver,
result = handle_tools_call(
{
"name": _sv.TOOL_EGRESS_ALLOW,
"arguments": {"routes_yaml": self._ROUTES, "justification": "x"},
},
ServerConfig(response_timeout_seconds=0.05),
ServerConfig(bottle_slug="dev", response_timeout_seconds=0.05),
)
pid = _sv.list_pending_proposals("dev")[0].id
self.assertIn(pid, result["content"][0]["text"]) # type: ignore[index]
# 2. operator decides out-of-band
_sv.write_response("dev", _sv.Response(proposal_id=pid, status=_sv.STATUS_APPROVED, notes="ok"))
# 3. agent resumes by polling — no re-proposing
poll = _check(self.resolver, {"arguments": {"proposal_id": pid}})
poll = self._check(pid)
self.assertFalse(poll["isError"])
self.assertIn("status: approved", poll["content"][0]["text"]) # type: ignore[index]
self.assertEqual([], _sv.list_pending_proposals("dev")) # resolved (response exists)
self.assertEqual([], _sv.list_pending_proposals("dev")) # resolved + archived
if __name__ == "__main__":