Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 01cee056be |
+11
-24
@@ -1,5 +1,4 @@
|
||||
# Run the project's test suite when package or runtime inputs change on a PR
|
||||
# or on push to main.
|
||||
# Run the project's test suite on every PR push and on push to main.
|
||||
#
|
||||
# The suite uses stdlib `unittest` discovery — no external Python
|
||||
# dependencies are required to execute it. Tests are split by directory:
|
||||
@@ -24,34 +23,22 @@ on:
|
||||
branches:
|
||||
- main
|
||||
paths:
|
||||
- 'bot_bottle/**'
|
||||
- 'tests/**/*.py'
|
||||
- 'cli.py'
|
||||
- 'scripts/coverage.sh'
|
||||
- 'scripts/critical-modules.txt'
|
||||
- 'scripts/diff_coverage.py'
|
||||
- 'scripts/tracker_policy.py'
|
||||
- 'scripts/firecracker-netpool.sh'
|
||||
- '**.py'
|
||||
- '.gitea/workflows/**.yml'
|
||||
- 'scripts/**'
|
||||
- 'README.md'
|
||||
# Dockerfiles and pyproject.toml are baked into the infra rootfs; a
|
||||
# change here alters what the integration/coverage jobs build locally.
|
||||
- 'Dockerfile*'
|
||||
- 'pyproject.toml'
|
||||
- 'requirements-dev.txt'
|
||||
- '.coveragerc'
|
||||
- '.dockerignore'
|
||||
pull_request:
|
||||
paths:
|
||||
- 'bot_bottle/**'
|
||||
- 'tests/**/*.py'
|
||||
- 'cli.py'
|
||||
- 'scripts/coverage.sh'
|
||||
- 'scripts/critical-modules.txt'
|
||||
- 'scripts/diff_coverage.py'
|
||||
- 'scripts/tracker_policy.py'
|
||||
- 'scripts/firecracker-netpool.sh'
|
||||
- '**.py'
|
||||
- '.gitea/workflows/**.yml'
|
||||
- 'scripts/**'
|
||||
- 'README.md'
|
||||
- 'Dockerfile*'
|
||||
- 'pyproject.toml'
|
||||
- 'requirements-dev.txt'
|
||||
- '.coveragerc'
|
||||
- '.dockerignore'
|
||||
workflow_dispatch:
|
||||
|
||||
jobs:
|
||||
|
||||
@@ -21,7 +21,7 @@ jobs:
|
||||
- uses: actions/checkout@v3
|
||||
with:
|
||||
fetch-depth: 0
|
||||
token: ${{ secrets.BADGE_PUSH_TOKEN }}
|
||||
token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
# No actions/setup-python: the runner image ships Python 3.12 and older
|
||||
# act_runner engines mishandle setup-python's PATH. Install into the
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
# bot-bottle
|
||||
|
||||
[](https://gitea.dideric.is/didericis/bot-bottle/actions?workflow=test.yml)
|
||||
[](https://coverage.readthedocs.io/)
|
||||
[](https://coverage.readthedocs.io/)
|
||||
[](https://gitea.dideric.is/didericis/bot-bottle/src/branch/main/docs/decisions/0004-coverage-policy.md)
|
||||
|
||||
**Problem:** Developer wants to run a coding agent without supervision, but they don't want a prompt injected or misbehaving agent wrecking their environment or exfiltrating sensitive data.
|
||||
|
||||
@@ -26,7 +26,6 @@ def provision_bottle(
|
||||
*,
|
||||
image_ref: str = "",
|
||||
tokens: dict[str, str] | None = None,
|
||||
env_var_secret: str | None = None,
|
||||
) -> RegisteredBottle:
|
||||
"""Register the bottle and provision its git-gate state. Rolls back the
|
||||
registration if provisioning fails so no orphan is left.
|
||||
@@ -36,7 +35,7 @@ def provision_bottle(
|
||||
``RegisteredBottle`` so callers can inject it into the agent container's
|
||||
environment."""
|
||||
inputs = registration_inputs(egress_plan)
|
||||
env_var_secret = env_var_secret or new_env_var_secret()
|
||||
env_var_secret = new_env_var_secret()
|
||||
reg = client.register_bottle(
|
||||
source_ip, image_ref=image_ref, policy=inputs.policy,
|
||||
metadata=inputs.metadata, tokens=tokens, env_var_secret=env_var_secret,
|
||||
|
||||
@@ -23,7 +23,6 @@ from ...orchestrator.client import OrchestratorClient
|
||||
from ...orchestrator.gateway import GATEWAY_NETWORK
|
||||
from ...orchestrator.lifecycle import INFRA_NAME, OrchestratorService
|
||||
from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME
|
||||
from ...orchestrator.reprovision import reprovision_bottles
|
||||
from ..consolidated_util import provision_bottle
|
||||
from ..consolidated_util import teardown_consolidated as _teardown_util
|
||||
from .gateway_provision import DockerGatewayTransport
|
||||
@@ -103,17 +102,17 @@ def _reprovision_running_bottles(
|
||||
when the orchestrator already has all tokens loaded. Best-effort: a single
|
||||
container exec failure never blocks a new bottle launch."""
|
||||
client = OrchestratorClient(orchestrator_url)
|
||||
bottles = client.list_bottles()
|
||||
if not bottles:
|
||||
return
|
||||
|
||||
# Build {source_ip: container_name} from live containers on the gateway
|
||||
# network, excluding the infra container itself.
|
||||
try:
|
||||
proc = run_docker([
|
||||
"docker", "network", "inspect",
|
||||
"--format", "{{range .Containers}}{{.Name}} {{.IPv4Address}}\n{{end}}",
|
||||
network,
|
||||
])
|
||||
except OSError as exc:
|
||||
log.info(f"egress token reprovision skipped: {exc}")
|
||||
return
|
||||
proc = run_docker([
|
||||
"docker", "network", "inspect",
|
||||
"--format", "{{range .Containers}}{{.Name}} {{.IPv4Address}}\n{{end}}",
|
||||
network,
|
||||
])
|
||||
ip_to_container: dict[str, str] = {}
|
||||
for line in proc.stdout.splitlines():
|
||||
parts = line.strip().split()
|
||||
@@ -122,15 +121,26 @@ def _reprovision_running_bottles(
|
||||
if ip:
|
||||
ip_to_container[ip] = parts[0]
|
||||
|
||||
secrets_by_ip: dict[str, str] = {}
|
||||
for source_ip, container_name in ip_to_container.items():
|
||||
reprovisioned = 0
|
||||
for bottle in bottles:
|
||||
bottle_id = bottle.get("bottle_id")
|
||||
source_ip = bottle.get("source_ip")
|
||||
if not isinstance(bottle_id, str) or not isinstance(source_ip, str):
|
||||
continue
|
||||
container_name = ip_to_container.get(source_ip)
|
||||
if not container_name:
|
||||
continue
|
||||
proc = run_docker(
|
||||
["docker", "exec", container_name, "printenv", ENV_VAR_SECRET_NAME]
|
||||
)
|
||||
if proc.returncode == 0 and proc.stdout.strip():
|
||||
secrets_by_ip[source_ip] = proc.stdout.strip()
|
||||
if proc.returncode != 0 or not proc.stdout.strip():
|
||||
continue
|
||||
try:
|
||||
if client.reprovision_gateway(bottle_id, proc.stdout.strip()):
|
||||
reprovisioned += 1
|
||||
except Exception: # noqa: BLE001 — best-effort, never block a launch
|
||||
pass
|
||||
|
||||
reprovisioned = reprovision_bottles(client, secrets_by_ip)
|
||||
if reprovisioned:
|
||||
log.info(
|
||||
"reprovisioned egress tokens",
|
||||
|
||||
@@ -22,9 +22,6 @@ class FirecrackerBottlePlan(BottlePlan):
|
||||
# (egress proxy credentials, git-gate/supervise headers); set by launch
|
||||
# from the orchestrator registration. Empty pre-registration.
|
||||
identity_token: str = ""
|
||||
# Applied to every agent SSH exec and mirrored into /run inside the VM so
|
||||
# the host can recover it after the infra VM restarts.
|
||||
env_var_secret: str = ""
|
||||
|
||||
@property
|
||||
def container_name(self) -> str:
|
||||
|
||||
@@ -1,23 +1,8 @@
|
||||
"""Cleanup for the Firecracker backend.
|
||||
|
||||
Reaps *orphans* only — resources with no live VM behind them:
|
||||
|
||||
* orphan run dirs: a per-bottle run dir (holding the ~1G rootfs.ext4)
|
||||
whose firecracker process has exited. These leak when a launch is
|
||||
hard-killed before its teardown runs (host OOM/crash, a cancelled CI
|
||||
job, `kill -9`); the clean-exit path already removes its own dir in
|
||||
launch.py.
|
||||
* orphan VM pids: a firecracker process whose run dir is already gone
|
||||
— a VMM left lingering after its dir was removed.
|
||||
|
||||
A run dir with a *live* firecracker process is a running bottle and is
|
||||
left strictly alone: it is neither killed nor removed. (The backend's
|
||||
`enumerate_active` registry is still a stub — #354 — so a live process
|
||||
is the only reliable "this bottle is in use" signal we have. Once the
|
||||
registry lands, registry-orphaned-but-running VMs can be reaped too.)
|
||||
|
||||
TAP slots free themselves (the flock drops when the launcher exits), so
|
||||
there is nothing to reclaim there.
|
||||
Orphans are: firecracker VMM processes whose config lives under our run
|
||||
dir, and the per-bottle run dirs. TAP slots free themselves (the flock
|
||||
drops when the launcher exits), so there is nothing to reclaim there.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -37,79 +22,38 @@ def _run_root() -> Path:
|
||||
return util.cache_dir() / "run"
|
||||
|
||||
|
||||
def _run_dir_of(cmd: str, run_root: Path) -> Path | None:
|
||||
"""The bottle run dir a firecracker cmdline belongs to, or None.
|
||||
|
||||
A bottle VM is launched with `--config-file <run_root>/<slug>/config.json`,
|
||||
so the run dir is the config file's parent when it sits directly under
|
||||
the run root. Anything else (a builder VM, the infra VM elsewhere) is
|
||||
not ours to reap here.
|
||||
"""
|
||||
toks = cmd.split()
|
||||
for i, tok in enumerate(toks):
|
||||
if tok == "--config-file" and i + 1 < len(toks):
|
||||
parent = Path(toks[i + 1]).parent
|
||||
if parent.parent == run_root:
|
||||
return parent
|
||||
return None
|
||||
|
||||
|
||||
def _scan_processes(run_root: Path) -> tuple[set[str], list[int]]:
|
||||
"""Inspect running firecracker VMs under ``run_root``.
|
||||
|
||||
Returns ``(live_run_dirs, orphan_pids)``:
|
||||
* ``live_run_dirs`` — run dirs backed by a running VM (never reaped);
|
||||
* ``orphan_pids`` — firecracker pids whose run dir no longer exists
|
||||
(a lingering VMM to kill).
|
||||
"""
|
||||
def _orphan_vm_pids() -> list[int]:
|
||||
"""firecracker processes whose --config-file is under our run dir."""
|
||||
run_root = str(_run_root())
|
||||
result = subprocess.run(
|
||||
["pgrep", "-a", "firecracker"],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
if result.returncode != 0:
|
||||
return set(), []
|
||||
live: set[str] = set()
|
||||
orphan_pids: list[int] = []
|
||||
return []
|
||||
pids: list[int] = []
|
||||
for line in result.stdout.splitlines():
|
||||
parts = line.split(None, 1)
|
||||
if len(parts) != 2:
|
||||
if len(parts) != 2 or run_root not in parts[1]:
|
||||
continue
|
||||
try:
|
||||
pid = int(parts[0])
|
||||
pids.append(int(parts[0]))
|
||||
except ValueError:
|
||||
continue
|
||||
run_dir = _run_dir_of(parts[1], run_root)
|
||||
if run_dir is None:
|
||||
continue
|
||||
if run_dir.is_dir():
|
||||
live.add(str(run_dir))
|
||||
else:
|
||||
orphan_pids.append(pid)
|
||||
return live, orphan_pids
|
||||
return pids
|
||||
|
||||
|
||||
def live_run_dirs() -> tuple[Path, ...]:
|
||||
"""Run directories backed by currently running agent microVMs."""
|
||||
live, _ = _scan_processes(_run_root())
|
||||
return tuple(Path(path) for path in sorted(live))
|
||||
|
||||
|
||||
def _orphan_run_dirs(run_root: Path, live: set[str]) -> list[str]:
|
||||
"""Run dirs with no live VM behind them — the leaked ones to remove."""
|
||||
def _run_dirs() -> list[str]:
|
||||
run_root = _run_root()
|
||||
if not run_root.is_dir():
|
||||
return []
|
||||
return sorted(
|
||||
str(p) for p in run_root.iterdir()
|
||||
if p.is_dir() and str(p) not in live
|
||||
)
|
||||
return sorted(str(p) for p in run_root.iterdir() if p.is_dir())
|
||||
|
||||
|
||||
def prepare_cleanup() -> FirecrackerBottleCleanupPlan:
|
||||
run_root = _run_root()
|
||||
live, orphan_pids = _scan_processes(run_root)
|
||||
return FirecrackerBottleCleanupPlan(
|
||||
vm_pids=tuple(orphan_pids),
|
||||
run_dirs=tuple(_orphan_run_dirs(run_root, live)),
|
||||
vm_pids=tuple(_orphan_vm_pids()),
|
||||
run_dirs=tuple(_run_dirs()),
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -25,27 +25,16 @@ The TAP slot allocation, rootfs build, and VM boot are the caller's job.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import subprocess
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
|
||||
from ...egress import EgressPlan
|
||||
from ...git_gate import GitGatePlan
|
||||
from ...log import info
|
||||
from ...orchestrator.client import OrchestratorClient, OrchestratorClientError
|
||||
from ...orchestrator.client import OrchestratorClient
|
||||
from ...orchestrator.lifecycle import (
|
||||
OrchestratorStartError, # re-exported so callers can catch it
|
||||
)
|
||||
from ...orchestrator.reprovision import reprovision_bottles
|
||||
from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME
|
||||
from ..consolidated_util import (
|
||||
provision_bottle,
|
||||
teardown_consolidated as _teardown_util,
|
||||
)
|
||||
from . import cleanup, infra_vm, util
|
||||
|
||||
_ENV_VAR_SECRET_PATH = "/run/bot-bottle/env-var-secret"
|
||||
from ..consolidated_util import provision_bottle, teardown_consolidated as _teardown_util
|
||||
from . import infra_vm
|
||||
|
||||
|
||||
class ConsolidatedLaunchError(RuntimeError):
|
||||
@@ -64,54 +53,6 @@ class LaunchContext:
|
||||
env_var_secret: str = "" # encryption key injected into the agent's env
|
||||
|
||||
|
||||
def _guest_ip_from_config(config_path: Path) -> str:
|
||||
"""Read the kernel's configured guest IP from a Firecracker config."""
|
||||
try:
|
||||
config = json.loads(config_path.read_text())
|
||||
args = config["boot-source"]["boot_args"]
|
||||
ip_arg = next(part for part in args.split() if part.startswith("ip="))
|
||||
return ip_arg.removeprefix("ip=").split(":", 1)[0]
|
||||
except (OSError, ValueError, KeyError, TypeError, StopIteration):
|
||||
return ""
|
||||
|
||||
|
||||
def persist_env_var_secret(private_key: Path, guest_ip: str, secret: str) -> None:
|
||||
"""Mirror the exec-time key into guest tmpfs for restart recovery."""
|
||||
proc = subprocess.run(
|
||||
util.ssh_base_argv(private_key, guest_ip)
|
||||
+ [f"umask 077; mkdir -p /run/bot-bottle; cat > {_ENV_VAR_SECRET_PATH}"],
|
||||
input=secret, capture_output=True, text=True, check=False,
|
||||
)
|
||||
if proc.returncode != 0:
|
||||
raise ConsolidatedLaunchError(
|
||||
f"failed to persist {ENV_VAR_SECRET_NAME} in agent VM: "
|
||||
f"{proc.stderr.strip() or '<no stderr>'}"
|
||||
)
|
||||
|
||||
|
||||
def _reprovision_running_bottles(client: OrchestratorClient) -> None:
|
||||
"""Read keys from live agent VMs and restore the restarted gateway."""
|
||||
try:
|
||||
secrets_by_ip: dict[str, str] = {}
|
||||
for run_dir in cleanup.live_run_dirs():
|
||||
guest_ip = _guest_ip_from_config(run_dir / "config.json")
|
||||
private_key = run_dir / "bottle_id_ed25519"
|
||||
if not guest_ip or not private_key.is_file():
|
||||
continue
|
||||
proc = subprocess.run(
|
||||
util.ssh_base_argv(private_key, guest_ip)
|
||||
+ [f"cat {_ENV_VAR_SECRET_PATH}"],
|
||||
capture_output=True, text=True, check=False,
|
||||
)
|
||||
if proc.returncode == 0 and proc.stdout.strip():
|
||||
secrets_by_ip[guest_ip] = proc.stdout.strip()
|
||||
count = reprovision_bottles(client, secrets_by_ip)
|
||||
if count:
|
||||
info(f"reprovisioned egress tokens for {count} Firecracker bottle(s)")
|
||||
except (OSError, OrchestratorClientError) as exc:
|
||||
info(f"egress token reprovision skipped: {exc}")
|
||||
|
||||
|
||||
def launch_consolidated(
|
||||
egress_plan: EgressPlan,
|
||||
git_gate_plan: GitGatePlan,
|
||||
@@ -126,7 +67,6 @@ def launch_consolidated(
|
||||
infra = infra_vm.ensure_running()
|
||||
url = infra.control_plane_url
|
||||
client = OrchestratorClient(url)
|
||||
_reprovision_running_bottles(client)
|
||||
|
||||
transport = infra_vm.gateway_transport()
|
||||
reg = provision_bottle(
|
||||
|
||||
@@ -26,7 +26,6 @@ from __future__ import annotations
|
||||
|
||||
import dataclasses
|
||||
import os
|
||||
import shutil
|
||||
from contextlib import ExitStack, contextmanager
|
||||
from pathlib import Path
|
||||
from typing import Callable, Generator
|
||||
@@ -55,10 +54,8 @@ from . import firecracker_vm, image_builder, isolation_probe, netpool, util
|
||||
from .bottle import FirecrackerBottle
|
||||
from .bottle_plan import FirecrackerBottlePlan
|
||||
from ...orchestrator.config_store import resolve_teardown_timeout
|
||||
from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME
|
||||
from .consolidated_launch import (
|
||||
launch_consolidated,
|
||||
persist_env_var_secret,
|
||||
teardown_consolidated,
|
||||
)
|
||||
|
||||
@@ -155,7 +152,6 @@ def launch(
|
||||
git_gate_plan=git_gate_plan,
|
||||
egress_plan=egress_plan,
|
||||
identity_token=ctx.identity_token,
|
||||
env_var_secret=ctx.env_var_secret,
|
||||
# Deliver the identity token as egress proxy credentials — clients
|
||||
# honor `HTTPS_PROXY=http://id:token@gw` without app changes; the
|
||||
# gateway reads Proxy-Authorization, validates the (source_ip,
|
||||
@@ -171,10 +167,6 @@ def launch(
|
||||
# Step 6: build the per-bottle rootfs + SSH key, then boot.
|
||||
run_dir = util.cache_dir() / "run" / plan.slug
|
||||
run_dir.mkdir(parents=True, exist_ok=True)
|
||||
# Remove the run dir on teardown so the per-bottle rootfs.ext4 (~1G)
|
||||
# doesn't leak. Registered before vm.terminate below so it runs *after*
|
||||
# it (ExitStack is LIFO): the VM is gone before we rm its rootfs.
|
||||
stack.callback(lambda: shutil.rmtree(run_dir, ignore_errors=True))
|
||||
rootfs = run_dir / "rootfs.ext4"
|
||||
util.build_rootfs_ext4(agent_base, rootfs)
|
||||
private_key, pubkey = util.generate_keypair(run_dir)
|
||||
@@ -190,7 +182,6 @@ def launch(
|
||||
)
|
||||
stack.callback(vm.terminate)
|
||||
firecracker_vm.wait_for_ssh(vm, private_key)
|
||||
persist_env_var_secret(private_key, slot.guest_ip, ctx.env_var_secret)
|
||||
|
||||
# Authoritative fail-closed egress-boundary check, before the agent
|
||||
# runs: prove the VM cannot reach the host directly.
|
||||
@@ -285,8 +276,6 @@ def _agent_guest_env(plan: FirecrackerBottlePlan, host_ip: str) -> dict[str, str
|
||||
env["GIT_GATE_URL"] = plan.agent_git_gate_url
|
||||
if plan.agent_supervise_url:
|
||||
env["MCP_SUPERVISE_URL"] = plan.agent_supervise_url
|
||||
if plan.env_var_secret:
|
||||
env[ENV_VAR_SECRET_NAME] = plan.env_var_secret
|
||||
for entry in egress_agent_env_entries(plan.egress_plan):
|
||||
key, _, value = entry.partition("=")
|
||||
env[key] = value
|
||||
|
||||
@@ -399,9 +399,6 @@ fi
|
||||
chown -R 0:0 /root 2>/dev/null || true
|
||||
|
||||
mkdir -p /etc/dropbear /run
|
||||
# Keep restart-recovery key material memory-backed, separate from both the
|
||||
# agent rootfs and the infra VM's persistent registry volume.
|
||||
mount -t tmpfs -o mode=0755 tmpfs /run 2>/dev/null || true
|
||||
# -R: generate host keys on demand. -E: log auth failures to stderr,
|
||||
# captured in the host-side console.log for debugging.
|
||||
/bb-dropbear -R -E -p 22 &
|
||||
|
||||
@@ -20,9 +20,6 @@ class MacosContainerBottlePlan(BottlePlan):
|
||||
# bottle is registered. See launch.py's stamp for why it lives here and not
|
||||
# only in the exec-time proxy env.
|
||||
identity_token: str = ""
|
||||
# Generated before `container run` so it becomes part of the container's
|
||||
# configured environment and can be read back after an infra restart.
|
||||
env_var_secret: str = ""
|
||||
|
||||
@property
|
||||
def container_name(self) -> str:
|
||||
|
||||
@@ -38,12 +38,7 @@ from ...egress import EgressPlan
|
||||
from ...git_gate import GitGatePlan
|
||||
from ...log import info
|
||||
from ...orchestrator.client import OrchestratorClient, OrchestratorClientError
|
||||
from ...orchestrator.reprovision import reprovision_bottles
|
||||
from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME
|
||||
from ..consolidated_util import (
|
||||
provision_bottle,
|
||||
teardown_consolidated as _teardown_util,
|
||||
)
|
||||
from ..consolidated_util import provision_bottle, teardown_consolidated as _teardown_util
|
||||
from . import util as container_mod
|
||||
from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active
|
||||
from .gateway import GATEWAY_NETWORK
|
||||
@@ -89,35 +84,12 @@ def ensure_gateway(
|
||||
needs `gateway_ip` at run time."""
|
||||
service = service or MacosInfraService()
|
||||
infra = service.ensure_running()
|
||||
endpoint = GatewayEndpoint(
|
||||
return GatewayEndpoint(
|
||||
orchestrator_url=infra.control_plane_url,
|
||||
gateway_ip=infra.gateway_ip,
|
||||
gateway_ca_pem=service.ca_cert_pem(),
|
||||
network=service.network,
|
||||
)
|
||||
_reprovision_running_bottles(endpoint)
|
||||
return endpoint
|
||||
|
||||
|
||||
def _reprovision_running_bottles(endpoint: GatewayEndpoint) -> None:
|
||||
"""Recover keys from live Apple containers and restore gateway tokens."""
|
||||
try:
|
||||
secrets_by_ip: dict[str, str] = {}
|
||||
for agent in enumerate_active():
|
||||
name = f"{CONTAINER_NAME_PREFIX}{agent.slug}"
|
||||
source_ip = container_mod.inspect_container_network_ip(name, endpoint.network)
|
||||
if not source_ip:
|
||||
continue
|
||||
secret = container_mod.read_container_env(name, ENV_VAR_SECRET_NAME)
|
||||
if secret:
|
||||
secrets_by_ip[source_ip] = secret
|
||||
count = reprovision_bottles(
|
||||
OrchestratorClient(endpoint.orchestrator_url), secrets_by_ip,
|
||||
)
|
||||
if count:
|
||||
info(f"reprovisioned egress tokens for {count} macOS bottle(s)")
|
||||
except (OrchestratorClientError, EnumerationError, OSError) as exc:
|
||||
info(f"egress token reprovision skipped: {exc}")
|
||||
|
||||
|
||||
def live_source_ips(network: str) -> list[str]:
|
||||
@@ -154,7 +126,6 @@ def register_agent(
|
||||
endpoint: GatewayEndpoint,
|
||||
image_ref: str = "",
|
||||
tokens: dict[str, str] | None = None,
|
||||
env_var_secret: str | None = None,
|
||||
) -> LaunchContext:
|
||||
"""Register the (already running) agent by its address and provision its
|
||||
git-gate state into the gateway. `source_ip` must be read from the live
|
||||
@@ -174,7 +145,6 @@ def register_agent(
|
||||
reg = provision_bottle(
|
||||
client, source_ip, egress_plan, git_gate_plan, AppleGatewayTransport(),
|
||||
image_ref=image_ref, tokens=tokens,
|
||||
env_var_secret=env_var_secret,
|
||||
)
|
||||
return LaunchContext(
|
||||
bottle_id=reg.bottle_id,
|
||||
|
||||
@@ -68,7 +68,6 @@ from .gateway_hosts import (
|
||||
)
|
||||
from .bottle_plan import MacosContainerBottlePlan
|
||||
from ...orchestrator.config_store import resolve_teardown_timeout
|
||||
from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME, new_env_var_secret
|
||||
from .consolidated_launch import (
|
||||
GatewayEndpoint,
|
||||
ensure_gateway,
|
||||
@@ -143,7 +142,6 @@ def launch(
|
||||
plan = _provision_git_gate_keys(plan)
|
||||
plan = _install_gateway_ca(plan, endpoint)
|
||||
plan = _stamp_agent_urls(plan, endpoint)
|
||||
plan = dataclasses.replace(plan, env_var_secret=new_env_var_secret())
|
||||
|
||||
# Step 3: run the agent. It has no identity token yet — registration
|
||||
# needs the address this run assigns.
|
||||
@@ -178,7 +176,6 @@ def launch(
|
||||
endpoint=endpoint,
|
||||
image_ref=plan.image,
|
||||
tokens=token_values,
|
||||
env_var_secret=plan.env_var_secret,
|
||||
)
|
||||
stack.callback(
|
||||
teardown_consolidated, ctx.bottle_id,
|
||||
@@ -409,8 +406,6 @@ def _agent_env_entries(
|
||||
env.append(f"GIT_GATE_URL={plan.agent_git_gate_url}")
|
||||
if plan.agent_supervise_url:
|
||||
env.append(f"MCP_SUPERVISE_URL={plan.agent_supervise_url}")
|
||||
if getattr(plan, "env_var_secret", ""):
|
||||
env.append(f"{ENV_VAR_SECRET_NAME}={plan.env_var_secret}")
|
||||
for name, value in sorted(plan.agent_provision.guest_env.items()):
|
||||
env.append(f"{name}={value}")
|
||||
# Forwarded vars: bare name → inherits from the `container run` process env
|
||||
|
||||
@@ -361,12 +361,6 @@ def exec_container(name: str, argv: list[str]) -> None:
|
||||
)
|
||||
|
||||
|
||||
def read_container_env(name: str, env_name: str) -> str:
|
||||
"""Read one configured env value from a running container, or ``""``."""
|
||||
result = _run_container_op([_CONTAINER, "exec", name, "printenv", env_name])
|
||||
return result.stdout.strip() if result.returncode == 0 else ""
|
||||
|
||||
|
||||
def exec_container_as_root(name: str, argv: list[str]) -> None:
|
||||
"""`exec_container`, but as uid 0 inside the container.
|
||||
|
||||
|
||||
@@ -1,35 +0,0 @@
|
||||
"""Shared host-side join for backend-discovered bottle encryption keys."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from .client import OrchestratorClient, OrchestratorClientError
|
||||
|
||||
|
||||
def reprovision_bottles(
|
||||
client: OrchestratorClient,
|
||||
secrets_by_source_ip: dict[str, str],
|
||||
) -> int:
|
||||
"""Restore tokens for registered bottles whose backend exposes a key.
|
||||
|
||||
Backends own discovery because containers and microVMs have different
|
||||
enumeration primitives. This helper owns the shared registry join and
|
||||
intentionally tolerates one bad/missing key without blocking a launch.
|
||||
"""
|
||||
restored = 0
|
||||
for bottle in client.list_bottles():
|
||||
bottle_id = bottle.get("bottle_id")
|
||||
source_ip = bottle.get("source_ip")
|
||||
if not isinstance(bottle_id, str) or not isinstance(source_ip, str):
|
||||
continue
|
||||
secret = secrets_by_source_ip.get(source_ip, "").strip()
|
||||
if not secret:
|
||||
continue
|
||||
try:
|
||||
if client.reprovision_gateway(bottle_id, secret):
|
||||
restored += 1
|
||||
except OrchestratorClientError:
|
||||
continue
|
||||
return restored
|
||||
|
||||
|
||||
__all__ = ["reprovision_bottles"]
|
||||
@@ -1,131 +0,0 @@
|
||||
# PRD prd-new: Encrypted at-rest egress secrets (SecretProvider, interim slice)
|
||||
|
||||
- **Status:** Draft
|
||||
- **Author:** didericis
|
||||
- **Created:** 2026-07-21
|
||||
- **Issue:** #355
|
||||
|
||||
## Summary
|
||||
|
||||
An interim step toward the generic `SecretProvider` (#355) that stops short
|
||||
of per-request minting. Today the orchestrator holds each bottle's egress
|
||||
auth tokens **in process memory only**, so any infra-container recreation
|
||||
silently strips every already-running bottle of its upstream credentials.
|
||||
This PRD makes those secrets survive a gateway restart by persisting them
|
||||
**encrypted**, under a key that is not itself sitting next to the
|
||||
ciphertext.
|
||||
|
||||
The end state in #355 — short-lived, scoped credentials minted per request
|
||||
— removes the need to store anything durable at all. That is a larger
|
||||
change gated on per-upstream minting support. This slice buys back
|
||||
restart-survivability now without regressing to plaintext secrets at rest.
|
||||
|
||||
## Problem
|
||||
|
||||
`Orchestrator._tokens` (`bot_bottle/orchestrator/service.py:74-79`) is a
|
||||
plain in-memory dict, deliberately never written to the registry DB:
|
||||
|
||||
> Held **in memory only** — never written to the registry DB — so the
|
||||
> gateway can inject each bottle's upstream credential without secrets at
|
||||
> rest. Lost on restart (re-launch re-registers them); the future
|
||||
> SecretProvider (#355) replaces this with per-request minting.
|
||||
|
||||
The registry itself *is* durable (SQLite on a container-only volume), and
|
||||
so is the gateway CA since #450 / `2cd44cf7`. The tokens are now the only
|
||||
piece of gateway state that does not survive a restart, which makes the
|
||||
failure mode both silent and confusing.
|
||||
|
||||
### Observed failure
|
||||
|
||||
Checking out a branch that touches `bot_bottle/**/*.py` changes
|
||||
`source_hash()` (`bot_bottle/orchestrator/lifecycle.py:88-99`).
|
||||
`MacosInfraService._source_current()`
|
||||
(`bot_bottle/backend/macos_container/infra.py:159-169`) sees the mismatch
|
||||
and `ensure_running()` force-removes and recreates the infra container
|
||||
(`infra.py:198-208`). The registry rows survive on the DB volume; the CA
|
||||
survives on its host bind-mount; `_tokens` comes back empty.
|
||||
|
||||
Every already-running bottle then fails closed, mid-session, on its next
|
||||
outbound request:
|
||||
|
||||
- `/resolve` succeeds — the bottle is still `active` in
|
||||
`orchestrator_bottles` and its policy blob is served intact, including
|
||||
`- host: "api.anthropic.com"` with `auth_scheme: Bearer` /
|
||||
`token_env: EGRESS_TOKEN_0`.
|
||||
- `tokens_for()` returns `{}`, so the resolved env overlay has no
|
||||
`EGRESS_TOKEN_0`.
|
||||
- `decide()` (`bot_bottle/egress_addon_core.py:644-652`) blocks with
|
||||
`egress: route for 'api.anthropic.com' declared auth but env var
|
||||
'EGRESS_TOKEN_0' is unset` — an 89-byte `403` on every request.
|
||||
|
||||
Confirmed live on the macOS backend on 2026-07-21: two bottles running
|
||||
since 20:29/20:30 were still registered `active` with valid policy after
|
||||
the 23:05 infra recreation, and both took 89-byte `403`s from then on,
|
||||
while a bottle launched *after* the recreation egressed normally. The
|
||||
recovery today is to relaunch every affected bottle.
|
||||
|
||||
Note this is a re-attachment blocker distinct from #443/#445 and from #450
|
||||
— the CA and the gateway address were both fine. It is specifically the
|
||||
credential wipe.
|
||||
|
||||
## Goals / Success criteria
|
||||
|
||||
1. A bottle's egress auth tokens survive infra-container recreation: an
|
||||
already-running bottle keeps egressing across a gateway restart with no
|
||||
relaunch and no operator action.
|
||||
2. Secrets are **never** at rest in plaintext, and never at rest next to a
|
||||
key that trivially decrypts them.
|
||||
3. Compromise of the registry DB file alone does not yield usable
|
||||
upstream credentials.
|
||||
4. The stored form is revocable and rotatable without relaunching bottles
|
||||
that are not affected.
|
||||
5. When the retained bottle record reaches the lifecycle status `removed`,
|
||||
its stored secrets are destroyed. Status/event persistence and the removal
|
||||
transition land separately; this interim slice intentionally retains the
|
||||
ciphertext across today's teardown/reconcile calls until that lifecycle is
|
||||
available.
|
||||
6. Migration is transparent: existing bottles keep working, no manifest
|
||||
changes required.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- **Per-request minting** of short-lived scoped credentials. That is the
|
||||
#355 end state; this PRD is explicitly the interim slice and should not
|
||||
foreclose it.
|
||||
- Generalizing `DeployKeyProvisioner` into the full `SecretProvider` ABC,
|
||||
or the manifest-level `{ provider: <name> }` reference surface.
|
||||
- User-extensible provider discovery (`~/.bot-bottle/contrib/<name>/`).
|
||||
- Changing the `/resolve` contract's shape (it already carries `tokens`).
|
||||
- Fixing the *trigger* — `source_hash` churn on branch switch. Recreating
|
||||
infra is legitimate; it just must not cost running bottles their
|
||||
credentials. A separate guard that refuses recreation while bottles are
|
||||
active is complementary and out of scope here.
|
||||
|
||||
## Design
|
||||
|
||||
> **TODO (didericis):** the encryption flow goes here — key custody, where
|
||||
> the key material lives relative to the ciphertext, the wrap/unwrap path
|
||||
> at register and at `/resolve`, and what an attacker who holds only the
|
||||
> DB (or only the host, or only the infra container) can recover.
|
||||
|
||||
Constraints the design has to satisfy, for reference while drafting:
|
||||
|
||||
- The gateway's `PolicyResolver` needs the cleartext at request time, on
|
||||
the data-plane path, so unwrap has to be cheap enough to sit in a
|
||||
per-flow `/resolve` (or be cached in memory after first unwrap).
|
||||
- The infra container is recreated routinely and unattended. Anything
|
||||
requiring an interactive unlock on every recreation defeats the goal.
|
||||
- The DB lives on a container-only volume that the host does not mount, so
|
||||
host-side and guest-side components see different filesystems — that
|
||||
asymmetry is available as a place to split custody.
|
||||
- The agent must never be able to reach the key material. It is a separate
|
||||
container with no control-plane token, which is the existing boundary.
|
||||
|
||||
## Open questions
|
||||
|
||||
- Where does the unwrap key live, and what recreates/re-derives it when the
|
||||
infra container is rebuilt?
|
||||
- Is the cleartext cached in memory after first unwrap, or unwrapped per
|
||||
request? (Latency vs. exposure window.)
|
||||
- What is the rotation story — re-wrap in place, or force re-registration?
|
||||
- Does this land behind a flag, or replace `_tokens` outright?
|
||||
@@ -1,178 +0,0 @@
|
||||
# Firecracker Image Remote Store
|
||||
|
||||
**Date:** 2026-07-22
|
||||
**Context:** PR #459 (run-dir leak fix) surfaced that committed Firecracker snapshots
|
||||
currently live only on the host machine. Once a host is wiped or a run-dir is
|
||||
evicted, a user's preserved bottle is gone. This note investigates a secure remote
|
||||
store so committed images survive host turnover and can be restored without the
|
||||
user having to pre-flag which sessions to keep.
|
||||
|
||||
## Verdict
|
||||
|
||||
Backblaze B2 + Cloudflare CDN is the cost-optimal choice for most deployments.
|
||||
Cloudflare R2 is the simpler zero-config option at slightly higher storage cost.
|
||||
Self-hosted MinIO is the right call for air-gapped or on-premises installs.
|
||||
|
||||
On the image-size front, zstd-compressing the committed tar before upload
|
||||
produces roughly a 60–70% reduction with negligible impact on restore latency.
|
||||
OverlayFS (for in-flight working rootfs, not for the committed artifact) cuts
|
||||
per-instance disk use to ~10–50 MB per extra bottle sharing the same base.
|
||||
|
||||
---
|
||||
|
||||
## What Gets Stored
|
||||
|
||||
The Firecracker backend produces two artifact types:
|
||||
|
||||
| Artifact | Created by | Size | Lifetime |
|
||||
|----------|-----------|------|---------|
|
||||
| `rootfs/agent-<digest>/` (dir) | `image_builder.py` → `mke2fs -d` | ~1 GB as ext4 | Cached per Dockerfile hash; evictable |
|
||||
| `committed/<slug>/rootfs.tar` | `FirecrackerFreezer._freeze` via SSH tar | 500 MB–1 GB | User-preserved; must survive host wipe |
|
||||
|
||||
The committed artifact is a tar of the guest's live filesystem streamed out over
|
||||
SSH (`freezer.py:58–90`). At resume time `launch.py` calls `mke2fs -d` to
|
||||
rebuild a fresh ext4 from this tar. The tar — not the ext4 — is what needs to be
|
||||
pushed to remote storage and pulled back at restore time.
|
||||
|
||||
The base image cache (`rootfs/agent-<digest>/`) is derivable from the Dockerfile
|
||||
and can be rebuilt on demand; it is lower priority for remote storage.
|
||||
|
||||
---
|
||||
|
||||
## Storage Candidates
|
||||
|
||||
### Object storage
|
||||
|
||||
| Provider | Storage | Egress | Notes |
|
||||
|----------|---------|--------|-------|
|
||||
| **Backblaze B2** | $0.006/GB | Free (via Cloudflare Bandwidth Alliance) | Cheapest storage; pairs with Cloudflare CDN to eliminate egress |
|
||||
| **Cloudflare R2** | $0.015/GB | $0 always | Zero-config egress; no lifecycle transitions (limitation) |
|
||||
| **Wasabi** | $0.0069/GB | Free (1:1 ratio) | 90-day minimum retention; good for archival; lifecycle evaluated daily |
|
||||
| **AWS S3** | $0.023/GB | $0.09/GB | Richest lifecycle support; expensive at scale; avoid unless already in AWS |
|
||||
| **MinIO** (self-hosted) | Host cost only | None | S3-compatible; best for private/on-prem deployments |
|
||||
|
||||
**B2 + Cloudflare CDN** is effectively $0.006/GB with zero egress — about 18×
|
||||
cheaper than S3 for restore-heavy workloads. **R2** is the zero-config choice
|
||||
($0 egress by default, no Bandwidth Alliance pairing needed) at a slightly
|
||||
higher storage rate.
|
||||
|
||||
**Cloudflare R2's missing lifecycle support** is the main caveat: auto-eviction
|
||||
rules (evict images older than N days) cannot currently be expressed natively in
|
||||
R2. Wasabi and S3 both support declarative lifecycle policies.
|
||||
|
||||
### Retention policy recommendation
|
||||
|
||||
The comment proposes:
|
||||
- Retain images for ~1 week by default
|
||||
- Warn when approaching a capacity threshold
|
||||
- Auto-evict oldest images once threshold is exceeded
|
||||
|
||||
This maps cleanly to an application-level policy (not a provider lifecycle rule),
|
||||
which avoids the R2 limitation and works consistently across providers:
|
||||
|
||||
1. On `commit`: upload tar, record `(slug, size_bytes, uploaded_at)` in a local
|
||||
or remote manifest file.
|
||||
2. On startup / on `list`: scan the manifest, warn if total stored size exceeds
|
||||
e.g. 80% of the configured threshold.
|
||||
3. On eviction run (CLI or cron): delete objects older than `retention_days`
|
||||
(default 7) that push total over `max_capacity`; oldest-first.
|
||||
|
||||
This keeps the policy logic in bot-bottle and the storage provider as a dumb
|
||||
object store — no vendor-specific lifecycle API required.
|
||||
|
||||
---
|
||||
|
||||
## Image Size Reduction
|
||||
|
||||
### Current artifact sizes
|
||||
|
||||
A typical committed tar for a Claude Code agent image is 500 MB–1 GB uncompressed.
|
||||
The per-run ext4 (copy of the base, written at `start`) adds another ~1 GB of
|
||||
local disk. The leak fix in this PR addresses the ext4 copies; the remote store
|
||||
addresses the committed tars.
|
||||
|
||||
### Compression
|
||||
|
||||
zstd compression of the committed tar before upload is the highest-leverage
|
||||
single change:
|
||||
|
||||
| Codec | Typical size (1 GB rootfs) | Compress speed | Decompress speed |
|
||||
|-------|---------------------------|---------------|-----------------|
|
||||
| gzip | 350–430 MB | ~100 MB/s | ~500 MB/s |
|
||||
| **zstd (default)** | **330–360 MB** | **~400 MB/s** | **~2 GB/s** |
|
||||
| xz | 290–320 MB | ~20 MB/s | ~200 MB/s |
|
||||
|
||||
**zstd is the best trade-off**: 65–67% size reduction, near-instantaneous
|
||||
decompression. The `tar` call in `freezer.py` could pipe through `zstd` before
|
||||
writing to disk and to the remote; `resume` decompresses on the way back. A
|
||||
`.tar.zst` suffix marks compressed artifacts so old tars remain restorable
|
||||
without the codec.
|
||||
|
||||
### SquashFS for the base image cache
|
||||
|
||||
The `rootfs/agent-<digest>/` directory (the buildah-exported tree) is rebuilt by
|
||||
`image_builder.py` and turned into per-run ext4 by `mke2fs -d`. Storing the
|
||||
base as a SquashFS image instead of a flat directory tree would reduce it from
|
||||
~1 GB to ~330–360 MB and make the cache remote-friendly. Firecracker does not
|
||||
directly boot SquashFS, but the existing `mke2fs -d` path reads a directory tree
|
||||
— a SquashFS mount could serve as the source. This is a larger change and lower
|
||||
priority than tar compression.
|
||||
|
||||
### OverlayFS for per-run rootfs
|
||||
|
||||
Multiple simultaneous bottles sharing the same agent image today each get a full
|
||||
`mke2fs -d` copy (~1 GB). OverlayFS (read-only base + writable sparse overlay)
|
||||
would reduce this to ~10–50 MB per instance beyond the first:
|
||||
|
||||
- Mount the base image directory as read-only lower layer
|
||||
- Attach a sparse ext4 or tmpfs writable layer per bottle
|
||||
- Pass the merged overlay to Firecracker as the block device
|
||||
|
||||
E2B's public write-up on Firecracker + OverlayFS confirms this approach works
|
||||
at scale. The `launch.py` changes would be non-trivial (device mapper or
|
||||
`fuse-overlayfs` plumbing), so this is a follow-up rather than a prerequisite
|
||||
for the remote store.
|
||||
|
||||
---
|
||||
|
||||
## Recommended Approach
|
||||
|
||||
**Phase 1 — remote store with zstd (tight scope, actionable now)**
|
||||
|
||||
1. Add `--zstd` to the `tar` call in `FirecrackerFreezer._freeze`; name the
|
||||
artifact `rootfs.tar.zst`. Keep uncompressed restore path for legacy tars.
|
||||
2. Add a `bb firecracker upload <slug>` / `bb firecracker pull <slug>` pair that
|
||||
pushes/fetches the compressed tar to the configured object store (S3-compatible
|
||||
API, so B2, R2, MinIO, and Wasabi all work with the same client).
|
||||
3. Store a `manifest.json` in the bucket (or a local mirror) tracking slug →
|
||||
`{size, uploaded_at}`. Use it for threshold warnings and eviction.
|
||||
4. Default retention: 7 days, configurable via `firecracker.image_retention_days`
|
||||
in `~/.config/bot-bottle/config.toml` (or equivalent).
|
||||
5. Warn at 80% of `max_capacity` (default e.g. 50 GB); evict oldest on commit
|
||||
once at 100%.
|
||||
|
||||
**Storage recommendation:** Cloudflare R2 for hosted deployments (zero egress,
|
||||
zero config), MinIO for private/on-premises.
|
||||
|
||||
**Phase 2 — base image cache compression**
|
||||
|
||||
Compress the `agent-<digest>` cache dir as a `.tar.zst` to save ~65% on repeated
|
||||
image uploads. Low urgency since the base image is rebuildable.
|
||||
|
||||
**Phase 3 — OverlayFS per-run disk**
|
||||
|
||||
Replace the full per-run ext4 copy with an OverlayFS sparse layer. Largest disk
|
||||
impact (~90–95% savings per concurrent bottle) but highest implementation
|
||||
complexity. Track as a separate PRD.
|
||||
|
||||
---
|
||||
|
||||
## Open Questions
|
||||
|
||||
- Does the host have a configured object-store credential path, or should the
|
||||
remote store be an opt-in with an explicit `bb config set image-store.url ...`?
|
||||
- Should `commit` automatically upload, or should upload be an explicit step to
|
||||
avoid surprise egress?
|
||||
- What is the acceptable cold-start latency for a restore from remote? A 330 MB
|
||||
zstd tar at 100 Mbit/s takes ~26 s; at 1 Gbit/s, ~2.6 s. This bounds the
|
||||
retention strategy (evict from local after successful upload vs keep local copy).
|
||||
@@ -1,160 +0,0 @@
|
||||
"""Backend-agnostic and backend-specific encrypted-secret recovery tests."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import subprocess
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
from bot_bottle.orchestrator.reprovision import reprovision_bottles
|
||||
from bot_bottle.backend.firecracker import consolidated_launch as fc
|
||||
from bot_bottle.backend.macos_container import consolidated_launch as mac
|
||||
from bot_bottle.backend.docker import consolidated_launch as docker
|
||||
from bot_bottle.orchestrator.client import OrchestratorClientError
|
||||
|
||||
|
||||
def _proc(returncode: int = 0, stdout: str = "", stderr: str = ""):
|
||||
return subprocess.CompletedProcess([], returncode, stdout=stdout, stderr=stderr)
|
||||
|
||||
|
||||
class TestSharedReprovision(unittest.TestCase):
|
||||
def test_joins_registry_records_by_source_ip(self) -> None:
|
||||
client = Mock()
|
||||
client.list_bottles.return_value = [
|
||||
{"bottle_id": "b1", "source_ip": "10.0.0.1"},
|
||||
{"bottle_id": "b2", "source_ip": "10.0.0.2"},
|
||||
{"bottle_id": 3, "source_ip": "10.0.0.3"},
|
||||
]
|
||||
client.reprovision_gateway.side_effect = [True, False]
|
||||
count = reprovision_bottles(
|
||||
client, {"10.0.0.1": " key-1\n", "10.0.0.2": "key-2"},
|
||||
)
|
||||
self.assertEqual(1, count)
|
||||
self.assertEqual(
|
||||
[("b1", "key-1"), ("b2", "key-2")],
|
||||
[call.args for call in client.reprovision_gateway.call_args_list],
|
||||
)
|
||||
|
||||
def test_one_failure_does_not_block_other_bottles(self) -> None:
|
||||
client = Mock()
|
||||
client.list_bottles.return_value = [
|
||||
{"bottle_id": "b1", "source_ip": "10.0.0.1"},
|
||||
{"bottle_id": "b2", "source_ip": "10.0.0.2"},
|
||||
]
|
||||
client.reprovision_gateway.side_effect = [
|
||||
OrchestratorClientError("bad key"), True,
|
||||
]
|
||||
self.assertEqual(
|
||||
1,
|
||||
reprovision_bottles(
|
||||
client, {"10.0.0.1": "key-1", "10.0.0.2": "key-2"},
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
class TestMacosReprovision(unittest.TestCase):
|
||||
def test_reads_configured_container_env_and_reprovisions(self) -> None:
|
||||
endpoint = mac.GatewayEndpoint("http://orch", "10.0.0.9", "PEM", "net")
|
||||
agent = SimpleNamespace(slug="demo")
|
||||
client = Mock()
|
||||
with patch.object(mac, "enumerate_active", return_value=[agent]), \
|
||||
patch.object(mac.container_mod, "inspect_container_network_ip",
|
||||
return_value="10.0.0.1"), \
|
||||
patch.object(mac.container_mod, "read_container_env", return_value="key"), \
|
||||
patch.object(mac, "OrchestratorClient", return_value=client), \
|
||||
patch.object(mac, "reprovision_bottles", return_value=1) as restore, \
|
||||
patch.object(mac, "info"):
|
||||
mac._reprovision_running_bottles(endpoint)
|
||||
restore.assert_called_once_with(client, {"10.0.0.1": "key"})
|
||||
|
||||
def test_enumeration_failure_is_best_effort(self) -> None:
|
||||
endpoint = mac.GatewayEndpoint("http://orch", "10.0.0.9", "PEM", "net")
|
||||
with patch.object(mac, "enumerate_active",
|
||||
side_effect=mac.EnumerationError("failed")), \
|
||||
patch.object(mac, "info") as info:
|
||||
mac._reprovision_running_bottles(endpoint)
|
||||
self.assertIn("skipped", info.call_args.args[0])
|
||||
|
||||
|
||||
class TestDockerReprovision(unittest.TestCase):
|
||||
def test_maps_network_containers_to_keys(self) -> None:
|
||||
inspect = _proc(stdout=(
|
||||
"bot-bottle-infra 172.18.0.2/16\n"
|
||||
"bot-bottle-a 172.18.0.3/16\n"
|
||||
"malformed\n"
|
||||
))
|
||||
key = _proc(stdout="secret\n")
|
||||
client = Mock()
|
||||
with patch.object(docker, "OrchestratorClient", return_value=client), \
|
||||
patch.object(docker, "run_docker", side_effect=[inspect, key]), \
|
||||
patch.object(docker, "reprovision_bottles", return_value=1) as restore, \
|
||||
patch.object(docker.log, "info"):
|
||||
docker._reprovision_running_bottles("http://orch")
|
||||
restore.assert_called_once_with(client, {"172.18.0.3": "secret"})
|
||||
|
||||
def test_missing_docker_is_best_effort(self) -> None:
|
||||
with patch.object(docker, "run_docker", side_effect=FileNotFoundError("docker")), \
|
||||
patch.object(docker.log, "info") as info:
|
||||
docker._reprovision_running_bottles("http://orch")
|
||||
self.assertIn("skipped", info.call_args.args[0])
|
||||
|
||||
|
||||
class TestFirecrackerReprovision(unittest.TestCase):
|
||||
def _run_dir(self, root: Path, ip: str = "10.243.0.3") -> Path:
|
||||
run_dir = root / "demo"
|
||||
run_dir.mkdir()
|
||||
(run_dir / "bottle_id_ed25519").write_text("key")
|
||||
(run_dir / "config.json").write_text(json.dumps({
|
||||
"boot-source": {"boot_args": f"root=/dev/vda ip={ip}::gw:mask::eth0:off"}
|
||||
}))
|
||||
return run_dir
|
||||
|
||||
def test_extracts_guest_ip_from_config(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
run_dir = self._run_dir(Path(tmp))
|
||||
self.assertEqual("10.243.0.3", fc._guest_ip_from_config(run_dir / "config.json"))
|
||||
self.assertEqual("", fc._guest_ip_from_config(run_dir / "missing.json"))
|
||||
|
||||
def test_persists_key_over_stdin_not_argv(self) -> None:
|
||||
with patch.object(fc.util, "ssh_base_argv", return_value=["ssh", "guest"]), \
|
||||
patch.object(fc.subprocess, "run", return_value=_proc()) as run:
|
||||
fc.persist_env_var_secret(Path("/key"), "10.0.0.1", "super-secret")
|
||||
self.assertEqual("super-secret", run.call_args.kwargs["input"])
|
||||
self.assertNotIn("super-secret", " ".join(run.call_args.args[0]))
|
||||
|
||||
def test_persist_failure_is_fatal_to_launch(self) -> None:
|
||||
with patch.object(fc.util, "ssh_base_argv", return_value=["ssh", "guest"]), \
|
||||
patch.object(fc.subprocess, "run", return_value=_proc(1, stderr="denied")):
|
||||
with self.assertRaisesRegex(fc.ConsolidatedLaunchError, "denied"):
|
||||
fc.persist_env_var_secret(Path("/key"), "10.0.0.1", "secret")
|
||||
|
||||
def test_reads_live_vm_key_and_reprovisions(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
run_dir = self._run_dir(Path(tmp))
|
||||
client = Mock()
|
||||
with patch.object(fc.cleanup, "live_run_dirs", return_value=(run_dir,)), \
|
||||
patch.object(fc.util, "ssh_base_argv", return_value=["ssh", "guest"]), \
|
||||
patch.object(fc.subprocess, "run", return_value=_proc(stdout="secret\n")), \
|
||||
patch.object(fc, "reprovision_bottles", return_value=1) as restore, \
|
||||
patch.object(fc, "info"):
|
||||
fc._reprovision_running_bottles(client)
|
||||
restore.assert_called_once_with(client, {"10.243.0.3": "secret"})
|
||||
|
||||
def test_unreadable_vm_is_skipped(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
run_dir = self._run_dir(Path(tmp))
|
||||
client = Mock()
|
||||
with patch.object(fc.cleanup, "live_run_dirs", return_value=(run_dir,)), \
|
||||
patch.object(fc.util, "ssh_base_argv", return_value=["ssh", "guest"]), \
|
||||
patch.object(fc.subprocess, "run", return_value=_proc(1)), \
|
||||
patch.object(fc, "reprovision_bottles", return_value=0) as restore:
|
||||
fc._reprovision_running_bottles(client)
|
||||
restore.assert_called_once_with(client, {})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -2,7 +2,6 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import dataclasses
|
||||
import unittest
|
||||
|
||||
from bot_bottle.backend.docker.consolidated_compose import consolidated_agent_compose
|
||||
@@ -47,15 +46,6 @@ class TestConsolidatedAgentCompose(unittest.TestCase):
|
||||
# forwarded secrets are bare names (value inherited from process env).
|
||||
self.assertIn("CLAUDE_CODE_OAUTH_TOKEN", env)
|
||||
|
||||
def test_env_var_secret_stays_a_bare_name(self) -> None:
|
||||
plan = _plan(with_egress=True, supervise=True, with_git=True)
|
||||
plan = dataclasses.replace(plan, env_var_secret="secret-value")
|
||||
env = consolidated_agent_compose(
|
||||
plan, gateway_ip=_GW, source_ip=_IP, network=_NET,
|
||||
)["services"]["agent"]["environment"]
|
||||
self.assertIn("ENV_VAR_SECRET", env)
|
||||
self.assertNotIn("ENV_VAR_SECRET=secret-value", env)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -91,7 +91,6 @@ class TestTeardownWarning(unittest.TestCase):
|
||||
bottle_id="b1", identity_token="t", source_ip="172.20.0.4",
|
||||
network="bot-bottle-gateway", gateway_ip="172.20.0.2",
|
||||
orchestrator_url="http://orch:8099",
|
||||
env_var_secret="encryption-key",
|
||||
)
|
||||
|
||||
images = BottleImages(agent="bot-bottle-claude:latest", sidecar="bot-bottle-sidecars:latest")
|
||||
@@ -108,7 +107,7 @@ class TestTeardownWarning(unittest.TestCase):
|
||||
mock.patch.object(
|
||||
launch_mod, "write_compose_file", return_value=Path("/tmp/compose.yml"),
|
||||
), \
|
||||
mock.patch.object(launch_mod, "compose_up") as compose_up, \
|
||||
mock.patch.object(launch_mod, "compose_up"), \
|
||||
mock.patch.object(launch_mod, "compose_dump_logs"), \
|
||||
mock.patch.object(
|
||||
launch_mod, "compose_down",
|
||||
@@ -123,9 +122,6 @@ class TestTeardownWarning(unittest.TestCase):
|
||||
self.assertIn("bot-bottle: warning:", output)
|
||||
self.assertIn("bot-bottle-test-teardown-abc", output)
|
||||
self.assertIn("compose-down", output)
|
||||
self.assertEqual(
|
||||
"encryption-key", compose_up.call_args.kwargs["env"]["ENV_VAR_SECRET"],
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -66,7 +66,6 @@ class TestNetpoolRenderers(unittest.TestCase):
|
||||
|
||||
self.assertIn("chown node:node /home/node", util._GUEST_INIT)
|
||||
self.assertIn("chmod 755 /home/node", util._GUEST_INIT)
|
||||
self.assertIn("mount -t tmpfs -o mode=0755 tmpfs /run", util._GUEST_INIT)
|
||||
|
||||
def test_nixos_module_is_non_invasive(self):
|
||||
# The NixOS module must NOT flip the host firewall backend or
|
||||
|
||||
@@ -10,9 +10,7 @@ classmethods forward to their module.
|
||||
from __future__ import annotations
|
||||
|
||||
import subprocess
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from unittest.mock import patch
|
||||
|
||||
from bot_bottle.backend.firecracker import cleanup as fc_cleanup
|
||||
@@ -25,77 +23,32 @@ def _proc(stdout: str = "", returncode: int = 0) -> "subprocess.CompletedProcess
|
||||
return subprocess.CompletedProcess([], returncode, stdout=stdout, stderr="")
|
||||
|
||||
|
||||
class TestProcessScan(unittest.TestCase):
|
||||
def test_run_dir_of_matches_only_direct_children(self):
|
||||
run_root = Path("/cache/run")
|
||||
self.assertEqual(
|
||||
Path("/cache/run/dev-a"),
|
||||
fc_cleanup._run_dir_of(
|
||||
f"firecracker --config-file {run_root}/dev-a/config.json", run_root
|
||||
),
|
||||
)
|
||||
# infra/builder VMs elsewhere, or nested paths, are not ours.
|
||||
self.assertIsNone(
|
||||
fc_cleanup._run_dir_of("firecracker --config-file /elsewhere/config.json", run_root)
|
||||
)
|
||||
self.assertIsNone(
|
||||
fc_cleanup._run_dir_of("firecracker --no-config", run_root)
|
||||
class TestOrphanEnumeration(unittest.TestCase):
|
||||
def test_orphan_vm_pids_filters_by_run_dir(self):
|
||||
run_root = str(fc_cleanup._run_root())
|
||||
out = (
|
||||
f"111 firecracker --config-file {run_root}/dev-a/config.json\n"
|
||||
"222 firecracker --config-file /somewhere/else/config.json\n"
|
||||
"notanint firecracker --config-file " + run_root + "/x\n"
|
||||
)
|
||||
with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(out)):
|
||||
self.assertEqual([111], fc_cleanup._orphan_vm_pids())
|
||||
|
||||
def test_scan_splits_live_dirs_from_orphan_pids(self):
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
run_root = Path(tmp)
|
||||
(run_root / "live-a").mkdir() # dir present -> live VM, protected
|
||||
# "gone-b" dir intentionally absent -> lingering VMM, orphan pid
|
||||
out = (
|
||||
f"111 firecracker --config-file {run_root}/live-a/config.json\n"
|
||||
f"222 firecracker --config-file {run_root}/gone-b/config.json\n"
|
||||
"333 firecracker --config-file /elsewhere/config.json\n"
|
||||
"notanint firecracker --config-file x\n"
|
||||
)
|
||||
with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(out)):
|
||||
live, orphan_pids = fc_cleanup._scan_processes(run_root)
|
||||
self.assertEqual({str(run_root / "live-a")}, live)
|
||||
self.assertEqual([222], orphan_pids)
|
||||
|
||||
def test_scan_empty_when_pgrep_fails(self):
|
||||
def test_orphan_vm_pids_empty_when_pgrep_fails(self):
|
||||
with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(returncode=1)):
|
||||
self.assertEqual((set(), []), fc_cleanup._scan_processes(Path("/x")))
|
||||
self.assertEqual([], fc_cleanup._orphan_vm_pids())
|
||||
|
||||
def test_live_run_dirs_returns_paths_in_stable_order(self):
|
||||
with patch.object(fc_cleanup, "_run_root", return_value=Path("/run")), \
|
||||
patch.object(fc_cleanup, "_scan_processes",
|
||||
return_value=({"/run/b", "/run/a"}, [])):
|
||||
self.assertEqual(
|
||||
(Path("/run/a"), Path("/run/b")), fc_cleanup.live_run_dirs(),
|
||||
)
|
||||
def test_run_dirs_empty_when_absent(self):
|
||||
with patch.object(fc_cleanup.util, "cache_dir") as cache:
|
||||
cache.return_value.__truediv__.return_value.is_dir.return_value = False
|
||||
self.assertEqual([], fc_cleanup._run_dirs())
|
||||
|
||||
def test_orphan_run_dirs_excludes_live_and_missing_root(self):
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
run_root = Path(tmp)
|
||||
(run_root / "live-a").mkdir()
|
||||
(run_root / "dead-b").mkdir()
|
||||
live = {str(run_root / "live-a")}
|
||||
self.assertEqual(
|
||||
[str(run_root / "dead-b")],
|
||||
fc_cleanup._orphan_run_dirs(run_root, live),
|
||||
)
|
||||
# absent run root -> nothing to reap
|
||||
self.assertEqual([], fc_cleanup._orphan_run_dirs(Path("/nope/run"), set()))
|
||||
|
||||
def test_prepare_cleanup_reaps_orphans_only(self):
|
||||
"""The live VM's dir is never in the plan; the dead one is."""
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
run_root = Path(tmp)
|
||||
(run_root / "live-a").mkdir()
|
||||
(run_root / "dead-b").mkdir()
|
||||
out = f"111 firecracker --config-file {run_root}/live-a/config.json\n"
|
||||
with patch.object(fc_cleanup, "_run_root", return_value=run_root), \
|
||||
patch.object(fc_cleanup.subprocess, "run", return_value=_proc(out)):
|
||||
plan = fc_cleanup.prepare_cleanup()
|
||||
self.assertEqual((), plan.vm_pids)
|
||||
self.assertEqual((str(run_root / "dead-b"),), plan.run_dirs)
|
||||
self.assertNotIn(str(run_root / "live-a"), plan.run_dirs)
|
||||
def test_prepare_cleanup_assembles_plan(self):
|
||||
with patch.object(fc_cleanup, "_orphan_vm_pids", return_value=[7]), \
|
||||
patch.object(fc_cleanup, "_run_dirs", return_value=["/run/x"]):
|
||||
plan = fc_cleanup.prepare_cleanup()
|
||||
self.assertEqual((7,), plan.vm_pids)
|
||||
self.assertEqual(("/run/x",), plan.run_dirs)
|
||||
|
||||
|
||||
class TestCleanupRemoval(unittest.TestCase):
|
||||
|
||||
@@ -49,7 +49,6 @@ def _plan(
|
||||
agent_git_gate_url: str = "",
|
||||
agent_supervise_url: str = "",
|
||||
image_policy: str = "fresh",
|
||||
env_var_secret: str = "",
|
||||
) -> MacosContainerBottlePlan:
|
||||
routes_path = stage_dir / "routes.yaml"
|
||||
routes_path.write_text("routes: []\n", encoding="utf-8")
|
||||
@@ -81,7 +80,6 @@ def _plan(
|
||||
),
|
||||
agent_git_gate_url=agent_git_gate_url,
|
||||
agent_supervise_url=agent_supervise_url,
|
||||
env_var_secret=env_var_secret,
|
||||
))
|
||||
|
||||
|
||||
@@ -187,12 +185,6 @@ class TestAgentRunArgv(unittest.TestCase):
|
||||
"bot-bottle-mac-gateway", self.argv[self.argv.index("--network") + 1],
|
||||
)
|
||||
|
||||
def test_env_var_secret_is_in_configured_container_environment(self) -> None:
|
||||
argv = _agent_run_argv(
|
||||
_plan(Path(self._tmp.name), env_var_secret="key-material"), _endpoint(),
|
||||
)
|
||||
self.assertIn("ENV_VAR_SECRET=key-material", argv)
|
||||
|
||||
def test_never_pins_an_ip(self) -> None:
|
||||
"""Apple Container 1.0.0 has no --ip: the address is DHCP-assigned and
|
||||
read back after start."""
|
||||
|
||||
@@ -28,21 +28,6 @@ class TestMacosContainerAvailability(unittest.TestCase):
|
||||
|
||||
|
||||
class TestMacosContainerCommands(unittest.TestCase):
|
||||
def test_read_container_env(self):
|
||||
completed = util.subprocess.CompletedProcess(
|
||||
args=[], returncode=0, stdout="secret\n", stderr="",
|
||||
)
|
||||
with patch.object(util, "_run_container_op", return_value=completed) as run:
|
||||
self.assertEqual("secret", util.read_container_env("bottle", "KEY"))
|
||||
run.assert_called_once_with(["container", "exec", "bottle", "printenv", "KEY"])
|
||||
|
||||
def test_read_container_env_returns_empty_on_failure(self):
|
||||
completed = util.subprocess.CompletedProcess(
|
||||
args=[], returncode=1, stdout="", stderr="missing",
|
||||
)
|
||||
with patch.object(util, "_run_container_op", return_value=completed):
|
||||
self.assertEqual("", util.read_container_env("bottle", "KEY"))
|
||||
|
||||
def test_dns_server_prefers_direct_host_ipv4_resolver(self):
|
||||
scutil = util.subprocess.CompletedProcess(
|
||||
args=[],
|
||||
|
||||
@@ -76,27 +76,6 @@ class TestTeardown(unittest.TestCase):
|
||||
self.assertEqual("DELETE", m.call_args.args[0].get_method())
|
||||
|
||||
|
||||
class TestReprovisionGateway(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.c = OrchestratorClient("http://orch:8080")
|
||||
|
||||
def test_success_posts_key(self) -> None:
|
||||
with patch(_URLOPEN, return_value=_resp(200, {"reprovisioned": True})) as opened:
|
||||
self.assertTrue(self.c.reprovision_gateway("b1", "key"))
|
||||
request = opened.call_args.args[0]
|
||||
self.assertEqual("POST", request.get_method())
|
||||
self.assertEqual({"env_var_secret": "key"}, json.loads(request.data))
|
||||
|
||||
def test_missing_stored_secret_is_false(self) -> None:
|
||||
with patch(_URLOPEN, side_effect=_http_error(404)):
|
||||
self.assertFalse(self.c.reprovision_gateway("b1", "key"))
|
||||
|
||||
def test_other_status_raises(self) -> None:
|
||||
with patch(_URLOPEN, side_effect=_http_error(400)):
|
||||
with self.assertRaises(OrchestratorClientError):
|
||||
self.c.reprovision_gateway("b1", "key")
|
||||
|
||||
|
||||
class TestHealthAndPolicy(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.c = OrchestratorClient("http://orch:8080")
|
||||
|
||||
@@ -6,7 +6,6 @@ server tests), plus one real-socket round-trip to prove the handler wiring.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import json
|
||||
import secrets
|
||||
import sqlite3
|
||||
@@ -64,43 +63,6 @@ class TestDispatch(unittest.TestCase):
|
||||
self.assertTrue(payload["bottle_id"])
|
||||
self.assertTrue(payload["identity_token"])
|
||||
|
||||
def test_register_and_reprovision_encrypted_tokens(self) -> None:
|
||||
key = base64.urlsafe_b64encode(b"unit-test-key").rstrip(b"=").decode()
|
||||
status, payload = dispatch(
|
||||
self.orch, "POST", "/bottles", _body({
|
||||
"source_ip": "10.243.0.11",
|
||||
"tokens": {"EGRESS_TOKEN_0": "upstream-secret"},
|
||||
"env_var_secret": key,
|
||||
}),
|
||||
)
|
||||
self.assertEqual(201, status)
|
||||
bottle_id = payload["bottle_id"]
|
||||
assert isinstance(bottle_id, str)
|
||||
self.orch._tokens.clear()
|
||||
status, response = dispatch(
|
||||
self.orch, "POST", f"/bottles/{bottle_id}/reprovision_gateway",
|
||||
_body({"env_var_secret": key}),
|
||||
)
|
||||
self.assertEqual((200, {"reprovisioned": True}), (status, response))
|
||||
self.assertEqual(
|
||||
{"EGRESS_TOKEN_0": "upstream-secret"}, self.orch.tokens_for(bottle_id),
|
||||
)
|
||||
|
||||
def test_reprovision_validates_request_and_missing_rows(self) -> None:
|
||||
status, _ = dispatch(
|
||||
self.orch, "POST", "/bottles/b1/reprovision_gateway", b"not-json",
|
||||
)
|
||||
self.assertEqual(400, status)
|
||||
status, _ = dispatch(
|
||||
self.orch, "POST", "/bottles/b1/reprovision_gateway", _body({}),
|
||||
)
|
||||
self.assertEqual(400, status)
|
||||
status, _ = dispatch(
|
||||
self.orch, "POST", "/bottles/b1/reprovision_gateway",
|
||||
_body({"env_var_secret": "key"}),
|
||||
)
|
||||
self.assertEqual(404, status)
|
||||
|
||||
def test_register_requires_source_ip(self) -> None:
|
||||
status, _ = dispatch(self.orch, "POST", "/bottles", _body({}))
|
||||
self.assertEqual(400, status)
|
||||
|
||||
@@ -173,58 +173,6 @@ if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
|
||||
class TestAgentSecrets(unittest.TestCase):
|
||||
"""store/get/delete for the bottled_agent_secrets table."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = Path(self._tmp.name) / "registry.db"
|
||||
self.store = RegistryStore(self.db)
|
||||
self.store.migrate()
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def test_store_and_get_roundtrip(self) -> None:
|
||||
self.store.store_agent_secrets("bottle-1", {"EGRESS_TOKEN_1": "enc-val-a"})
|
||||
got = self.store.get_agent_secrets("bottle-1")
|
||||
self.assertEqual({"EGRESS_TOKEN_1": "enc-val-a"}, got)
|
||||
|
||||
def test_get_returns_empty_when_none_stored(self) -> None:
|
||||
self.assertEqual({}, self.store.get_agent_secrets("no-such-bottle"))
|
||||
|
||||
def test_store_replaces_existing_rows(self) -> None:
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "old"})
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "new", "K2": "v2"})
|
||||
got = self.store.get_agent_secrets("bottle-1")
|
||||
self.assertEqual({"K": "new", "K2": "v2"}, got)
|
||||
|
||||
def test_delete_removes_secrets(self) -> None:
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "v"})
|
||||
self.store.delete_agent_secrets("bottle-1")
|
||||
self.assertEqual({}, self.store.get_agent_secrets("bottle-1"))
|
||||
|
||||
def test_delete_is_idempotent_on_missing(self) -> None:
|
||||
self.store.delete_agent_secrets("no-such-bottle") # must not raise
|
||||
|
||||
def test_secrets_are_isolated_by_bottle_id(self) -> None:
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "for-1"})
|
||||
self.store.store_agent_secrets("bottle-2", {"K": "for-2"})
|
||||
self.assertEqual({"K": "for-1"}, self.store.get_agent_secrets("bottle-1"))
|
||||
self.assertEqual({"K": "for-2"}, self.store.get_agent_secrets("bottle-2"))
|
||||
|
||||
def test_secrets_isolated_by_type(self) -> None:
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "injected"}, secret_type="injected_env_var")
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "other"}, secret_type="other_type")
|
||||
self.assertEqual({"K": "injected"}, self.store.get_agent_secrets("bottle-1"))
|
||||
self.assertEqual({"K": "other"}, self.store.get_agent_secrets("bottle-1", secret_type="other_type"))
|
||||
|
||||
def test_secrets_persist_across_reopen(self) -> None:
|
||||
self.store.store_agent_secrets("bottle-1", {"K": "v"})
|
||||
reopened = RegistryStore(self.db)
|
||||
self.assertEqual({"K": "v"}, reopened.get_agent_secrets("bottle-1"))
|
||||
|
||||
|
||||
class TestReapAbsent(unittest.TestCase):
|
||||
"""`reap_absent` — the self-heal for rows whose bottle is gone.
|
||||
|
||||
|
||||
@@ -1,96 +0,0 @@
|
||||
"""Unit tests for per-bottle egress secret encryption (PRD prd-new-secret-provider)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import unittest
|
||||
|
||||
from bot_bottle.orchestrator.secret_store import (
|
||||
ENV_VAR_SECRET_NAME,
|
||||
decrypt_value,
|
||||
encrypt_value,
|
||||
new_env_var_secret,
|
||||
)
|
||||
|
||||
|
||||
class TestNewEnvVarSecret(unittest.TestCase):
|
||||
def test_returns_non_empty_string(self) -> None:
|
||||
s = new_env_var_secret()
|
||||
self.assertIsInstance(s, str)
|
||||
self.assertTrue(len(s) > 0)
|
||||
|
||||
def test_secrets_are_unique(self) -> None:
|
||||
keys = {new_env_var_secret() for _ in range(50)}
|
||||
self.assertEqual(50, len(keys))
|
||||
|
||||
def test_no_padding_characters(self) -> None:
|
||||
# URL-safe base64, padding stripped — should round-trip cleanly
|
||||
for _ in range(20):
|
||||
self.assertNotIn("=", new_env_var_secret())
|
||||
|
||||
|
||||
class TestEncryptDecryptRoundtrip(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.secret = new_env_var_secret()
|
||||
|
||||
def _rt(self, plaintext: str) -> str:
|
||||
return decrypt_value(self.secret, encrypt_value(self.secret, plaintext))
|
||||
|
||||
def test_roundtrip_short_value(self) -> None:
|
||||
self.assertEqual("sk-abc123", self._rt("sk-abc123"))
|
||||
|
||||
def test_roundtrip_empty_string(self) -> None:
|
||||
self.assertEqual("", self._rt(""))
|
||||
|
||||
def test_roundtrip_long_value_crosses_block_boundary(self) -> None:
|
||||
# 32 bytes is exactly one HMAC-SHA256 block; 65 bytes crosses two.
|
||||
plaintext = "x" * 65
|
||||
self.assertEqual(plaintext, self._rt(plaintext))
|
||||
|
||||
def test_roundtrip_unicode(self) -> None:
|
||||
self.assertEqual("héllo wörld", self._rt("héllo wörld"))
|
||||
|
||||
def test_encrypt_produces_different_ciphertexts_each_call(self) -> None:
|
||||
ct1 = encrypt_value(self.secret, "same")
|
||||
ct2 = encrypt_value(self.secret, "same")
|
||||
self.assertNotEqual(ct1, ct2) # fresh nonce each call
|
||||
|
||||
def test_ciphertext_is_url_safe_base64(self) -> None:
|
||||
ct = encrypt_value(self.secret, "hello")
|
||||
# no '+', '/', '=' — URL-safe and padding-stripped
|
||||
for ch in ("+", "/", "="):
|
||||
self.assertNotIn(ch, ct)
|
||||
|
||||
|
||||
class TestDecryptErrors(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.secret = new_env_var_secret()
|
||||
|
||||
def test_wrong_key_raises_value_error(self) -> None:
|
||||
ct = encrypt_value(self.secret, "secret-token")
|
||||
other_key = new_env_var_secret()
|
||||
# Wrong key produces garbage bytes; decrypt_value raises ValueError
|
||||
# when the result is non-UTF-8 (which is very likely for 12-char data).
|
||||
# We allow it to succeed only if garbage happens to be valid UTF-8, but
|
||||
# the plaintext must not match.
|
||||
try:
|
||||
result = decrypt_value(other_key, ct)
|
||||
self.assertNotEqual("secret-token", result)
|
||||
except ValueError:
|
||||
pass
|
||||
|
||||
def test_truncated_blob_raises_value_error(self) -> None:
|
||||
with self.assertRaises(ValueError):
|
||||
decrypt_value(self.secret, "dG9vc2hvcnQ") # "tooshort" — under 16 nonce bytes
|
||||
|
||||
def test_invalid_base64_raises_value_error(self) -> None:
|
||||
with self.assertRaises(ValueError):
|
||||
decrypt_value(self.secret, "!!not-base64!!")
|
||||
|
||||
|
||||
class TestConstant(unittest.TestCase):
|
||||
def test_env_var_secret_name(self) -> None:
|
||||
self.assertEqual("ENV_VAR_SECRET", ENV_VAR_SECRET_NAME)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -15,7 +15,6 @@ from bot_bottle.orchestrator.broker import LaunchBroker, LaunchRequest, StubBrok
|
||||
from bot_bottle.orchestrator.registry import RegistryStore
|
||||
from bot_bottle.orchestrator.service import Orchestrator
|
||||
from bot_bottle.orchestrator.gateway import Gateway
|
||||
from bot_bottle.orchestrator.secret_store import new_env_var_secret
|
||||
from bot_bottle.store_manager import StoreManager
|
||||
from bot_bottle.supervise import (
|
||||
Proposal,
|
||||
@@ -118,25 +117,6 @@ class TestOrchestrator(unittest.TestCase):
|
||||
rec = self.orch.launch_bottle("10.243.0.6")
|
||||
self.assertEqual({}, self.orch.tokens_for(rec.bottle_id))
|
||||
|
||||
def test_encrypted_tokens_can_be_reprovisioned_after_memory_loss(self) -> None:
|
||||
key = new_env_var_secret()
|
||||
rec = self.orch.launch_bottle(
|
||||
"10.243.0.12", tokens={"EGRESS_TOKEN_0": "secret"},
|
||||
env_var_secret=key,
|
||||
)
|
||||
self.assertNotEqual({}, self.store.get_agent_secrets(rec.bottle_id))
|
||||
self.orch._tokens.clear()
|
||||
self.assertTrue(self.orch.reprovision_from_secret(rec.bottle_id, key))
|
||||
self.assertEqual({"EGRESS_TOKEN_0": "secret"}, self.orch.tokens_for(rec.bottle_id))
|
||||
|
||||
def test_reprovision_rejects_missing_rows_and_wrong_key(self) -> None:
|
||||
self.assertFalse(self.orch.reprovision_from_secret("missing", new_env_var_secret()))
|
||||
rec = self.orch.launch_bottle(
|
||||
"10.243.0.13", tokens={"K": "value"},
|
||||
env_var_secret=new_env_var_secret(),
|
||||
)
|
||||
self.assertFalse(self.orch.reprovision_from_secret(rec.bottle_id, new_env_var_secret()))
|
||||
|
||||
def test_set_policy_live_reload(self) -> None:
|
||||
rec = self.orch.launch_bottle("10.243.0.3")
|
||||
self.assertTrue(self.orch.set_policy(rec.bottle_id, '{"x":1}'))
|
||||
|
||||
Reference in New Issue
Block a user