Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 74f79cd690 | |||
| 0a3ac27f63 | |||
| 89058fbaec | |||
| c435e9088e | |||
| 3b6f68d26d | |||
| 572904df44 | |||
| 0f98d75eff |
@@ -7,10 +7,13 @@ imports it rather than re-implementing it.
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import dataclasses
|
||||||
|
|
||||||
from ..egress import EgressPlan
|
from ..egress import EgressPlan
|
||||||
from ..git_gate import GitGatePlan
|
from ..git_gate import GitGatePlan
|
||||||
from ..orchestrator.client import OrchestratorClient
|
from ..orchestrator.client import OrchestratorClient, RegisteredBottle
|
||||||
from ..orchestrator.registration import registration_inputs
|
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
|
from .docker.gateway_provision import GatewayTransport, deprovision_git_gate, provision_git_gate
|
||||||
|
|
||||||
|
|
||||||
@@ -23,21 +26,27 @@ def provision_bottle(
|
|||||||
*,
|
*,
|
||||||
image_ref: str = "",
|
image_ref: str = "",
|
||||||
tokens: dict[str, str] | None = None,
|
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
|
"""Register the bottle and provision its git-gate state. Rolls back the
|
||||||
registration if provisioning fails so no orphan is left. Returns the
|
registration if provisioning fails so no orphan is left.
|
||||||
`RegisteredBottle` from the orchestrator."""
|
|
||||||
|
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)
|
inputs = registration_inputs(egress_plan)
|
||||||
|
env_var_secret = env_var_secret or new_env_var_secret()
|
||||||
reg = client.register_bottle(
|
reg = client.register_bottle(
|
||||||
source_ip, image_ref=image_ref, policy=inputs.policy,
|
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:
|
try:
|
||||||
provision_git_gate(transport, reg.bottle_id, git_gate_plan)
|
provision_git_gate(transport, reg.bottle_id, git_gate_plan)
|
||||||
except Exception:
|
except Exception:
|
||||||
client.teardown_bottle(reg.bottle_id)
|
client.teardown_bottle(reg.bottle_id)
|
||||||
raise
|
raise
|
||||||
return reg
|
return dataclasses.replace(reg, env_var_secret=env_var_secret)
|
||||||
|
|
||||||
|
|
||||||
def teardown_consolidated(
|
def teardown_consolidated(
|
||||||
|
|||||||
@@ -39,6 +39,10 @@ class DockerBottlePlan(BottlePlan):
|
|||||||
# (egress proxy credentials, git-gate/supervise headers); set by launch
|
# (egress proxy credentials, git-gate/supervise headers); set by launch
|
||||||
# from the orchestrator registration. Empty pre-registration.
|
# from the orchestrator registration. Empty pre-registration.
|
||||||
identity_token: str = ""
|
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
|
@property
|
||||||
def container_name(self) -> str:
|
def container_name(self) -> str:
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ from __future__ import annotations
|
|||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from ...egress import egress_agent_env_entries
|
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 ..util import AGENT_CA_BUNDLE, AGENT_CA_PATH
|
||||||
from .bottle_plan import DockerBottlePlan
|
from .bottle_plan import DockerBottlePlan
|
||||||
from .egress import EGRESS_PORT
|
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.
|
# the secret value never lands on argv or in the compose file.
|
||||||
for name in sorted(plan.forwarded_env.keys()):
|
for name in sorted(plan.forwarded_env.keys()):
|
||||||
env.append(name)
|
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))
|
env.extend(egress_agent_env_entries(plan.egress_plan))
|
||||||
|
|
||||||
service: dict[str, Any] = {
|
service: dict[str, Any] = {
|
||||||
|
|||||||
@@ -15,12 +15,15 @@ from __future__ import annotations
|
|||||||
|
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
from ... import log
|
||||||
from ...docker_cmd import run_docker
|
from ...docker_cmd import run_docker
|
||||||
from ...egress import EgressPlan
|
from ...egress import EgressPlan
|
||||||
from ...git_gate import GitGatePlan
|
from ...git_gate import GitGatePlan
|
||||||
from ...orchestrator.client import OrchestratorClient
|
from ...orchestrator.client import OrchestratorClient
|
||||||
from ...orchestrator.gateway import GATEWAY_NETWORK
|
from ...orchestrator.gateway import GATEWAY_NETWORK
|
||||||
from ...orchestrator.lifecycle import INFRA_NAME, OrchestratorService
|
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 provision_bottle
|
||||||
from ..consolidated_util import teardown_consolidated as _teardown_util
|
from ..consolidated_util import teardown_consolidated as _teardown_util
|
||||||
from .gateway_provision import DockerGatewayTransport
|
from .gateway_provision import DockerGatewayTransport
|
||||||
@@ -41,6 +44,7 @@ class LaunchContext:
|
|||||||
network: str # the shared gateway network to attach to
|
network: str # the shared gateway network to attach to
|
||||||
gateway_ip: str # the gateway's address — the agent's proxy target
|
gateway_ip: str # the gateway's address — the agent's proxy target
|
||||||
orchestrator_url: str
|
orchestrator_url: str
|
||||||
|
env_var_secret: str = "" # encryption key injected into the agent's env
|
||||||
|
|
||||||
|
|
||||||
def _network_cidr(network: str) -> str:
|
def _network_cidr(network: str) -> str:
|
||||||
@@ -85,6 +89,55 @@ def _network_container_ips(network: str) -> list[str]:
|
|||||||
return ips
|
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(
|
def launch_consolidated(
|
||||||
egress_plan: EgressPlan,
|
egress_plan: EgressPlan,
|
||||||
git_gate_plan: GitGatePlan,
|
git_gate_plan: GitGatePlan,
|
||||||
@@ -96,9 +149,14 @@ def launch_consolidated(
|
|||||||
network: str = GATEWAY_NETWORK,
|
network: str = GATEWAY_NETWORK,
|
||||||
) -> LaunchContext:
|
) -> LaunchContext:
|
||||||
"""Ensure the infra container is up, allocate + register the bottle, and
|
"""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()
|
service = service or OrchestratorService()
|
||||||
url = service.ensure_running()
|
url = service.ensure_running()
|
||||||
|
_reprovision_running_bottles(url, network=network, infra_name=infra_name)
|
||||||
client = OrchestratorClient(url)
|
client = OrchestratorClient(url)
|
||||||
|
|
||||||
cidr = _network_cidr(network)
|
cidr = _network_cidr(network)
|
||||||
@@ -117,6 +175,7 @@ def launch_consolidated(
|
|||||||
network=network,
|
network=network,
|
||||||
gateway_ip=gateway_ip,
|
gateway_ip=gateway_ip,
|
||||||
orchestrator_url=url,
|
orchestrator_url=url,
|
||||||
|
env_var_secret=reg.env_var_secret,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -186,6 +186,7 @@ def launch(
|
|||||||
agent_git_gate_url=git_gate_url,
|
agent_git_gate_url=git_gate_url,
|
||||||
agent_supervise_url=supervise_url,
|
agent_supervise_url=supervise_url,
|
||||||
identity_token=ctx.identity_token,
|
identity_token=ctx.identity_token,
|
||||||
|
env_var_secret=ctx.env_var_secret,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Step 5: render + up the agent-only compose, pinned on the shared
|
# Step 5: render + up the agent-only compose, pinned on the shared
|
||||||
@@ -198,7 +199,12 @@ def launch(
|
|||||||
project = compose_project_name(plan.slug)
|
project = compose_project_name(plan.slug)
|
||||||
# Forwarded vars (OAuth token, host interpolations) flow through the
|
# Forwarded vars (OAuth token, host interpolations) flow through the
|
||||||
# subprocess env as bare names so values never land in the file.
|
# 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}
|
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(
|
info(
|
||||||
f"docker compose up -d (project {project}, agent on shared "
|
f"docker compose up -d (project {project}, agent on shared "
|
||||||
f"gateway {ctx.gateway_ip}, ip {ctx.source_ip})"
|
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
|
# (egress proxy credentials, git-gate/supervise headers); set by launch
|
||||||
# from the orchestrator registration. Empty pre-registration.
|
# from the orchestrator registration. Empty pre-registration.
|
||||||
identity_token: str = ""
|
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
|
@property
|
||||||
def container_name(self) -> str:
|
def container_name(self) -> str:
|
||||||
|
|||||||
@@ -88,6 +88,12 @@ def _scan_processes(run_root: Path) -> tuple[set[str], list[int]]:
|
|||||||
return live, orphan_pids
|
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]:
|
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."""
|
"""Run dirs with no live VM behind them — the leaked ones to remove."""
|
||||||
if not run_root.is_dir():
|
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
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
import subprocess
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
from ...egress import EgressPlan
|
from ...egress import EgressPlan
|
||||||
from ...git_gate import GitGatePlan
|
from ...git_gate import GitGatePlan
|
||||||
from ...orchestrator.client import OrchestratorClient
|
from ...log import info
|
||||||
|
from ...orchestrator.client import OrchestratorClient, OrchestratorClientError
|
||||||
from ...orchestrator.lifecycle import (
|
from ...orchestrator.lifecycle import (
|
||||||
OrchestratorStartError, # re-exported so callers can catch it
|
OrchestratorStartError, # re-exported so callers can catch it
|
||||||
)
|
)
|
||||||
from ..consolidated_util import provision_bottle, teardown_consolidated as _teardown_util
|
from ...orchestrator.reprovision import reprovision_bottles
|
||||||
from . import infra_vm
|
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):
|
class ConsolidatedLaunchError(RuntimeError):
|
||||||
@@ -50,6 +61,55 @@ class LaunchContext:
|
|||||||
source_ip: str # the VM's guest IP — the attribution key
|
source_ip: str # the VM's guest IP — the attribution key
|
||||||
gateway_ca_pem: str # the shared gateway CA the provisioner installs
|
gateway_ca_pem: str # the shared gateway CA the provisioner installs
|
||||||
orchestrator_url: str
|
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(
|
def launch_consolidated(
|
||||||
@@ -66,6 +126,7 @@ def launch_consolidated(
|
|||||||
infra = infra_vm.ensure_running()
|
infra = infra_vm.ensure_running()
|
||||||
url = infra.control_plane_url
|
url = infra.control_plane_url
|
||||||
client = OrchestratorClient(url)
|
client = OrchestratorClient(url)
|
||||||
|
_reprovision_running_bottles(client)
|
||||||
|
|
||||||
transport = infra_vm.gateway_transport()
|
transport = infra_vm.gateway_transport()
|
||||||
reg = provision_bottle(
|
reg = provision_bottle(
|
||||||
@@ -80,6 +141,7 @@ def launch_consolidated(
|
|||||||
source_ip=guest_ip,
|
source_ip=guest_ip,
|
||||||
gateway_ca_pem=infra.gateway_ca_pem(),
|
gateway_ca_pem=infra.gateway_ca_pem(),
|
||||||
orchestrator_url=url,
|
orchestrator_url=url,
|
||||||
|
env_var_secret=reg.env_var_secret,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -55,8 +55,10 @@ from . import firecracker_vm, image_builder, isolation_probe, netpool, util
|
|||||||
from .bottle import FirecrackerBottle
|
from .bottle import FirecrackerBottle
|
||||||
from .bottle_plan import FirecrackerBottlePlan
|
from .bottle_plan import FirecrackerBottlePlan
|
||||||
from ...orchestrator.config_store import resolve_teardown_timeout
|
from ...orchestrator.config_store import resolve_teardown_timeout
|
||||||
|
from ...orchestrator.secret_store import ENV_VAR_SECRET_NAME
|
||||||
from .consolidated_launch import (
|
from .consolidated_launch import (
|
||||||
launch_consolidated,
|
launch_consolidated,
|
||||||
|
persist_env_var_secret,
|
||||||
teardown_consolidated,
|
teardown_consolidated,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -153,6 +155,7 @@ def launch(
|
|||||||
git_gate_plan=git_gate_plan,
|
git_gate_plan=git_gate_plan,
|
||||||
egress_plan=egress_plan,
|
egress_plan=egress_plan,
|
||||||
identity_token=ctx.identity_token,
|
identity_token=ctx.identity_token,
|
||||||
|
env_var_secret=ctx.env_var_secret,
|
||||||
# Deliver the identity token as egress proxy credentials — clients
|
# Deliver the identity token as egress proxy credentials — clients
|
||||||
# honor `HTTPS_PROXY=http://id:token@gw` without app changes; the
|
# honor `HTTPS_PROXY=http://id:token@gw` without app changes; the
|
||||||
# gateway reads Proxy-Authorization, validates the (source_ip,
|
# gateway reads Proxy-Authorization, validates the (source_ip,
|
||||||
@@ -187,6 +190,7 @@ def launch(
|
|||||||
)
|
)
|
||||||
stack.callback(vm.terminate)
|
stack.callback(vm.terminate)
|
||||||
firecracker_vm.wait_for_ssh(vm, private_key)
|
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
|
# Authoritative fail-closed egress-boundary check, before the agent
|
||||||
# runs: prove the VM cannot reach the host directly.
|
# 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
|
env["GIT_GATE_URL"] = plan.agent_git_gate_url
|
||||||
if plan.agent_supervise_url:
|
if plan.agent_supervise_url:
|
||||||
env["MCP_SUPERVISE_URL"] = 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):
|
for entry in egress_agent_env_entries(plan.egress_plan):
|
||||||
key, _, value = entry.partition("=")
|
key, _, value = entry.partition("=")
|
||||||
env[key] = value
|
env[key] = value
|
||||||
|
|||||||
@@ -399,6 +399,9 @@ fi
|
|||||||
chown -R 0:0 /root 2>/dev/null || true
|
chown -R 0:0 /root 2>/dev/null || true
|
||||||
|
|
||||||
mkdir -p /etc/dropbear /run
|
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,
|
# -R: generate host keys on demand. -E: log auth failures to stderr,
|
||||||
# captured in the host-side console.log for debugging.
|
# captured in the host-side console.log for debugging.
|
||||||
/bb-dropbear -R -E -p 22 &
|
/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
|
# bottle is registered. See launch.py's stamp for why it lives here and not
|
||||||
# only in the exec-time proxy env.
|
# only in the exec-time proxy env.
|
||||||
identity_token: str = ""
|
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
|
@property
|
||||||
def container_name(self) -> str:
|
def container_name(self) -> str:
|
||||||
|
|||||||
@@ -38,7 +38,12 @@ from ...egress import EgressPlan
|
|||||||
from ...git_gate import GitGatePlan
|
from ...git_gate import GitGatePlan
|
||||||
from ...log import info
|
from ...log import info
|
||||||
from ...orchestrator.client import OrchestratorClient, OrchestratorClientError
|
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 . import util as container_mod
|
||||||
from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active
|
from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active
|
||||||
from .gateway import GATEWAY_NETWORK
|
from .gateway import GATEWAY_NETWORK
|
||||||
@@ -72,6 +77,7 @@ class LaunchContext:
|
|||||||
gateway_ip: str
|
gateway_ip: str
|
||||||
network: str
|
network: str
|
||||||
orchestrator_url: str
|
orchestrator_url: str
|
||||||
|
env_var_secret: str = "" # encryption key injected into the agent's env
|
||||||
|
|
||||||
|
|
||||||
def ensure_gateway(
|
def ensure_gateway(
|
||||||
@@ -83,12 +89,35 @@ def ensure_gateway(
|
|||||||
needs `gateway_ip` at run time."""
|
needs `gateway_ip` at run time."""
|
||||||
service = service or MacosInfraService()
|
service = service or MacosInfraService()
|
||||||
infra = service.ensure_running()
|
infra = service.ensure_running()
|
||||||
return GatewayEndpoint(
|
endpoint = GatewayEndpoint(
|
||||||
orchestrator_url=infra.control_plane_url,
|
orchestrator_url=infra.control_plane_url,
|
||||||
gateway_ip=infra.gateway_ip,
|
gateway_ip=infra.gateway_ip,
|
||||||
gateway_ca_pem=service.ca_cert_pem(),
|
gateway_ca_pem=service.ca_cert_pem(),
|
||||||
network=service.network,
|
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]:
|
def live_source_ips(network: str) -> list[str]:
|
||||||
@@ -125,6 +154,7 @@ def register_agent(
|
|||||||
endpoint: GatewayEndpoint,
|
endpoint: GatewayEndpoint,
|
||||||
image_ref: str = "",
|
image_ref: str = "",
|
||||||
tokens: dict[str, str] | None = None,
|
tokens: dict[str, str] | None = None,
|
||||||
|
env_var_secret: str | None = None,
|
||||||
) -> LaunchContext:
|
) -> LaunchContext:
|
||||||
"""Register the (already running) agent by its address and provision its
|
"""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
|
git-gate state into the gateway. `source_ip` must be read from the live
|
||||||
@@ -144,6 +174,7 @@ def register_agent(
|
|||||||
reg = provision_bottle(
|
reg = provision_bottle(
|
||||||
client, source_ip, egress_plan, git_gate_plan, AppleGatewayTransport(),
|
client, source_ip, egress_plan, git_gate_plan, AppleGatewayTransport(),
|
||||||
image_ref=image_ref, tokens=tokens,
|
image_ref=image_ref, tokens=tokens,
|
||||||
|
env_var_secret=env_var_secret,
|
||||||
)
|
)
|
||||||
return LaunchContext(
|
return LaunchContext(
|
||||||
bottle_id=reg.bottle_id,
|
bottle_id=reg.bottle_id,
|
||||||
@@ -152,6 +183,7 @@ def register_agent(
|
|||||||
gateway_ip=endpoint.gateway_ip,
|
gateway_ip=endpoint.gateway_ip,
|
||||||
network=endpoint.network,
|
network=endpoint.network,
|
||||||
orchestrator_url=endpoint.orchestrator_url,
|
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 .bottle_plan import MacosContainerBottlePlan
|
||||||
from ...orchestrator.config_store import resolve_teardown_timeout
|
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 (
|
from .consolidated_launch import (
|
||||||
GatewayEndpoint,
|
GatewayEndpoint,
|
||||||
ensure_gateway,
|
ensure_gateway,
|
||||||
@@ -142,6 +143,7 @@ def launch(
|
|||||||
plan = _provision_git_gate_keys(plan)
|
plan = _provision_git_gate_keys(plan)
|
||||||
plan = _install_gateway_ca(plan, endpoint)
|
plan = _install_gateway_ca(plan, endpoint)
|
||||||
plan = _stamp_agent_urls(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
|
# Step 3: run the agent. It has no identity token yet — registration
|
||||||
# needs the address this run assigns.
|
# needs the address this run assigns.
|
||||||
@@ -176,6 +178,7 @@ def launch(
|
|||||||
endpoint=endpoint,
|
endpoint=endpoint,
|
||||||
image_ref=plan.image,
|
image_ref=plan.image,
|
||||||
tokens=token_values,
|
tokens=token_values,
|
||||||
|
env_var_secret=plan.env_var_secret,
|
||||||
)
|
)
|
||||||
stack.callback(
|
stack.callback(
|
||||||
teardown_consolidated, ctx.bottle_id,
|
teardown_consolidated, ctx.bottle_id,
|
||||||
@@ -406,6 +409,8 @@ def _agent_env_entries(
|
|||||||
env.append(f"GIT_GATE_URL={plan.agent_git_gate_url}")
|
env.append(f"GIT_GATE_URL={plan.agent_git_gate_url}")
|
||||||
if plan.agent_supervise_url:
|
if plan.agent_supervise_url:
|
||||||
env.append(f"MCP_SUPERVISE_URL={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()):
|
for name, value in sorted(plan.agent_provision.guest_env.items()):
|
||||||
env.append(f"{name}={value}")
|
env.append(f"{name}={value}")
|
||||||
# Forwarded vars: bare name → inherits from the `container run` process env
|
# 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:
|
def exec_container_as_root(name: str, argv: list[str]) -> None:
|
||||||
"""`exec_container`, but as uid 0 inside the container.
|
"""`exec_container`, but as uid 0 inside the container.
|
||||||
|
|
||||||
|
|||||||
@@ -41,10 +41,13 @@ class OrchestratorClientError(RuntimeError):
|
|||||||
@dataclass(frozen=True)
|
@dataclass(frozen=True)
|
||||||
class RegisteredBottle:
|
class RegisteredBottle:
|
||||||
"""What `POST /bottles` returns: the minted bottle id and the per-bottle
|
"""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
|
bottle_id: str
|
||||||
identity_token: str
|
identity_token: str
|
||||||
|
env_var_secret: str = ""
|
||||||
|
|
||||||
|
|
||||||
class OrchestratorClient:
|
class OrchestratorClient:
|
||||||
@@ -120,17 +123,21 @@ class OrchestratorClient:
|
|||||||
metadata: str = "",
|
metadata: str = "",
|
||||||
policy: str = "",
|
policy: str = "",
|
||||||
tokens: dict[str, str] | None = None,
|
tokens: dict[str, str] | None = None,
|
||||||
|
env_var_secret: str = "",
|
||||||
) -> RegisteredBottle:
|
) -> RegisteredBottle:
|
||||||
"""Register a bottle and broker its launch (`POST /bottles`). `tokens`
|
"""Register a bottle and broker its launch (`POST /bottles`). `tokens`
|
||||||
are the per-bottle egress auth values (env_name -> value) the
|
are the per-bottle egress auth values (env_name -> value) the
|
||||||
orchestrator holds in memory for the gateway to inject. Returns the
|
orchestrator holds in memory for the gateway to inject. When
|
||||||
minted id + identity token."""
|
*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", {
|
payload = self._ok("POST", "/bottles", {
|
||||||
"source_ip": source_ip,
|
"source_ip": source_ip,
|
||||||
"image_ref": image_ref,
|
"image_ref": image_ref,
|
||||||
"metadata": metadata,
|
"metadata": metadata,
|
||||||
"policy": policy,
|
"policy": policy,
|
||||||
"tokens": tokens or {},
|
"tokens": tokens or {},
|
||||||
|
"env_var_secret": env_var_secret,
|
||||||
})
|
})
|
||||||
bottle_id = payload.get("bottle_id")
|
bottle_id = payload.get("bottle_id")
|
||||||
token = payload.get("identity_token")
|
token = payload.get("identity_token")
|
||||||
@@ -138,6 +145,24 @@ class OrchestratorClient:
|
|||||||
raise OrchestratorClientError("register: response missing bottle_id/identity_token")
|
raise OrchestratorClientError("register: response missing bottle_id/identity_token")
|
||||||
return RegisteredBottle(bottle_id=bottle_id, identity_token=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:
|
def teardown_bottle(self, bottle_id: str) -> bool:
|
||||||
"""Tear a bottle down (`DELETE /bottles/<id>`). False if the
|
"""Tear a bottle down (`DELETE /bottles/<id>`). False if the
|
||||||
orchestrator didn't know it (404) — idempotent for cleanup paths."""
|
orchestrator didn't know it (404) — idempotent for cleanup paths."""
|
||||||
|
|||||||
@@ -9,9 +9,13 @@ vsock / unix-socket portability caveats):
|
|||||||
GET /bottles -> 200 {"bottles": [ <redacted record>, ...]}
|
GET /bottles -> 200 {"bottles": [ <redacted record>, ...]}
|
||||||
POST /bottles -> 201 {"bottle_id","identity_token"} (launch)
|
POST /bottles -> 201 {"bottle_id","identity_token"} (launch)
|
||||||
body: {"source_ip", ["image_ref"],
|
body: {"source_ip", ["image_ref"],
|
||||||
["metadata"], ["policy"]}
|
["metadata"], ["policy"],
|
||||||
|
["tokens"], ["env_var_secret"]}
|
||||||
PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload)
|
PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload)
|
||||||
body: {"policy"}
|
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)
|
DELETE /bottles/<bottle_id> -> 200 {"torn_down": true} | 404 (teardown)
|
||||||
POST /reconcile -> 200 {"reaped": [bottle_id, ...]}
|
POST /reconcile -> 200 {"reaped": [bottle_id, ...]}
|
||||||
body: {"live_source_ips": [...],
|
body: {"live_source_ips": [...],
|
||||||
@@ -116,12 +120,14 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
|
|||||||
tokens = {
|
tokens = {
|
||||||
k: v for k, v in raw_tokens.items() if isinstance(k, str) and isinstance(v, str)
|
k: v for k, v in raw_tokens.items() if isinstance(k, str) and isinstance(v, str)
|
||||||
} if isinstance(raw_tokens, dict) else {}
|
} if isinstance(raw_tokens, dict) else {}
|
||||||
|
env_var_secret = data.get("env_var_secret", "")
|
||||||
rec = orch.launch_bottle(
|
rec = orch.launch_bottle(
|
||||||
source_ip,
|
source_ip,
|
||||||
image_ref=image_ref if isinstance(image_ref, str) else "",
|
image_ref=image_ref if isinstance(image_ref, str) else "",
|
||||||
metadata=metadata if isinstance(metadata, str) else "",
|
metadata=metadata if isinstance(metadata, str) else "",
|
||||||
policy=policy if isinstance(policy, str) else "",
|
policy=policy if isinstance(policy, str) else "",
|
||||||
tokens=tokens,
|
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}
|
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 200, {"updated": True}
|
||||||
return 404, {"error": "no such bottle"}
|
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/"):
|
if method == "DELETE" and route.startswith("/bottles/"):
|
||||||
bottle_id = route[len("/bottles/"):]
|
bottle_id = route[len("/bottles/"):]
|
||||||
if orch.teardown_bottle(bottle_id):
|
if orch.teardown_bottle(bottle_id):
|
||||||
|
|||||||
@@ -113,6 +113,22 @@ _MIGRATIONS = TableMigrations(
|
|||||||
# egress allowlist / routes / git config selected by source IP. The
|
# egress allowlist / routes / git config selected by source IP. The
|
||||||
# multi-tenant gateway resolves it per request via `attribute`.
|
# multi-tenant gateway resolves it per request via `attribute`.
|
||||||
"ALTER TABLE orchestrator_bottles ADD COLUMN policy TEXT NOT NULL DEFAULT ''",
|
"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 None
|
||||||
return rec
|
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__ = [
|
__all__ = [
|
||||||
"BottleRecord",
|
"BottleRecord",
|
||||||
|
|||||||
@@ -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"]
|
||||||
@@ -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"]
|
||||||
@@ -87,13 +87,22 @@ class Orchestrator:
|
|||||||
metadata: str = "",
|
metadata: str = "",
|
||||||
policy: str = "",
|
policy: str = "",
|
||||||
tokens: dict[str, str] | None = None,
|
tokens: dict[str, str] | None = None,
|
||||||
|
env_var_secret: str = "",
|
||||||
) -> BottleRecord:
|
) -> BottleRecord:
|
||||||
"""Register a bottle (with its gateway policy + in-memory egress auth
|
"""Register a bottle (with its gateway policy + in-memory egress auth
|
||||||
tokens) and broker its launch. Rolls the registry entry back if the
|
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)
|
rec = self.registry.register(source_ip, metadata=metadata, policy=policy)
|
||||||
if tokens:
|
if tokens:
|
||||||
self._tokens[rec.bottle_id] = dict(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(
|
req = LaunchRequest(
|
||||||
op="launch",
|
op="launch",
|
||||||
bottle_id=rec.bottle_id,
|
bottle_id=rec.bottle_id,
|
||||||
@@ -284,6 +293,26 @@ class Orchestrator:
|
|||||||
))
|
))
|
||||||
return True, ""
|
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 ----------------------------------------------
|
# --- consolidated gateway ----------------------------------------------
|
||||||
|
|
||||||
def ensure_gateway(self) -> None:
|
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()
|
||||||
@@ -2,6 +2,7 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import dataclasses
|
||||||
import unittest
|
import unittest
|
||||||
|
|
||||||
from bot_bottle.backend.docker.consolidated_compose import consolidated_agent_compose
|
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).
|
# forwarded secrets are bare names (value inherited from process env).
|
||||||
self.assertIn("CLAUDE_CODE_OAUTH_TOKEN", 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__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
@@ -91,6 +91,7 @@ class TestTeardownWarning(unittest.TestCase):
|
|||||||
bottle_id="b1", identity_token="t", source_ip="172.20.0.4",
|
bottle_id="b1", identity_token="t", source_ip="172.20.0.4",
|
||||||
network="bot-bottle-gateway", gateway_ip="172.20.0.2",
|
network="bot-bottle-gateway", gateway_ip="172.20.0.2",
|
||||||
orchestrator_url="http://orch:8099",
|
orchestrator_url="http://orch:8099",
|
||||||
|
env_var_secret="encryption-key",
|
||||||
)
|
)
|
||||||
|
|
||||||
images = BottleImages(agent="bot-bottle-claude:latest", sidecar="bot-bottle-sidecars:latest")
|
images = BottleImages(agent="bot-bottle-claude:latest", sidecar="bot-bottle-sidecars:latest")
|
||||||
@@ -107,7 +108,7 @@ class TestTeardownWarning(unittest.TestCase):
|
|||||||
mock.patch.object(
|
mock.patch.object(
|
||||||
launch_mod, "write_compose_file", return_value=Path("/tmp/compose.yml"),
|
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_dump_logs"), \
|
||||||
mock.patch.object(
|
mock.patch.object(
|
||||||
launch_mod, "compose_down",
|
launch_mod, "compose_down",
|
||||||
@@ -122,6 +123,9 @@ class TestTeardownWarning(unittest.TestCase):
|
|||||||
self.assertIn("bot-bottle: warning:", output)
|
self.assertIn("bot-bottle: warning:", output)
|
||||||
self.assertIn("bot-bottle-test-teardown-abc", output)
|
self.assertIn("bot-bottle-test-teardown-abc", output)
|
||||||
self.assertIn("compose-down", output)
|
self.assertIn("compose-down", output)
|
||||||
|
self.assertEqual(
|
||||||
|
"encryption-key", compose_up.call_args.kwargs["env"]["ENV_VAR_SECRET"],
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
|
|||||||
@@ -66,6 +66,7 @@ class TestNetpoolRenderers(unittest.TestCase):
|
|||||||
|
|
||||||
self.assertIn("chown node:node /home/node", util._GUEST_INIT)
|
self.assertIn("chown node:node /home/node", util._GUEST_INIT)
|
||||||
self.assertIn("chmod 755 /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):
|
def test_nixos_module_is_non_invasive(self):
|
||||||
# The NixOS module must NOT flip the host firewall backend or
|
# The NixOS module must NOT flip the host firewall backend or
|
||||||
|
|||||||
@@ -62,6 +62,14 @@ class TestProcessScan(unittest.TestCase):
|
|||||||
with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(returncode=1)):
|
with patch.object(fc_cleanup.subprocess, "run", return_value=_proc(returncode=1)):
|
||||||
self.assertEqual((set(), []), fc_cleanup._scan_processes(Path("/x")))
|
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):
|
def test_orphan_run_dirs_excludes_live_and_missing_root(self):
|
||||||
with tempfile.TemporaryDirectory() as tmp:
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
run_root = Path(tmp)
|
run_root = Path(tmp)
|
||||||
|
|||||||
@@ -49,6 +49,7 @@ def _plan(
|
|||||||
agent_git_gate_url: str = "",
|
agent_git_gate_url: str = "",
|
||||||
agent_supervise_url: str = "",
|
agent_supervise_url: str = "",
|
||||||
image_policy: str = "fresh",
|
image_policy: str = "fresh",
|
||||||
|
env_var_secret: str = "",
|
||||||
) -> MacosContainerBottlePlan:
|
) -> MacosContainerBottlePlan:
|
||||||
routes_path = stage_dir / "routes.yaml"
|
routes_path = stage_dir / "routes.yaml"
|
||||||
routes_path.write_text("routes: []\n", encoding="utf-8")
|
routes_path.write_text("routes: []\n", encoding="utf-8")
|
||||||
@@ -80,6 +81,7 @@ def _plan(
|
|||||||
),
|
),
|
||||||
agent_git_gate_url=agent_git_gate_url,
|
agent_git_gate_url=agent_git_gate_url,
|
||||||
agent_supervise_url=agent_supervise_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],
|
"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:
|
def test_never_pins_an_ip(self) -> None:
|
||||||
"""Apple Container 1.0.0 has no --ip: the address is DHCP-assigned and
|
"""Apple Container 1.0.0 has no --ip: the address is DHCP-assigned and
|
||||||
read back after start."""
|
read back after start."""
|
||||||
|
|||||||
@@ -28,6 +28,21 @@ class TestMacosContainerAvailability(unittest.TestCase):
|
|||||||
|
|
||||||
|
|
||||||
class TestMacosContainerCommands(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):
|
def test_dns_server_prefers_direct_host_ipv4_resolver(self):
|
||||||
scutil = util.subprocess.CompletedProcess(
|
scutil = util.subprocess.CompletedProcess(
|
||||||
args=[],
|
args=[],
|
||||||
|
|||||||
@@ -76,6 +76,27 @@ class TestTeardown(unittest.TestCase):
|
|||||||
self.assertEqual("DELETE", m.call_args.args[0].get_method())
|
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):
|
class TestHealthAndPolicy(unittest.TestCase):
|
||||||
def setUp(self) -> None:
|
def setUp(self) -> None:
|
||||||
self.c = OrchestratorClient("http://orch:8080")
|
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
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import base64
|
||||||
import json
|
import json
|
||||||
import secrets
|
import secrets
|
||||||
import sqlite3
|
import sqlite3
|
||||||
@@ -63,6 +64,43 @@ class TestDispatch(unittest.TestCase):
|
|||||||
self.assertTrue(payload["bottle_id"])
|
self.assertTrue(payload["bottle_id"])
|
||||||
self.assertTrue(payload["identity_token"])
|
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:
|
def test_register_requires_source_ip(self) -> None:
|
||||||
status, _ = dispatch(self.orch, "POST", "/bottles", _body({}))
|
status, _ = dispatch(self.orch, "POST", "/bottles", _body({}))
|
||||||
self.assertEqual(400, status)
|
self.assertEqual(400, status)
|
||||||
|
|||||||
@@ -173,6 +173,58 @@ if __name__ == "__main__":
|
|||||||
unittest.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):
|
class TestReapAbsent(unittest.TestCase):
|
||||||
"""`reap_absent` — the self-heal for rows whose bottle is gone.
|
"""`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()
|
||||||
@@ -15,6 +15,7 @@ from bot_bottle.orchestrator.broker import LaunchBroker, LaunchRequest, StubBrok
|
|||||||
from bot_bottle.orchestrator.registry import RegistryStore
|
from bot_bottle.orchestrator.registry import RegistryStore
|
||||||
from bot_bottle.orchestrator.service import Orchestrator
|
from bot_bottle.orchestrator.service import Orchestrator
|
||||||
from bot_bottle.orchestrator.gateway import Gateway
|
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.store_manager import StoreManager
|
||||||
from bot_bottle.supervise import (
|
from bot_bottle.supervise import (
|
||||||
Proposal,
|
Proposal,
|
||||||
@@ -117,6 +118,25 @@ class TestOrchestrator(unittest.TestCase):
|
|||||||
rec = self.orch.launch_bottle("10.243.0.6")
|
rec = self.orch.launch_bottle("10.243.0.6")
|
||||||
self.assertEqual({}, self.orch.tokens_for(rec.bottle_id))
|
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:
|
def test_set_policy_live_reload(self) -> None:
|
||||||
rec = self.orch.launch_bottle("10.243.0.3")
|
rec = self.orch.launch_bottle("10.243.0.3")
|
||||||
self.assertTrue(self.orch.set_policy(rec.bottle_id, '{"x":1}'))
|
self.assertTrue(self.orch.set_policy(rec.bottle_id, '{"x":1}'))
|
||||||
|
|||||||
Reference in New Issue
Block a user