Compare commits

..

7 Commits

Author SHA1 Message Date
didericis-codex 74f79cd690 refactor(secrets): keep ciphertext format minimal
test / integration-docker (pull_request) Successful in 23s
lint / lint (push) Successful in 1m4s
test / unit (pull_request) Successful in 1m57s
test / integration-firecracker (pull_request) Successful in 3m44s
test / coverage (pull_request) Successful in 22s
test / publish-infra (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Failing after 14m28s
2026-07-22 18:05:26 +00:00
didericis-codex 0a3ac27f63 test(secrets): satisfy static coverage checks
test / integration-docker (pull_request) Successful in 20s
test / unit (pull_request) Successful in 46s
lint / lint (push) Successful in 2m46s
test / integration-firecracker (pull_request) Successful in 3m26s
test / coverage (pull_request) Successful in 22s
test / publish-infra (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Failing after 14m24s
2026-07-22 17:35:27 +00:00
didericis-codex 89058fbaec fix(secrets): recover encrypted tokens on all backends
test / integration-docker (pull_request) Successful in 24s
lint / lint (push) Failing after 1m0s
test / unit (pull_request) Successful in 2m3s
test / integration-firecracker (pull_request) Successful in 4m48s
test / coverage (pull_request) Successful in 20s
test / publish-infra (pull_request) Has been skipped
tracker-policy-pr / check-pr (pull_request) Failing after 11m36s
2026-07-22 17:33:16 +00:00
didericis-claude c435e9088e fix(types): remove unused new_env_var_secret import
tracker-policy-pr / check-pr (pull_request) Successful in 12s
lint / lint (push) Successful in 2m55s
test / integration-docker (pull_request) Successful in 15s
test / unit (pull_request) Successful in 39s
test / integration-firecracker (pull_request) Successful in 3m29s
test / coverage (pull_request) Failing after 18s
test / publish-infra (pull_request) Has been skipped
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-22 03:13:00 +00:00
didericis 3b6f68d26d test(secrets): cover per-bottle egress secret encryption
tracker-policy-pr / check-pr (pull_request) Successful in 13s
test / integration-docker (pull_request) Successful in 19s
lint / lint (push) Failing after 57s
test / unit (pull_request) Successful in 1m44s
test / integration-firecracker (pull_request) Successful in 3m18s
test / coverage (pull_request) Failing after 28s
test / publish-infra (pull_request) Has been skipped
Add unit coverage for the encrypted at-rest egress secrets: the
secret_store round-trip (new_env_var_secret/encrypt_value/decrypt_value)
and the registry/service wiring that injects ENV_VAR_SECRET.

Recovered from the bot-bottle-claude-agent-1 VM after the CI runner's
disk filled and forced its rootfs read-only; committed work was already
on origin, these were the agent's uncommitted changes.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 22:57:30 -04:00
didericis-claude 572904df44 feat(secrets): encrypt egress tokens at rest with per-bottle ENV_VAR_SECRET
Implements the interim secret-provider design (PRD prd-new-secret-provider):
each agent receives a random ENV_VAR_SECRET injected into its container env
at launch. The host uses this key to encrypt each egress auth token value
(HMAC-SHA256 CTR mode, stdlib-only) and store it in a new
bottled_agent_secrets table (one row per env var, key column plaintext for
auditing). The key never touches the DB.

On infra container restart the in-memory token map is lost. launch_consolidated
now calls _reprovision_running_bottles after ensure_running: for each
registered bottle still alive on the gateway network it execs
`printenv ENV_VAR_SECRET` into the agent container and posts the result to the
new POST /bottles/<id>/reprovision_gateway control-plane endpoint, which
decrypts the stored rows and restores _tokens — no manual intervention needed.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-21 22:57:30 -04:00
didericis 0f98d75eff docs(prd): draft encrypted at-rest egress secrets (SecretProvider interim)
The orchestrator holds each bottle's egress auth tokens in process memory
only, so recreating the infra container strips every already-running
bottle of its upstream credentials. The registry row and the gateway CA
both survive; the tokens do not, so /resolve serves an intact policy with
an empty token map and the addon fails closed on `token_env unset`.

Drafts the interim slice of #355: persist the tokens encrypted so they
survive a restart, without regressing to plaintext at rest and without
foreclosing the per-request minting end state. Design section is left for
the author to fill in.
2026-07-21 22:57:30 -04:00
33 changed files with 1064 additions and 18 deletions
+15 -6
View File
@@ -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(
+4
View File
@@ -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:
@@ -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] = {
@@ -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/<id>/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,
)
+6
View File
@@ -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})"
@@ -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:
@@ -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():
@@ -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 '<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(
@@ -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,
)
+6
View File
@@ -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
+3
View File
@@ -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 &
@@ -20,6 +20,9 @@ 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,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,
)
@@ -68,6 +68,7 @@ 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,
@@ -142,6 +143,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.
@@ -176,6 +178,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,
@@ -406,6 +409,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
@@ -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.
+28 -3
View File
@@ -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/<id>/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/<id>`). False if the
orchestrator didn't know it (404) — idempotent for cleanup paths."""
+24 -1
View File
@@ -9,9 +9,13 @@ vsock / unix-socket portability caveats):
GET /bottles -> 200 {"bottles": [ <redacted record>, ...]}
POST /bottles -> 201 {"bottle_id","identity_token"} (launch)
body: {"source_ip", ["image_ref"],
["metadata"], ["policy"]}
["metadata"], ["policy"],
["tokens"], ["env_var_secret"]}
PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload)
body: {"policy"}
POST /bottles/<bottle_id>/reprovision_gateway
-> 200 {"reprovisioned": true} | 404
body: {"env_var_secret"}
DELETE /bottles/<bottle_id> -> 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):
+67
View File
@@ -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",
+35
View File
@@ -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"]
+94
View File
@@ -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/<id>/reprovision_gateway``; the orchestrator decrypts the
stored rows and re-populates ``_tokens``.
Encryption scheme: HMAC-SHA256 used as a PRF in CTR mode (stdlib-only,
no external deps). Each value is encrypted independently. The output blob is
``nonce (16 bytes) || ciphertext`` encoded as URL-safe base64 (no padding).
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"]
+30 -1
View File
@@ -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:
@@ -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: <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?
@@ -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()
+10
View File
@@ -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()
+5 -1
View File
@@ -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__":
+1
View File
@@ -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
+8
View File
@@ -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)
@@ -49,6 +49,7 @@ 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")
@@ -80,6 +81,7 @@ def _plan(
),
agent_git_gate_url=agent_git_gate_url,
agent_supervise_url=agent_supervise_url,
env_var_secret=env_var_secret,
))
@@ -185,6 +187,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."""
+15
View File
@@ -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=[],
+21
View File
@@ -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")
@@ -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)
+52
View File
@@ -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.
@@ -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()
+20
View File
@@ -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}'))