diff --git a/bot_bottle/backend/consolidated_util.py b/bot_bottle/backend/consolidated_util.py index 354212e..520ca0a 100644 --- a/bot_bottle/backend/consolidated_util.py +++ b/bot_bottle/backend/consolidated_util.py @@ -7,10 +7,13 @@ imports it rather than re-implementing it. from __future__ import annotations +import dataclasses + from ..egress import EgressPlan from ..git_gate import GitGatePlan -from ..orchestrator.client import OrchestratorClient +from ..orchestrator.client import OrchestratorClient, RegisteredBottle from ..orchestrator.registration import registration_inputs +from ..orchestrator.secret_store import new_env_var_secret from .docker.gateway_provision import GatewayTransport, deprovision_git_gate, provision_git_gate @@ -23,21 +26,27 @@ 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. Returns the - `RegisteredBottle` from the orchestrator.""" + registration if provisioning fails so no orphan is left. + + Generates a fresh ENV_VAR_SECRET, passes it to the orchestrator so it can + encrypt the token values at rest, and stamps the secret onto the returned + ``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() reg = client.register_bottle( source_ip, image_ref=image_ref, policy=inputs.policy, - metadata=inputs.metadata, tokens=tokens, + metadata=inputs.metadata, tokens=tokens, env_var_secret=env_var_secret, ) try: provision_git_gate(transport, reg.bottle_id, git_gate_plan) except Exception: client.teardown_bottle(reg.bottle_id) raise - return reg + return dataclasses.replace(reg, env_var_secret=env_var_secret) def teardown_consolidated( diff --git a/bot_bottle/backend/docker/bottle_plan.py b/bot_bottle/backend/docker/bottle_plan.py index 55cee06..a327137 100644 --- a/bot_bottle/backend/docker/bottle_plan.py +++ b/bot_bottle/backend/docker/bottle_plan.py @@ -39,6 +39,10 @@ class DockerBottlePlan(BottlePlan): # (egress proxy credentials, git-gate/supervise headers); set by launch # from the orchestrator registration. Empty pre-registration. identity_token: str = "" + # Encryption key for the agent's stored egress secrets; injected into the + # agent container as ENV_VAR_SECRET via the compose subprocess env (bare + # name — value never written to the compose file). Empty pre-registration. + env_var_secret: str = "" @property def container_name(self) -> str: diff --git a/bot_bottle/backend/docker/consolidated_compose.py b/bot_bottle/backend/docker/consolidated_compose.py index fde30c9..0065224 100644 --- a/bot_bottle/backend/docker/consolidated_compose.py +++ b/bot_bottle/backend/docker/consolidated_compose.py @@ -17,6 +17,7 @@ from __future__ import annotations from typing import Any from ...egress import egress_agent_env_entries +from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME from ..util import AGENT_CA_BUNDLE, AGENT_CA_PATH from .bottle_plan import DockerBottlePlan from .egress import EGRESS_PORT @@ -58,6 +59,10 @@ def consolidated_agent_compose( # the secret value never lands on argv or in the compose file. for name in sorted(plan.forwarded_env.keys()): env.append(name) + # ENV_VAR_SECRET: bare name so the value comes from the compose subprocess + # env (set in launch.py) and is never written to the compose file on disk. + if getattr(plan, "env_var_secret", ""): + env.append(ENV_VAR_SECRET_NAME) env.extend(egress_agent_env_entries(plan.egress_plan)) service: dict[str, Any] = { diff --git a/bot_bottle/backend/docker/consolidated_launch.py b/bot_bottle/backend/docker/consolidated_launch.py index 15555b5..9f9768a 100644 --- a/bot_bottle/backend/docker/consolidated_launch.py +++ b/bot_bottle/backend/docker/consolidated_launch.py @@ -15,12 +15,15 @@ from __future__ import annotations from dataclasses import dataclass +from ... import log from ...docker_cmd import run_docker from ...egress import EgressPlan from ...git_gate import GitGatePlan 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 @@ -41,6 +44,7 @@ class LaunchContext: network: str # the shared gateway network to attach to gateway_ip: str # the gateway's address — the agent's proxy target orchestrator_url: str + env_var_secret: str = "" # encryption key injected into the agent's env def _network_cidr(network: str) -> str: @@ -85,6 +89,55 @@ def _network_container_ips(network: str) -> list[str]: return ips +def _reprovision_running_bottles( + orchestrator_url: str, + network: str = GATEWAY_NETWORK, + infra_name: str = INFRA_NAME, +) -> None: + """Re-inject egress tokens for any registered bottles that lost their + in-memory tokens (e.g., after an infra container restart). + + For each registered bottle whose source IP maps to a live container on the + gateway network, reads ENV_VAR_SECRET via ``docker exec … printenv`` and + calls ``POST /bottles//reprovision_gateway``. Idempotent — a no-op + 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) + # 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 + ip_to_container: dict[str, str] = {} + for line in proc.stdout.splitlines(): + parts = line.strip().split() + if len(parts) >= 2 and parts[0] != infra_name: + ip = parts[1].split("/", 1)[0] + if ip: + ip_to_container[ip] = parts[0] + + secrets_by_ip: dict[str, str] = {} + for source_ip, container_name in ip_to_container.items(): + 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() + + reprovisioned = reprovision_bottles(client, secrets_by_ip) + if reprovisioned: + log.info( + "reprovisioned egress tokens", + context={"count": reprovisioned}, + ) + + def launch_consolidated( egress_plan: EgressPlan, git_gate_plan: GitGatePlan, @@ -96,9 +149,14 @@ def launch_consolidated( network: str = GATEWAY_NETWORK, ) -> LaunchContext: """Ensure the infra container is up, allocate + register the bottle, and - provision its git-gate state. Returns the agent's attach context.""" + provision its git-gate state. Returns the agent's attach context. + + Also reprovisiones egress tokens for any already-running bottles that lost + their in-memory credentials (e.g. after an infra container restart), so + they regain egress access before the new bottle is registered.""" service = service or OrchestratorService() url = service.ensure_running() + _reprovision_running_bottles(url, network=network, infra_name=infra_name) client = OrchestratorClient(url) cidr = _network_cidr(network) @@ -117,6 +175,7 @@ def launch_consolidated( network=network, gateway_ip=gateway_ip, orchestrator_url=url, + env_var_secret=reg.env_var_secret, ) diff --git a/bot_bottle/backend/docker/launch.py b/bot_bottle/backend/docker/launch.py index 8f08fbe..e767d92 100644 --- a/bot_bottle/backend/docker/launch.py +++ b/bot_bottle/backend/docker/launch.py @@ -186,6 +186,7 @@ def launch( agent_git_gate_url=git_gate_url, agent_supervise_url=supervise_url, identity_token=ctx.identity_token, + env_var_secret=ctx.env_var_secret, ) # Step 5: render + up the agent-only compose, pinned on the shared @@ -198,7 +199,12 @@ def launch( project = compose_project_name(plan.slug) # Forwarded vars (OAuth token, host interpolations) flow through the # subprocess env as bare names so values never land in the file. + # ENV_VAR_SECRET follows the same pattern: bare name in the compose + # spec, value only in the subprocess env so it is never written to disk. compose_env: dict[str, str] = {**os.environ, **plan.forwarded_env} + if plan.env_var_secret: + from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME + compose_env[ENV_VAR_SECRET_NAME] = plan.env_var_secret info( f"docker compose up -d (project {project}, agent on shared " f"gateway {ctx.gateway_ip}, ip {ctx.source_ip})" diff --git a/bot_bottle/backend/firecracker/bottle_plan.py b/bot_bottle/backend/firecracker/bottle_plan.py index 92c17c4..5a9603a 100644 --- a/bot_bottle/backend/firecracker/bottle_plan.py +++ b/bot_bottle/backend/firecracker/bottle_plan.py @@ -22,6 +22,9 @@ 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: diff --git a/bot_bottle/backend/firecracker/cleanup.py b/bot_bottle/backend/firecracker/cleanup.py index 6e9fe2d..bce6c9d 100644 --- a/bot_bottle/backend/firecracker/cleanup.py +++ b/bot_bottle/backend/firecracker/cleanup.py @@ -88,6 +88,12 @@ def _scan_processes(run_root: Path) -> tuple[set[str], list[int]]: return live, orphan_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.""" if not run_root.is_dir(): diff --git a/bot_bottle/backend/firecracker/consolidated_launch.py b/bot_bottle/backend/firecracker/consolidated_launch.py index 3e25917..cb72690 100644 --- a/bot_bottle/backend/firecracker/consolidated_launch.py +++ b/bot_bottle/backend/firecracker/consolidated_launch.py @@ -25,16 +25,27 @@ 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 ...orchestrator.client import OrchestratorClient +from ...log import info +from ...orchestrator.client import OrchestratorClient, OrchestratorClientError from ...orchestrator.lifecycle import ( OrchestratorStartError, # re-exported so callers can catch it ) -from ..consolidated_util import provision_bottle, teardown_consolidated as _teardown_util -from . import infra_vm +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" class ConsolidatedLaunchError(RuntimeError): @@ -50,6 +61,55 @@ class LaunchContext: source_ip: str # the VM's guest IP — the attribution key gateway_ca_pem: str # the shared gateway CA the provisioner installs orchestrator_url: str + 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 ''}" + ) + + +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( @@ -66,6 +126,7 @@ 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( @@ -80,6 +141,7 @@ def launch_consolidated( source_ip=guest_ip, gateway_ca_pem=infra.gateway_ca_pem(), orchestrator_url=url, + env_var_secret=reg.env_var_secret, ) diff --git a/bot_bottle/backend/firecracker/launch.py b/bot_bottle/backend/firecracker/launch.py index a76f9ef..d72f40c 100644 --- a/bot_bottle/backend/firecracker/launch.py +++ b/bot_bottle/backend/firecracker/launch.py @@ -55,8 +55,10 @@ 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, ) @@ -153,6 +155,7 @@ 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, @@ -187,6 +190,7 @@ 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. @@ -281,6 +285,8 @@ 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 diff --git a/bot_bottle/backend/firecracker/util.py b/bot_bottle/backend/firecracker/util.py index 0d9aad8..465c348 100644 --- a/bot_bottle/backend/firecracker/util.py +++ b/bot_bottle/backend/firecracker/util.py @@ -399,6 +399,9 @@ 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 & diff --git a/bot_bottle/backend/macos_container/bottle_plan.py b/bot_bottle/backend/macos_container/bottle_plan.py index b167685..96fc781 100644 --- a/bot_bottle/backend/macos_container/bottle_plan.py +++ b/bot_bottle/backend/macos_container/bottle_plan.py @@ -23,6 +23,9 @@ class MacosContainerBottlePlan(BottlePlan): # Guest-local container engine (issue #392). Gates the derived image, the # device-mode relaxation, and the resident podman service. nested_containers: bool = False + # 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: diff --git a/bot_bottle/backend/macos_container/consolidated_launch.py b/bot_bottle/backend/macos_container/consolidated_launch.py index ebc97cb..6c25f5e 100644 --- a/bot_bottle/backend/macos_container/consolidated_launch.py +++ b/bot_bottle/backend/macos_container/consolidated_launch.py @@ -38,7 +38,12 @@ from ...egress import EgressPlan from ...git_gate import GitGatePlan from ...log import info from ...orchestrator.client import OrchestratorClient, OrchestratorClientError -from ..consolidated_util import provision_bottle, teardown_consolidated as _teardown_util +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 util as container_mod from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active from .gateway import GATEWAY_NETWORK @@ -72,6 +77,7 @@ class LaunchContext: gateway_ip: str network: str orchestrator_url: str + env_var_secret: str = "" # encryption key injected into the agent's env def ensure_gateway( @@ -83,12 +89,35 @@ def ensure_gateway( needs `gateway_ip` at run time.""" service = service or MacosInfraService() infra = service.ensure_running() - return GatewayEndpoint( + endpoint = 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]: @@ -125,6 +154,7 @@ 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 @@ -144,6 +174,7 @@ 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, @@ -152,6 +183,7 @@ def register_agent( gateway_ip=endpoint.gateway_ip, network=endpoint.network, orchestrator_url=endpoint.orchestrator_url, + env_var_secret=reg.env_var_secret, ) diff --git a/bot_bottle/backend/macos_container/launch.py b/bot_bottle/backend/macos_container/launch.py index 45b0a76..8902266 100644 --- a/bot_bottle/backend/macos_container/launch.py +++ b/bot_bottle/backend/macos_container/launch.py @@ -69,6 +69,7 @@ from .gateway_hosts import ( from . import nested_containers as nested_containers_mod 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, @@ -169,6 +170,7 @@ 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. @@ -203,6 +205,7 @@ 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, @@ -443,6 +446,8 @@ 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 diff --git a/bot_bottle/backend/macos_container/util.py b/bot_bottle/backend/macos_container/util.py index 404d9c3..23a6403 100644 --- a/bot_bottle/backend/macos_container/util.py +++ b/bot_bottle/backend/macos_container/util.py @@ -361,6 +361,12 @@ 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. diff --git a/bot_bottle/orchestrator/client.py b/bot_bottle/orchestrator/client.py index 86c919c..b3128bf 100644 --- a/bot_bottle/orchestrator/client.py +++ b/bot_bottle/orchestrator/client.py @@ -41,10 +41,13 @@ class OrchestratorClientError(RuntimeError): @dataclass(frozen=True) class RegisteredBottle: """What `POST /bottles` returns: the minted bottle id and the per-bottle - identity token the agent presents for app-layer attribution.""" + identity token the agent presents for app-layer attribution. `env_var_secret` + is set by the caller (not from the server response) and carries the + encryption key so it can be injected into the agent container's env.""" bottle_id: str identity_token: str + env_var_secret: str = "" class OrchestratorClient: @@ -120,17 +123,21 @@ class OrchestratorClient: metadata: str = "", policy: str = "", tokens: dict[str, str] | None = None, + env_var_secret: str = "", ) -> RegisteredBottle: """Register a bottle and broker its launch (`POST /bottles`). `tokens` are the per-bottle egress auth values (env_name -> value) the - orchestrator holds in memory for the gateway to inject. Returns the - minted id + identity token.""" + orchestrator holds in memory for the gateway to inject. When + *env_var_secret* is provided, the orchestrator also encrypts the token + values and stores them in ``bottled_agent_secrets`` for restart + recovery. Returns the minted id + identity token.""" payload = self._ok("POST", "/bottles", { "source_ip": source_ip, "image_ref": image_ref, "metadata": metadata, "policy": policy, "tokens": tokens or {}, + "env_var_secret": env_var_secret, }) bottle_id = payload.get("bottle_id") token = payload.get("identity_token") @@ -138,6 +145,24 @@ class OrchestratorClient: raise OrchestratorClientError("register: response missing bottle_id/identity_token") return RegisteredBottle(bottle_id=bottle_id, identity_token=token) + def reprovision_gateway(self, bottle_id: str, env_var_secret: str) -> bool: + """Re-inject a bottle's egress tokens from its ENV_VAR_SECRET + (`POST /bottles//reprovision_gateway`). Returns True when the + orchestrator successfully decrypted and restored the tokens, False + when it had no stored secrets for this bottle (404).""" + status, _ = self._request( + "POST", + f"/bottles/{bottle_id}/reprovision_gateway", + {"env_var_secret": env_var_secret}, + ) + if status == 404: + return False + if not 200 <= status < 300: + raise OrchestratorClientError( + f"reprovision_gateway {bottle_id}: HTTP {status}" + ) + return True + def teardown_bottle(self, bottle_id: str) -> bool: """Tear a bottle down (`DELETE /bottles/`). False if the orchestrator didn't know it (404) — idempotent for cleanup paths.""" diff --git a/bot_bottle/orchestrator/control_plane.py b/bot_bottle/orchestrator/control_plane.py index 06ad8c7..1625cec 100644 --- a/bot_bottle/orchestrator/control_plane.py +++ b/bot_bottle/orchestrator/control_plane.py @@ -9,9 +9,13 @@ vsock / unix-socket portability caveats): GET /bottles -> 200 {"bottles": [ , ...]} POST /bottles -> 201 {"bottle_id","identity_token"} (launch) body: {"source_ip", ["image_ref"], - ["metadata"], ["policy"]} + ["metadata"], ["policy"], + ["tokens"], ["env_var_secret"]} PUT /bottles//policy -> 200 {"updated": true} | 404 (live reload) body: {"policy"} + POST /bottles//reprovision_gateway + -> 200 {"reprovisioned": true} | 404 + body: {"env_var_secret"} DELETE /bottles/ -> 200 {"torn_down": true} | 404 (teardown) POST /reconcile -> 200 {"reaped": [bottle_id, ...]} body: {"live_source_ips": [...], @@ -116,12 +120,14 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches tokens = { k: v for k, v in raw_tokens.items() if isinstance(k, str) and isinstance(v, str) } if isinstance(raw_tokens, dict) else {} + env_var_secret = data.get("env_var_secret", "") rec = orch.launch_bottle( source_ip, image_ref=image_ref if isinstance(image_ref, str) else "", metadata=metadata if isinstance(metadata, str) else "", policy=policy if isinstance(policy, str) else "", tokens=tokens, + env_var_secret=env_var_secret if isinstance(env_var_secret, str) else "", ) return 201, {"bottle_id": rec.bottle_id, "identity_token": rec.identity_token} @@ -138,6 +144,23 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches return 200, {"updated": True} return 404, {"error": "no such bottle"} + if ( + method == "POST" + and route.startswith("/bottles/") + and route.endswith("/reprovision_gateway") + ): + bottle_id = route[len("/bottles/") : -len("/reprovision_gateway")] + try: + data = _parse_json_object(body) + except ValueError as e: + return 400, {"error": f"invalid JSON: {e}"} + env_var_secret = data.get("env_var_secret") + if not isinstance(env_var_secret, str) or not env_var_secret: + return 400, {"error": "env_var_secret (string) is required"} + if orch.reprovision_from_secret(bottle_id, env_var_secret): + return 200, {"reprovisioned": True} + return 404, {"error": "no stored secrets for this bottle"} + if method == "DELETE" and route.startswith("/bottles/"): bottle_id = route[len("/bottles/"):] if orch.teardown_bottle(bottle_id): diff --git a/bot_bottle/orchestrator/registry.py b/bot_bottle/orchestrator/registry.py index 06c6d74..a7f34b1 100644 --- a/bot_bottle/orchestrator/registry.py +++ b/bot_bottle/orchestrator/registry.py @@ -113,6 +113,22 @@ _MIGRATIONS = TableMigrations( # egress allowlist / routes / git config selected by source IP. The # multi-tenant gateway resolves it per request via `attribute`. "ALTER TABLE orchestrator_bottles ADD COLUMN policy TEXT NOT NULL DEFAULT ''", + # v4 — per-bottle encrypted egress secrets (PRD prd-new-secret-provider). + # One row per env-var: key (env-var name) is plaintext for auditing; + # value is the encrypted token string. The encryption key (ENV_VAR_SECRET) + # lives only in the agent's environment — a row alone cannot recover the + # credential. + """ + CREATE TABLE IF NOT EXISTS bottled_agent_secrets ( + bottled_agent_id TEXT NOT NULL, + key TEXT NOT NULL, + value TEXT NOT NULL, + type TEXT NOT NULL DEFAULT 'injected_env_var' + ) + """, + # v5 — index for fast per-bottle lookups and bulk DELETE on teardown. + "CREATE INDEX IF NOT EXISTS idx_bottled_agent_secrets_id " + "ON bottled_agent_secrets (bottled_agent_id, type)", ], ) @@ -326,6 +342,57 @@ class RegistryStore(DbStore): return None return rec + # --- encrypted egress secret store ------------------------------------ + + def store_agent_secrets( + self, + bottle_id: str, + encrypted_values: dict[str, str], + secret_type: str = "injected_env_var", + ) -> None: + """Replace all stored secrets for *bottle_id* with *encrypted_values* + (env-var name → encrypted ciphertext). Deletes then re-inserts so a + re-registration is always consistent with the current token set.""" + with self._connection() as conn: + conn.execute( + "DELETE FROM bottled_agent_secrets " + "WHERE bottled_agent_id = ? AND type = ?", + (bottle_id, secret_type), + ) + conn.executemany( + "INSERT INTO bottled_agent_secrets " + "(bottled_agent_id, key, value, type) VALUES (?, ?, ?, ?)", + [(bottle_id, k, v, secret_type) for k, v in encrypted_values.items()], + ) + self._chmod() + + def get_agent_secrets( + self, + bottle_id: str, + secret_type: str = "injected_env_var", + ) -> dict[str, str]: + """Return {env_var_name: encrypted_value} for *bottle_id*, or {} if none.""" + with self._connection() as conn: + rows = conn.execute( + "SELECT key, value FROM bottled_agent_secrets " + "WHERE bottled_agent_id = ? AND type = ?", + (bottle_id, secret_type), + ).fetchall() + return {row[0]: row[1] for row in rows} + + def delete_agent_secrets( + self, + bottle_id: str, + secret_type: str = "injected_env_var", + ) -> None: + """Remove all stored secrets for *bottle_id* (e.g. on teardown).""" + with self._connection() as conn: + conn.execute( + "DELETE FROM bottled_agent_secrets " + "WHERE bottled_agent_id = ? AND type = ?", + (bottle_id, secret_type), + ) + __all__ = [ "BottleRecord", diff --git a/bot_bottle/orchestrator/reprovision.py b/bot_bottle/orchestrator/reprovision.py new file mode 100644 index 0000000..e55ebea --- /dev/null +++ b/bot_bottle/orchestrator/reprovision.py @@ -0,0 +1,35 @@ +"""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"] diff --git a/bot_bottle/orchestrator/secret_store.py b/bot_bottle/orchestrator/secret_store.py new file mode 100644 index 0000000..a7a6e51 --- /dev/null +++ b/bot_bottle/orchestrator/secret_store.py @@ -0,0 +1,94 @@ +"""Symmetric encryption for per-bottle egress secrets (PRD prd-new-secret-provider). + +Each agent receives a random ENV_VAR_SECRET at startup — passed as an env var, +never logged or persisted. The host uses this key to encrypt each egress auth +token value before writing it to the bottled_agent_secrets table; the DB rows +(ciphertext, plaintext env-var name) without the key are insufficient to +recover the credentials. + +On orchestrator restart the in-memory token map is lost. The host-side +reattachment path reads ENV_VAR_SECRET from the running agent container via +``docker exec … printenv ENV_VAR_SECRET`` and posts it to +``POST /bottles//reprovision_gateway``; the orchestrator decrypts the +stored rows and re-populates ``_tokens``. + +Encryption scheme: HMAC-SHA256 used as a PRF in CTR mode (stdlib-only, +no external deps). Each value is encrypted independently. The output blob is +``nonce (16 bytes) || ciphertext`` encoded as URL-safe base64 (no padding). + + keystream_block_i = HMAC-SHA256(key, nonce || i.to_bytes(4, "big")) + ciphertext_i = plaintext_i XOR keystream_block_i[:len(plaintext_i)] +""" + +from __future__ import annotations + +import base64 +import hashlib +import hmac +import secrets + +_KEY_BYTES = 32 # 256-bit key from ENV_VAR_SECRET +_NONCE_BYTES = 16 # 128-bit random nonce per encrypt call +_BLOCK = 32 # HMAC-SHA256 output width == one keystream block + +# Env-var name the agent container receives at startup. +ENV_VAR_SECRET_NAME = "ENV_VAR_SECRET" + + +def new_env_var_secret() -> str: + """Generate a fresh ENV_VAR_SECRET: 32 random bytes as URL-safe base64.""" + return base64.urlsafe_b64encode(secrets.token_bytes(_KEY_BYTES)).rstrip(b"=").decode() + + +def _b64dec(s: str) -> bytes: + return base64.urlsafe_b64decode(s + "=" * (-len(s) % 4)) + + +def _keystream(key: bytes, nonce: bytes, block_index: int) -> bytes: + return hmac.new( + key, nonce + block_index.to_bytes(4, "big"), hashlib.sha256 + ).digest() + + +def encrypt_value(secret_b64: str, plaintext: str) -> str: + """Encrypt a single string value with *secret_b64* (the ENV_VAR_SECRET). + + Returns a URL-safe base64 blob ``nonce || ciphertext`` suitable for + the ``bottled_agent_secrets.value`` column.""" + key = _b64dec(secret_b64) + pt = plaintext.encode() + nonce = secrets.token_bytes(_NONCE_BYTES) + ct = bytearray() + for i in range(0, len(pt), _BLOCK): + chunk = pt[i : i + _BLOCK] + ks = _keystream(key, nonce, i)[: len(chunk)] + ct.extend(p ^ k for p, k in zip(chunk, ks)) + return base64.urlsafe_b64encode(nonce + bytes(ct)).rstrip(b"=").decode() + + +def decrypt_value(secret_b64: str, blob_b64: str) -> str: + """Decrypt a blob produced by :func:`encrypt_value`. + + Returns the original plaintext string. Raises ``ValueError`` for malformed + input or a key mismatch (wrong key produces garbage, not an error, unless + the plaintext is non-UTF-8 — treat all such failures as wrong key).""" + key = _b64dec(secret_b64) + try: + blob = _b64dec(blob_b64) + except Exception as exc: + raise ValueError(f"invalid ciphertext blob: {exc}") from exc + if len(blob) < _NONCE_BYTES: + raise ValueError("ciphertext blob too short") + nonce, ciphertext = blob[:_NONCE_BYTES], blob[_NONCE_BYTES:] + pt = bytearray() + for i in range(0, len(ciphertext), _BLOCK): + chunk = ciphertext[i : i + _BLOCK] + ks = _keystream(key, nonce, i)[: len(chunk)] + pt.extend(c ^ k for c, k in zip(chunk, ks)) + try: + return bytes(pt).decode() + except UnicodeDecodeError as exc: + raise ValueError(f"decryption produced non-UTF-8 output (wrong key?): {exc}") from exc + + +__all__ = ["ENV_VAR_SECRET_NAME", "new_env_var_secret", "encrypt_value", "decrypt_value"] diff --git a/bot_bottle/orchestrator/service.py b/bot_bottle/orchestrator/service.py index b449782..f1d2963 100644 --- a/bot_bottle/orchestrator/service.py +++ b/bot_bottle/orchestrator/service.py @@ -87,13 +87,22 @@ class Orchestrator: metadata: str = "", policy: str = "", tokens: dict[str, str] | None = None, + env_var_secret: str = "", ) -> BottleRecord: """Register a bottle (with its gateway policy + in-memory egress auth tokens) and broker its launch. Rolls the registry entry back if the - launch doesn't take, so a failure leaves no orphan.""" + launch doesn't take, so a failure leaves no orphan. + + When *env_var_secret* is provided alongside *tokens*, the token values + are also encrypted and written to ``bottled_agent_secrets`` so they can + survive an orchestrator restart (see ``reprovision_from_secret``).""" rec = self.registry.register(source_ip, metadata=metadata, policy=policy) if tokens: self._tokens[rec.bottle_id] = dict(tokens) + if env_var_secret: + from .secret_store import encrypt_value + encrypted = {k: encrypt_value(env_var_secret, v) for k, v in tokens.items()} + self.registry.store_agent_secrets(rec.bottle_id, encrypted) req = LaunchRequest( op="launch", bottle_id=rec.bottle_id, @@ -284,6 +293,26 @@ class Orchestrator: )) return True, "" + # --- secret reprovision ----------------------------------------------- + + def reprovision_from_secret(self, bottle_id: str, env_var_secret: str) -> bool: + """Re-inject a bottle's egress tokens from its ENV_VAR_SECRET. + + Reads the encrypted rows from ``bottled_agent_secrets``, decrypts each + value with *env_var_secret*, and restores ``_tokens[bottle_id]``. + Returns True on success, False when no stored secrets exist for this + bottle or decryption fails (wrong key / corrupt data).""" + from .secret_store import decrypt_value + encrypted = self.registry.get_agent_secrets(bottle_id) + if not encrypted: + return False + try: + self._tokens[bottle_id] = {k: decrypt_value(env_var_secret, v) + for k, v in encrypted.items()} + except ValueError: + return False + return True + # --- consolidated gateway ---------------------------------------------- def ensure_gateway(self) -> None: diff --git a/docs/prds/prd-new-secret-provider-encrypted-env.md b/docs/prds/prd-new-secret-provider-encrypted-env.md new file mode 100644 index 0000000..f01d451 --- /dev/null +++ b/docs/prds/prd-new-secret-provider-encrypted-env.md @@ -0,0 +1,131 @@ +# 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: }` reference surface. +- User-extensible provider discovery (`~/.bot-bottle/contrib//`). +- 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? diff --git a/tests/unit/test_backend_secret_reprovision.py b/tests/unit/test_backend_secret_reprovision.py new file mode 100644 index 0000000..33e9dc5 --- /dev/null +++ b/tests/unit/test_backend_secret_reprovision.py @@ -0,0 +1,160 @@ +"""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() diff --git a/tests/unit/test_consolidated_compose.py b/tests/unit/test_consolidated_compose.py index 33495e1..6a86867 100644 --- a/tests/unit/test_consolidated_compose.py +++ b/tests/unit/test_consolidated_compose.py @@ -2,6 +2,7 @@ from __future__ import annotations +import dataclasses import unittest from bot_bottle.backend.docker.consolidated_compose import consolidated_agent_compose @@ -46,6 +47,15 @@ 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() diff --git a/tests/unit/test_docker_launch_teardown.py b/tests/unit/test_docker_launch_teardown.py index 55474a5..38b8a6e 100644 --- a/tests/unit/test_docker_launch_teardown.py +++ b/tests/unit/test_docker_launch_teardown.py @@ -91,6 +91,7 @@ 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") @@ -107,7 +108,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"), \ + mock.patch.object(launch_mod, "compose_up") as compose_up, \ mock.patch.object(launch_mod, "compose_dump_logs"), \ mock.patch.object( launch_mod, "compose_down", @@ -122,6 +123,9 @@ 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__": diff --git a/tests/unit/test_firecracker_backend.py b/tests/unit/test_firecracker_backend.py index d72d795..3a1387b 100644 --- a/tests/unit/test_firecracker_backend.py +++ b/tests/unit/test_firecracker_backend.py @@ -66,6 +66,7 @@ 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 diff --git a/tests/unit/test_firecracker_cleanup.py b/tests/unit/test_firecracker_cleanup.py index 8670783..a66660d 100644 --- a/tests/unit/test_firecracker_cleanup.py +++ b/tests/unit/test_firecracker_cleanup.py @@ -62,6 +62,14 @@ class TestProcessScan(unittest.TestCase): with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(returncode=1)): self.assertEqual((set(), []), fc_cleanup._scan_processes(Path("/x"))) + 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_orphan_run_dirs_excludes_live_and_missing_root(self): with tempfile.TemporaryDirectory() as tmp: run_root = Path(tmp) diff --git a/tests/unit/test_macos_container_launch_wiring.py b/tests/unit/test_macos_container_launch_wiring.py index 5432565..81a427b 100644 --- a/tests/unit/test_macos_container_launch_wiring.py +++ b/tests/unit/test_macos_container_launch_wiring.py @@ -50,6 +50,7 @@ def _plan( agent_supervise_url: str = "", image_policy: str = "fresh", nested_containers: bool = False, + env_var_secret: str = "", ) -> MacosContainerBottlePlan: routes_path = stage_dir / "routes.yaml" routes_path.write_text("routes: []\n", encoding="utf-8") @@ -82,6 +83,7 @@ def _plan( ), agent_git_gate_url=agent_git_gate_url, agent_supervise_url=agent_supervise_url, + env_var_secret=env_var_secret, )) @@ -187,6 +189,12 @@ 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.""" diff --git a/tests/unit/test_macos_container_util.py b/tests/unit/test_macos_container_util.py index 5024fe4..3a419c8 100644 --- a/tests/unit/test_macos_container_util.py +++ b/tests/unit/test_macos_container_util.py @@ -28,6 +28,21 @@ 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=[], diff --git a/tests/unit/test_orchestrator_client.py b/tests/unit/test_orchestrator_client.py index d6e07bb..88a060b 100644 --- a/tests/unit/test_orchestrator_client.py +++ b/tests/unit/test_orchestrator_client.py @@ -76,6 +76,27 @@ 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") diff --git a/tests/unit/test_orchestrator_control_plane.py b/tests/unit/test_orchestrator_control_plane.py index ea75d5c..9ebd8b6 100644 --- a/tests/unit/test_orchestrator_control_plane.py +++ b/tests/unit/test_orchestrator_control_plane.py @@ -6,6 +6,7 @@ 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 @@ -63,6 +64,43 @@ 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) diff --git a/tests/unit/test_orchestrator_registry.py b/tests/unit/test_orchestrator_registry.py index 3328a71..90d0ff4 100644 --- a/tests/unit/test_orchestrator_registry.py +++ b/tests/unit/test_orchestrator_registry.py @@ -173,6 +173,58 @@ 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. diff --git a/tests/unit/test_orchestrator_secret_store.py b/tests/unit/test_orchestrator_secret_store.py new file mode 100644 index 0000000..64dd453 --- /dev/null +++ b/tests/unit/test_orchestrator_secret_store.py @@ -0,0 +1,96 @@ +"""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() diff --git a/tests/unit/test_orchestrator_service.py b/tests/unit/test_orchestrator_service.py index efee477..61ed48d 100644 --- a/tests/unit/test_orchestrator_service.py +++ b/tests/unit/test_orchestrator_service.py @@ -15,6 +15,7 @@ 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, @@ -117,6 +118,25 @@ 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}'))