Compare commits

..

2 Commits

Author SHA1 Message Date
didericis-claude c527841d55 fix(ci): use absolute github.workspace paths for coverage artifact upload/download
test / integration-docker (pull_request) Successful in 13s
tracker-policy-pr / check-pr (pull_request) Successful in 15s
test / unit (pull_request) Successful in 31s
test / integration-firecracker (pull_request) Successful in 3m11s
test / coverage (pull_request) Failing after 1m44s
test / publish-infra (pull_request) Has been skipped
The delphi-ci runner resolves relative paths in upload-artifact and
download-artifact from a different CWD than run: shell steps, so
'.coverage.unit' etc. were never found. Using ${{ github.workspace }}
gives an absolute path that does not depend on the JS action's CWD.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-21 04:53:33 +00:00
didericis-claude 5940b75bb7 ci: artifact-based coverage and local Firecracker candidate flow
test / integration-docker (pull_request) Successful in 12s
tracker-policy-pr / check-pr (pull_request) Successful in 8s
test / unit (pull_request) Successful in 33s
test / integration-firecracker (pull_request) Successful in 3m13s
test / coverage (pull_request) Failing after 1m45s
test / publish-infra (pull_request) Has been skipped
Each test job now runs once under coverage and uploads a small .coverage.*
artifact. The coverage job combines them on ubuntu-latest — no test reruns,
no KVM dependency. The infra candidate is built directly on the KVM runner,
eliminating the build-infra job and the ~70 s upload + ~83 s combined
download. For PRs, no rootfs artifact is transferred at all. Main-branch
pushes upload the tested rootfs and matching dropbear so publish-infra
publishes the byte-identical artifact. relative_files = True in .coveragerc
lets coverage files from different runners combine without path remapping.

Closes #446
2026-07-21 04:04:30 +00:00
102 changed files with 547 additions and 5147 deletions
+8 -39
View File
@@ -57,27 +57,16 @@ jobs:
run: python3 -m pip install --break-system-packages -r requirements-dev.txt
- name: Run unit tests with coverage
env:
COVERAGE_FILE: ${{ github.workspace }}/.coverage.unit
run: python3 -m coverage run -m unittest discover -t . -s tests/unit -v
run: python3 -m coverage run --data-file=.coverage.unit -m unittest discover -t . -s tests/unit -v
- name: Report unit coverage
env:
COVERAGE_FILE: ${{ github.workspace }}/.coverage.unit
run: python3 -m coverage report -m
# upload-artifact@v3's glob skips dotfiles, so a bare `.coverage.unit`
# silently uploads nothing ("No files were found"). Stage it under a
# non-dot name; the coverage job renames it back before `coverage
# combine`. `cp` also fails loudly if coverage never wrote the file.
- name: Stage unit coverage for upload
run: cp .coverage.unit coverage-unit.dat
run: python3 -m coverage report --data-file=.coverage.unit -m
- name: Upload unit coverage artifact
uses: actions/upload-artifact@v3
with:
name: coverage-unit
path: coverage-unit.dat
path: ${{ github.workspace }}/.coverage.unit
integration-docker:
runs-on: ubuntu-latest
@@ -102,18 +91,13 @@ jobs:
- name: Run integration tests (docker) with coverage
env:
BOT_BOTTLE_BACKEND: docker
COVERAGE_FILE: ${{ github.workspace }}/.coverage.docker
run: python3 -m coverage run -m unittest discover -t . -s tests/integration -v
# Non-dot name so upload-artifact's dotfile-skipping glob picks it up.
- name: Stage docker coverage for upload
run: cp .coverage.docker coverage-docker.dat
run: python3 -m coverage run --data-file=.coverage.docker -m unittest discover -t . -s tests/integration -v
- name: Upload docker coverage artifact
uses: actions/upload-artifact@v3
with:
name: coverage-docker
path: coverage-docker.dat
path: ${{ github.workspace }}/.coverage.docker
# Integration tests against the Firecracker backend. Runs on a self-hosted
# KVM runner (label `kvm`) where /dev/kvm and the TAP/nft pool are available.
@@ -155,7 +139,7 @@ jobs:
- name: Build infra candidate from this checkout
env:
BOT_BOTTLE_FC_DROPBEAR: /var/cache/bot-bottle-fc/dropbear
run: python3 -m bot_bottle.backend.firecracker.publish_infra --output infra-candidate --reuse-published
run: python3 -m bot_bottle.backend.firecracker.publish_infra --output infra-candidate
- name: Replace the persistent infra VM with the candidate
run: python3 -c 'from bot_bottle.backend.firecracker import infra_vm; infra_vm.stop()'
@@ -167,18 +151,13 @@ jobs:
env:
BOT_BOTTLE_BACKEND: firecracker
BOT_BOTTLE_INFRA_ARTIFACT_DIR: ${{ github.workspace }}/infra-candidate
COVERAGE_FILE: ${{ github.workspace }}/.coverage.firecracker
run: python3 -m coverage run -m unittest discover -t . -s tests/integration -v
# Non-dot name so upload-artifact's dotfile-skipping glob picks it up.
- name: Stage firecracker coverage for upload
run: cp .coverage.firecracker coverage-firecracker.dat
run: python3 -m coverage run --data-file=.coverage.firecracker -m unittest discover -t . -s tests/integration -v
- name: Upload firecracker coverage artifact
uses: actions/upload-artifact@v3
with:
name: coverage-firecracker
path: coverage-firecracker.dat
path: ${{ github.workspace }}/.coverage.firecracker
# Only upload the large rootfs artifact on main-branch pushes;
# PRs avoid the ~194 MB transfer. publish-infra only runs on main
@@ -202,8 +181,6 @@ jobs:
#
# Runs on ubuntu-latest — no KVM needed, no test reruns. Coverage files use
# relative_files = True (.coveragerc) so they combine cleanly across runners.
# Each test job sets COVERAGE_FILE to an absolute path so coverage.py writes
# to a known location that upload-artifact can find regardless of runner env.
#
# Restricted to the same events as integration-firecracker: it depends on
# that job's coverage artifact and skips for fork PRs alongside it.
@@ -243,14 +220,6 @@ jobs:
name: coverage-firecracker
path: ${{ github.workspace }}
# Rename the non-dot upload names back to the .coverage.* files that
# `coverage combine` discovers (see the staging steps in each test job).
- name: Reassemble coverage data files
run: |
mv coverage-unit.dat .coverage.unit
mv coverage-docker.dat .coverage.docker
mv coverage-firecracker.dat .coverage.firecracker
- name: Combined coverage (unit + integration, incl. firecracker)
run: PYTHON=python3 bash scripts/coverage.sh aggregate critical
-3
View File
@@ -14,9 +14,6 @@ on:
jobs:
update-badges:
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@v3
with:
+36 -13
View File
@@ -1,22 +1,45 @@
# Shared infra image: gateway data plane + orchestrator control plane.
# Firecracker single infra-VM image (PRD 0070 Stage B).
#
# Used directly by the Docker backend (run as one `bot-bottle-infra`
# container, replacing the prior two-container split). The Firecracker
# backend extends this via Dockerfile.infra.fc, adding buildah/crun/
# netavark for in-VM agent-image building.
#
# Dockerfile.orchestrator is the single definition of the orchestrator
# content (the lean `bot_bottle` package on python:3.12-slim). Both this
# image and Dockerfile.infra.fc pull it in via `COPY --from`.
# The per-host infra VM runs the orchestrator control plane, the gateway
# data plane, AND builds agent images (buildah) — all in one microVM (see
# backend/firecracker/infra_vm.py). It composes:
# * FROM the gateway image (mitmproxy / git / gitleaks / supervise + the
# flat daemon modules) — now trixie-based, so buildah 1.39 is available;
# * `COPY --from` the orchestrator image's content (the single definition
# of the control-plane payload — see Dockerfile.orchestrator), so this
# VM and the docker backend share one orchestrator definition; and
# * buildah, installed HERE only (the docker orchestrator/gateway images
# never carry it).
#
# multi-`FROM` can't union two bases (that's multi-stage, not multiple
# inheritance), so the orchestrator content is pulled in via `COPY --from`
# rather than a second base. Both images share the trixie `python:3.12-slim`
# base, so the copy is clean (same python; future installed deps copy too).
#
# The docker backend keeps orchestrator + gateway as separate images; this
# combined image exists only for the Firecracker single-VM cut. Splitting a
# service back into its own VM later is a routing change, not a repackaging
# (PRD 0070's "secret concentration"; a disposable builder can boot from
# this same image on its own TAP).
FROM bot-bottle-gateway:latest
# The orchestrator content, from its single definition. The gateway image
# already has the flat daemon modules under /app; this adds the full
# `bot_bottle` package so `python3 -m bot_bottle.orchestrator` resolves —
# used by gateway_init when BOT_BOTTLE_GATEWAY_DAEMONS includes `orchestrator`.
# --- in-VM agent-image builder (PRD 0069 Stage 3) -------------------
# The Firecracker backend builds users' agent Dockerfiles *inside this VM*
# with buildah (rootless, daemonless) instead of on the host — no host
# Docker daemon, no root-equivalent `docker` group. `crun` is the OCI
# runtime; `netavark` + `aardvark-dns` are the network backend for `FROM`
# pulls + `RUN` egress. Requires the trixie base (buildah 1.39: bookworm's
# 1.28 can't parse Dockerfile heredocs that agent images use).
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
buildah crun netavark aardvark-dns \
&& rm -rf /var/lib/apt/lists/*
# vfs + chroot: buildah works as root in the bare microVM (no
# fuse-overlayfs / overlay module / subuid maps). Matches image_builder.
ENV STORAGE_DRIVER=vfs \
BUILDAH_ISOLATION=chroot
# The orchestrator content, pulled from its single definition. The gateway
# image already has the flat daemon modules under /app; this adds the full
# `bot_bottle` package so `python3 -m bot_bottle.orchestrator` resolves.
COPY --from=bot-bottle-orchestrator:latest /app/bot_bottle /app/bot_bottle
-23
View File
@@ -1,23 +0,0 @@
# Firecracker infra VM image (PRD 0070 Stage B).
#
# Extends the shared infra base (Dockerfile.infra: gateway + orchestrator
# control plane) with the in-VM agent-image builder. The Firecracker backend
# builds users' agent Dockerfiles *inside this VM* with buildah (rootless,
# daemonless) instead of on the host — no host Docker daemon, no
# root-equivalent `docker` group.
#
# Requires the trixie base from bot-bottle-gateway (buildah 1.39: bookworm's
# 1.28 can't parse Dockerfile heredocs that agent images use).
#
# `crun` is the OCI runtime; `netavark` + `aardvark-dns` are the network
# backend for `FROM` pulls + `RUN` egress. `vfs` + `chroot`: buildah works
# as root in the bare microVM (no fuse-overlayfs / overlay module / subuid
# maps). Matches image_builder.
FROM bot-bottle-infra:latest
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
buildah crun netavark aardvark-dns \
&& rm -rf /var/lib/apt/lists/*
ENV STORAGE_DRIVER=vfs \
BUILDAH_ISOLATION=chroot
+2 -43
View File
@@ -45,10 +45,6 @@ PROVIDER_TEMPLATES = frozenset({PROVIDER_CLAUDE, PROVIDER_CODEX, PROVIDER_PI})
# forward_host_credentials is enabled. Pipelock must pass these through
# (no TLS MITM) or its header DLP blocks the injected JWT.
CODEX_HOST_CREDENTIAL_HOSTS = ("api.openai.com", "chatgpt.com")
# Host that egress injects the host Claude bearer on when Claude
# forward_host_credentials is enabled.
CLAUDE_HOST_CREDENTIAL_HOSTS = ("api.anthropic.com",)
PromptMode = Literal[
"append_file",
"read_prompt_file",
@@ -261,28 +257,7 @@ class AgentProvider(ABC):
Default: Debian/node — writes the git-gate insteadOf gitconfig
and sets user.name/email as node. Workspace copy runs through
BottleBackend.provision_workspace against the running bottle."""
from .log import die, info
# Firecracker exports image rootfs files through an unprivileged host
# tar extraction, so image-time ownership of XDG directories is not
# preserved. Git consults ~/.config/git even when the actual config
# is ~/.gitconfig; an unreadable directory there can prevent the
# git-gate insteadOf rules below from taking effect. Repair this at
# runtime, after every backend's copy/export path has completed.
git_xdg_dir = f"{plan.guest_home}/.config/git"
repair = bottle.exec(
f"chown node:node {shlex.quote(plan.guest_home)} && "
f"chmod 755 {shlex.quote(plan.guest_home)} && "
f"mkdir -p {shlex.quote(git_xdg_dir)} && "
f"chown -R node:node {shlex.quote(f'{plan.guest_home}/.config')} && "
f"chmod -R u+rwX,go+rX {shlex.quote(f'{plan.guest_home}/.config')}",
user="root",
)
if repair.returncode != 0:
die(
"git provisioning: could not make the runtime Git config "
f"directory readable: {(repair.stderr or repair.stdout).strip()}"
)
from .log import info
manifest_bottle = plan.manifest.bottle
if manifest_bottle.git:
@@ -305,27 +280,11 @@ class AgentProvider(ABC):
f"{len(manifest_bottle.git)} insteadOf rule(s)"
)
bottle.cp_in(str(config_file), guest_gitconfig)
permissions = bottle.exec(
bottle.exec(
f"chown node:node {shlex.quote(guest_gitconfig)} && "
f"chmod 644 {shlex.quote(guest_gitconfig)}",
user="root",
)
if permissions.returncode != 0:
die(
"git provisioning: could not set ownership on "
f"{guest_gitconfig}: "
f"{(permissions.stderr or permissions.stdout).strip()}"
)
configured = bottle.exec(
"git config --global --get-regexp '^url\\..*\\.insteadof$'",
user="node",
)
if configured.returncode != 0:
die(
"git provisioning: the runtime user cannot read the "
f"git-gate insteadOf rules from {guest_gitconfig}: "
f"{(configured.stderr or configured.stdout).strip()}"
)
gu = manifest_bottle.git_user
if not gu.is_empty():
+4 -37
View File
@@ -37,10 +37,10 @@ import os
import shlex
import sys
from abc import ABC, abstractmethod
from contextlib import AbstractContextManager, contextmanager
from contextlib import AbstractContextManager
from dataclasses import dataclass
from pathlib import Path
from typing import TYPE_CHECKING, Any, Generator, Generic, Sequence, TypeVar
from typing import TYPE_CHECKING, Any, Generic, Sequence, TypeVar
from ..agent_provider import AgentProvisionPlan, get_provider, build_agent_provision_plan
from ..egress import EgressPlan
@@ -83,9 +83,6 @@ class BottleSpec:
# True when launched via --headless (no TTY, no interactive prompts).
# The git-gate host-key preflight uses this to error rather than prompt.
headless: bool = False
# Image startup policy. "fresh" preserves the normal build path;
# "cached" reuses the current local image/artifact without rebuilding.
image_policy: str = "fresh"
@dataclass(frozen=True)
@@ -281,18 +278,6 @@ PlanT = TypeVar("PlanT", bound=BottlePlan)
CleanupT = TypeVar("CleanupT", bound=BottleCleanupPlan)
@dataclass(frozen=True)
class BottleImages:
"""Resolved image references (or artifact paths) for a bottle launch.
For Docker/macOS-container backends, `agent` and `sidecar` are string
image refs. For the smolmachines backend they are Path objects pointing
to pre-built `.smolmachine` artifacts."""
agent: str | Path
sidecar: str | Path = ""
class BottleBackend(ABC, Generic[PlanT, CleanupT]):
"""Abstract base for selectable bottle backends. Concrete subclasses
(e.g. DockerBottleBackend) own their own prepare/launch impls.
@@ -452,27 +437,9 @@ class BottleBackend(ABC, Generic[PlanT, CleanupT]):
prompt file, Dockerfile path, and guest home all live on
`agent_provision_plan` — the source of truth."""
def prelaunch_checks(self, plan: PlanT) -> None:
"""Raise StaleImageError if any cached image used by this plan is stale.
No-op default; backends override to call the shared check_stale*
helpers on their image/artifact timestamps. Called by the CLI before
launch so the operator can be prompted outside the launch context."""
@contextmanager
def launch(self, plan: PlanT) -> Generator[Bottle, None, None]:
"""Template: build or load images, then delegate to _launch_impl."""
images = self._build_or_load_images(plan)
with self._launch_impl(plan, images) as bottle:
yield bottle
@abstractmethod
def _build_or_load_images(self, plan: PlanT) -> BottleImages:
"""Return the agent and sidecar image references (or artifact paths)
for this plan, building fresh images when the policy requires it."""
@abstractmethod
def _launch_impl(self, plan: PlanT, images: BottleImages) -> AbstractContextManager[Bottle]:
"""Bring up the bottle using pre-resolved images; yield a handle; tear down on exit."""
def launch(self, plan: PlanT) -> AbstractContextManager[Bottle]:
"""Build/run the bottle and yield a handle; tear down on exit."""
def provision(self, plan: PlanT, bottle: "Bottle") -> str | None:
"""Copy host-side files (CA cert, prompt, skills, .git) into
-68
View File
@@ -1,68 +0,0 @@
"""Shared helpers for the consolidated launch sequence (PRD 0070).
Logic that was duplicated across the docker, macos_container, and
firecracker consolidated_launch modules — extracted so each backend
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, 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
def provision_bottle(
client: OrchestratorClient,
source_ip: str,
egress_plan: EgressPlan,
git_gate_plan: GitGatePlan,
transport: GatewayTransport,
*,
image_ref: str = "",
tokens: dict[str, 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.
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 = new_env_var_secret()
reg = client.register_bottle(
source_ip, image_ref=image_ref, policy=inputs.policy,
metadata=inputs.metadata, tokens=tokens, env_var_secret=env_var_secret,
)
try:
provision_git_gate(transport, reg.bottle_id, git_gate_plan)
except Exception:
client.teardown_bottle(reg.bottle_id)
raise
return dataclasses.replace(reg, env_var_secret=env_var_secret)
def teardown_consolidated(
bottle_id: str,
transport: GatewayTransport,
*,
orchestrator_url: str,
timeout: float | None = None,
) -> None:
"""Deregister the bottle and remove its git-gate state. Both steps are
idempotent so this is safe from a cleanup trap."""
from ..orchestrator.config_store import DEFAULT_TEARDOWN_TIMEOUT_SECONDS
OrchestratorClient(
orchestrator_url,
timeout=timeout if timeout is not None else DEFAULT_TEARDOWN_TIMEOUT_SECONDS,
).teardown_bottle(bottle_id)
deprovision_git_gate(transport, bottle_id)
__all__ = ["provision_bottle", "teardown_consolidated"]
+3 -9
View File
@@ -31,7 +31,7 @@ from ...env import ResolvedEnv
from ...git_gate import GitGatePlan
from ...supervise import SupervisePlan
from ...manifest import Manifest
from .. import ActiveAgent, BottleBackend, BottleImages, BottleSpec
from .. import ActiveAgent, BottleBackend, BottleSpec
from . import cleanup as _cleanup
from . import enumerate as _enumerate
from . import launch as _launch
@@ -100,15 +100,9 @@ class DockerBottleBackend(BottleBackend["DockerBottlePlan", "DockerBottleCleanup
stage_dir=stage_dir,
)
def prelaunch_checks(self, plan: DockerBottlePlan) -> None:
_launch.stale_checks(plan)
def _build_or_load_images(self, plan: DockerBottlePlan) -> BottleImages:
return _launch.build_or_load_images(plan)
@contextmanager
def _launch_impl(self, plan: DockerBottlePlan, images: BottleImages) -> Generator[DockerBottle, None, None]:
with _launch.launch(plan, images, provision=self.provision) as bottle:
def launch(self, plan: DockerBottlePlan) -> Generator[DockerBottle, None, None]:
with _launch.launch(plan, provision=self.provision) as bottle:
yield bottle
def ensure_orchestrator(self) -> str:
-4
View File
@@ -39,10 +39,6 @@ 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,7 +17,6 @@ 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
@@ -59,10 +58,6 @@ 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] = {
@@ -1,13 +1,19 @@
"""Consolidated bottle launch sequence for the docker backend (PRD 0070).
Composes the orchestrator primitives into the register/teardown sequence:
Composes the orchestrator primitives into the register/teardown sequence that
replaces the per-bottle gateway:
1. ensure the single infra container (control plane + gateway) is up;
2. allocate the bottle a pinned source IP on the gateway network;
3. register it and provision its git-gate repos/creds into the gateway.
1. ensure the orchestrator control plane + shared gateway are up;
2. allocate the bottle a pinned source IP on the gateway network (the
attribution key), skipping the gateway's own address + live bottles;
3. register it (egress policy blob + slug metadata) → bottle id + identity
token;
4. provision its git-gate repos/creds into the running gateway.
Returns a `LaunchContext` with everything the agent container needs to
attach. The agent `docker run` itself is the backend's job; this owns the
It returns a `LaunchContext` with everything the agent container needs to
attach — network, pinned IP, the gateway's address (its proxy target), the
orchestrator URL, and the identity token. The agent `docker run` itself is
the backend's job (it owns provider provisioning); this owns the
orchestrator-facing wiring so that sequence stays testable in isolation.
"""
@@ -15,18 +21,19 @@ 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 ..consolidated_util import provision_bottle
from ..consolidated_util import teardown_consolidated as _teardown_util
from .gateway_provision import DockerGatewayTransport
from ...orchestrator.gateway import GATEWAY_NAME, GATEWAY_NETWORK
from ...orchestrator.lifecycle import OrchestratorService
from ...orchestrator.registration import registration_inputs
from .gateway_net import next_free_ip
from .gateway_provision import (
DockerGatewayTransport,
deprovision_git_gate,
provision_git_gate,
)
class ConsolidatedLaunchError(RuntimeError):
@@ -43,7 +50,6 @@ 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:
@@ -69,85 +75,28 @@ def _container_ip(name: str, network: str) -> str:
ip = proc.stdout.strip()
if proc.returncode != 0 or not ip:
raise ConsolidatedLaunchError(
f"container {name} has no address on {network}: {proc.stderr.strip()}"
f"gateway {name} has no address on {network}: {proc.stderr.strip()}"
)
return ip
def _network_container_ips(network: str) -> list[str]:
"""Every address currently assigned on the gateway network — the ground
truth for "in use": the infra container and every live agent. Read from
the network so a new bottle can't collide with anything actually attached."""
truth for "in use": the gateway + orchestrator infrastructure containers
and every live agent. Read from the network so a new bottle can't collide
with anything actually attached (a registry-only view would miss the
orchestrator/gateway containers)."""
proc = run_docker([
"docker", "network", "inspect", "--format",
"{{range .Containers}}{{.IPv4Address}} {{end}}", network,
])
ips: list[str] = []
for entry in proc.stdout.split():
# entries look like "172.20.0.2/16" — keep the address.
ips.append(entry.split("/", 1)[0])
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)
bottles = client.list_bottles()
if not bottles:
return
# Build {source_ip: container_name} from live containers on the gateway
# network, excluding the infra container itself.
proc = run_docker([
"docker", "network", "inspect",
"--format", "{{range .Containers}}{{.Name}} {{.IPv4Address}}\n{{end}}",
network,
])
ip_to_container: dict[str, str] = {}
for line in proc.stdout.splitlines():
parts = line.strip().split()
if len(parts) >= 2 and parts[0] != infra_name:
ip = parts[1].split("/", 1)[0]
if ip:
ip_to_container[ip] = parts[0]
reprovisioned = 0
for bottle in bottles:
bottle_id = bottle.get("bottle_id")
source_ip = bottle.get("source_ip")
if not isinstance(bottle_id, str) or not isinstance(source_ip, str):
continue
container_name = ip_to_container.get(source_ip)
if not container_name:
continue
proc = run_docker(
["docker", "exec", container_name, "printenv", ENV_VAR_SECRET_NAME]
)
if proc.returncode != 0 or not proc.stdout.strip():
continue
try:
if client.reprovision_gateway(bottle_id, proc.stdout.strip()):
reprovisioned += 1
except Exception: # noqa: BLE001 — best-effort, never block a launch
pass
if reprovisioned:
log.info(
"reprovisioned egress tokens",
context={"count": reprovisioned},
)
def launch_consolidated(
egress_plan: EgressPlan,
git_gate_plan: GitGatePlan,
@@ -155,29 +104,33 @@ def launch_consolidated(
image_ref: str = "",
tokens: dict[str, str] | None = None,
service: OrchestratorService | None = None,
infra_name: str = INFRA_NAME,
gateway_name: str = GATEWAY_NAME,
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.
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."""
"""Ensure the orchestrator + gateway are up, allocate + register the
bottle, and provision its git-gate state. Returns the agent's attach
context. Raises `ConsolidatedLaunchError` (or the primitives' own errors)
if any step fails — the caller tears down on failure."""
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)
gateway_ip = _container_ip(infra_name, network)
gateway_ip = _container_ip(gateway_name, network)
source_ip = next_free_ip(cidr, _network_container_ips(network))
transport = DockerGatewayTransport(infra_name)
reg = provision_bottle(
client, source_ip, egress_plan, git_gate_plan, transport,
image_ref=image_ref, tokens=tokens,
inputs = registration_inputs(egress_plan)
reg = client.register_bottle(
source_ip, image_ref=image_ref, policy=inputs.policy,
metadata=inputs.metadata, tokens=tokens,
)
try:
provision_git_gate(
DockerGatewayTransport(gateway_name), reg.bottle_id, git_gate_plan)
except Exception:
# Roll the registration back so a provisioning failure leaves no orphan.
client.teardown_bottle(reg.bottle_id)
raise
return LaunchContext(
bottle_id=reg.bottle_id,
identity_token=reg.identity_token,
@@ -185,17 +138,16 @@ def launch_consolidated(
network=network,
gateway_ip=gateway_ip,
orchestrator_url=url,
env_var_secret=reg.env_var_secret,
)
def teardown_consolidated(
bottle_id: str, *, orchestrator_url: str, infra_name: str = INFRA_NAME,
timeout: float | None = None,
bottle_id: str, *, orchestrator_url: str, gateway_name: str = GATEWAY_NAME,
) -> None:
"""Deregister the bottle and remove its git-gate state. Idempotent."""
_teardown_util(bottle_id, DockerGatewayTransport(infra_name),
orchestrator_url=orchestrator_url, timeout=timeout)
"""Deregister the bottle and remove its git-gate state from the gateway.
Both steps are idempotent so this is safe from a cleanup trap."""
OrchestratorClient(orchestrator_url).teardown_bottle(bottle_id)
deprovision_git_gate(DockerGatewayTransport(gateway_name), bottle_id)
__all__ = [
+23 -65
View File
@@ -42,9 +42,7 @@ from ...git_gate import (
provision_git_gate_dynamic_keys,
revoke_git_gate_provisioned_keys,
)
from ...image_cache import check_stale
from ...log import die, info, warn
from .. import BottleImages
from ...log import info, warn
from . import util as docker_mod
from .bottle import DockerBottle
from .bottle_plan import DockerBottlePlan
@@ -64,7 +62,6 @@ from .compose import (
write_compose_file,
)
from .consolidated_compose import consolidated_agent_compose
from ...orchestrator.config_store import resolve_teardown_timeout
from .consolidated_launch import launch_consolidated, teardown_consolidated
from ...orchestrator.gateway import DockerGateway
@@ -73,47 +70,16 @@ from ...orchestrator.gateway import DockerGateway
_REPO_DIR = str(Path(__file__).resolve().parent.parent.parent.parent)
def build_or_load_images(plan: DockerBottlePlan) -> BottleImages:
"""Resolve the agent image ref for this plan.
Returns the committed snapshot if one exists, the cached image when the
policy is 'cached', or builds a fresh image and returns that."""
committed = read_committed_image(plan.slug)
if committed and docker_mod.image_exists(committed):
info(f"using committed image {committed!r}")
return BottleImages(agent=committed)
if plan.spec.image_policy == "cached":
if not docker_mod.image_exists(plan.image):
die(
f"cached agent image {plan.image!r} not found; "
"run without --cached-images to build it"
)
info(f"using cached agent image {plan.image!r}")
return BottleImages(agent=plan.image)
docker_mod.build_image(plan.image, _REPO_DIR, dockerfile=plan.dockerfile_path)
docker_mod.verify_agent_image(
plan.image, runtime_for(plan.agent_provider_template).smoke_test,
)
return BottleImages(agent=plan.image)
@contextmanager
def launch(
plan: DockerBottlePlan,
images: BottleImages,
*,
provision: Callable[[DockerBottlePlan, "DockerBottle"], str | None],
) -> Generator[DockerBottle, None, None]:
"""Launch and provision a Docker bottle via compose. Teardown on exit."""
"""Build, launch, and provision a Docker bottle via compose.
Teardown on exit."""
stack = ExitStack()
# Stamp the resolved agent image ref into the plan so compose rendering
# picks up the right image (may be a committed snapshot or cached ref).
plan = dataclasses.replace(
plan,
agent_provision=dataclasses.replace(plan.agent_provision, image=str(images.agent)),
)
_bottle_for_revoke = plan.manifest.bottle
_git_gate_dir_for_revoke = git_gate_state_dir(plan.slug)
@@ -130,6 +96,25 @@ def launch(
)
try:
# Step 1: agent image. Use a committed snapshot when one exists
# and is present in the local daemon; otherwise build from the
# Dockerfile. (The gateway image is built by the orchestrator.)
committed = read_committed_image(plan.slug)
if committed and docker_mod.image_exists(committed):
info(f"using committed image {committed!r}")
plan = dataclasses.replace(
plan,
agent_provision=dataclasses.replace(plan.agent_provision, image=committed),
)
else:
docker_mod.build_image(
plan.image, _REPO_DIR,
dockerfile=plan.dockerfile_path,
)
docker_mod.verify_agent_image(
plan.image, runtime_for(plan.agent_provider_template).smoke_test,
)
# Step 2: mint the git-gate dynamic (gitea) deploy keys, if any, before
# provisioning the bottle's repos into the shared gateway.
git_gate_plan = plan.git_gate_plan
@@ -148,14 +133,11 @@ def launch(
token_values = egress_resolve_token_values(
plan.egress_plan.token_env_map, effective_env,
)
teardown_timeout = resolve_teardown_timeout()
ctx = launch_consolidated(
plan.egress_plan, git_gate_plan, image_ref=plan.image, tokens=token_values,
)
stack.callback(
teardown_consolidated, ctx.bottle_id,
orchestrator_url=ctx.orchestrator_url,
timeout=teardown_timeout,
teardown_consolidated, ctx.bottle_id, orchestrator_url=ctx.orchestrator_url,
)
# Step 4: install the SHARED gateway CA into the agent (replaces the
@@ -186,7 +168,6 @@ 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
@@ -199,12 +180,7 @@ 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})"
@@ -231,21 +207,3 @@ def launch(
yield bottle
finally:
teardown()
def stale_checks(plan: DockerBottlePlan) -> None:
"""Raise StaleImageError if a cached image is older than the configured
threshold. Only runs when image_policy is 'cached'. Called by the backend
class's _image_stale_checks before _launch_impl starts any resources."""
if plan.spec.image_policy != "cached":
return
committed = read_committed_image(plan.slug)
if committed and docker_mod.image_exists(committed):
ts = docker_mod.image_created_at(committed)
if ts is not None:
check_stale(f"agent image {committed!r}", ts)
return
if docker_mod.image_exists(plan.image):
ts = docker_mod.image_created_at(plan.image)
if ts is not None:
check_stale(f"agent image {plan.image!r}", ts)
-44
View File
@@ -5,7 +5,6 @@ existence, and building images."""
from __future__ import annotations
import os
from datetime import datetime, timezone
import re
import shutil
import subprocess
@@ -198,46 +197,3 @@ def commit_container(container_name: str, image_tag: str) -> None:
f"{(result.stderr or '').strip() or '<no stderr>'}"
)
info(f"committed {container_name!r}{image_tag!r}")
def image_created_at(ref: str) -> datetime | None:
"""Return Docker's image Created timestamp as an aware UTC datetime, or
None when the field is absent or unparseable. Callers should skip the
stale check when None is returned."""
r = subprocess.run(
["docker", "image", "inspect", "--format", "{{.Created}}", ref],
capture_output=True,
text=True,
check=False,
)
if r.returncode != 0:
die(
f"docker image inspect for {ref!r} failed: "
f"{(r.stderr or '').strip() or '<no stderr>'}"
)
raw = r.stdout.strip()
if not raw:
return None
try:
return _parse_docker_timestamp(raw)
except ValueError:
return None
def _parse_docker_timestamp(raw: str) -> datetime:
text = raw.strip()
if text.endswith("Z"):
text = text[:-1] + "+00:00"
dot = text.find(".")
if dot != -1:
tz_plus = text.find("+", dot)
tz_minus = text.find("-", dot)
tz_candidates = [pos for pos in (tz_plus, tz_minus) if pos != -1]
if tz_candidates:
tz_pos = min(tz_candidates)
frac = text[dot + 1:tz_pos]
text = text[:dot + 1] + frac[:6].ljust(6, "0") + text[tz_pos:]
dt = datetime.fromisoformat(text)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
+4 -11
View File
@@ -18,7 +18,7 @@ from ...env import ResolvedEnv
from ...git_gate import GitGatePlan
from ...manifest import Manifest
from ...supervise import SupervisePlan
from .. import ActiveAgent, BottleBackend, BottleImages, BottleSpec
from .. import ActiveAgent, BottleBackend, BottleSpec
from . import cleanup as _cleanup
from . import enumerate as _enumerate
from . import launch as _launch
@@ -92,18 +92,11 @@ class FirecrackerBottleBackend(
stage_dir=stage_dir,
)
def _build_or_load_images(self, plan: FirecrackerBottlePlan) -> BottleImages:
return BottleImages(agent=_launch.build_or_load_agent_base(plan))
def prelaunch_checks(self, plan: FirecrackerBottlePlan) -> None:
_launch.stale_checks(plan)
@contextmanager
def _launch_impl(
self, plan: FirecrackerBottlePlan, images: BottleImages,
def launch(
self, plan: FirecrackerBottlePlan
) -> Generator[FirecrackerBottle, None, None]:
assert isinstance(images.agent, Path)
with _launch.launch(plan, images.agent, provision=self.provision) as bottle:
with _launch.launch(plan, provision=self.provision) as bottle:
yield bottle
def prepare_cleanup(self) -> FirecrackerBottleCleanupPlan:
@@ -33,7 +33,8 @@ from ...orchestrator.client import OrchestratorClient
from ...orchestrator.lifecycle import (
OrchestratorStartError, # re-exported so callers can catch it
)
from ..consolidated_util import provision_bottle, teardown_consolidated as _teardown_util
from ...orchestrator.registration import registration_inputs
from ..docker.gateway_provision import deprovision_git_gate, provision_git_gate
from . import infra_vm
@@ -50,7 +51,6 @@ 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 launch_consolidated(
@@ -68,11 +68,18 @@ def launch_consolidated(
url = infra.control_plane_url
client = OrchestratorClient(url)
transport = infra_vm.gateway_transport()
reg = provision_bottle(
client, guest_ip, egress_plan, git_gate_plan, transport,
image_ref=image_ref, tokens=tokens,
inputs = registration_inputs(egress_plan)
reg = client.register_bottle(
guest_ip, image_ref=image_ref, policy=inputs.policy,
metadata=inputs.metadata, tokens=tokens,
)
try:
provision_git_gate(
infra_vm.gateway_transport(), reg.bottle_id, git_gate_plan)
except Exception:
client.teardown_bottle(reg.bottle_id)
raise
# The shared gateway CA every agent on this host trusts for TLS
# interception — fetched from the infra VM over SSH.
return LaunchContext(
@@ -81,19 +88,16 @@ def launch_consolidated(
source_ip=guest_ip,
gateway_ca_pem=infra.gateway_ca_pem(),
orchestrator_url=url,
env_var_secret=reg.env_var_secret,
)
def teardown_consolidated(
bottle_id: str, *, orchestrator_url: str, timeout: float | None = None,
) -> None:
def teardown_consolidated(bottle_id: str, *, orchestrator_url: str) -> None:
"""Deregister the bottle and remove its git-gate state from the gateway
VM. Both steps are idempotent so this is safe from a cleanup trap. Does
NOT stop the infra VM — it's a persistent per-host singleton shared by
every bottle."""
_teardown_util(bottle_id, infra_vm.gateway_transport(),
orchestrator_url=orchestrator_url, timeout=timeout)
OrchestratorClient(orchestrator_url).teardown_bottle(bottle_id)
deprovision_git_gate(infra_vm.gateway_transport(), bottle_id)
__all__ = [
+1 -1
View File
@@ -5,7 +5,7 @@ backend — we stream the guest root filesystem out over the control
channel (SSH here). Unlike the other backends this needs no Docker: the
tar *is* the resumable artifact. `resume` extracts it and rebuilds a
fresh per-bottle ext4 with `mke2fs -d` (see `util.build_committed_rootfs_dir`
and `launch.build_or_load_agent_base`). The bottle keeps running after the
and `launch._build_agent_base`). The bottle keeps running after the
snapshot.
"""
@@ -58,12 +58,6 @@ def _rootfs_digest(dockerfile: Path) -> str:
return h.hexdigest()[:16]
def cached_agent_rootfs_dir(dockerfile: Path) -> Path | None:
"""Return the ready cached rootfs for ``dockerfile``, if one exists."""
base = util.cache_dir() / "rootfs" / f"agent-{_rootfs_digest(dockerfile)}"
return base if (base / ".bb-ready").is_file() else None
def build_agent_rootfs_dir(
dockerfile: Path, *, image_tag: str, smoke_test: tuple[str, ...] = (),
) -> Path:
@@ -78,7 +72,7 @@ def build_agent_rootfs_dir(
silent-failure image at build time rather than at first agent use."""
digest = _rootfs_digest(dockerfile)
base = util.cache_dir() / "rootfs" / f"agent-{digest}"
if cached_agent_rootfs_dir(dockerfile) is not None:
if (base / ".bb-ready").is_file():
info(f"using cached agent rootfs {base.name}")
return base
@@ -41,7 +41,7 @@ from . import util
_ARTIFACT_FORMAT = "1"
_REPO_ROOT = Path(__file__).resolve().parents[3]
_DOCKERFILES = ("Dockerfile.orchestrator", "Dockerfile.gateway", "Dockerfile.infra", "Dockerfile.infra.fc")
_DOCKERFILES = ("Dockerfile.orchestrator", "Dockerfile.gateway", "Dockerfile.infra")
_DEFAULT_BASE = "https://gitea.dideric.is"
_DEFAULT_OWNER = "didericis"
+13 -16
View File
@@ -33,7 +33,6 @@ from pathlib import Path
from typing import Generator
from ...log import die, info
from .. import util as backend_util
from ..docker import util as docker_mod
from ..docker.gateway_provision import GatewayProvisionError
from . import firecracker_vm, infra_artifact, netpool, util
@@ -94,18 +93,19 @@ class InfraVm:
"""The gateway's mitmproxy CA (PEM) that agents install to trust its
TLS interception. Generated a moment after boot, so this polls over
SSH until it appears (mirrors DockerGateway.ca_cert_pem)."""
def _fetch() -> str | None:
deadline = time.monotonic() + timeout
while True:
proc = subprocess.run(
util.ssh_base_argv(self.private_key, self.guest_ip)
+ [f"cat {_GATEWAY_CA_PATH}"],
capture_output=True, text=True, timeout=15, check=False,
)
ok = proc.returncode == 0 and "BEGIN CERTIFICATE" in proc.stdout
return proc.stdout if ok else None
try:
return backend_util.poll_ca_cert(_fetch, timeout=timeout)
except TimeoutError as exc:
die(str(exc))
if proc.returncode == 0 and "BEGIN CERTIFICATE" in proc.stdout:
return proc.stdout
if time.monotonic() >= deadline:
die(f"gateway CA not available after {timeout:g}s: "
f"{proc.stderr.strip() or 'empty'}")
time.sleep(_HEALTH_POLL_SECONDS)
def ensure_built() -> None:
@@ -125,19 +125,16 @@ def ensure_built() -> None:
def build_infra_images_with_docker() -> None:
"""Build the four fixed images from source with host Docker: orchestrator,
gateway, the shared infra base (Dockerfile.infra), then the Firecracker
infra image (Dockerfile.infra.fc: FROM infra + buildah). The launch host
uses this only in `BOT_BOTTLE_INFRA_BUILD=local` mode; `publish_infra`
uses it off-host to produce the published artifact."""
"""Build the three fixed images from source with host Docker: orchestrator,
gateway, then the combined infra image (`COPY --from` orchestrator, `FROM`
gateway). The launch host uses this only in `BOT_BOTTLE_INFRA_BUILD=local`
mode; `publish_infra` uses it off-host to produce the published artifact."""
docker_mod.build_image(
_ORCHESTRATOR_IMAGE, str(_REPO_ROOT), dockerfile="Dockerfile.orchestrator")
docker_mod.build_image(
_GATEWAY_IMAGE, str(_REPO_ROOT), dockerfile="Dockerfile.gateway")
docker_mod.build_image(
"bot-bottle-infra:latest", str(_REPO_ROOT), dockerfile="Dockerfile.infra")
docker_mod.build_image(
_INFRA_IMAGE, str(_REPO_ROOT), dockerfile="Dockerfile.infra.fc")
_INFRA_IMAGE, str(_REPO_ROOT), dockerfile="Dockerfile.infra")
def build_infra_rootfs_dir() -> Path:
+13 -42
View File
@@ -45,15 +45,13 @@ from ...git_gate import (
provision_git_gate_dynamic_keys,
revoke_git_gate_provisioned_keys,
)
from ...image_cache import check_stale_path
from ...log import die, info, warn
from ...log import info, warn
from ...supervise import SUPERVISE_PORT
from ..docker.egress import EGRESS_PORT
from ..util import AGENT_CA_BUNDLE, AGENT_CA_PATH
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 .consolidated_launch import (
launch_consolidated,
teardown_consolidated,
@@ -66,7 +64,6 @@ _GIT_HTTP_PORT = 9420
@contextmanager
def launch(
plan: FirecrackerBottlePlan,
agent_base: Path,
*,
provision: Callable[[FirecrackerBottlePlan, "FirecrackerBottle"], str | None],
) -> Generator[FirecrackerBottle, None, None]:
@@ -88,9 +85,11 @@ def launch(
raise teardown_exc
try:
# Step 1 (rootfs resolution/build) runs in BottleBackend.launch before
# this context starts resources. ``agent_base`` is the selected cache,
# fresh build, or committed snapshot.
# Step 1: agent rootfs. Built from the Dockerfile inside a Firecracker
# builder VM (buildah, no host docker); a committed snapshot is reused
# when present. Returns the base dir the per-bottle ext4 is made from.
plan, agent_base = _build_agent_base(plan)
# Step 2: mint the git-gate dynamic (gitea) deploy keys, if any.
git_gate_plan = plan.git_gate_plan
if git_gate_plan.upstreams:
@@ -113,7 +112,6 @@ def launch(
token_values = egress_resolve_token_values(
plan.egress_plan.token_env_map, effective_env,
)
teardown_timeout = resolve_teardown_timeout()
ctx = launch_consolidated(
plan.egress_plan, git_gate_plan,
guest_ip=slot.guest_ip,
@@ -123,7 +121,6 @@ def launch(
stack.callback(
teardown_consolidated, ctx.bottle_id,
orchestrator_url=ctx.orchestrator_url,
timeout=teardown_timeout,
)
# Step 5: install the SHARED gateway CA (replaces the per-bottle CA).
@@ -209,7 +206,9 @@ def launch(
teardown()
def build_or_load_agent_base(plan: FirecrackerBottlePlan) -> Path:
def _build_agent_base(
plan: FirecrackerBottlePlan,
) -> tuple[FirecrackerBottlePlan, Path]:
"""Produce the agent's base rootfs dir. Primary path: build the Dockerfile
inside a Firecracker builder VM (buildah, no host docker), smoke-testing
the image before export. A committed snapshot (freeze/migrate) is resumed
@@ -218,36 +217,13 @@ def build_or_load_agent_base(plan: FirecrackerBottlePlan) -> Path:
committed_tar = committed_rootfs_path(plan.slug)
if committed and committed_tar.is_file():
info(f"resuming from committed rootfs {committed_tar}")
return util.build_committed_rootfs_dir(committed_tar)
dockerfile = Path(plan.dockerfile_path)
if plan.spec.image_policy == "cached":
cached = image_builder.cached_agent_rootfs_dir(dockerfile)
if cached is None:
die(
f"cached agent rootfs for {plan.image!r} not found; "
"run without --cached-images to build it"
)
info(f"using cached agent rootfs {cached.name}")
return cached
return image_builder.build_agent_rootfs_dir(
dockerfile,
return plan, util.build_committed_rootfs_dir(committed_tar)
base = image_builder.build_agent_rootfs_dir(
Path(plan.dockerfile_path),
image_tag=plan.image,
smoke_test=runtime_for(plan.agent_provider_template).smoke_test,
)
def stale_checks(plan: FirecrackerBottlePlan) -> None:
"""Raise when the cached rootfs selected by this plan is stale."""
if plan.spec.image_policy != "cached":
return
committed = read_committed_image(plan.slug)
committed_tar = committed_rootfs_path(plan.slug)
if committed and committed_tar.is_file():
check_stale_path(f"agent rootfs {committed_tar}", committed_tar)
return
cached = image_builder.cached_agent_rootfs_dir(Path(plan.dockerfile_path))
if cached is not None:
check_stale_path(f"agent rootfs {cached}", cached / ".bb-ready")
return plan, base
# --- agent guest env -------------------------------------------------
@@ -263,11 +239,6 @@ def _agent_guest_env(plan: FirecrackerBottlePlan, host_ip: str) -> dict[str, str
"HTTPS_PROXY": proxy_url, "HTTP_PROXY": proxy_url,
"https_proxy": proxy_url, "http_proxy": proxy_url,
"NO_PROXY": no_proxy, "no_proxy": no_proxy,
# Rootfs export can leave Git's implicit XDG paths unreadable even
# after the runtime repair. Bypass that discovery and name the
# provisioned global config explicitly so insteadOf can never fall
# through to the credential-bearing upstream URL.
"GIT_CONFIG_GLOBAL": f"{plan.guest_home}/.gitconfig",
"NODE_EXTRA_CA_CERTS": AGENT_CA_PATH,
"SSL_CERT_FILE": AGENT_CA_BUNDLE,
"REQUESTS_CA_BUNDLE": AGENT_CA_BUNDLE,
@@ -131,29 +131,6 @@ def build_artifact(out_dir: Path) -> tuple[str, Path, Path]:
return version, gz, sha
def _try_download_published(out_dir: Path) -> tuple[str, Path, Path] | None:
"""If this version's artifact is already in the registry, download the gz
and sha to out_dir and return (version, gz_path, sha_path). Returns None
when not yet published."""
version = infra_artifact.infra_artifact_version(infra_vm._infra_init())
sha_url = infra_artifact.artifact_url(version, "rootfs.ext4.gz.sha256")
try:
with urllib.request.urlopen(infra_artifact._open(sha_url)):
pass
except urllib.error.HTTPError as e:
if e.code == 404:
return None
raise SystemExit(f"registry check failed (HTTP {e.code}): {sha_url}")
except urllib.error.URLError as e:
raise SystemExit(f"registry unreachable: {sha_url} ({e.reason})")
print(f"infra rootfs {version} already published — downloading instead of building")
gz = out_dir / "rootfs.ext4.gz"
sha = out_dir / "rootfs.ext4.gz.sha256"
infra_artifact._download(infra_artifact.artifact_url(version, "rootfs.ext4.gz"), gz)
infra_artifact._download(infra_artifact.artifact_url(version, "rootfs.ext4.gz.sha256"), sha)
return version, gz, sha
def _publish_bundle(root: Path, token: str) -> str:
version_file = root / "version.txt"
# Guard the read so a missing version.txt is a clean error, not a raw
@@ -210,8 +187,6 @@ def main(argv: list[str] | None = None) -> int:
help="build a candidate bundle in DIR without publishing")
mode.add_argument("--publish-dir", type=Path,
help="publish an already-built and tested candidate bundle")
parser.add_argument("--reuse-published", action="store_true",
help="with --output: download from registry if already published instead of building")
args = parser.parse_args(argv)
_, _, token = infra_artifact._config()
@@ -222,14 +197,6 @@ def main(argv: list[str] | None = None) -> int:
if args.output is not None:
args.output.mkdir(parents=True, exist_ok=True)
reused = None
if args.reuse_published:
reused = _try_download_published(args.output)
if reused is not None:
version, _, _ = reused
(args.output / "version.txt").write_text(version + "\n", encoding="utf-8")
print(f"reused published infra rootfs candidate {version}")
return 0
version, _gz, _sha = build_artifact(args.output)
(args.output / "version.txt").write_text(version + "\n", encoding="utf-8")
print(f"built infra rootfs candidate {version}")
-6
View File
@@ -373,12 +373,6 @@ mount -o remount,rw / 2>/dev/null
# scratch dirs there — git worktrees, build temp, `git init /tmp/...`, etc.
mkdir -p /tmp && chmod 1777 /tmp
# Rootfs export also maps the image's original owners to the unprivileged
# host build uid. That uid is not guaranteed to be node's uid in the guest;
# restore the home-directory boundary before any SSH provisioning runs.
chown node:node /home/node 2>/dev/null || true
chmod 755 /home/node 2>/dev/null || true
# Install the per-bottle SSH pubkey from the kernel cmdline.
KEY=$(sed -n 's/.*bb_pubkey=\([^ ]*\).*/\1/p' /proc/cmdline | base64 -d 2>/dev/null)
if [ -n "$KEY" ]; then
+4 -10
View File
@@ -12,7 +12,7 @@ from ...env import ResolvedEnv
from ...git_gate import GitGatePlan
from ...supervise import SupervisePlan
from ...manifest import Manifest
from .. import ActiveAgent, BottleBackend, BottleImages, BottleSpec
from .. import ActiveAgent, BottleBackend, BottleSpec
from . import cleanup as _cleanup
from . import enumerate as _enumerate
from . import launch as _launch
@@ -82,17 +82,11 @@ class MacosContainerBottleBackend(
stage_dir=stage_dir,
)
def prelaunch_checks(self, plan: MacosContainerBottlePlan) -> None:
_launch.stale_checks(plan)
def _build_or_load_images(self, plan: MacosContainerBottlePlan) -> BottleImages:
return _launch.build_or_load_images(plan)
@contextmanager
def _launch_impl(
self, plan: MacosContainerBottlePlan, images: BottleImages
def launch(
self, plan: MacosContainerBottlePlan
) -> Generator[MacosContainerBottle, None, None]:
with _launch.launch(plan, images, provision=self.provision) as bottle:
with _launch.launch(plan, provision=self.provision) as bottle:
yield bottle
def ensure_orchestrator(self) -> str:
@@ -36,11 +36,9 @@ from dataclasses import dataclass
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 . import util as container_mod
from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active
from ...orchestrator.client import OrchestratorClient
from ...orchestrator.registration import registration_inputs
from ..docker.gateway_provision import deprovision_git_gate, provision_git_gate
from .gateway import GATEWAY_NETWORK
from .gateway_provision import AppleGatewayTransport
from .infra import MacosInfraService, OrchestratorStartError
@@ -72,7 +70,6 @@ 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(
@@ -92,32 +89,6 @@ def ensure_gateway(
)
def live_source_ips(network: str) -> list[str]:
"""Every running agent container's address on `network`.
The reconciliation input: the orchestrator lives inside the infra
container and cannot enumerate the host's containers, so the host has to
tell it which bottles are actually up. Containers that have not been
assigned an address yet contribute nothing the reap's grace window, not
this list, is what protects an in-flight launch.
Raises `EnumerationError` when the live set cannot be determined
authoritatively: either the container listing fails or any individual
inspect fails. Callers must skip reconciliation in that case to avoid
unregistering healthy bottles."""
ips: list[str] = []
for agent in enumerate_active():
name = f"{CONTAINER_NAME_PREFIX}{agent.slug}"
ip = container_mod.inspect_container_network_ip(name, network)
if ip is None:
raise EnumerationError(
f"container inspect {name!r} failed; live set is not authoritative"
)
if ip:
ips.append(ip)
return ips
def register_agent(
egress_plan: EgressPlan,
git_gate_plan: GitGatePlan,
@@ -132,20 +103,17 @@ def register_agent(
container it is the attribution key the gateway resolves policy by.
Raises on failure; the caller tears down."""
client = OrchestratorClient(endpoint.orchestrator_url)
# Self-heal before registering: a launcher that died hard (SIGKILL, closed
# terminal, host sleep) never ran its teardown callback, leaving an active
# row with no container. vmnet recycles addresses, so such a row can
# collide with this bottle's — and `by_source_ip` fail-closes on ambiguity,
# which would resolve no policy at all and deny every host. Best-effort: a
# reconciliation failure must not block an otherwise-fine launch.
try:
client.reconcile(live_source_ips(endpoint.network))
except (OrchestratorClientError, EnumerationError) as e:
info(f"registry reconciliation skipped: {e}")
reg = provision_bottle(
client, source_ip, egress_plan, git_gate_plan, AppleGatewayTransport(),
image_ref=image_ref, tokens=tokens,
inputs = registration_inputs(egress_plan)
reg = client.register_bottle(
source_ip, image_ref=image_ref, policy=inputs.policy,
metadata=inputs.metadata, tokens=tokens,
)
try:
provision_git_gate(AppleGatewayTransport(), reg.bottle_id, git_gate_plan)
except Exception:
# Roll the registration back so a provisioning failure leaves no orphan.
client.teardown_bottle(reg.bottle_id)
raise
return LaunchContext(
bottle_id=reg.bottle_id,
identity_token=reg.identity_token,
@@ -153,25 +121,21 @@ def register_agent(
gateway_ip=endpoint.gateway_ip,
network=endpoint.network,
orchestrator_url=endpoint.orchestrator_url,
env_var_secret=reg.env_var_secret,
)
def teardown_consolidated(
bottle_id: str, *, orchestrator_url: str, timeout: float | None = None,
) -> None:
def teardown_consolidated(bottle_id: str, *, orchestrator_url: str) -> None:
"""Deregister the bottle and remove its git-gate state from the gateway.
Both steps are idempotent so this is safe from a cleanup trap. Does NOT
stop the gateway it's a persistent per-host singleton."""
_teardown_util(bottle_id, AppleGatewayTransport(),
orchestrator_url=orchestrator_url, timeout=timeout)
OrchestratorClient(orchestrator_url).teardown_bottle(bottle_id)
deprovision_git_gate(AppleGatewayTransport(), bottle_id)
__all__ = [
"GatewayEndpoint",
"LaunchContext",
"ensure_gateway",
"live_source_ips",
"register_agent",
"teardown_consolidated",
"ConsolidatedLaunchError",
@@ -19,10 +19,6 @@ CONTAINER_NAME_PREFIX = "bot-bottle-"
_INFRA_NAMES = frozenset({INFRA_NAME})
class EnumerationError(RuntimeError):
"""container list failed; the resulting live set is not authoritative."""
def enumerate_active() -> list[ActiveAgent]:
result = subprocess.run(
["container", "list", "--quiet"],
@@ -31,10 +27,7 @@ def enumerate_active() -> list[ActiveAgent]:
check=False,
)
if result.returncode != 0:
raise EnumerationError(
f"container list failed: "
f"{(result.stderr or '').strip() or '<no stderr>'}"
)
return []
out: list[ActiveAgent] = []
for name in sorted(line.strip() for line in result.stdout.splitlines()):
if not name.startswith(CONTAINER_NAME_PREFIX) or name in _INFRA_NAMES:
+14 -22
View File
@@ -41,7 +41,7 @@ from dataclasses import dataclass
from pathlib import Path
from ... import log
from ...orchestrator.gateway import GATEWAY_CA_CERT, MITMPROXY_HOME
from ...orchestrator.gateway import GATEWAY_CA_CERT
from ...orchestrator.lifecycle import (
DEFAULT_PORT,
DEFAULT_STARTUP_TIMEOUT_SECONDS,
@@ -52,9 +52,7 @@ from ...paths import (
CONTROL_PLANE_TOKEN_ENV,
HOST_DB_FILENAME,
host_control_plane_token,
host_gateway_ca_dir,
)
from .. import util as backend_util
from . import util as container_mod
from .gateway import (
DEFAULT_CA_TIMEOUT_SECONDS,
@@ -219,14 +217,6 @@ class MacosInfraService:
# Container-only DB volume: one kernel writes bot-bottle.db, never
# shared with the host or another guest.
"--volume", f"{self._db_volume}:{_DB_ROOT_IN_CONTAINER}",
# The DB needs a container-only ext4 volume for coherent SQLite
# locking, but the CA has no such constraint. Keep it in the host
# app-data root so infra-container recreation and Apple Container
# volume pruning cannot silently rotate every bottle's trust
# anchor (issue #450).
"--mount",
container_mod.bind_mount_spec(
str(host_gateway_ca_dir()), MITMPROXY_HOME),
# Bind-mount the control-plane source (read-only); a code change
# takes effect on relaunch with no image rebuild.
"--mount",
@@ -270,19 +260,21 @@ class MacosInfraService:
def ca_cert_pem(self, *, timeout: float = DEFAULT_CA_TIMEOUT_SECONDS) -> str:
"""The gateway's mitmproxy CA (PEM) agents install to trust its TLS
interception. Read through the container path backed by the persistent
host CA directory; polls because mitmproxy writes it a beat after
start."""
def _fetch() -> str | None:
interception. Read out of the container (the CA lives on a
container-internal path, not a host mount); polls because mitmproxy
writes it a beat after start."""
deadline = time.monotonic() + timeout
while True:
result = container_mod.run_container_argv(
["container", "exec", self._name, "cat", GATEWAY_CA_CERT])
return result.stdout if result.returncode == 0 and result.stdout.strip() else None
try:
return backend_util.poll_ca_cert(_fetch, timeout=timeout)
except TimeoutError as exc:
raise GatewayError(
f"gateway CA not available in {self._name} after {timeout:g}s"
) from exc
if result.returncode == 0 and result.stdout.strip():
return result.stdout
if time.monotonic() >= deadline:
raise GatewayError(
f"gateway CA not available in {self._name} after {timeout:g}s: "
f"{(result.stderr or '').strip() or 'empty'}"
)
time.sleep(_CA_POLL_SECONDS)
def stop(self) -> None:
"""Remove the infra container (idempotent). The DB volume persists."""
+18 -47
View File
@@ -53,9 +53,7 @@ from ...git_gate import (
revoke_git_gate_provisioned_keys,
)
from ...git_http_backend import DEFAULT_PORT as _GIT_HTTP_PORT
from ...image_cache import check_stale
from ...log import die, info, warn
from .. import BottleImages
from ...supervise import SUPERVISE_PORT
from ..docker.egress import EGRESS_PORT
from ..util import AGENT_CA_BUNDLE, AGENT_CA_PATH
@@ -67,7 +65,6 @@ from .gateway_hosts import (
set_gateway_host,
)
from .bottle_plan import MacosContainerBottlePlan
from ...orchestrator.config_store import resolve_teardown_timeout
from .consolidated_launch import (
GatewayEndpoint,
ensure_gateway,
@@ -79,43 +76,18 @@ _REPO_DIR = str(Path(__file__).resolve().parent.parent.parent.parent)
_AGENT_SLEEP_SECONDS = "2147483647"
def build_or_load_images(plan: MacosContainerBottlePlan) -> BottleImages:
"""Resolve the agent image ref for this plan. The gateway's own image is
built by `ensure_gateway` it belongs to the shared singleton."""
committed = read_committed_image(plan.slug)
if committed and container_mod.image_exists(committed):
info(f"using committed image {committed!r}")
return BottleImages(agent=committed)
if plan.spec.image_policy == "cached":
if not container_mod.image_exists(plan.image):
die(
f"cached agent image {plan.image!r} not found; "
"run without --cached-images to build it"
)
info(f"using cached agent image {plan.image!r}")
return BottleImages(agent=plan.image)
container_mod.build_image(plan.image, _REPO_DIR, dockerfile=plan.dockerfile_path)
return BottleImages(agent=plan.image)
@contextmanager
def launch(
plan: MacosContainerBottlePlan,
images: BottleImages,
*,
provision: Callable[[MacosContainerBottlePlan, "MacosContainerBottle"], str | None],
) -> Generator[MacosContainerBottle, None, None]:
"""Run, register, provision, and yield an Apple Container bottle on the
shared per-host gateway."""
"""Build, run, register, provision, and yield an Apple Container bottle on
the shared per-host gateway."""
stack = ExitStack()
bottle_for_revoke = plan.manifest.bottle
git_gate_dir_for_revoke = git_gate_state_dir(plan.slug)
plan = dataclasses.replace(
plan,
agent_provision=dataclasses.replace(plan.agent_provision, image=str(images.agent)),
)
def teardown() -> None:
teardown_exc: BaseException | None = None
try:
@@ -128,6 +100,8 @@ def launch(
raise teardown_exc
try:
plan = _build_images(plan)
# Step 1: the per-host singletons. Must precede the agent run — its
# proxy env needs the gateway's address at `container run` time.
endpoint = ensure_gateway()
@@ -168,7 +142,6 @@ def launch(
token_values = egress_resolve_token_values(
plan.egress_plan.token_env_map, effective_env,
)
teardown_timeout = resolve_teardown_timeout()
ctx = register_agent(
plan.egress_plan,
plan.git_gate_plan,
@@ -180,7 +153,6 @@ def launch(
stack.callback(
teardown_consolidated, ctx.bottle_id,
orchestrator_url=ctx.orchestrator_url,
timeout=teardown_timeout,
)
info(
f"agent {plan.container_name} registered "
@@ -218,23 +190,22 @@ def launch(
teardown()
def stale_checks(plan: MacosContainerBottlePlan) -> None:
"""Raise StaleImageError if a cached image is older than the configured
threshold. Only runs when image_policy is 'cached'. Called by the backend
class's _image_stale_checks before _launch_impl starts any resources."""
if plan.spec.image_policy != "cached":
return
def _build_images(plan: MacosContainerBottlePlan) -> MacosContainerBottlePlan:
"""Build the agent image. The gateway's own image is built by
`ensure_gateway` it belongs to the shared singleton, not to a bottle."""
committed = read_committed_image(plan.slug)
if committed and container_mod.image_exists(committed):
ts = container_mod.image_created_at(committed)
if ts is not None:
check_stale(f"agent image {committed!r}", ts)
return
if container_mod.image_exists(plan.image):
ts = container_mod.image_created_at(plan.image)
if ts is not None:
check_stale(f"agent image {plan.image!r}", ts)
info(f"using committed image {committed!r}")
return dataclasses.replace(
plan,
agent_provision=dataclasses.replace(
plan.agent_provision, image=committed,
),
)
container_mod.build_image(
plan.image, _REPO_DIR, dockerfile=plan.dockerfile_path,
)
return plan
def _provision_git_gate_keys(
@@ -10,7 +10,6 @@ import shutil
import subprocess
import tempfile
import time
from datetime import datetime, timezone
from typing import Iterable
from ...log import die, info
@@ -573,41 +572,6 @@ def try_container_ipv4_on_network(name: str, network: str) -> str:
return ""
def inspect_container_network_ip(name: str, network: str) -> str | None:
"""IP of `name` on `network`, distinguishing inspect failure from "not yet".
Returns:
- the IP string when the container has one on `network`
- "" when inspect succeeds but no address is assigned yet (in-flight DHCP)
- None when the inspect command itself fails (authoritative list impossible)
"""
result = subprocess.run(
[_CONTAINER, "inspect", name],
capture_output=True, text=True, check=False,
)
if result.returncode != 0:
return None
try:
data = json.loads(result.stdout or "[]")
except json.JSONDecodeError:
return None
if isinstance(data, list):
data = data[0] if data else {}
if not isinstance(data, dict):
return None
status = data.get("status")
networks = status.get("networks") if isinstance(status, dict) else None
if not isinstance(networks, list):
return ""
for entry in networks:
if not isinstance(entry, dict) or entry.get("network") != network:
continue
raw = entry.get("ipv4Address")
if isinstance(raw, str) and raw:
return raw.split("/", 1)[0]
return ""
def wait_container_ipv4_on_network(
name: str, network: str, *, timeout: float = 15.0, poll: float = 0.25,
) -> str:
@@ -662,39 +626,6 @@ def image_id(ref: str) -> str:
raise AssertionError("unreachable")
def image_created_at(ref: str) -> datetime | None:
"""Return the image creation timestamp as an aware UTC datetime, or None
when the field is absent or unparseable (e.g. FROM-scratch images, images
pulled from registries that omit the field). Callers should skip the stale
check when None is returned rather than treating it as an error."""
result = subprocess.run(
[_CONTAINER, "image", "inspect", ref],
capture_output=True,
text=True,
check=False,
)
if result.returncode != 0:
die(
f"container image inspect for {ref!r} failed: "
f"{(result.stderr or '').strip() or '<no stderr>'}"
)
try:
data = json.loads(result.stdout or "{}")
except json.JSONDecodeError as exc:
die(f"container image inspect for {ref!r} returned malformed JSON: {exc}")
if isinstance(data, list) and data:
data = data[0]
if isinstance(data, dict):
value = data.get("created") or data.get("Created")
if isinstance(value, str) and value:
try:
ts = value.rstrip("Z")
return datetime.fromisoformat(ts).replace(tzinfo=timezone.utc)
except ValueError:
pass
return None
def save(ref: str, output: str) -> None:
subprocess.run([_CONTAINER, "image", "save", ref, "-o", output], check=True)
-20
View File
@@ -7,8 +7,6 @@ from __future__ import annotations
import hashlib
import os
import ssl
import time
from collections.abc import Callable
from pathlib import Path
from typing import TYPE_CHECKING
@@ -17,24 +15,6 @@ from ..log import die, info
if TYPE_CHECKING:
from ..egress import EgressPlan
_CA_POLL_INTERVAL = 0.5
def poll_ca_cert(fetch: Callable[[], str | None], *, timeout: float) -> str:
"""Poll `fetch` until it returns a non-empty PEM string or `timeout` expires.
`fetch` should return the PEM on success and `None` (or empty string) when
the cert is not yet available. Raises `TimeoutError` if the cert never
appears within `timeout` seconds."""
deadline = time.monotonic() + timeout
while True:
result = fetch()
if result:
return result
if time.monotonic() >= deadline:
raise TimeoutError(f"CA cert not available after {timeout:g}s")
time.sleep(_CA_POLL_INTERVAL)
# Debian-family CA layout, shared by every backend (all guest images
# are Debian-family). AGENT_CA_PATH is the source path that
+5 -4
View File
@@ -19,7 +19,6 @@ from .commit import cmd_commit
from .edit import cmd_edit
from .info import cmd_info
from .init import cmd_init
from .login import cmd_login
from .resume import cmd_resume
from .start import cmd_start
from .supervise import cmd_supervise
@@ -34,7 +33,6 @@ COMMANDS = {
"info": cmd_info,
"init": cmd_init,
"list": cmd_list,
"login": cmd_login,
"resume": cmd_resume,
"start": cmd_start,
"supervise": cmd_supervise,
@@ -45,7 +43,7 @@ COMMANDS = {
# the host (TAP pool, /dev/kvm, firecracker) and never opens the store, so
# gating it on the schema breaks preflight on a fresh CI runner where stdin
# isn't a TTY and the migration prompt can't be answered.
NO_MIGRATION_COMMANDS = frozenset({"backend", "login"})
NO_MIGRATION_COMMANDS = frozenset({"backend"})
def usage() -> None:
@@ -58,7 +56,6 @@ def usage() -> None:
sys.stderr.write(" info print env, skills, and prompt details for a named agent\n")
sys.stderr.write(" init interactively create a new agent and add it to bot-bottle.json\n")
sys.stderr.write(" list list available agents or active containers\n")
sys.stderr.write(" login register this host with a bot-bottle console\n")
sys.stderr.write(
" resume re-launch a bottle by its identity "
"(continues state from PRD 0016)\n"
@@ -114,3 +111,7 @@ def main(argv: list[str] | None = None) -> int:
return e.code if isinstance(e.code, int) else 1
except KeyboardInterrupt:
return 130
if __name__ == "__main__":
sys.exit(main())
-15
View File
@@ -1,15 +0,0 @@
"""Entry point for `python -m bot_bottle.cli`.
`cli.py` at the repo root is the usual way in; this makes the package
runnable too, so the CLI works from an installed copy where there is no
`cli.py` on disk to point at.
"""
from __future__ import annotations
import sys
from . import main
if __name__ == "__main__":
sys.exit(main())
-168
View File
@@ -1,168 +0,0 @@
"""bb login — register this host with a bot-bottle console.
Opens a device-authorization flow against the target console, waits for the
operator to approve, then writes access and refresh tokens to
~/.bot-bottle/console.json (or $BOT_BOTTLE_ROOT/console.json).
Usage:
bb login [--console-url URL] [--label LABEL]
Flags:
--console-url URL Target console URL (overrides BB_CONSOLE_URL env var)
--label LABEL Host label shown in the console (default: hostname)
"""
from __future__ import annotations
import json
import os
import socket
import sys
import tempfile
import time
import urllib.error
import urllib.request
from pathlib import Path
from typing import Any
from ..paths import bot_bottle_root
_CONSOLE_URL_ENV = "BB_CONSOLE_URL"
_POLL_SLEEP = 2 # seconds between polls; matches console's poll_interval default
def _usage() -> None:
sys.stderr.write(
"usage: bb login [--console-url URL] [--label LABEL]\n"
"\n"
"Options:\n"
" --console-url URL Console base URL (or BB_CONSOLE_URL env var)\n"
" --label LABEL Host label shown in the console (default: hostname)\n"
)
def _flag(argv: list[str], name: str) -> str | None:
for i, arg in enumerate(argv):
if arg == name and i + 1 < len(argv):
return argv[i + 1]
if arg.startswith(f"{name}="):
return arg[len(name) + 1:]
return None
def _post(url: str, payload: dict[str, Any]) -> dict[str, Any]:
data = json.dumps(payload).encode()
req = urllib.request.Request(
url, data=data, headers={"Content-Type": "application/json"}
)
with urllib.request.urlopen(req, timeout=10) as resp:
return json.loads(resp.read())
def _get(url: str) -> tuple[int, dict[str, Any]]:
req = urllib.request.Request(url)
try:
with urllib.request.urlopen(req, timeout=10) as resp:
return resp.status, json.loads(resp.read())
except urllib.error.HTTPError as e:
return e.code, {}
def _save_credentials(
console_url: str, host_id: str, access_token: str, refresh_token: str
) -> Path:
path = bot_bottle_root() / "console.json"
path.parent.mkdir(parents=True, exist_ok=True)
content = (
json.dumps(
{
"url": console_url,
"host_id": host_id,
"access_token": access_token,
"refresh_token": refresh_token,
},
indent=2,
)
+ "\n"
)
fd, tmp_path_str = tempfile.mkstemp(dir=path.parent, prefix=".console-")
tmp = Path(tmp_path_str)
try:
tmp.chmod(0o600)
with os.fdopen(fd, "w") as f:
f.write(content)
os.replace(tmp, path)
except OSError:
try:
tmp.unlink()
except OSError:
pass
raise
return path
def cmd_login(argv: list[str]) -> int:
if "--help" in argv or "-h" in argv:
_usage()
return 0
console_url = _flag(argv, "--console-url") or os.environ.get(_CONSOLE_URL_ENV)
if not console_url:
sys.stderr.write(
"bb login: --console-url or BB_CONSOLE_URL is required\n"
)
return 1
console_url = console_url.rstrip("/")
label = _flag(argv, "--label") or socket.gethostname()
try:
resp = _post(f"{console_url}/api/v1/hosts/authorize", {"label": label})
except (OSError, ValueError) as exc:
sys.stderr.write(f"bb login: failed to start authorization: {exc}\n")
return 1
device_code = resp["device_code"]
user_code = resp["user_code"]
expires_in = resp.get("expires_in", 300)
poll_sleep = max(1, min(int(resp.get("poll_interval", _POLL_SLEEP)), 60))
sys.stderr.write(
f"\nOpen this URL in your browser to authorize this host:\n\n"
f" {console_url}/hosts/authorize?code={user_code}\n\n"
f"Waiting for approval"
)
deadline = time.monotonic() + expires_in
while time.monotonic() < deadline:
sys.stderr.write(".")
sys.stderr.flush()
time.sleep(poll_sleep)
try:
code, result = _get(
f"{console_url}/api/v1/hosts/authorize/{device_code}"
)
except (OSError, ValueError):
continue
if code == 410:
break
st = result.get("status")
if st == "approved":
sys.stderr.write("\n\nApproved.\n")
path = _save_credentials(
console_url,
result["host_id"],
result["access_token"],
result["refresh_token"],
)
sys.stderr.write(f"Credentials saved to {path}\n")
return 0
if st == "denied":
sys.stderr.write("\n\nDenied by operator.\n")
return 1
sys.stderr.write("\n\nAuthorization timed out.\n")
return 1
+4 -33
View File
@@ -35,7 +35,6 @@ from ..bottle_state import (
is_preserved,
mark_preserved,
)
from ..image_cache import StaleImageError
from ..log import info, die
from ..manifest import Manifest, ManifestIndex
from ._common import PROG, USER_CWD, read_tty_line
@@ -65,14 +64,6 @@ def cmd_start(argv: list[str]) -> int:
"skip all prompts. For orchestrators, CI, and webhooks."
),
)
parser.add_argument(
"--cached-images",
action="store_true",
help=(
"quickstart with existing local agent and sidecar images; "
"only valid with --headless"
),
)
parser.add_argument(
"--bottle",
action="append",
@@ -105,8 +96,6 @@ def cmd_start(argv: list[str]) -> int:
help="agent name defined in bot-bottle.json (omit to pick interactively)",
)
args = parser.parse_args(argv)
if args.cached_images and not args.headless:
die("--cached-images is only supported with --headless")
dry_run = args.dry_run or os.environ.get("BOT_BOTTLE_DRY_RUN") == "1"
if args.no_cache or os.environ.get("BOT_BOTTLE_NO_CACHE") == "1":
@@ -158,10 +147,6 @@ def cmd_start(argv: list[str]) -> int:
label, color = tui.name_color_modal(default_label=agent_name)
label, color = _resolve_unique_label(label, color)
image_policy = _select_image_policy()
if image_policy is None:
return 0
spec = BottleSpec(
manifest=manifest,
agent_name=agent_name,
@@ -170,7 +155,6 @@ def cmd_start(argv: list[str]) -> int:
label=label,
color=color,
bottle_names=bottle_names,
image_policy=image_policy,
)
return _launch_bottle(
spec,
@@ -229,7 +213,6 @@ def _start_headless(
color=args.color or "",
bottle_names=bottle_names,
headless=True,
image_policy="cached" if args.cached_images else "fresh",
)
return _launch_bottle(
spec,
@@ -412,13 +395,6 @@ def _text_prompt_yes() -> bool:
return reply in ("y", "Y", "yes", "YES")
def _select_image_policy() -> str | None:
return tui.filter_select(
["fresh", "cached"],
title="Select image startup mode",
)
def _text_render_preflight():
def _render(plan: DockerBottlePlan, backend_name: str) -> None:
print(file=sys.stderr)
@@ -561,15 +537,6 @@ def _launch_bottle(
return 0
backend = get_bottle_backend(backend_name)
try:
backend.prelaunch_checks(plan)
except StaleImageError as exc:
if assume_yes:
die(str(exc))
sys.stderr.write(f"bot-bottle: {exc}\nLaunch anyway? [y/N] ")
sys.stderr.flush()
if read_tty_line() not in ("y", "Y", "yes", "YES"):
return 0
with backend.launch(plan) as bottle:
agent_provider_template = getattr(plan, "agent_provider_template", "claude")
extra_args: tuple[str, ...] = ()
@@ -588,6 +555,10 @@ def _launch_bottle(
f"session ended (exit {exit_code}); "
f"container {bottle.name} will be removed"
)
# While the container is still alive: always snapshot the
# transcript and — if the agent exited non-zero — mark
# the state for preservation. This picks up crashes /
# Ctrl-Cs / OOM kills before cleanup removes the state dir.
if agent_provider_template == "claude":
capture_claude_session_state(identity, exit_code)
return 0
-71
View File
@@ -1,71 +0,0 @@
"""SQLite-backed bot-bottle configuration store."""
from __future__ import annotations
from pathlib import Path
try:
from .db_store import DbStore
from .migrations import TableMigrations
from .paths import host_db_path
except ImportError:
from db_store import DbStore # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from migrations import TableMigrations # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from paths import host_db_path # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS = 1
class ConfigStore(DbStore):
"""SQLite configuration for host-side bot-bottle settings."""
def __init__(self, db_path: Path | None = None) -> None:
migrations = TableMigrations("config_store", [
# v1 — host-side bot-bottle settings
"""
CREATE TABLE IF NOT EXISTS bot_bottle_config (
id INTEGER PRIMARY KEY CHECK (id = 1),
cached_image_stale_warning_days INTEGER NOT NULL DEFAULT 1
)
""",
])
super().__init__(db_path or host_db_path(), migrations)
def cached_image_stale_warning_days(self) -> int:
if not self.db_path.is_file():
return DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS
with self._connect() as conn:
row = conn.execute(
"""
SELECT cached_image_stale_warning_days
FROM bot_bottle_config
WHERE id = 1
""",
).fetchone()
if row is None:
return DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS
try:
return int(row["cached_image_stale_warning_days"])
except (TypeError, ValueError):
return DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS
def set_cached_image_stale_warning_days(self, days: int) -> Path:
with self._connect() as conn:
conn.execute(
"""
INSERT INTO bot_bottle_config (id, cached_image_stale_warning_days)
VALUES (1, ?)
ON CONFLICT(id) DO UPDATE SET
cached_image_stale_warning_days = excluded.cached_image_stale_warning_days
""",
(days,),
)
self._chmod()
return self.db_path
__all__ = [
"DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS",
"ConfigStore",
]
+2 -15
View File
@@ -10,7 +10,7 @@
# Current Node LTS; slim variant keeps the image small while still
# providing apt-get for any future additions.
FROM node:22-trixie-slim
FROM node:22-slim
# Install runtime system deps. claude-code shells out to git for several
# features (status checks, commits, PR creation) — without git in the
@@ -21,15 +21,7 @@ FROM node:22-trixie-slim
# to it) works against egress's bumped TLS without the agent needing
# local DNS.
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
git \
ca-certificates \
curl \
openssh-client \
podman \
ripgrep \
iproute2 \
dnsutils \
&& apt-get install -y --no-install-recommends git ca-certificates curl ripgrep iproute2 dnsutils \
&& rm -rf /var/lib/apt/lists/*
# App-specific deps. Python isn't required by claude-code itself
@@ -47,11 +39,6 @@ RUN apt-get update \
RUN npm install -g --no-fund --no-audit @anthropic-ai/claude-code@2.1.172 \
&& npm cache clean --force
# Git reads both ~/.gitconfig and ~/.config/git/config. Keep its XDG config
# path traversable by the non-root runtime user so permission errors do not
# suppress bot-bottle's git-gate insteadOf rules.
RUN install -d -o node -g node -m 755 /home/node/.config /home/node/.config/git
# Run as a non-root user. The node image already provides a `node` user
# (uid 1000) with a home directory, which is where claude-code will write
# its session state.
+5 -17
View File
@@ -23,9 +23,8 @@ from ...agent_provider import (
provider_startup_args,
)
from ...backend.docker import util as docker_mod
from ...egress import CLAUDE_HOST_CREDENTIAL_TOKEN_REF, EgressRoute
from ...egress import EgressRoute
from ...log import die, info, warn
from .claude_auth import claude_host_access_token
if TYPE_CHECKING:
@@ -119,6 +118,7 @@ class ClaudeAgentProvider(AgentProvider):
color: str = "",
provider_settings: dict[str, object] | None = None,
) -> AgentProvisionPlan:
del forward_host_credentials, host_env
resolved_guest_env = dict(guest_env or {})
startup_args = provider_startup_args(provider_settings)
guest_home = self.guest_home
@@ -180,24 +180,13 @@ class ClaudeAgentProvider(AgentProvider):
claude_settings,
f"{guest_home}/.claude/settings.json",
))
provisioned_env: dict[str, str] = {}
if forward_host_credentials:
_host_env = host_env or dict(os.environ)
provisioned_env[CLAUDE_HOST_CREDENTIAL_TOKEN_REF] = (
claude_host_access_token(_host_env)
)
cred_token_ref = (
CLAUDE_HOST_CREDENTIAL_TOKEN_REF if forward_host_credentials
else auth_token
)
egress_routes = (EgressRoute(
host="api.anthropic.com",
auth_scheme="Bearer" if (auth_token or forward_host_credentials) else "",
token_ref=cred_token_ref,
auth_scheme="Bearer" if auth_token else "",
token_ref=auth_token,
),)
hidden_env_names: frozenset[str] = frozenset()
if auth_token or forward_host_credentials:
if auth_token:
env_vars["CLAUDE_CODE_OAUTH_TOKEN"] = "egress-placeholder"
hidden_env_names = frozenset({"CLAUDE_CODE_OAUTH_TOKEN"})
@@ -219,7 +208,6 @@ class ClaudeAgentProvider(AgentProvider):
files=tuple(files),
egress_routes=egress_routes,
hidden_env_names=hidden_env_names,
provisioned_env=provisioned_env,
)
def provision_skills(self, plan: "BottlePlan", bottle: "Bottle") -> None:
-114
View File
@@ -1,114 +0,0 @@
"""Host Claude auth helpers.
Reads the host's Claude Code credentials and returns only the access
token needed by egress. Does not expose refresh tokens or raw payloads.
Credential storage by platform:
Linux ~/.claude/.credentials.json
macOS macOS Keychain, service "Claude Code-credentials"
(file path is tried first; Keychain is the fallback)
"""
from __future__ import annotations
import json
import os
import subprocess
import sys
from datetime import datetime, timezone
from pathlib import Path
from ...log import die
_KEYCHAIN_SERVICE = "Claude Code-credentials"
def claude_auth_path(host_env: dict[str, str] | None = None) -> Path:
env = os.environ if host_env is None else host_env
home = env.get("HOME")
if home:
return Path(home) / ".claude" / ".credentials.json"
return Path.home() / ".claude" / ".credentials.json"
def _read_keychain() -> dict[str, object] | None:
"""Try the macOS Keychain. Returns parsed JSON dict or None."""
if sys.platform != "darwin":
return None
try:
result = subprocess.run(
["security", "find-generic-password", "-s", _KEYCHAIN_SERVICE, "-w"],
capture_output=True,
text=True,
timeout=10,
)
except (FileNotFoundError, subprocess.TimeoutExpired):
return None
if result.returncode != 0 or not result.stdout.strip():
return None
try:
raw = json.loads(result.stdout.strip())
except json.JSONDecodeError:
return None
return raw if isinstance(raw, dict) else None
def claude_host_access_token(
host_env: dict[str, str] | None = None,
*,
now: datetime | None = None,
) -> str:
path = claude_auth_path(host_env)
raw: dict[str, object] | None = None
if path.is_file():
try:
raw = json.loads(path.read_text())
except (OSError, json.JSONDecodeError) as e:
die(f"claude host credentials: could not read valid JSON at {path}: {e}")
if not isinstance(raw, dict):
die(f"claude host credentials: {path} must contain a JSON object")
else:
raw = _read_keychain()
if raw is None:
die(
f"claude host credentials: auth file missing at {path} and "
f"macOS Keychain lookup for '{_KEYCHAIN_SERVICE}' failed. "
"Run `claude login` on the host or disable "
"agent_provider.forward_host_credentials."
)
oauth = raw.get("claudeAiOauth")
if not isinstance(oauth, dict):
die(
"claude host credentials: claudeAiOauth is missing from credentials. "
"Run `claude login` on the host or disable "
"agent_provider.forward_host_credentials."
)
access_token = oauth.get("accessToken")
if not isinstance(access_token, str) or not access_token:
die(
"claude host credentials: claudeAiOauth.accessToken is missing or empty. "
"Run `claude login` on the host and restart the bottle."
)
# expiresAt is in milliseconds
expires_at = oauth.get("expiresAt")
if isinstance(expires_at, (int, float)):
check_now = now or datetime.now(timezone.utc)
exp_dt = datetime.fromtimestamp(float(expires_at) / 1000.0, timezone.utc)
if exp_dt <= check_now:
die(
"claude host credentials: host Claude access token is expired. "
"Run `claude login` on the host and restart the bottle."
)
return access_token
__all__ = [
"claude_auth_path",
"claude_host_access_token",
]
+2 -11
View File
@@ -3,17 +3,10 @@
# Mirrors the default Claude image shape: Node LTS, git/network tooling,
# non-root node user, and the provider CLI installed for that user.
FROM node:22-trixie-slim
FROM node:22-slim
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
git \
ca-certificates \
curl \
openssh-client \
podman \
procps \
ripgrep \
&& apt-get install -y --no-install-recommends git ca-certificates curl procps ripgrep \
&& rm -rf /var/lib/apt/lists/*
# App-specific deps. Python isn't required by codex itself
@@ -24,8 +17,6 @@ RUN apt-get update \
&& apt-get install -y --no-install-recommends python3 python3-pip python3-venv \
&& rm -rf /var/lib/apt/lists/*
RUN install -d -o node -g node -m 755 /home/node/.config /home/node/.config/git
USER node
WORKDIR /home/node
+2 -5
View File
@@ -2,7 +2,7 @@
#
# Node LTS, git/network tooling, and the Pi coding-agent CLI installed globally.
FROM node:22-trixie-slim
FROM node:22-slim
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
@@ -10,8 +10,6 @@ RUN apt-get update \
ca-certificates \
curl \
fd-find \
openssh-client \
podman \
ripgrep \
&& ln -s /usr/bin/fdfind /usr/local/bin/fd \
&& rm -rf /var/lib/apt/lists/*
@@ -23,8 +21,7 @@ RUN apt-get update \
RUN npm install -g --ignore-scripts --no-fund --no-audit @earendil-works/pi-coding-agent \
&& npm cache clean --force
RUN install -d -o node -g node -m 755 /home/node/.config /home/node/.config/git \
&& mkdir -p /home/node/.pi/agent \
RUN mkdir -p /home/node/.pi/agent \
/home/node/.pi/context-mode/sessions \
/tmp/pi-subagents-uid-1000 \
&& chown -R node:node /home/node/.pi /tmp \
-3
View File
@@ -30,7 +30,6 @@ if TYPE_CHECKING:
from .manifest import ManifestBottle
CODEX_HOST_CREDENTIAL_TOKEN_REF = "BOT_BOTTLE_CODEX_HOST_ACCESS_TOKEN"
CLAUDE_HOST_CREDENTIAL_TOKEN_REF = "BOT_BOTTLE_CLAUDE_HOST_ACCESS_TOKEN"
EGRESS_HOSTNAME = "egress"
@@ -146,7 +145,6 @@ def egress_manifest_routes(
outbound_detectors=r.OutboundDetectors,
inbound_detectors=r.InboundDetectors,
outbound_on_match=r.OutboundOnMatch,
preserve_auth=r.PreserveAuth,
))
return tuple(out)
@@ -402,7 +400,6 @@ class Egress(ABC):
)
__all__ = [
"CLAUDE_HOST_CREDENTIAL_TOKEN_REF",
"CODEX_HOST_CREDENTIAL_TOKEN_REF",
"EGRESS_HOSTNAME",
"EGRESS_ROUTES_FILENAME",
+1 -5
View File
@@ -367,10 +367,7 @@ class EgressAddon:
# Strip agent-set Authorization after DLP scan so smuggled tokens
# are caught above; the route may inject gateway-owned auth below.
# Routes with preserve_auth=True pass the header through as-is so the
# agent's own credentials (e.g. registry bearer tokens) reach the upstream.
if route is None or not route.preserve_auth:
flow.request.headers.pop("authorization", None)
flow.request.headers.pop("authorization", None)
# Build headers mapping for match evaluation
req_headers = {k.lower(): v for k, v in flow.request.headers.items()}
@@ -382,7 +379,6 @@ class EgressAddon:
env,
request_method=flow.request.method,
request_headers=req_headers,
deny_reason=config.deny_reason,
)
if decision.action == "block":
+8 -58
View File
@@ -78,7 +78,6 @@ class Route:
inbound_detectors: tuple[str, ...] | None = None
# "" means unset → DEFAULT_OUTBOUND_ON_MATCH. See OUTBOUND_ON_MATCH_VALUES.
outbound_on_match: str = ""
preserve_auth: bool = False
LOG_OFF = 0 # no logging
@@ -90,14 +89,6 @@ LOG_FULL = 2 # log block/warn events + full request and response bodies
class Config:
routes: tuple[Route, ...]
log: int = LOG_OFF
# Why this Config is a deny-all, when it is one for a reason *other* than
# the bottle's own policy genuinely not listing the host. A deny-all is
# indistinguishable from "policy loaded, host not allowed" at the decision
# point — both are simply "no matching route" — so without this the
# operator sees `host X is not in the allowlist` and goes hunting for a
# missing route that was never the problem. Empty for a normally-parsed
# policy; `decide` prefers it over the allowlist wording when set.
deny_reason: str = ""
@dataclass(frozen=True)
@@ -309,18 +300,11 @@ def _parse_one(idx: int, raw: object) -> Route:
idx, host, raw_dict,
)
preserve_auth_raw = raw_dict.get("preserve_auth", False)
if preserve_auth_raw is not True and preserve_auth_raw is not False:
raise ValueError(
f"{label} ({host}): 'preserve_auth' must be a boolean"
)
preserve_auth: bool = preserve_auth_raw
for k in raw_dict:
if k not in ("host", "matches", "auth_scheme", "token_env", "dlp", "git", "preserve_auth"):
if k not in ("host", "matches", "auth_scheme", "token_env", "dlp", "git"):
raise ValueError(
f"{label} ({host}): unknown key {k!r}; accepted keys "
f"are 'host', 'matches', 'auth_scheme', 'token_env', 'dlp', 'git', 'preserve_auth'"
f"are 'host', 'matches', 'auth_scheme', 'token_env', 'dlp', 'git'"
)
return Route(
@@ -332,7 +316,6 @@ def _parse_one(idx: int, raw: object) -> Route:
outbound_detectors=outbound_detectors,
inbound_detectors=inbound_detectors,
outbound_on_match=outbound_on_match,
preserve_auth=preserve_auth,
)
@@ -385,8 +368,6 @@ def route_to_yaml_dict(r: Route) -> dict[str, object]:
dlp["outbound_on_match"] = r.outbound_on_match
if dlp:
d["dlp"] = dlp
if r.preserve_auth:
d["preserve_auth"] = True
return d
@@ -424,40 +405,16 @@ class PolicyResolverLike(typing.Protocol):
...
# Deny-all explanations. Each names the *actual* failure so an operator isn't
# sent looking for a missing egress route when the bottle never had a policy
# to begin with — the failure mode that made a bricked registration read like
# a misconfigured allowlist.
DENY_UNATTRIBUTED = (
"egress: this request was not attributed to any bottle, so no egress "
"policy applies and every host is denied. Either the bottle's registry "
"row is missing/ambiguous (torn down, or another bottle claimed its "
"source IP), or the request carried no matching identity token — check "
"that the caller's proxy URL includes it. This is not an allowlist problem."
)
DENY_UNPARSEABLE = (
"egress: this bottle's egress policy could not be parsed, so it is being "
"treated as deny-all. Fix the bottle's egress.routes; every host is denied "
"until it loads."
)
DENY_RESOLVER_ERROR = (
"egress: the orchestrator could not be reached to resolve this bottle's "
"egress policy, so every host is denied (fail-closed). Check that the "
"control plane is up; this is not an allowlist problem."
)
def _config_from_policy(policy: "str | None") -> "Config":
"""Parse a resolved policy blob into a Config, fail-closed: None / empty /
unparseable all become a deny-all Config (no routes every request
blocked). Each deny-all carries the reason it is one, so the block message
names the real fault instead of blaming the allowlist."""
blocked)."""
if not policy:
return Config(routes=(), deny_reason=DENY_UNATTRIBUTED)
return Config(routes=()) # unattributed or empty → deny-all
try:
return load_config(policy)
except ValueError:
return Config(routes=(), deny_reason=DENY_UNPARSEABLE)
return Config(routes=()) # unparseable policy → deny
def resolve_client_config(
@@ -471,7 +428,7 @@ def resolve_client_config(
try:
policy = resolver.resolve(client_ip, identity_token)
except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught
return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR)
return Config(routes=()) # orchestrator unreachable/errored → deny
return _config_from_policy(policy)
@@ -500,7 +457,7 @@ def resolve_client_context(
client_ip, identity_token,
)
except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught
return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR), "", {}
return Config(routes=()), "", {} # orchestrator unreachable/errored → deny
return _config_from_policy(policy), (bottle_id or ""), tokens
@@ -615,16 +572,12 @@ def decide(
*,
request_method: str = "GET",
request_headers: typing.Mapping[str, str] | None = None,
deny_reason: str = "",
) -> Decision:
"""`deny_reason` is `Config.deny_reason`: when the deny-all came from a
missing/unparseable policy rather than the bottle's own allowlist, report
that instead of implying a route is merely absent."""
route = match_route(routes, request_host)
if route is None:
return Decision(
action="block",
reason=deny_reason or (
reason=(
f"egress: host {request_host!r} is not in the "
f"bottle's egress.routes allowlist. Declare a "
f"route for it or remove the request."
@@ -899,9 +852,6 @@ __all__ = [
"is_git_push_request",
"is_git_fetch_request",
"load_config",
"DENY_UNATTRIBUTED",
"DENY_UNPARSEABLE",
"DENY_RESOLVER_ERROR",
"resolve_client_config",
"resolve_client_context",
"PolicyResolverLike",
+15 -35
View File
@@ -61,11 +61,6 @@ class _DaemonSpec:
_EGRESS_ONLY_ENV_PREFIXES: tuple[str, ...] = ("EGRESS_TOKEN_",)
_READY_GATED_DAEMONS: tuple[str, ...] = ("git-gate", "git-http")
# Daemons that must be requested explicitly via BOT_BOTTLE_GATEWAY_DAEMONS
# and are NOT started in the default (env-var-unset) case. The orchestrator
# only runs in the combined infra container, never in a standalone gateway.
_OPT_IN_DAEMONS: frozenset[str] = frozenset({"orchestrator"})
def _env_for_daemon(name: str, base_env: dict[str, str]) -> dict[str, str]:
"""Egress sees the full bundle env. Everyone else gets a copy
@@ -80,14 +75,7 @@ def _env_for_daemon(name: str, base_env: dict[str, str]) -> dict[str, str]:
}
# The orchestrator is listed first so it starts before the gateway daemons,
# giving the control plane a head start to accept /resolve calls. The gateway
# daemons tolerate early /resolve failures and retry per-request.
_DAEMONS: tuple[_DaemonSpec, ...] = (
_DaemonSpec("orchestrator", (
"python3", "-m", "bot_bottle.orchestrator",
"--host", "0.0.0.0", "--port", "8099", "--broker", "stub",
)),
_DaemonSpec("egress", ("/bin/sh", "/app/egress-entrypoint.sh")),
_DaemonSpec("git-gate", ("/bin/sh", "/git-gate-entrypoint.sh")),
_DaemonSpec("git-http", ("python3", "-m", "bot_bottle.git_http_backend")),
@@ -115,20 +103,18 @@ def _selected_daemons(
env: dict[str, str],
all_daemons: Sequence[_DaemonSpec] | None = None,
) -> tuple[_DaemonSpec, ...]:
"""Filter the daemon set by the BOT_BOTTLE_GATEWAY_DAEMONS env var.
"""Filter the daemon set by the BOT_BOTTLE_GATEWAY_DAEMONS env
var. Unknown names in the list are ignored the caller is the
source of truth for which daemons are wired.
When the var is unset/empty, return all non-opt-in daemons (the
standard gateway subset). Opt-in daemons (e.g. `orchestrator`) only
run when explicitly named they never start in a plain gateway
container that doesn't set the env var. Unknown names are ignored.
`all_daemons` defaults to `_DAEMONS` resolved at call time (not at
definition time), so tests can pass a custom list."""
`all_daemons` defaults to `_DAEMONS` resolved at call time (not
at definition time), so tests can monkey-patch the module-level
`_DAEMONS` and have the new value take effect."""
if all_daemons is None:
all_daemons = _DAEMONS
raw = env.get("BOT_BOTTLE_GATEWAY_DAEMONS", "").strip()
if not raw:
return tuple(d for d in all_daemons if d.name not in _OPT_IN_DAEMONS)
return tuple(all_daemons)
wanted = {n.strip() for n in raw.split(",") if n.strip()}
return tuple(d for d in all_daemons if d.name in wanted)
@@ -150,7 +136,7 @@ def _pump(name: str, stream: IO[bytes]) -> None:
def _spawn(spec: _DaemonSpec) -> subprocess.Popen[bytes]:
env = _env_for_daemon(spec.name, dict(os.environ))
proc = subprocess.Popen( # pylint: disable=consider-using-with
proc = subprocess.Popen(
_argv_for_daemon(spec.name, spec.argv, env),
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
@@ -197,14 +183,6 @@ class _Supervisor:
except ProcessLookupError:
pass
def _sigkill_all(self) -> None:
for _, p in self.procs:
if p.poll() is None:
try:
p.kill()
except ProcessLookupError:
pass
def request_restart(self, daemon_name: str) -> bool:
"""Queue a daemon restart for the main loop to process.
@@ -257,7 +235,12 @@ class _Supervisor:
f"grace ({_GRACE_SECONDS:.0f}s) elapsed; SIGKILL on "
f"{', '.join(still_running)}"
)
self._sigkill_all()
for _, p in self.procs:
if p.poll() is None:
try:
p.kill()
except ProcessLookupError:
pass
done = all(p.poll() is not None for _, p in self.procs)
if done:
@@ -378,10 +361,7 @@ def main(argv: Sequence[str] | None = None) -> int:
# --signal HUP <bundle>` after writing routes.yaml. The kernel
# delivers SIGHUP to PID 1 (this supervisor); forward it to
# mitmdump so it reloads its addon.
signal.signal(
signal.SIGHUP,
lambda *_: sup.forward_signal(signal.SIGHUP, "egress"), # type: ignore[misc]
)
signal.signal(signal.SIGHUP, lambda *_: sup.forward_signal(signal.SIGHUP, "egress")) # type: ignore
while not sup.tick():
time.sleep(_POLL_INTERVAL)
-42
View File
@@ -1,42 +0,0 @@
"""Shared helpers for cached-image quickstart stale checks."""
from __future__ import annotations
from datetime import datetime, timezone
from pathlib import Path
try:
from .config_store import ConfigStore
except ImportError:
from config_store import ConfigStore # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
class StaleImageError(Exception):
"""Raised when a cached image or artifact exceeds the configured staleness
threshold. Callers can catch this to prompt interactively; headless paths
let it propagate as a fatal error."""
def check_stale(label: str, created_at: datetime) -> None:
"""Raise StaleImageError if `created_at` is older than the configured
stale-warning threshold. Negative threshold disables the check."""
threshold_days = ConfigStore().cached_image_stale_warning_days()
if threshold_days < 0:
return
now = datetime.now(timezone.utc)
created = created_at.astimezone(timezone.utc)
age = now - created
if age.total_seconds() <= threshold_days * 86400:
return
raise StaleImageError(
f"cached {label} is {age.days} day(s) old; "
"quickstart does not verify it matches the current Dockerfile/context"
)
def check_stale_path(label: str, path: Path) -> None:
"""Raise StaleImageError if `path`'s mtime exceeds the staleness threshold."""
check_stale(label, datetime.fromtimestamp(path.stat().st_mtime, tz=timezone.utc))
__all__ = ["StaleImageError", "check_stale", "check_stale_path"]
+4 -10
View File
@@ -25,9 +25,8 @@ class ManifestAgentProvider:
header, and sets a placeholder CLAUDE_CODE_OAUTH_TOKEN in the agent
so the Claude Code CLI starts.
`forward_host_credentials` forwards the host provider auth token into
the egress sidecar (Codex and Claude). For Codex this reads
`~/.codex/auth.json`; for Claude it reads `~/.claude/.credentials.json`.
`forward_host_credentials` forwards the host Codex auth token into
the egress daemon (Codex only).
"""
template: str = "claude"
@@ -93,15 +92,10 @@ class ManifestAgentProvider:
f"is only supported for built-in templates "
f"({', '.join(sorted(PROVIDER_TEMPLATES))})"
)
if forward_host_credentials and template not in {"codex", "claude"}:
if forward_host_credentials and template != "codex":
raise ManifestError(
f"bottle '{bottle_name}' agent_provider.forward_host_credentials "
"is only supported for templates 'codex' and 'claude'"
)
if forward_host_credentials and auth_token:
raise ManifestError(
f"bottle '{bottle_name}' agent_provider.forward_host_credentials "
"and auth_token both set; use one or the other"
"is currently only supported for template 'codex'"
)
settings = _parse_provider_settings(bottle_name, template, d.get("settings"))
return cls(
+2 -15
View File
@@ -71,7 +71,6 @@ class ManifestEgressRoute:
OutboundDetectors: tuple[str, ...] | None = None
InboundDetectors: tuple[str, ...] | None = None
OutboundOnMatch: str = ""
PreserveAuth: bool = False
@classmethod
def from_dict(cls, bottle_name: str, idx: int, raw: object) -> "ManifestEgressRoute":
@@ -191,22 +190,11 @@ class ManifestEgressRoute:
f"only 'fetch' is accepted"
)
# --- preserve_auth ---
preserve_auth = False
if "preserve_auth" in d:
raw_preserve_auth = d.get("preserve_auth")
if not isinstance(raw_preserve_auth, bool):
raise ManifestError(
f"{label} preserve_auth must be a boolean "
f"(was {type(raw_preserve_auth).__name__})"
)
preserve_auth = raw_preserve_auth
for k in d:
if k not in ("host", "matches", "auth", "role", "dlp", "git", "preserve_auth"):
if k not in ("host", "matches", "auth", "role", "dlp", "git"):
raise ManifestError(
f"{label} has unknown key {k!r}; accepted keys are "
f"'host', 'matches', 'auth', 'role', 'dlp', 'git', 'preserve_auth'"
f"'host', 'matches', 'auth', 'role', 'dlp', 'git'"
)
return cls(
@@ -219,7 +207,6 @@ class ManifestEgressRoute:
OutboundDetectors=outbound_detectors,
InboundDetectors=inbound_detectors,
OutboundOnMatch=outbound_on_match,
PreserveAuth=preserve_auth,
)
+3 -43
View File
@@ -15,7 +15,6 @@ from __future__ import annotations
import json
import urllib.error
import urllib.request
from collections.abc import Iterable
from dataclasses import dataclass
from ..paths import host_control_plane_token
@@ -41,13 +40,10 @@ 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. `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."""
identity token the agent presents for app-layer attribution."""
bottle_id: str
identity_token: str
env_var_secret: str = ""
class OrchestratorClient:
@@ -123,21 +119,17 @@ 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. 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."""
orchestrator holds in memory for the gateway to inject. 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")
@@ -145,24 +137,6 @@ 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."""
@@ -173,20 +147,6 @@ class OrchestratorClient:
raise OrchestratorClientError(f"teardown {bottle_id}: HTTP {status}")
return True
def reconcile(
self, live_source_ips: Iterable[str], *, grace_seconds: float | None = None,
) -> list[str]:
"""Drop registry rows for bottles that are no longer running
(`POST /reconcile`), returning the reaped bottle ids. `live_source_ips`
is the caller's enumeration of its live bottles — the orchestrator
can't see the backend from inside the infra container."""
body: dict[str, object] = {"live_source_ips": list(live_source_ips)}
if grace_seconds is not None:
body["grace_seconds"] = grace_seconds
payload = self._ok("POST", "/reconcile", body)
reaped = payload.get("reaped")
return [r for r in reaped if isinstance(r, str)] if isinstance(reaped, list) else []
def set_policy(self, bottle_id: str, policy: str) -> bool:
"""Live-reload a bottle's policy (`PUT /bottles/<id>/policy`). False on
404 (unknown bottle)."""
-107
View File
@@ -1,107 +0,0 @@
"""Per-host orchestrator configuration store (settings in bot-bottle.db).
Co-tenants the shared `bot-bottle.db` via the `DbStore` framework. Settings
are readable by the host launch path directly (no HTTP round-trip to the
orchestrator), so they take effect even before the orchestrator is reachable.
"""
from __future__ import annotations
import os
import sqlite3
from pathlib import Path
from ..db_store import DbStore
from ..migrations import TableMigrations
from ..paths import host_db_path
TEARDOWN_TIMEOUT_ENV = "BOT_BOTTLE_ORCHESTRATOR_TEARDOWN_TIMEOUT_SECONDS"
DEFAULT_TEARDOWN_TIMEOUT_SECONDS = 30.0
_MIGRATIONS = TableMigrations(
"orchestrator_config",
[
"""
CREATE TABLE IF NOT EXISTS orchestrator_config (
id INTEGER PRIMARY KEY CHECK (id = 1),
teardown_timeout_seconds REAL
)
""",
],
)
class OrchestratorConfigStore(DbStore):
"""Orchestrator settings in the shared host DB."""
def __init__(self, db_path: Path | None = None) -> None:
super().__init__(db_path or host_db_path(), _MIGRATIONS)
def _connect(self) -> sqlite3.Connection:
conn = super()._connect()
conn.execute("PRAGMA busy_timeout=5000")
return conn
def get_teardown_timeout_seconds(self) -> float | None:
"""Return the configured teardown timeout, or None if not set."""
try:
with self._connection() as conn:
row = conn.execute(
"SELECT teardown_timeout_seconds FROM orchestrator_config WHERE id = 1"
).fetchone()
except sqlite3.OperationalError:
return None
return row["teardown_timeout_seconds"] if row else None
def set_teardown_timeout_seconds(self, value: float) -> None:
"""Persist the teardown timeout."""
with self._connection() as conn:
conn.execute(
"INSERT OR REPLACE INTO orchestrator_config"
" (id, teardown_timeout_seconds) VALUES (1, ?)",
(value,),
)
self._chmod()
def delete_teardown_timeout_seconds(self) -> bool:
"""Clear the stored teardown timeout. Returns True if a value existed."""
with self._connection() as conn:
cur = conn.execute(
"UPDATE orchestrator_config SET teardown_timeout_seconds = NULL"
" WHERE id = 1 AND teardown_timeout_seconds IS NOT NULL"
)
return cur.rowcount > 0
def resolve_teardown_timeout(db_path: Path | None = None) -> float:
"""Return the teardown timeout to use, in priority order:
1. ``BOT_BOTTLE_ORCHESTRATOR_TEARDOWN_TIMEOUT_SECONDS`` env var
2. ``teardown_timeout_seconds`` in the orchestrator config DB
3. ``DEFAULT_TEARDOWN_TIMEOUT_SECONDS`` (30 s)
"""
raw = os.environ.get(TEARDOWN_TIMEOUT_ENV, "").strip()
if raw:
try:
value = float(raw)
if value > 0:
return value
except ValueError:
pass
store = OrchestratorConfigStore(db_path)
if not store.is_migrated():
store.migrate()
db_value = store.get_teardown_timeout_seconds()
if db_value is not None and db_value > 0:
return db_value
return DEFAULT_TEARDOWN_TIMEOUT_SECONDS
__all__ = [
"OrchestratorConfigStore",
"resolve_teardown_timeout",
"TEARDOWN_TIMEOUT_ENV",
"DEFAULT_TEARDOWN_TIMEOUT_SECONDS",
]
+1 -48
View File
@@ -9,17 +9,10 @@ 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"],
["tokens"], ["env_var_secret"]}
["metadata"], ["policy"]}
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": [...],
["grace_seconds"]}
POST /attribute -> 200 {"bottle_id"} | 403
POST /resolve -> 200 {"bottle_id","policy"} | 403
body: {"source_ip","identity_token"}
@@ -120,14 +113,12 @@ 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}
@@ -144,50 +135,12 @@ 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):
return 200, {"torn_down": True}
return 404, {"error": "no such bottle"}
if method == "POST" and route == "/reconcile":
# Host-driven self-heal: the caller enumerates its live bottles (only
# the host can see the backend) and the orchestrator drops rows for
# every other active bottle. Trusted-caller only — an agent that could
# reach this would be able to unregister its neighbours.
try:
data = _parse_json_object(body)
except ValueError as e:
return 400, {"error": f"invalid JSON: {e}"}
raw_ips = data.get("live_source_ips")
if not isinstance(raw_ips, list):
return 400, {"error": "live_source_ips (list of strings) is required"}
live = [ip for ip in raw_ips if isinstance(ip, str) and ip]
grace = data.get("grace_seconds")
kwargs = (
{"grace_seconds": float(grace)}
if isinstance(grace, (int, float)) and not isinstance(grace, bool)
else {}
)
return 200, {"reaped": orch.reconcile(live, **kwargs)}
if method == "POST" and route == "/attribute":
try:
data = _parse_json_object(body)
+10 -41
View File
@@ -27,7 +27,6 @@ from ..paths import (
CONTROL_PLANE_TOKEN_ENV,
host_control_plane_token,
host_db_path,
host_gateway_ca_dir,
)
from ..supervise import DB_PATH_IN_CONTAINER
@@ -49,23 +48,14 @@ GATEWAY_LABEL = "bot-bottle-orch-gateway=1"
# the source IP the gateway attributes by is the address on this network.
GATEWAY_NETWORK = "bot-bottle-gateway"
# mitmproxy's CA dir in the bundle. The host's gateway-CA dir (see
# `host_gateway_ca_dir`) is bind-mounted here so the gateway's self-generated
# CA stays STABLE across container recreation — every agent installs this one
# CA to trust the shared gateway's TLS interception, so it must not rotate when
# the gateway restarts. A host bind-mount rather than a named volume: a named
# volume is silently wiped by `docker volume prune`, minting a fresh CA that
# breaks every running bottle (issue #450).
# mitmproxy's CA dir in the bundle. A persistent named volume here keeps the
# gateway's self-generated CA STABLE across container recreation — every agent
# installs this one CA to trust the shared gateway's TLS interception, so it
# must not rotate when the gateway restarts.
MITMPROXY_HOME = "/home/mitmproxy/.mitmproxy"
GATEWAY_CA_VOLUME = "bot-bottle-gateway-mitmproxy"
GATEWAY_CA_CERT = f"{MITMPROXY_HOME}/mitmproxy-ca-cert.pem"
# The CA material mitmproxy writes into its confdir. mitmproxy reuses these on
# startup when present and generates them only on first run, so persisting them
# is what makes the CA stable; deleting them (see `rotate_gateway_ca`) forces a
# fresh CA on the next start. `mitmproxy-ca.pem` (cert + private key) is the
# signing identity; the rest are derived encodings agents/clients consume.
GATEWAY_CA_GLOB = "mitmproxy-ca*"
# The gateway data-plane image + its Dockerfile. Kept as a local constant
# rather than imported from the backend layer, which would drag
# the whole backend layer into the lean orchestrator (see #359); unify when
@@ -83,26 +73,6 @@ def _host_db_dir() -> str:
return str(db_dir)
def rotate_gateway_ca(ca_dir: Path | None = None) -> list[Path]:
"""Delete the persisted mitmproxy CA so the next gateway start mints a
fresh one the explicit, deliberate CA-rollover path (issue #450).
Persistence keeps the CA stable across restarts precisely because mitmproxy
reuses the on-disk CA; rotation is therefore just removing that material.
Returns the files removed (empty when there was no CA yet); idempotent.
This only clears the on-disk CA. It does NOT stop the running gateway (whose
mitmproxy still holds the old CA in memory) or re-provision agents the
caller recreates the gateway to mint the new CA and re-attaches bottles.
`rotate-ca` on the orchestrator CLI wires those steps together."""
ca_dir = ca_dir if ca_dir is not None else host_gateway_ca_dir()
removed: list[Path] = []
for path in sorted(ca_dir.glob(GATEWAY_CA_GLOB)):
path.unlink()
removed.append(path)
return removed
class GatewayError(Exception):
"""The shared gateway failed to build/start/stop (non-zero `docker` exit)."""
@@ -251,10 +221,9 @@ class DockerGateway(Gateway):
"--name", self.name,
"--label", GATEWAY_LABEL,
"--network", self.network,
# Persist the self-generated CA on the host so it survives both
# container recreation AND docker volume pruning (agents trust it)
# — see host_gateway_ca_dir / issue #450.
"--volume", f"{host_gateway_ca_dir()}:{MITMPROXY_HOME}",
# Persist the self-generated CA so it survives restarts (agents
# trust it) — see GATEWAY_CA_VOLUME.
"--volume", f"{GATEWAY_CA_VOLUME}:{MITMPROXY_HOME}",
# Share the one host DB: the supervise daemon queues proposals
# into the same file the orchestrator (and the operator, over
# HTTP) reads — no second, disconnected DB in the container.
@@ -303,7 +272,7 @@ class DockerGateway(Gateway):
__all__ = [
"Gateway", "DockerGateway", "GatewayError", "rotate_gateway_ca",
"Gateway", "DockerGateway", "GatewayError",
"GATEWAY_NAME", "GATEWAY_LABEL", "GATEWAY_IMAGE", "GATEWAY_NETWORK",
"GATEWAY_CA_CERT", "GATEWAY_CA_GLOB",
"GATEWAY_CA_VOLUME", "GATEWAY_CA_CERT",
]
+164 -159
View File
@@ -1,15 +1,17 @@
"""Orchestrator + gateway lifecycle (PRD 0070, docker slice).
Runs both the orchestrator control plane and the gateway data plane inside
a single `bot-bottle-infra` container on the shared gateway network
matching the structure already used by the macOS and Firecracker backends.
`gateway_init` is PID 1 and supervises both; the infra container is an
idempotent per-host singleton.
Runs the orchestrator control plane **as a container** on the shared gateway
network, alongside the gateway container. This is the PRD's "virtualize the
orchestrator": container↔container between the gateway and the orchestrator
avoids the host firewall (which drops containerhost traffic), and the gateway
reaches the control plane by container name over docker DNS. The host CLI
reaches it via a published loopback port.
The combined container replaces the prior two-container split
(bot-bottle-orchestrator + bot-bottle-orch-gateway). The host CLI reaches
the control plane via a published loopback port; gateway daemons reach it
over 127.0.0.1 (same container).
The orchestrator runs with the **register-only broker** the *backend*
launches agent containers (compose), so the orchestrator needs no docker
socket. That keeps this control-plane container unprivileged; the host manages
both containers. `ensure_running` is an idempotent singleton (fixed container
names + the published port).
"""
from __future__ import annotations
@@ -23,75 +25,50 @@ from pathlib import Path
from .. import log
from ..docker_cmd import run_docker
from ..paths import (
CONTROL_PLANE_TOKEN_ENV,
bot_bottle_root,
host_control_plane_token,
host_gateway_ca_dir,
)
from ..supervise import DB_PATH_IN_CONTAINER
from .gateway import (
GATEWAY_DOCKERFILE,
GATEWAY_IMAGE,
GATEWAY_NETWORK,
GatewayError,
MITMPROXY_HOME,
_host_db_dir,
)
from ..paths import CONTROL_PLANE_TOKEN_ENV, bot_bottle_root, host_control_plane_token
from .gateway import GATEWAY_IMAGE, GATEWAY_NAME, GATEWAY_NETWORK, DockerGateway, GatewayError
DEFAULT_PORT = 8099
DEFAULT_STARTUP_TIMEOUT_SECONDS = 45.0
INFRA_NAME = "bot-bottle-infra"
INFRA_LABEL = "bot-bottle-infra=1"
# The combined infra image: gateway data plane + orchestrator content.
# Built from Dockerfile.infra (FROM gateway + COPY --from orchestrator).
INFRA_IMAGE = os.environ.get("BOT_BOTTLE_INFRA_IMAGE", "bot-bottle-infra:latest")
INFRA_DOCKERFILE = "Dockerfile.infra"
# Baked as a container label so `ensure_running` can detect whether the
# running container is executing the current bind-mounted source.
INFRA_SOURCE_HASH_LABEL = "bot-bottle-infra-source-hash"
# Orchestrator image: the single canonical definition of the control-plane
# content (lean: python:3.12-slim + bot_bottle package, no mitmproxy/git).
# Used as a build intermediate: `Dockerfile.infra` COPY --from this image.
ORCHESTRATOR_NAME = "bot-bottle-orchestrator"
ORCHESTRATOR_LABEL = "bot-bottle-orchestrator=1"
# The control-plane's own runtime image — lean (python + the stdlib-only
# `bot_bottle` package, bind-mounted at run time), distinct from the heavy
# gateway data-plane image it used to borrow (#384). Env override for
# operators pinning a published build.
ORCHESTRATOR_IMAGE = os.environ.get(
"BOT_BOTTLE_ORCHESTRATOR_IMAGE", "bot-bottle-orchestrator:latest"
)
ORCHESTRATOR_DOCKERFILE = "Dockerfile.orchestrator"
# Baked onto the container as a label so `ensure_running` can tell whether the
# running process is executing the *current* bind-mounted source — see
# `source_hash`.
ORCHESTRATOR_SOURCE_HASH_LABEL = "bot-bottle-orchestrator-source-hash"
# The gateway daemons + orchestrator the infra container runs.
# BOT_BOTTLE_GATEWAY_DAEMONS listing `orchestrator` opts it in to
# gateway_init's supervise tree (see gateway_init._OPT_IN_DAEMONS).
_INFRA_DAEMONS = "egress,git-http,supervise,orchestrator"
# The bind-mount path for the live control-plane source inside the
# container. Separate from /app so the gateway's baked scripts
# (egress_addon.py, egress-entrypoint.sh) are not overlaid.
_SRC_IN_CONTAINER = "/bot-bottle-src"
# Bot-bottle host-root bind-mount inside the container (DB + state).
# The repo root is bind-mounted into the control-plane container so
# `python -m bot_bottle.orchestrator` resolves the package (the orchestrator
# is stdlib-only, so the lean orchestrator image's python is enough).
_REPO_ROOT = Path(__file__).resolve().parents[2]
_APP_DIR = "/app"
_ROOT_IN_CONTAINER = "/bot-bottle-root"
# The supervise daemon writes proposals into the host DB directory.
_SUPERVISE_DB_DIR_IN_CONTAINER = os.path.dirname(DB_PATH_IN_CONTAINER)
_HEALTH_POLL_SECONDS = 0.25
DEFAULT_STARTUP_TIMEOUT_SECONDS = 45.0
_HEALTH_REQUEST_TIMEOUT_SECONDS = 1.0
_REPO_ROOT = Path(__file__).resolve().parents[2]
class OrchestratorStartError(RuntimeError):
"""The infra container did not become healthy within the timeout."""
"""The orchestrator container did not become healthy within the timeout."""
def source_hash(repo_root: Path) -> str:
"""Content hash of the orchestrator's bind-mounted Python source (the
`bot_bottle` package the control-plane process imports). Changes only
when the code that would actually run changes `ensure_running`
recreates the container on a mismatch so a code change takes effect,
but leaves a healthy up-to-date container alone to preserve in-memory
egress tokens."""
`bot_bottle` package the control-plane process imports). This only
changes when the code that would actually run inside the container
changes `ensure_running` recreates the container on a mismatch and
otherwise leaves a healthy one alone, so a bottle launch that isn't
accompanied by a code change doesn't restart the process and drop every
*other* active bottle's in-memory egress tokens (`Orchestrator._tokens`
in `service.py`, never persisted to disk by design)."""
h = hashlib.sha256()
for path in sorted((repo_root / "bot_bottle").rglob("*.py")):
h.update(str(path.relative_to(repo_root)).encode())
@@ -100,37 +77,57 @@ def source_hash(repo_root: Path) -> str:
class OrchestratorService:
"""Manages the single per-host infra container (control plane + gateway).
"""Manages the orchestrator control-plane container + the shared gateway.
Callers only need `ensure_running()` + `url`.
`infra_name` / `infra_label` let backends run independent infra containers
on the same host without name collisions (e.g. isolated integration tests
that can't share the production INFRA_NAME singleton)."""
`orchestrator_name` / `orchestrator_label` let backends run independent
orchestrators on the same host without name collisions (e.g. the
Firecracker backend uses `bot-bottle-fc-orchestrator` alongside the Docker
backend's `bot-bottle-orchestrator`); `gateway_name` gives the paired
gateway container the same treatment (e.g. isolated integration tests
that can't share the production `GATEWAY_NAME` singleton). Subclass and
override `_gateway()` for anything `_gateway_image`/`gateway_name` can't
express (a genuinely backend-specific gateway variant)."""
def __init__(
self,
*,
port: int = DEFAULT_PORT,
network: str = GATEWAY_NETWORK,
image: str = INFRA_IMAGE,
image: str = ORCHESTRATOR_IMAGE,
gateway_image: str = GATEWAY_IMAGE,
gateway_name: str = GATEWAY_NAME,
repo_root: Path = _REPO_ROOT,
host_root: Path | None = None,
infra_name: str = INFRA_NAME,
infra_label: str = INFRA_LABEL,
orchestrator_name: str = ORCHESTRATOR_NAME,
orchestrator_label: str = ORCHESTRATOR_LABEL,
) -> None:
self.port = port
self.network = network
# Two distinct images (#384): `image` is the lean control-plane
# runtime this container runs; `_gateway_image` is the heavy egress /
# git-gate / supervise data plane the gateway container runs. They
# were one conflated image before the split.
self.image = image
self._gateway_image = gateway_image
self._gateway_name = gateway_name
self._repo_root = repo_root
self._host_root = host_root or bot_bottle_root()
self._infra_name = infra_name
self._infra_label = infra_label
self._orchestrator_name = orchestrator_name
self._orchestrator_label = orchestrator_label
@property
def url(self) -> str:
"""Host-side control-plane URL (published loopback port)."""
return f"http://127.0.0.1:{self.port}"
@property
def internal_url(self) -> str:
"""Control-plane URL as the gateway container reaches it — by name over
docker DNS on the shared network. This is the gateway's
BOT_BOTTLE_ORCHESTRATOR_URL."""
return f"http://{self._orchestrator_name}:{self.port}"
def is_healthy(self, *, timeout: float = _HEALTH_REQUEST_TIMEOUT_SECONDS) -> bool:
try:
with urllib.request.urlopen(f"{self.url}/health", timeout=timeout) as resp:
@@ -142,131 +139,139 @@ class OrchestratorService:
proc = run_docker(["docker", "ps", "--filter", f"name=^/{name}$", "--format", "{{.Names}}"])
return name in proc.stdout.split()
def _infra_source_current(self, current_hash: str) -> bool:
"""True iff the running infra container was started from the current
bind-mounted source. Mirrors the macOS backend's `_source_current`."""
if not self._container_running(self._infra_name):
return False
proc = run_docker([
"docker", "inspect", "--format",
"{{ index .Config.Labels \"" + INFRA_SOURCE_HASH_LABEL + "\" }}",
self._infra_name,
])
if proc.returncode != 0:
return True # can't compare → don't churn a working container
return proc.stdout.strip() == current_hash
def _ensure_network(self) -> None:
if run_docker(["docker", "network", "inspect", self.network]).returncode == 0:
return
proc = run_docker(["docker", "network", "create", self.network])
if proc.returncode != 0 and "already exists" not in proc.stderr:
raise GatewayError(
f"gateway network {self.network} failed to create: {proc.stderr.strip()}"
)
def _build_images(self) -> None:
"""Build the gateway base, the orchestrator intermediate, then the
infra image. All are cache-aware: a no-op when nothing changed."""
for tag, dockerfile in (
(GATEWAY_IMAGE, GATEWAY_DOCKERFILE),
(ORCHESTRATOR_IMAGE, ORCHESTRATOR_DOCKERFILE),
(self.image, INFRA_DOCKERFILE),
):
argv = ["docker", "build", "-t", tag,
"-f", str(self._repo_root / dockerfile),
str(self._repo_root)]
if os.environ.get("BOT_BOTTLE_NO_CACHE"):
argv.insert(2, "--no-cache")
proc = run_docker(argv)
if proc.returncode != 0:
raise GatewayError(f"{dockerfile} build failed: {proc.stderr.strip()}")
def _run_infra_container(self, current_hash: str) -> None:
"""Start the combined infra container (idempotent: clears a stale
fixed-name container first). Labels the container with `current_hash`
so a later `ensure_running` can detect a real code change."""
self._ensure_network()
run_docker(["docker", "rm", "--force", self._infra_name])
def _run_orchestrator_container(self, current_hash: str) -> None:
"""Start the control-plane container (idempotent: clears a stale
fixed-name container first). Register-only broker no docker socket.
Labels the container with `current_hash` so a later `ensure_running`
can detect a real code change (see `source_hash`)."""
run_docker(["docker", "rm", "--force", self._orchestrator_name])
proc = run_docker([
"docker", "run", "--detach",
"--name", self._infra_name,
"--label", self._infra_label,
"--label", f"{INFRA_SOURCE_HASH_LABEL}={current_hash}",
"--name", self._orchestrator_name,
"--label", self._orchestrator_label,
"--label", f"{ORCHESTRATOR_SOURCE_HASH_LABEL}={current_hash}",
"--network", self.network,
# Host CLI reaches the control plane here (loopback only).
# gateway_init always starts the orchestrator on DEFAULT_PORT (8099)
# inside the container; self.port is the host-side published port.
"--publish", f"127.0.0.1:{self.port}:{DEFAULT_PORT}",
# Persist the mitmproxy CA on the host so it survives container
# recreation AND docker volume pruning (issue #450): every agent
# trusts this one CA, so a fresh one would break all running bottles.
"--volume", f"{host_gateway_ca_dir()}:{MITMPROXY_HOME}",
# Shared supervise DB (same file the operator reads over HTTP).
"--volume", f"{_host_db_dir()}:{_SUPERVISE_DB_DIR_IN_CONTAINER}",
"--env", f"SUPERVISE_DB_PATH={DB_PATH_IN_CONTAINER}",
# Live control-plane source, mounted to a path that does not
# overlay the gateway's baked /app scripts.
"--volume", f"{self._repo_root}:{_SRC_IN_CONTAINER}:ro",
# PYTHONPATH lets the orchestrator (and other Python daemons)
# import the live source ahead of the installed package.
"--env", f"PYTHONPATH={_SRC_IN_CONTAINER}",
# Orchestrator registry DB on the host (sole writer: control plane).
# Host CLI reaches the control plane here; bound to loopback so it
# is not exposed on the host's external interfaces. NOTE: the
# container is still on `self.network` (the shared gateway network),
# so agents can reach it by container IP — which is exactly why the
# control plane requires the secret below rather than trusting the
# network boundary.
"--publish", f"127.0.0.1:{self.port}:{self.port}",
"--volume", f"{self._repo_root}:{_APP_DIR}:ro",
"--workdir", _APP_DIR,
# Persist the registry DB on the host (sole-owner: only the
# orchestrator opens bot-bottle.db).
"--volume", f"{self._host_root}:{_ROOT_IN_CONTAINER}",
"--env", f"BOT_BOTTLE_ROOT={_ROOT_IN_CONTAINER}",
# Control-plane secret: required by the orchestrator (to enforce)
# and by the gateway daemons (to present on /resolve calls).
# The control-plane secret it requires on every route but /health.
# Bare `--env NAME` → docker inherits the value from the run env
# below, so the secret never lands on argv / `docker inspect`.
"--env", CONTROL_PLANE_TOKEN_ENV,
# Gateway daemons reach the orchestrator over loopback at its
# fixed internal port (DEFAULT_PORT), independent of self.port.
"--env", f"BOT_BOTTLE_ORCHESTRATOR_URL=http://127.0.0.1:{DEFAULT_PORT}",
# Opt the orchestrator into gateway_init's supervise tree.
"--env", f"BOT_BOTTLE_GATEWAY_DAEMONS={_INFRA_DAEMONS}",
"--entrypoint", "python3",
self.image,
"-m", "bot_bottle.orchestrator",
"--host", "0.0.0.0", "--port", str(self.port), "--broker", "stub",
], env={**os.environ, CONTROL_PLANE_TOKEN_ENV: host_control_plane_token()})
if proc.returncode != 0:
raise OrchestratorStartError(
f"infra container failed to start: {proc.stderr.strip()}"
f"orchestrator container failed to start: {proc.stderr.strip()}"
)
def _gateway(self) -> DockerGateway:
return DockerGateway(
self._gateway_image,
name=self._gateway_name,
network=self.network,
orchestrator_url=self.internal_url,
)
def _ensure_orchestrator_image(self) -> None:
"""Build the lean control-plane image from `Dockerfile.orchestrator`
when it's missing (#384). Cheap — a `FROM python:*-slim` base with no
deps to install, so the layer cache makes rebuilds a no-op. Unlike the
gateway image this is build-if-missing, not build-every-time: the
control plane bind-mounts its source, so a code change is caught by the
source-hash recreate (below), not by an image rebuild."""
if run_docker(["docker", "image", "inspect", self.image]).returncode == 0:
return
argv = ["docker", "build", "-t", self.image,
"-f", str(self._repo_root / ORCHESTRATOR_DOCKERFILE),
str(self._repo_root)]
if os.environ.get("BOT_BOTTLE_NO_CACHE"):
argv.insert(2, "--no-cache")
proc = run_docker(argv)
if proc.returncode != 0:
raise GatewayError(
f"orchestrator image build failed: {proc.stderr.strip()}"
)
def _orchestrator_source_current(self, current_hash: str) -> bool:
"""True iff the running orchestrator container was created from the
*current* bind-mounted source. Mirrors `DockerGateway`'s
image-staleness check, but by content hash rather than image id since
the orchestrator runs bind-mounted source, not a built image."""
if not self._container_running(self._orchestrator_name):
return False
proc = run_docker([
"docker", "inspect", "--format",
"{{ index .Config.Labels \"" + ORCHESTRATOR_SOURCE_HASH_LABEL + "\" }}",
self._orchestrator_name,
])
if proc.returncode != 0:
return True # can't compare -> don't churn a working container
return proc.stdout.strip() == current_hash
def ensure_running(
self, *, startup_timeout: float = DEFAULT_STARTUP_TIMEOUT_SECONDS,
) -> str:
"""Ensure the infra container (control plane + gateway) is up; return
the host control-plane URL. Idempotent a healthy container on current
source is left untouched. Raises `OrchestratorStartError` on timeout."""
self._build_images()
"""Ensure the control plane + shared gateway are up; return the host
control-plane URL. Idempotent a healthy control plane running
current code and a running gateway are left untouched. Raises
`OrchestratorStartError` on timeout."""
gateway = self._gateway()
gateway.ensure_built() # rebuild the bundle image on a source change
gateway.ensure_running() # creates the shared network + (re)starts gateway
# Recreate the orchestrator container only when its bind-mounted
# source has actually changed since it started — its Python process
# loaded that code at startup and won't reload, so a stale container
# would keep running OLD control-plane code. Recreating on *every*
# launch (the prior behaviour) would drop every other active
# bottle's in-memory egress tokens each time a new bottle starts,
# since the orchestrator process holds them only in memory (#381).
current_hash = source_hash(self._repo_root)
if self.is_healthy() and self._infra_source_current(current_hash):
if self.is_healthy() and self._orchestrator_source_current(current_hash):
return self.url
log.info("starting infra container", context={"name": self._infra_name})
self._run_infra_container(current_hash)
self._ensure_orchestrator_image()
log.info(
"starting orchestrator container",
context={"name": self._orchestrator_name},
)
self._run_orchestrator_container(current_hash)
deadline = time.monotonic() + startup_timeout
while time.monotonic() < deadline:
if self.is_healthy():
log.info("infra container healthy", context={"url": self.url})
log.info("orchestrator healthy", context={"url": self.url})
return self.url
time.sleep(_HEALTH_POLL_SECONDS)
raise OrchestratorStartError(
f"infra container at {self.url} did not become healthy within {startup_timeout:g}s"
f"orchestrator at {self.url} did not become healthy within {startup_timeout:g}s"
)
def stop(self) -> None:
"""Remove the infra container (idempotent)."""
run_docker(["docker", "rm", "--force", self._infra_name])
"""Remove the orchestrator + gateway containers (idempotent)."""
run_docker(["docker", "rm", "--force", self._orchestrator_name])
self._gateway().stop()
__all__ = [
"OrchestratorService",
"OrchestratorStartError",
"INFRA_NAME",
"INFRA_IMAGE",
"INFRA_SOURCE_HASH_LABEL",
"ORCHESTRATOR_NAME",
"ORCHESTRATOR_IMAGE",
"DEFAULT_PORT",
"DEFAULT_STARTUP_TIMEOUT_SECONDS",
"source_hash",
]
-139
View File
@@ -32,7 +32,6 @@ import hmac
import secrets
import sqlite3
import time
from collections.abc import Iterable
from dataclasses import dataclass
from pathlib import Path
@@ -43,12 +42,6 @@ from ..paths import host_db_path
# 256 bits of urandom, URL-safe — unguessable per-bottle identity token.
IDENTITY_TOKEN_BYTES = 32
# How recently a row must have been registered to be exempt from
# `reap_absent`. Covers the window between `container run` and the address
# becoming visible to another launch's enumeration, so reconciliation never
# reaps a bottle that is still coming up.
DEFAULT_REAP_GRACE_SECONDS = 120.0
def new_identity_token() -> str:
"""A fresh per-bottle identity token (PRD 0070 attribution defence)."""
@@ -113,22 +106,6 @@ _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)",
],
)
@@ -248,70 +225,6 @@ class RegistryStore(DbStore):
).fetchall()
return [_row_to_record(r) for r in rows]
def reap_absent(
self,
live_source_ips: Iterable[str],
*,
grace_seconds: float = DEFAULT_REAP_GRACE_SECONDS,
now: float | None = None,
) -> list[BottleRecord]:
"""Delete active rows whose source IP is not held by a live bottle.
A row only ever leaves the registry two ways: an explicit
`teardown_bottle` (the launcher's cleanup callback) or the supersede
sweep in `register`. Neither runs when the launching CLI dies hard
SIGKILL, a closed terminal, a host sleep/crash so the row outlives
its container. That orphan is not inert: source IPs are recycled by
the backend's DHCP, and `by_source_ip` fail-closes on ambiguity, so a
leftover row at a reused address can brick the *next* bottle that
lands on it (no policy resolved -> every host denied, reported to the
agent as "not in the allowlist"). Reconciling against the live set at
launch keeps the registry from accumulating those landmines.
Restores the invariant the data plane needs: **at most one active row
per live address, and none at all for a dead one.** Two cases, because
a dead bottle's address may already have been handed to a live one:
* no live bottle holds the address every row there is an orphan;
* a live bottle holds it but several rows claim it the newest
registration is authoritative and the rest are orphans, the same
rule `register`'s same-IP supersede sweep applies. Without this
second case a recycled address stays ambiguous, which is exactly
the state that resolves no policy.
`grace_seconds` protects an in-flight launch: registration happens
moments after `container run`, and a concurrent launch's address may
not be visible to the caller's enumeration yet. Rows younger than the
grace window are never reaped, so reconciliation can't race a bottle
that is still coming up. Returns the deleted records."""
live = {ip for ip in live_source_ips if ip}
cutoff = (time.time() if now is None else now) - grace_seconds
with self._connection() as conn:
rows = conn.execute(
"SELECT * FROM orchestrator_bottles WHERE state = 'active'",
).fetchall()
by_ip: dict[str, list[BottleRecord]] = {}
for row in rows:
rec = _row_to_record(row)
by_ip.setdefault(rec.source_ip, []).append(rec)
candidates: list[BottleRecord] = []
for ip, recs in by_ip.items():
if ip not in live:
candidates.extend(recs)
continue
# Keep the newest claim on a live address; supersede the rest.
recs.sort(key=lambda r: r.created_at)
candidates.extend(recs[:-1])
doomed = [r for r in candidates if r.created_at <= cutoff]
for rec in doomed:
conn.execute(
"DELETE FROM orchestrator_bottles WHERE bottle_id = ?",
(rec.bottle_id,),
)
if doomed:
self._chmod()
return doomed
def by_source_ip(self, source_ip: str) -> BottleRecord | None:
"""Network-layer attribution: the single active bottle at this source
IP, or None if unknown or ambiguous (more than one a
@@ -342,57 +255,6 @@ 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",
@@ -400,5 +262,4 @@ __all__ = [
"new_identity_token",
"default_db_path",
"IDENTITY_TOKEN_BYTES",
"DEFAULT_REAP_GRACE_SECONDS",
]
-62
View File
@@ -1,62 +0,0 @@
"""Rotate the shared gateway's mitmproxy CA (issue #450).
python -m bot_bottle.orchestrator.rotate_ca
A deliberate CA rollover has two halves: drop the *persisted* CA so a fresh one
is minted, and drop the *running* gateway so its mitmproxy (which holds the old
CA in memory) is replaced. This one-shot command does both:
1. Delete the persisted CA under the host gateway-CA dir the next gateway
start generates a new one (mitmproxy reuses an existing CA, generates only
when absent).
2. Force-remove the infra / standalone-gateway containers so the stale
in-memory CA is gone; the next bottle launch's idempotent `ensure_running`
brings the gateway back up and mints the fresh CA.
It does NOT re-provision the new CA into already-running bottles those must be
re-attached so they install the new trust anchor. Rotation is thus an explicit,
operator-driven action with a brief egress interruption, not an automatic one.
"""
from __future__ import annotations
import sys
from pathlib import Path
from ..docker_cmd import run_docker
from ..paths import host_gateway_ca_dir
from .gateway import GATEWAY_NAME, rotate_gateway_ca
from .lifecycle import INFRA_NAME
# The containers whose mitmproxy would still be serving the old CA from memory:
# the consolidated infra container and the standalone per-host gateway.
_GATEWAY_CONTAINERS = (INFRA_NAME, GATEWAY_NAME)
def _out(msg: str) -> None:
sys.stdout.write(f"rotate-ca: {msg}\n")
def main(argv: list[str] | None = None) -> int:
del argv # no flags — a single deliberate action
ca_dir: Path = host_gateway_ca_dir()
removed = rotate_gateway_ca(ca_dir)
if removed:
_out(f"removed {len(removed)} CA file(s) from {ca_dir}")
else:
_out(f"no persisted CA under {ca_dir}; a fresh one is minted on next start")
# Drop any running gateway so its in-memory (now-stale) CA is replaced on
# the next launch. `rm --force` on an absent name is a tolerated no-op.
for name in _GATEWAY_CONTAINERS:
proc = run_docker(["docker", "rm", "--force", name])
if proc.returncode == 0 and proc.stdout.strip():
_out(f"removed running container {name}")
_out("done — the next bottle launch remints the CA; re-attach bottles to "
"install the new trust anchor")
return 0
if __name__ == "__main__":
raise SystemExit(main())
-94
View File
@@ -1,94 +0,0 @@
"""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"]
+2 -60
View File
@@ -13,20 +13,15 @@ Launch lifecycle:
and returns the record. If the broker rejects/fails, the registry entry
is rolled back so a failed launch leaves no orphan.
* `teardown_bottle` sends a signed teardown request, then deregisters.
* `reconcile` sweeps rows whose bottle is no longer running the
self-heal for the teardown paths that never got to run (a hard-killed
launcher), since an orphan row at a recycled source IP bricks the next
bottle that lands on it.
"""
from __future__ import annotations
import json
from collections.abc import Iterable
from datetime import datetime, timezone
from .broker import LaunchBroker, LaunchRequest, sign_request
from .registry import DEFAULT_REAP_GRACE_SECONDS, BottleRecord, RegistryStore
from .registry import BottleRecord, RegistryStore
from .gateway import Gateway
from ..supervise import (
AuditEntry,
@@ -87,22 +82,13 @@ 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.
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``)."""
launch doesn't take, so a failure leaves no orphan."""
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,
@@ -131,30 +117,6 @@ class Orchestrator:
self._tokens.pop(bottle_id, None)
return True
def reconcile(
self,
live_source_ips: Iterable[str],
*,
grace_seconds: float = DEFAULT_REAP_GRACE_SECONDS,
) -> list[str]:
"""Drop registry rows for bottles that are no longer running, and
forget their in-memory egress tokens. Returns the reaped bottle ids.
The caller supplies the live set because only the host can enumerate
its own containers the orchestrator runs *inside* the infra
container and has no view of the backend. Deliberately does not
broker a teardown: the container is already gone, so there is nothing
to stop, and a broker error must not stop the sweep from clearing
the row that would otherwise brick the next bottle at that address.
See `RegistryStore.reap_absent` for why orphans accumulate and why
they are harmful rather than merely untidy."""
reaped = self.registry.reap_absent(
live_source_ips, grace_seconds=grace_seconds)
for rec in reaped:
self._tokens.pop(rec.bottle_id, None)
return [rec.bottle_id for rec in reaped]
def tokens_for(self, bottle_id: str) -> dict[str, str]:
"""The bottle's in-memory egress auth tokens (env_name -> value), or
empty. The gateway injects these per request; they are never
@@ -293,26 +255,6 @@ 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:
-26
View File
@@ -33,13 +33,6 @@ HOST_DB_FILENAME = "bot-bottle.db"
CONTROL_PLANE_TOKEN_FILENAME = "control-plane-token"
CONTROL_PLANE_TOKEN_ENV = "BOT_BOTTLE_CONTROL_PLANE_TOKEN"
# The host directory holding the gateway's persistent mitmproxy CA. Bind-mounted
# into the infra/gateway container at mitmproxy's confdir so the self-generated
# CA survives container recreation — every agent installs this one CA to trust
# the shared gateway's TLS interception, so it must not rotate on restart. See
# host_gateway_ca_dir() for why this is a host bind-mount, not a named volume.
GATEWAY_CA_DIRNAME = "gateway-ca"
def bot_bottle_root() -> Path:
"""The app data root — `$BOT_BOTTLE_ROOT` if set, else `~/.bot-bottle`."""
@@ -66,23 +59,6 @@ def host_db_dir() -> Path:
return db_dir
def host_gateway_ca_dir() -> Path:
"""The directory holding the gateway's persistent mitmproxy CA, created if
missing. Backends bind-mount this into the infra/gateway container at
mitmproxy's confdir so the CA persists across container recreation.
A host bind-mount under the app-data root deliberately NOT a Docker
named volume. A named volume survives `docker rm` but is silently wiped by
`docker volume prune` / `docker system prune --volumes` during routine host
maintenance; the gateway then mints a fresh CA that every already-running
bottle distrusts, failing the TLS handshake even after it reconnects to the
moved gateway (issue #450). A path under the root docker never prunes it,
and it stays directly inspectable + rotatable from the host."""
ca_dir = bot_bottle_root() / GATEWAY_CA_DIRNAME
ca_dir.mkdir(parents=True, exist_ok=True)
return ca_dir
def host_control_plane_token() -> str:
"""The per-host control-plane secret, minted (256-bit, url-safe) and
persisted 0600 on first use, then reused.
@@ -118,10 +94,8 @@ __all__ = [
"HOST_DB_FILENAME",
"CONTROL_PLANE_TOKEN_FILENAME",
"CONTROL_PLANE_TOKEN_ENV",
"GATEWAY_CA_DIRNAME",
"bot_bottle_root",
"host_db_path",
"host_db_dir",
"host_gateway_ca_dir",
"host_control_plane_token",
]
-4
View File
@@ -6,11 +6,9 @@ from pathlib import Path
try:
from .audit_store import AuditStore
from .config_store import ConfigStore
from .queue_store import QueueStore
except ImportError:
from audit_store import AuditStore # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from config_store import ConfigStore # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from queue_store import QueueStore # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
_instance: StoreManager | None = None
@@ -49,13 +47,11 @@ class StoreManager:
return (
QueueStore("", self.db_path).is_migrated()
and AuditStore(self.db_path).is_migrated()
and ConfigStore(self.db_path).is_migrated()
)
def migrate(self) -> None:
QueueStore("", self.db_path).migrate()
AuditStore(self.db_path).migrate()
ConfigStore(self.db_path).migrate()
__all__ = ["StoreManager"]
-38
View File
@@ -312,44 +312,6 @@ reaches over the RPC rather than a shared mount into the VM. WAL on the
shared DB is therefore a deliberate, tested future change — not enabled ad
hoc. `sqlite3` itself is stdlib, so "the host needs SQLite" is a non-cost.
### Gateway CA: host-resident, like the DB
The shared gateway bumps TLS with a self-generated mitmproxy CA, and **every
bottle installs that CA** into its trust store to accept the bumped leaves. So
the CA is durable per-host state with the same rule as the DB: it must outlive
any single gateway container, or a restart mints a fresh CA that every
already-running bottle distrusts — the TLS handshake then fails even after the
bottle re-resolves and reconnects to the moved gateway (issue #450, a
re-attachment blocker distinct from #443/#445).
The CA lives on the **host filesystem** at `bot_bottle_root()/gateway-ca`
(`host_gateway_ca_dir()`), bind-mounted into the container at mitmproxy's
confdir. This is deliberately a host bind-mount, **not a container-runtime
named volume**: a named volume survives ordinary container removal but can be
silently wiped by Docker's or Apple Container's volume-prune commands during
routine host maintenance, which is exactly how the ephemeral-CA symptom shows
up in practice. A path under the app-data root is not managed or pruned by the
container runtime, and stays directly inspectable and rotatable from the host.
mitmproxy reuses an existing CA and generates one only on first run, so the
bind-mount alone gives
"adopt-existing, generate-on-first-run" for free.
The macOS backend uses the same host-resident CA directory and bind-mounts it
into the consolidated Apple infra container. Its `bot-bottle-mac-db` named
volume remains container-only because that prevents incoherent cross-kernel
SQLite locking, but the CA is deliberately not stored there: Apple Container
also has a `container volume prune` operation, and the named volume is
temporarily unreferenced while the infra container is recreated. Keeping the
CA on the host makes both ordinary recreation and volume pruning safe.
**Deliberate rollover** is the explicit inverse: `rotate_gateway_ca()` removes
the persisted CA material so the next start remints it, and the
`python -m bot_bottle.orchestrator.rotate_ca` one-shot wires that together with
dropping the running gateway container (whose mitmproxy still holds the old CA
in memory). Rotation does not auto-re-provision the new CA into running bottles
— those re-attach to install the new anchor — so it is an operator action with
a brief egress interruption, never an implicit one.
## Sequencing
Jump straight to the **virtualized** end state (not a host-daemon stepping
@@ -1,146 +0,0 @@
# PRD prd-new: Claude forward_host_credentials
- **Status:** Draft
- **Author:** claude
- **Created:** 2026-07-01
- **Issue:** #325
## Summary
Add `agent_provider.forward_host_credentials: true` support for the
`claude` template, mirroring the existing Codex flow. When enabled,
bot-bottle reads the host's Claude OAuth session key from
`~/.claude/.credentials.json` at launch, forwards it only to the egress sidecar,
and injects a placeholder `CLAUDE_CODE_OAUTH_TOKEN` into the agent so
Claude Code starts without ever seeing the real credential.
## Problem
Running a Claude agent in a container today requires the operator to
manually extract a long-lived OAuth token (`claude setup-token`), export
it as `BOT_BOTTLE_CLAUDE_OAUTH_TOKEN`, and reference it explicitly in
the manifest with `agent_provider.auth_token:
"BOT_BOTTLE_CLAUDE_OAUTH_TOKEN"`. This is a two-step manual ceremony
that is easy to skip or do incorrectly.
The host already stores a valid Claude session in `~/.claude/.credentials.json`
after `claude login`. Codex already automates an
equivalent extraction from `~/.codex/auth.json`. There is no reason
Claude bottles cannot do the same.
## Goals / Success Criteria
- A Claude bottle with `forward_host_credentials: true` in the manifest
uses the host's `~/.claude/.credentials.json` session key at launch with no
additional operator steps.
- The agent container receives only `CLAUDE_CODE_OAUTH_TOKEN=egress-placeholder`
— never the real token.
- The real session key lives only in the egress sidecar's environment.
- Missing, malformed, or expired host Claude auth fails launch with a
clear operator-facing message.
- Existing `auth_token` behavior is unchanged.
- `forward_host_credentials: true` is rejected in the manifest when both
`auth_token` and `forward_host_credentials` are set, since they serve
the same purpose.
## Non-goals
- Refreshing Claude OAuth tokens in the sidecar.
- Writing a dummy `~/.claude.json` auth state to the agent (unlike the
Codex flow, Claude Code reads its credential from `CLAUDE_CODE_OAUTH_TOKEN`
in env, not from an auth file — no guest-side auth marker is needed).
- Supporting `forward_host_credentials` for providers other than `codex`
and `claude`.
## Design
### Manifest schema
```yaml
agent_provider:
template: claude
forward_host_credentials: true
```
Rejects in manifest validation when:
- Template is not `codex` or `claude`.
- Both `auth_token` and `forward_host_credentials` are set.
### Host auth extraction (`contrib/claude/claude_auth.py`)
Claude Code credential storage varies by platform:
- **Linux**: `~/.claude/.credentials.json`
- **macOS**: macOS Keychain, service `"Claude Code-credentials"`
(the file path is tried first; Keychain is the fallback when the file
is absent)
`~/.claude.json` contains only UI state and profile metadata — no token.
The credentials JSON schema (same whether from file or Keychain):
```json
{
"claudeAiOauth": {
"accessToken": "<access-token>",
"refreshToken": "<refresh-token>",
"expiresAt": 1748276587173,
"scopes": ["user:inference", "user:profile"]
}
}
```
`expiresAt` is in **milliseconds** (not seconds).
At prepare/launch time, when `forward_host_credentials: true`:
1. Try `~/.claude/.credentials.json`; on macOS, if absent, run
`security find-generic-password -s "Claude Code-credentials" -w`
and parse its stdout as JSON.
2. Require a `claudeAiOauth` dict.
3. Require a non-empty `claudeAiOauth.accessToken` string.
4. If `claudeAiOauth.expiresAt` is present, divide by 1000 and require
the result to be in the future.
5. Return only the access token to the launch path.
Errors name the missing or invalid condition and point the operator at
`claude login`, without printing token values.
### Egress route
When `forward_host_credentials: true`:
- Provision the session key in `provisioned_env` under
`BOT_BOTTLE_CLAUDE_HOST_ACCESS_TOKEN` (new constant in `egress.py`).
- Set up the `api.anthropic.com` egress route with `auth_scheme: Bearer`
and `token_ref: BOT_BOTTLE_CLAUDE_HOST_ACCESS_TOKEN`.
- Set `CLAUDE_CODE_OAUTH_TOKEN=egress-placeholder` in the agent env and
add it to `hidden_env_names`.
No dummy auth file and no `verify` step are needed — Claude Code reads
the credential from the env var, not from a file.
### Constants
- `CLAUDE_HOST_CREDENTIAL_TOKEN_REF = "BOT_BOTTLE_CLAUDE_HOST_ACCESS_TOKEN"`
in `egress.py` (alongside the existing `CODEX_HOST_CREDENTIAL_TOKEN_REF`).
- `CLAUDE_HOST_CREDENTIAL_HOSTS = ("api.anthropic.com",)` in
`agent_provider.py` (alongside the existing `CODEX_HOST_CREDENTIAL_HOSTS`).
### Data flow
```
Host ~/.claude/.credentials.json → bot-bottle launch
├──► egress sidecar env (real token only)
└──► agent env: CLAUDE_CODE_OAUTH_TOKEN=egress-placeholder
Agent → HTTPS to api.anthropic.com (via egress)
Egress → injects Authorization: Bearer <real token>
Egress → forwards to api.anthropic.com
```
## Open questions
None — the Codex precedent makes the design clear.
@@ -1,156 +0,0 @@
# PRD prd-new: Consolidate infra backend for Docker
- **Status:** Active
- **Author:** Claude
- **Created:** 2026-07-20
- **Issue:** #431
## Summary
The Docker backend runs two containers — `bot-bottle-orch-gateway` (gateway
data plane) and `bot-bottle-orchestrator` (control plane) — where the
macOS and Firecracker backends already run a single combined infra
unit. This PRD collapses Docker to the same model: one `bot-bottle-infra`
container running both processes under the `gateway_init` supervise tree, a
restructured `Dockerfile.infra` as the shared gateway+orchestrator base,
and a handful of extracted shared utilities (CA cert polling, teardown
sequence, launch skeleton) that are currently duplicated across all three
`consolidated_launch.py` files.
## Goals / success criteria
- Docker backend starts exactly one infra container instead of two.
- `Dockerfile.infra` is the shared base image (gateway + orchestrator, no
buildah); the Firecracker image layers buildah on top of it.
- The orchestrator process runs under the `gateway_init` supervise tree
inside the combined container (one PID-1, one restart/health surface).
- CA cert polling, the teardown sequence, and the shared launch skeleton
(ensure-infra → register → provision → return context) live in a single
shared module; all three backends import from it.
- No functional change to macOS or Firecracker launch paths.
## Non-goals
- Changing per-bottle isolation — agents stay one-VM/container-each.
- Consolidating transport implementations (`DockerGatewayTransport`,
`AppleGatewayTransport`, `SshGatewayTransport`) — these are already the
right abstraction boundary.
- macOS DHCP-inversion of registration order — irreducible backend
difference, stays as-is.
- Any changes to the orchestrator RPC protocol or the attribution model.
## Design
### Dockerfile restructuring
**Current shape:**
- `Dockerfile.gateway` — data plane (mitmproxy, gitleaks, git, openssh,
supervise daemons)
- `Dockerfile.orchestrator` — control plane (python:3.12-slim + bot_bottle
package; stdlib-only, no third-party deps)
- `Dockerfile.infra` — Firecracker only: `FROM bot-bottle-gateway` +
buildah + `COPY --from bot-bottle-orchestrator`
**New shape:**
- `Dockerfile.gateway` — unchanged
- `Dockerfile.orchestrator` — unchanged (single definition of orchestrator
content; both Docker infra and Firecracker infra `COPY --from` it)
- `Dockerfile.infra`**shared base**: `FROM bot-bottle-gateway` + `COPY
--from bot-bottle-orchestrator` (no buildah — Docker infra image)
- `Dockerfile.infra.fc` — Firecracker only: `FROM bot-bottle-infra` +
buildah/crun/netavark/aardvark-dns (layered on the shared base, same net
result as today)
The comment in `Dockerfile.infra` that says "the docker backend keeps
orchestrator + gateway as separate images; this combined image exists only
for the Firecracker single-VM cut" is removed.
### Orchestrator in the supervise tree
`gateway_init` already supervises the data-plane daemons (egress, git-http,
supervise-MCP). The orchestrator control plane is added as another supervised
process: `python3 -m bot_bottle.orchestrator --host 0.0.0.0 --port <port>
--broker stub`.
The orchestrator source is bind-mounted (`/app` → repo root, as today) so
dev live-reload still works. `source_hash`-based container recreation in
`OrchestratorService.ensure_running` continues to apply — a code change
recreates the combined infra container, which bounces both gateway and
orchestrator. This is acceptable: the docker backend is a dev/legacy target
where in-flight egress connections across a code deploy are not a hard
requirement.
### `OrchestratorService` changes
`OrchestratorService` currently starts two containers in sequence: gateway
first (`DockerGateway.ensure_running`), then orchestrator. After this PRD:
- Single `docker run` of `bot-bottle-infra:latest`
- Container name: `bot-bottle-infra` (replaces `bot-bottle-orch-gateway` +
`bot-bottle-orchestrator`)
- Published ports: `127.0.0.1:{host_port}:8099` for the control plane
(`gateway_init` listens on a fixed internal port 8099; the caller-chosen
host port maps to it)
- Bind mounts: repo root + host root (same as today)
- `DockerGateway` becomes an implementation detail of `OrchestratorService`
rather than a separately started container; the gateway image name
(`GATEWAY_IMAGE`) is no longer referenced at runtime, only at build time
for the `Dockerfile.infra` base
The `_gateway()` / `ensure_running` two-step in `OrchestratorService` is
replaced by a single `_run_infra_container()`.
### Shared backend utilities
Three items are duplicated across
`backend/docker/consolidated_launch.py`,
`backend/macos_container/consolidated_launch.py`, and
`backend/firecracker/consolidated_launch.py`:
1. **CA cert polling loop** — `deadline = time.monotonic() + timeout; while
...: try fetch CA; sleep` — extracted to
`backend/consolidated_util.py:poll_ca_cert(transport, *, timeout)`.
2. **Teardown sequence**`OrchestratorClient(url).teardown_bottle(id)` +
`deprovision_git_gate(transport, id)` — extracted to
`backend/consolidated_util.py:teardown_consolidated(url, transport,
bottle_id)`.
3. **Launch skeleton** — all three follow: ensure-infra → allocate/register
→ provision git-gate → fetch CA cert → return launch context. The macOS
inversion (agent starts before registration, source IP from DHCP) is the
only deviation. Extract a shared `_provision_bottle(transport, bottle_id,
plan, orchestrator_url)` helper covering the register → provision →
return-token steps; the backends keep their own `launch_consolidated`
wrappers for the before/after (infra-ensure + agent-start + IP
allocation), calling the shared helper.
The new `backend/consolidated_util.py` module holds only backend-neutral,
transport-agnostic logic. All three backends import from it.
## Implementation chunks
1. **(this PR)** Dockerfile restructuring: rename current `Dockerfile.infra`
content to `Dockerfile.infra.fc`; write new `Dockerfile.infra` as
gateway+orchestrator base. Update Firecracker image-build references from
`Dockerfile.infra``Dockerfile.infra.fc`.
2. Add orchestrator process to `gateway_init` supervise tree.
3. Collapse `OrchestratorService` to a single-container start; rename
container from `bot-bottle-orch-gateway`/`bot-bottle-orchestrator`
`bot-bottle-infra`; update image name constant.
4. Extract `backend/consolidated_util.py` with `poll_ca_cert`,
`teardown_consolidated`, and `_provision_bottle`; update all three
`consolidated_launch.py` files to import from it.
5. Update tests that reference the old container names or two-container
startup sequence.
## Open questions
None — the supervise-tree approach and shared Dockerfile layering were
confirmed in issue #431.
@@ -1,47 +0,0 @@
# PRD prd-new: Modernize built-in agent images
- **Status:** Draft
- **Author:** Codex
- **Created:** 2026-07-21
- **Issue:** #451
## Summary
Keep every built-in agent provider on Debian's current stable release and make
Podman available inside each image. This gives agents a consistent, modern
userspace and an OCI container tool without requiring per-project setup.
## Problem
The Claude, Codex, and Pi images inherit the generic `node:22-slim` tag. That
tag does not state which Debian release the project supports and currently
leaves the images on the older Bookworm release. None of the built-in images
installs Podman, so tasks that need to inspect or build OCI images must first
modify the bottle or cannot run at all.
## Goals / success criteria
- Every Dockerfile under `bot_bottle/contrib/*/Dockerfile` explicitly inherits
`node:22-trixie-slim`, based on Debian 13 (the current stable release).
- Every built-in agent image installs Podman from Debian stable.
- Every built-in agent image retains an SSH client for Git-over-SSH workflows.
- The non-root agent user owns a traversable XDG Git configuration directory,
so Git can load bot-bottle's global git-gate rewrites without permission
errors.
- A shared test enforces both requirements for current and future built-in
providers.
## Non-goals
- Configuring privileged or nested-container execution for bottles.
- Pinning Podman outside Debian's stable package repository.
- Changing the Node.js or agent CLI release policy.
## Design
Use the explicit `node:22-trixie-slim` base rather than the floating `slim`
variant. Install the `podman` package with each image's existing `apt-get`
dependency layer, so package metadata and caches are still removed in the same
layer. Treat Debian stable as the Podman stability and update channel; this
keeps the images stdlib/distribution-first and avoids adding a third-party
package repository.
@@ -34,7 +34,6 @@ from tests._docker import skip_unless_docker
# image instead of leaking a new dangling tag on every invocation.
_TEST_ORCHESTRATOR_IMAGE = "bot-bottle-orchestrator:itest"
_TEST_GATEWAY_IMAGE = "bot-bottle-gateway:itest"
_TEST_INFRA_IMAGE = "bot-bottle-infra:itest"
@skip_unless_docker()
@@ -70,17 +69,20 @@ class TestDockerControlPlaneAuthIntegration(unittest.TestCase):
os.environ["BOT_BOTTLE_ROOT"] = cls._tmp.name
cls.addClassCleanup(_restore_root)
infra_name = f"bot-bottle-infra-itest-{suffix}"
orchestrator_name = f"bot-bottle-orch-itest-{suffix}"
gateway_name = f"bot-bottle-gw-itest-{suffix}"
network = f"bot-bottle-net-itest-{suffix}"
host_root = Path(cls._tmp.name)
cls.addClassCleanup(
cls._teardown_docker, infra_name, network, host_root
cls._teardown_docker, orchestrator_name, gateway_name, network, host_root
)
cls.svc = OrchestratorService(
infra_name=infra_name,
orchestrator_name=orchestrator_name,
gateway_name=gateway_name,
network=network,
image=_TEST_INFRA_IMAGE,
image=_TEST_ORCHESTRATOR_IMAGE,
gateway_image=_TEST_GATEWAY_IMAGE,
port=20000 + secrets.randbelow(10000),
host_root=host_root,
)
@@ -89,23 +91,23 @@ class TestDockerControlPlaneAuthIntegration(unittest.TestCase):
@staticmethod
def _teardown_docker(
infra_name: str, network: str, host_root: Path
orchestrator_name: str, gateway_name: str, network: str, host_root: Path
) -> None:
subprocess.run(
["docker", "rm", "--force", infra_name],
["docker", "rm", "--force", orchestrator_name, gateway_name],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False,
)
subprocess.run(
["docker", "network", "rm", network],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False,
)
# The infra container (no USER directive) wrote the registry
# The orchestrator container (no USER directive) wrote the registry
# DB as root into the throwaway host_root; chown it back so the
# (non-root) tempdir cleanup can remove it. Same workaround
# test_multitenant_isolation.py uses for the identical bind mount.
subprocess.run(
["docker", "run", "--rm", "-v", f"{host_root}:/r",
"--entrypoint", "chown", _TEST_INFRA_IMAGE, "-R",
"--entrypoint", "chown", _TEST_GATEWAY_IMAGE, "-R",
f"{os.getuid()}:{os.getgid()}", "/r"],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False,
)
+1 -1
View File
@@ -172,7 +172,7 @@ class TestSandboxEscape(unittest.TestCase):
# base image without producing five confusing
# command-not-found failures down the suite.
missing: list[str] = []
for tool in ("curl", "git", "dig", "ssh"):
for tool in ("curl", "git", "dig"):
r = cls._bottle.exec(f"command -v {tool} >/dev/null 2>&1")
if r.returncode != 0:
missing.append(tool)
+1 -66
View File
@@ -9,15 +9,11 @@ import unittest
from pathlib import Path
from bot_bottle.agent_provider import (
CLAUDE_HOST_CREDENTIAL_HOSTS,
CODEX_HOST_CREDENTIAL_HOSTS,
build_agent_provision_plan,
prompt_args,
)
from bot_bottle.egress import (
CLAUDE_HOST_CREDENTIAL_TOKEN_REF,
CODEX_HOST_CREDENTIAL_TOKEN_REF,
)
from bot_bottle.egress import CODEX_HOST_CREDENTIAL_TOKEN_REF
def _jwt(exp: int) -> str:
@@ -296,67 +292,6 @@ class TestAgentProviderRuntime(unittest.TestCase):
)
self.assertEqual({}, plan.provisioned_env)
def test_claude_forward_host_credentials_populates_egress_route(self):
access_token = "sk-ant-oat01-test-key" # gitleaks:allow
with tempfile.TemporaryDirectory(prefix="bb-provider.") as tmp:
home = Path(tmp) / "host-claude"
cred_dir = home / ".claude"
cred_dir.mkdir(parents=True)
(cred_dir / ".credentials.json").write_text(json.dumps({
"claudeAiOauth": {"accessToken": access_token},
}))
plan = build_agent_provision_plan(
template="claude",
dockerfile="",
state_dir=Path(tmp),
instance_name="bot-bottle-test",
prompt_file=Path(tmp) / "prompt.txt",
forward_host_credentials=True,
host_env={"HOME": str(home)},
)
self.assertEqual(1, len(plan.egress_routes))
route = plan.egress_routes[0]
self.assertIn(route.host, CLAUDE_HOST_CREDENTIAL_HOSTS)
self.assertEqual("Bearer", route.auth_scheme)
self.assertEqual(CLAUDE_HOST_CREDENTIAL_TOKEN_REF, route.token_ref)
self.assertEqual("egress-placeholder", plan.env_vars["CLAUDE_CODE_OAUTH_TOKEN"])
self.assertEqual(frozenset({"CLAUDE_CODE_OAUTH_TOKEN"}), plan.hidden_env_names)
def test_claude_forward_host_credentials_populates_provisioned_env(self):
access_token = "sk-ant-oat01-test-key" # gitleaks:allow
with tempfile.TemporaryDirectory(prefix="bb-provider.") as tmp:
home = Path(tmp) / "host-claude"
cred_dir = home / ".claude"
cred_dir.mkdir(parents=True)
(cred_dir / ".credentials.json").write_text(json.dumps({
"claudeAiOauth": {"accessToken": access_token},
}))
plan = build_agent_provision_plan(
template="claude",
dockerfile="",
state_dir=Path(tmp),
instance_name="bot-bottle-test",
prompt_file=Path(tmp) / "prompt.txt",
forward_host_credentials=True,
host_env={"HOME": str(home)},
)
self.assertEqual(
{CLAUDE_HOST_CREDENTIAL_TOKEN_REF: access_token},
plan.provisioned_env,
)
def test_claude_without_forward_host_credentials_has_empty_provisioned_env(self):
with tempfile.TemporaryDirectory(prefix="bb-provider.") as tmp:
plan = build_agent_provision_plan(
template="claude",
dockerfile="",
state_dir=Path(tmp),
instance_name="bot-bottle-test",
prompt_file=Path(tmp) / "prompt.txt",
forward_host_credentials=False,
)
self.assertEqual({}, plan.provisioned_env)
def test_pi_plan_writes_default_ollama_models(self):
with tempfile.TemporaryDirectory(prefix="bb-provider.") as tmp:
plan = build_agent_provision_plan(
-29
View File
@@ -1,29 +0,0 @@
"""Unit: shared cross-backend helpers in backend/util.py."""
from __future__ import annotations
import unittest
from unittest.mock import patch
from bot_bottle.backend import util as backend_util
class TestPollCaCert(unittest.TestCase):
def test_returns_pem_on_first_success(self) -> None:
result = backend_util.poll_ca_cert(lambda: "PEM", timeout=1.0)
self.assertEqual("PEM", result)
def test_raises_timeout_error_when_cert_never_appears(self) -> None:
with self.assertRaises(TimeoutError):
backend_util.poll_ca_cert(lambda: None, timeout=0.0)
def test_polls_until_cert_appears(self) -> None:
responses = iter([None, None, "-----BEGIN CERTIFICATE-----\n"])
with patch("bot_bottle.backend.util.time.sleep") as mock_sleep:
result = backend_util.poll_ca_cert(lambda: next(responses), timeout=5.0)
self.assertTrue(result.startswith("-----BEGIN CERTIFICATE-----"))
self.assertEqual(2, mock_sleep.call_count)
if __name__ == "__main__":
unittest.main()
-49
View File
@@ -1,49 +0,0 @@
"""Unit contracts shared by all built-in agent images."""
from __future__ import annotations
import re
import unittest
from pathlib import Path
_CONTRIB_DIR = Path(__file__).resolve().parents[2] / "bot_bottle/contrib"
_AGENT_DOCKERFILES = tuple(sorted(_CONTRIB_DIR.glob("*/Dockerfile")))
class TestBuiltinAgentImages(unittest.TestCase):
def test_all_use_debian_trixie_stable(self):
self.assertTrue(_AGENT_DOCKERFILES)
for dockerfile in _AGENT_DOCKERFILES:
with self.subTest(provider=dockerfile.parent.name):
self.assertRegex(
dockerfile.read_text(),
r"(?m)^FROM node:22-trixie-slim\s*$",
)
def test_all_install_podman(self):
for dockerfile in _AGENT_DOCKERFILES:
with self.subTest(provider=dockerfile.parent.name):
self.assertRegex(
dockerfile.read_text(),
re.compile(r"(?m)^\s*podman(?:\s|\\|$)"),
)
def test_all_install_ssh_client(self):
for dockerfile in _AGENT_DOCKERFILES:
with self.subTest(provider=dockerfile.parent.name):
self.assertRegex(
dockerfile.read_text(),
re.compile(r"(?m)^\s*openssh-client(?:\s|\\|$)"),
)
def test_all_prepare_node_git_config_directory(self):
for dockerfile in _AGENT_DOCKERFILES:
with self.subTest(provider=dockerfile.parent.name):
dockerfile_text = dockerfile.read_text()
self.assertIn("install -d -o node -g node", dockerfile_text)
self.assertIn("/home/node/.config/git", dockerfile_text)
if __name__ == "__main__":
unittest.main()
-268
View File
@@ -1,268 +0,0 @@
"""Unit tests for bb login command."""
from __future__ import annotations
import json
import os
import tempfile
import unittest
import urllib.error
from email.message import Message
from typing import Any
from unittest.mock import MagicMock, patch
class TestFlagParsing(unittest.TestCase):
def test_console_url_flag(self) -> None:
from bot_bottle.cli.login import _flag
self.assertEqual(_flag(["--console-url", "http://x"], "--console-url"), "http://x")
def test_console_url_equals_form(self) -> None:
from bot_bottle.cli.login import _flag
self.assertEqual(
_flag(["--console-url=http://x"], "--console-url"), "http://x"
)
def test_label_flag(self) -> None:
from bot_bottle.cli.login import _flag
self.assertEqual(_flag(["--label", "my-mac"], "--label"), "my-mac")
def test_missing_flag_returns_none(self) -> None:
from bot_bottle.cli.login import _flag
self.assertIsNone(_flag([], "--console-url"))
class TestHttpHelpers(unittest.TestCase):
def test_post_sends_json_and_decodes_response(self) -> None:
from bot_bottle.cli.login import _post
response = MagicMock()
response.__enter__.return_value.read.return_value = b'{"ok": true}'
with patch("urllib.request.urlopen", return_value=response) as urlopen:
self.assertEqual(
_post("http://console/start", {"label": "host"}), {"ok": True}
)
request = urlopen.call_args.args[0]
self.assertEqual(request.data, b'{"label": "host"}')
self.assertEqual(request.get_header("Content-type"), "application/json")
def test_get_decodes_success_response(self) -> None:
from bot_bottle.cli.login import _get
response = MagicMock()
response.__enter__.return_value.status = 200
response.__enter__.return_value.read.return_value = b'{"status": "pending"}'
with patch("urllib.request.urlopen", return_value=response):
self.assertEqual(
_get("http://console/status"), (200, {"status": "pending"})
)
def test_get_returns_http_error_status(self) -> None:
from bot_bottle.cli.login import _get
error = urllib.error.HTTPError(
"http://console/status", 410, "gone", Message(), None
)
with patch("urllib.request.urlopen", side_effect=error):
self.assertEqual(_get("http://console/status"), (410, {}))
class TestSaveCredentials(unittest.TestCase):
def test_writes_json_and_sets_perms(self) -> None:
from bot_bottle.cli.login import _save_credentials
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
path = _save_credentials("http://c", "hid", "at", "rt")
self.assertTrue(path.exists())
data = json.loads(path.read_text())
self.assertEqual(data["url"], "http://c")
self.assertEqual(data["host_id"], "hid")
self.assertEqual(data["access_token"], "at")
self.assertEqual(data["refresh_token"], "rt")
self.assertEqual(oct(path.stat().st_mode & 0o777), oct(0o600))
def test_temp_file_is_private_before_replace(self) -> None:
"""Temp file must be 0600 at the moment os.replace is called."""
from bot_bottle.cli.login import _save_credentials
from pathlib import Path as _Path
tmp_perms_at_replace: list[int] = []
real_replace = os.replace
def _spy_replace(
src: str | os.PathLike[str], dst: str | os.PathLike[str]
) -> None:
tmp_perms_at_replace.append(_Path(src).stat().st_mode & 0o777)
real_replace(src, dst)
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("os.replace", side_effect=_spy_replace):
path = _save_credentials("http://c", "hid", "at", "rt")
self.assertEqual(len(tmp_perms_at_replace), 1)
self.assertEqual(oct(tmp_perms_at_replace[0]), oct(0o600))
self.assertEqual(oct(path.stat().st_mode & 0o777), oct(0o600))
def test_cleanup_on_write_failure(self) -> None:
"""Temp file is removed and no credentials remain if replace fails."""
from bot_bottle.cli.login import _save_credentials
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("os.replace", side_effect=OSError("disk full")):
with self.assertRaises(OSError):
_save_credentials("http://c", "hid", "at", "rt")
leftovers = [f for f in os.listdir(tmp) if f.startswith(".console-")]
self.assertEqual(leftovers, [])
class TestCmdLoginMissingUrl(unittest.TestCase):
def test_help_returns_0(self) -> None:
from bot_bottle.cli.login import cmd_login
self.assertEqual(cmd_login(["--help"]), 0)
def test_returns_1_without_url(self) -> None:
from bot_bottle.cli.login import cmd_login
with patch.dict(os.environ, {}, clear=True):
os.environ.pop("BB_CONSOLE_URL", None)
result = cmd_login([])
self.assertEqual(result, 1)
def test_reads_env_var(self) -> None:
"""Exits 1 (network error) not because of missing URL when env var is set."""
from bot_bottle.cli.login import cmd_login
def _fail_post(_url: str, _payload: dict[str, Any]) -> dict[str, Any]:
raise OSError("connection refused")
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(
os.environ,
{
"BB_CONSOLE_URL": "http://localhost:9999",
"BOT_BOTTLE_ROOT": tmp,
},
):
with patch("bot_bottle.cli.login._post", side_effect=_fail_post):
result = cmd_login([])
self.assertEqual(result, 1)
class TestCmdLoginFlow(unittest.TestCase):
def _run_with_mocks(
self, poll_responses: list[dict[str, Any]], tmp: str
) -> int:
from bot_bottle.cli.login import cmd_login
start_resp = {
"device_code": "dc123",
"user_code": "ABC-DEF",
"expires_in": 300,
"poll_interval": 0,
}
poll_iter = iter(poll_responses)
def _fake_post(
_url: str, _payload: dict[str, Any]
) -> dict[str, Any]:
return start_resp
def _fake_get(_url: str) -> tuple[int, dict[str, Any]]:
return 200, next(poll_iter, {"status": "pending"})
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("bot_bottle.cli.login._post", side_effect=_fake_post):
with patch("bot_bottle.cli.login._get", side_effect=_fake_get):
with patch("time.sleep"):
return cmd_login(["--console-url", "http://console"])
def test_approved_flow_returns_0(self) -> None:
approved = {
"status": "approved",
"host_id": "hid",
"access_token": "at",
"refresh_token": "rt",
}
with tempfile.TemporaryDirectory() as tmp:
result = self._run_with_mocks(
[{"status": "pending"}, approved], tmp
)
self.assertEqual(result, 0)
with open(os.path.join(tmp, "console.json"), encoding="utf-8") as f:
creds = json.loads(f.read())
self.assertEqual(creds["host_id"], "hid")
def test_denied_flow_returns_1(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
result = self._run_with_mocks([{"status": "denied"}], tmp)
self.assertEqual(result, 1)
def test_timeout_returns_1(self) -> None:
from bot_bottle.cli.login import cmd_login
start_resp = {
"device_code": "dc",
"user_code": "ZZZ-ZZZ",
"expires_in": 0, # already expired; loop never runs
"poll_interval": 2,
}
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("bot_bottle.cli.login._post", return_value=start_resp):
result = cmd_login(["--console-url", "http://console"])
self.assertEqual(result, 1)
def test_poll_interval_from_server_is_used(self) -> None:
"""time.sleep must be called with the server-provided poll_interval."""
from bot_bottle.cli.login import cmd_login
server_interval = 7
start_resp = {
"device_code": "dc",
"user_code": "ABC-DEF",
"expires_in": 300,
"poll_interval": server_interval,
}
approved = {
"status": "approved",
"host_id": "hid",
"access_token": "at",
"refresh_token": "rt",
}
poll_iter = iter([{"status": "pending"}, approved])
def _fake_get(_url: str) -> tuple[int, dict[str, str]]:
return 200, next(poll_iter)
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("bot_bottle.cli.login._post", return_value=start_resp):
with patch(
"bot_bottle.cli.login._get", side_effect=_fake_get
):
with patch("time.sleep") as mock_sleep:
result = cmd_login(["--console-url", "http://console"])
self.assertEqual(result, 0)
self.assertTrue(mock_sleep.called)
for call in mock_sleep.call_args_list:
self.assertEqual(call.args[0], server_interval)
class TestDispatcherRegistration(unittest.TestCase):
def test_login_in_commands(self) -> None:
from bot_bottle.cli import COMMANDS
self.assertIn("login", COMMANDS)
def test_login_in_no_migration(self) -> None:
from bot_bottle.cli import NO_MIGRATION_COMMANDS
self.assertIn("login", NO_MIGRATION_COMMANDS)
if __name__ == "__main__":
unittest.main()
-40
View File
@@ -1,40 +0,0 @@
"""The CLI package is runnable as `python -m bot_bottle.cli`."""
from __future__ import annotations
import subprocess
import sys
import unittest
from pathlib import Path
_REPO_ROOT = Path(__file__).resolve().parents[2]
def _run(*args: str) -> subprocess.CompletedProcess[str]:
return subprocess.run(
[sys.executable, "-m", "bot_bottle.cli", *args],
cwd=_REPO_ROOT,
capture_output=True,
text=True,
check=False,
)
class TestModuleEntry(unittest.TestCase):
def test_help_exits_zero(self) -> None:
result = _run("--help")
self.assertEqual(result.returncode, 0, result.stderr)
self.assertIn("login", result.stderr)
def test_no_args_prints_usage(self) -> None:
# main() returns 2 with no command, matching the cli.py entry point.
self.assertEqual(_run().returncode, 2)
def test_subcommand_help_reaches_handler(self) -> None:
result = _run("login", "--help")
self.assertEqual(result.returncode, 0, result.stderr)
self.assertIn("--console-url", result.stderr)
if __name__ == "__main__":
unittest.main()
-12
View File
@@ -176,18 +176,6 @@ class TestCmdStartHeadless(unittest.TestCase):
self.assertEqual("researcher-2", self._spec().label)
def test_cached_images_sets_cached_policy(self):
start_mod.cmd_start(
["--headless", "--cached-images", "researcher", "--bottle", "claude",
"--prompt", "Do it"]
)
self.assertEqual("cached", self._spec().image_policy)
def test_cached_images_requires_headless(self):
with self.assertRaises(Die):
start_mod.cmd_start(["--cached-images", "researcher"])
self._launch_mock.assert_not_called()
class TestPrepareWithPreflight(unittest.TestCase):
"""prepare_with_preflight calls render_preflight with the plan and backend name."""
-21
View File
@@ -65,12 +65,6 @@ class TestCmdStartSelector(unittest.TestCase):
)
self._modal_patch.start()
self._image_policy_patch = patch(
"bot_bottle.cli.start._select_image_policy",
return_value="fresh",
)
self._image_policy_patch.start()
self._env_patch = patch.dict(os.environ, {}, clear=False)
self._env_patch.start()
os.environ.pop("BOT_BOTTLE_BACKEND", None)
@@ -81,7 +75,6 @@ class TestCmdStartSelector(unittest.TestCase):
self._agent_picker_patch.stop()
self._bottle_picker_patch.stop()
self._modal_patch.stop()
self._image_policy_patch.stop()
self._env_patch.stop()
# ------------------------------------------------------------------
@@ -140,19 +133,6 @@ class TestCmdStartSelector(unittest.TestCase):
spec = self._launch_mock.call_args[0][0]
self.assertEqual(("claude", "dev"), spec.bottle_names)
def test_image_policy_forwarded_to_spec(self):
with patch("bot_bottle.cli.start._select_image_policy", return_value="cached"):
start_mod.cmd_start(["researcher"])
self._launch_mock.assert_called_once()
spec = self._launch_mock.call_args[0][0]
self.assertEqual("cached", spec.image_policy)
def test_image_policy_cancel_returns_0(self):
with patch("bot_bottle.cli.start._select_image_policy", return_value=None):
rc = start_mod.cmd_start(["researcher"])
self.assertEqual(0, rc)
self._launch_mock.assert_not_called()
def test_empty_bottle_selection_forwarded(self):
self._bottle_picker_mock.return_value = []
start_mod.cmd_start(["researcher"])
@@ -234,7 +214,6 @@ class TestCmdStartLabelCollision(unittest.TestCase):
).start()
# Stub the bottle picker to always return a selection.
patch.object(tui_mod, "filter_multiselect", return_value=["claude"]).start()
patch("bot_bottle.cli.start._select_image_policy", return_value="fresh").start()
self.addCleanup(patch.stopall)
def test_no_collision_proceeds_without_reprompt(self):
-146
View File
@@ -1,146 +0,0 @@
"""Unit: _launch_bottle StaleImageError handling.
Exercises prelaunch_checks / backend.launch flow:
- headless mode die on stale
- interactive mode, user declines stop without launching
- interactive mode, user confirms skip stale check and launch once
"""
from __future__ import annotations
import io
import tempfile
import unittest
from types import SimpleNamespace
from typing import Any, cast
from unittest.mock import MagicMock, patch
from bot_bottle.image_cache import StaleImageError
from bot_bottle.log import Die
def _fake_plan() -> Any:
provision = SimpleNamespace(startup_args=())
return cast(Any, SimpleNamespace(
agent_provision=provision,
agent_provider_template="claude",
slug="dev-abc",
))
def _ok_cm(bottle: Any) -> MagicMock:
"""Return a context-manager mock that yields `bottle`."""
cm = MagicMock()
cm.__enter__ = MagicMock(return_value=bottle)
cm.__exit__ = MagicMock(return_value=False)
return cm
class TestLaunchBottleStaleHandling(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.mkdtemp(prefix="cli-stale-test.")
def _spec(self) -> Any:
from bot_bottle.backend import BottleSpec
from bot_bottle.manifest import ManifestIndex
idx = ManifestIndex.from_json_obj({
"bottles": {"dev": {}},
"agents": {"demo": {"skills": [], "prompt": "", "bottle": "dev"}},
})
return BottleSpec(
manifest=idx,
agent_name="demo",
copy_cwd=False,
user_cwd=self._tmp,
identity="dev-abc",
)
def _run_launch(self, **patch_kwargs: Any) -> int:
import bot_bottle.cli.start as start_mod
spec = self._spec()
with patch.object(start_mod, "prepare_with_preflight",
return_value=(_fake_plan(), "dev-abc")), \
patch.object(start_mod, "settle_state"), \
patch.object(start_mod, "info"):
return start_mod._launch_bottle(
spec,
dry_run=False,
backend_name="docker",
**patch_kwargs,
)
def test_headless_stale_calls_die(self) -> None:
"""In headless mode (assume_yes=True), a StaleImageError from prelaunch_checks must call die()."""
import bot_bottle.cli.start as start_mod
backend_mock = MagicMock()
backend_mock.prelaunch_checks.side_effect = StaleImageError("image is 5 day(s) old")
with patch.object(start_mod, "get_bottle_backend", return_value=backend_mock), \
patch.object(start_mod, "die", side_effect=Die()):
with self.assertRaises(Die):
self._run_launch(assume_yes=True)
backend_mock.launch.assert_not_called()
def test_interactive_user_declines_stops_before_launch(self) -> None:
"""Interactive user answering 'n' → launch is never called."""
import bot_bottle.cli.start as start_mod
backend_mock = MagicMock()
backend_mock.prelaunch_checks.side_effect = StaleImageError("image is 5 day(s) old")
with patch.object(start_mod, "get_bottle_backend", return_value=backend_mock), \
patch.object(start_mod, "read_tty_line", return_value="n"), \
patch("sys.stderr", new_callable=io.StringIO):
rc = self._run_launch(assume_yes=False)
self.assertEqual(0, rc)
backend_mock.launch.assert_not_called()
def test_interactive_user_confirms_launches_once(self) -> None:
"""Interactive user answering 'y' → prelaunch stale error is bypassed; launch called once."""
import bot_bottle.cli.start as start_mod
bottle_mock = MagicMock()
bottle_mock.name = "dev-abc"
backend_mock = MagicMock()
backend_mock.prelaunch_checks.side_effect = StaleImageError("image is 5 day(s) old")
backend_mock.launch.return_value = _ok_cm(bottle_mock)
with patch.object(start_mod, "get_bottle_backend", return_value=backend_mock), \
patch.object(start_mod, "read_tty_line", return_value="y"), \
patch.object(start_mod, "attach_agent", return_value=0), \
patch.object(start_mod, "capture_claude_session_state"), \
patch("sys.stderr", new_callable=io.StringIO):
rc = self._run_launch(assume_yes=False)
self.assertEqual(0, rc)
backend_mock.prelaunch_checks.assert_called_once()
backend_mock.launch.assert_called_once()
def test_interactive_yes_uppercase_also_accepted(self) -> None:
"""'Y' or 'YES' should also be accepted as confirmation."""
import bot_bottle.cli.start as start_mod
bottle_mock = MagicMock()
bottle_mock.name = "dev-abc"
backend_mock = MagicMock()
backend_mock.prelaunch_checks.side_effect = StaleImageError("image is 5 day(s) old")
backend_mock.launch.return_value = _ok_cm(bottle_mock)
with patch.object(start_mod, "get_bottle_backend", return_value=backend_mock), \
patch.object(start_mod, "read_tty_line", return_value="YES"), \
patch.object(start_mod, "attach_agent", return_value=0), \
patch.object(start_mod, "capture_claude_session_state"), \
patch("sys.stderr", new_callable=io.StringIO):
rc = self._run_launch(assume_yes=False)
self.assertEqual(0, rc)
backend_mock.launch.assert_called_once()
if __name__ == "__main__":
unittest.main()
-86
View File
@@ -1,86 +0,0 @@
"""Unit tests for the host-side configuration store."""
from __future__ import annotations
import sqlite3
import tempfile
import unittest
from pathlib import Path
from bot_bottle.config_store import (
DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS,
ConfigStore,
)
from bot_bottle.store_manager import StoreManager
class TestConfigStore(unittest.TestCase):
def test_cached_image_warning_days_defaults_to_one(self) -> None:
with tempfile.TemporaryDirectory(prefix="config-store.") as tmp:
store = ConfigStore(Path(tmp) / "bot-bottle.db")
store.migrate()
self.assertEqual(
DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS,
store.cached_image_stale_warning_days(),
)
def test_cached_image_warning_days_reads_value(self) -> None:
with tempfile.TemporaryDirectory(prefix="config-store.") as tmp:
store = ConfigStore(Path(tmp) / "bot-bottle.db")
store.migrate()
store.set_cached_image_stale_warning_days(7)
self.assertEqual(7, store.cached_image_stale_warning_days())
def test_config_schema_uses_explicit_settings_columns(self) -> None:
with tempfile.TemporaryDirectory(prefix="config-store.") as tmp:
store = ConfigStore(Path(tmp) / "bot-bottle.db")
store.migrate()
with sqlite3.connect(store.db_path) as conn:
conn.row_factory = sqlite3.Row
columns = [
row["name"]
for row in conn.execute("PRAGMA table_info(bot_bottle_config)")
]
self.assertEqual([
"id",
"cached_image_stale_warning_days",
], columns)
def test_store_manager_includes_config_store(self) -> None:
with tempfile.TemporaryDirectory(prefix="config-store.") as tmp:
db = Path(tmp) / "bot-bottle.db"
manager = StoreManager(db)
self.assertFalse(manager.is_migrated())
manager.migrate()
self.assertTrue(manager.is_migrated())
def test_cached_image_warning_days_returns_default_when_db_missing(self) -> None:
# When the db file doesn't exist yet (parent exists, file doesn't),
# the store returns the default without touching the file.
with tempfile.TemporaryDirectory(prefix="config-store.") as tmp:
store = ConfigStore(Path(tmp) / "missing.db")
self.assertEqual(
DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS,
store.cached_image_stale_warning_days(),
)
def test_cached_image_warning_days_returns_default_on_null_value(self) -> None:
# If the row exists but the value is NULL (or not castable to int),
# the store falls back to the default.
with tempfile.TemporaryDirectory(prefix="config-store.") as tmp:
db_path = Path(tmp) / "bot-bottle.db"
store = ConfigStore(db_path)
store.migrate()
# Write a NULL value directly.
with sqlite3.connect(db_path) as conn:
conn.execute(
"UPDATE bot_bottle_config SET cached_image_stale_warning_days = NULL WHERE id = 1"
)
self.assertEqual(
DEFAULT_CACHED_IMAGE_STALE_WARNING_DAYS,
store.cached_image_stale_warning_days(),
)
if __name__ == "__main__":
unittest.main()
+3 -4
View File
@@ -15,7 +15,6 @@ from bot_bottle.git_gate import GitGatePlan
from bot_bottle.orchestrator.client import RegisteredBottle
_MOD = "bot_bottle.backend.docker.consolidated_launch"
_UTIL = "bot_bottle.backend.consolidated_util"
def _egress_plan() -> EgressPlan:
@@ -50,7 +49,7 @@ class TestLaunchConsolidated(unittest.TestCase):
patch(f"{_MOD}._container_ip", return_value="172.18.0.2"), \
patch(f"{_MOD}._network_container_ips", return_value=list(on_network)), \
patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_UTIL}.provision_git_gate", provision or Mock()):
patch(f"{_MOD}.provision_git_gate", provision or Mock()):
return launch_consolidated(_egress_plan(), _git_plan(), service=service)
def test_allocates_ip_registers_and_provisions(self) -> None:
@@ -85,8 +84,8 @@ class TestLaunchConsolidated(unittest.TestCase):
class TestTeardownConsolidated(unittest.TestCase):
def test_deregisters_and_deprovisions(self) -> None:
client = Mock()
with patch(f"{_UTIL}.OrchestratorClient", return_value=client), \
patch(f"{_UTIL}.deprovision_git_gate") as deprov:
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_MOD}.deprovision_git_gate") as deprov:
teardown_consolidated("b1", orchestrator_url="http://orch:8080")
client.teardown_bottle.assert_called_once_with("b1")
deprov.assert_called_once()
-186
View File
@@ -1,186 +0,0 @@
"""Unit: host Claude auth extraction."""
from __future__ import annotations
import json
import tempfile
import unittest
from datetime import datetime, timezone
from pathlib import Path
from unittest.mock import MagicMock, patch
from bot_bottle.contrib.claude.claude_auth import (
claude_auth_path,
claude_host_access_token,
)
from bot_bottle.log import Die
def _cred_json(access_token: str, **extra: object) -> str:
payload: dict[str, object] = {"claudeAiOauth": {"accessToken": access_token, **extra}}
return json.dumps(payload)
class TestClaudeHostAccessToken(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.TemporaryDirectory(prefix="bb-claude-auth.")
self.home = Path(self.tmp.name)
self.cred_dir = self.home / ".claude"
self.cred_dir.mkdir()
self.auth_path = self.cred_dir / ".credentials.json"
def tearDown(self):
self.tmp.cleanup()
def _write(self, payload: dict) -> None: # type: ignore[no-untyped-def]
self.auth_path.write_text(json.dumps(payload))
def test_auth_path_uses_home_env(self):
self.assertEqual(
self.auth_path,
claude_auth_path({"HOME": str(self.home)}),
)
# --- file-based (Linux) ---
def test_file_returns_access_token(self):
key = "sk-ant-oat01-real-key" # gitleaks:allow
self._write({"claudeAiOauth": {"accessToken": key}})
out = claude_host_access_token({"HOME": str(self.home)})
self.assertEqual(key, out)
def test_file_missing_claude_ai_oauth_dies(self):
self._write({"hasCompletedOnboarding": True})
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(self.home)})
def test_file_missing_access_token_dies(self):
self._write({"claudeAiOauth": {"expiresAt": 2000000000000}})
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(self.home)})
def test_file_empty_access_token_dies(self):
self._write({"claudeAiOauth": {"accessToken": ""}})
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(self.home)})
def test_file_expired_token_dies(self):
# expiresAt is milliseconds; 1_000_000 ms is year 1970
self._write({
"claudeAiOauth": {"accessToken": "sk-ant-oat01-x", "expiresAt": 1_000_000}, # gitleaks:allow
})
with self.assertRaises(Die):
claude_host_access_token(
{"HOME": str(self.home)},
now=datetime(2026, 1, 1, tzinfo=timezone.utc),
)
def test_file_future_expiry_is_accepted(self):
key = "sk-ant-oat01-y" # gitleaks:allow
# 2_000_000_000_000 ms ≈ year 2033
self._write({
"claudeAiOauth": {"accessToken": key, "expiresAt": 2_000_000_000_000},
})
out = claude_host_access_token(
{"HOME": str(self.home)},
now=datetime(2026, 1, 1, tzinfo=timezone.utc),
)
self.assertEqual(key, out)
def test_file_absent_expiry_is_accepted(self):
key = "sk-ant-oat01-z" # gitleaks:allow
self._write({"claudeAiOauth": {"accessToken": key}})
out = claude_host_access_token({"HOME": str(self.home)})
self.assertEqual(key, out)
def test_file_non_json_dies(self):
self.auth_path.write_text("not json {{{")
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(self.home)})
def test_file_json_array_root_dies(self):
self.auth_path.write_text("[]")
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(self.home)})
def test_file_extra_fields_are_ignored(self):
key = "sk-ant-oat01-real" # gitleaks:allow
self._write({
"claudeAiOauth": {
"accessToken": key,
"refreshToken": "sk-ant-ort01-secret", # gitleaks:allow
"scopes": ["user:inference"],
"expiresAt": 2_000_000_000_000,
},
})
out = claude_host_access_token({"HOME": str(self.home)})
self.assertEqual(key, out)
# --- macOS Keychain fallback ---
def _home_without_creds(self) -> Path:
"""A home dir that has .claude/ but no .credentials.json."""
empty = self.home / "no-creds"
(empty / ".claude").mkdir(parents=True)
return empty
def _mock_keychain(self, stdout: str, returncode: int = 0) -> MagicMock:
mock = MagicMock()
mock.returncode = returncode
mock.stdout = stdout
return mock
def test_keychain_used_when_file_absent(self):
key = "sk-ant-oat01-keychain" # gitleaks:allow
home = self._home_without_creds()
with patch(
"bot_bottle.contrib.claude.claude_auth.subprocess.run",
return_value=self._mock_keychain(_cred_json(key)),
), patch(
"bot_bottle.contrib.claude.claude_auth.sys.platform", "darwin",
):
out = claude_host_access_token({"HOME": str(home)})
self.assertEqual(key, out)
def test_keychain_failure_when_file_absent_dies(self):
home = self._home_without_creds()
with patch(
"bot_bottle.contrib.claude.claude_auth.subprocess.run",
return_value=self._mock_keychain("", returncode=44),
), patch(
"bot_bottle.contrib.claude.claude_auth.sys.platform", "darwin",
):
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(home)})
def test_no_file_no_keychain_on_linux_dies(self):
home = self._home_without_creds()
with patch("bot_bottle.contrib.claude.claude_auth.sys.platform", "linux"):
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(home)})
def test_keychain_non_json_dies(self):
home = self._home_without_creds()
with patch(
"bot_bottle.contrib.claude.claude_auth.subprocess.run",
return_value=self._mock_keychain("not-json"),
), patch(
"bot_bottle.contrib.claude.claude_auth.sys.platform", "darwin",
):
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(home)})
def test_keychain_security_not_found_dies(self):
home = self._home_without_creds()
with patch(
"bot_bottle.contrib.claude.claude_auth.subprocess.run",
side_effect=FileNotFoundError,
), patch(
"bot_bottle.contrib.claude.claude_auth.sys.platform", "darwin",
):
with self.assertRaises(Die):
claude_host_access_token({"HOME": str(home)})
if __name__ == "__main__":
unittest.main()
@@ -3,7 +3,6 @@
from __future__ import annotations
import contextlib
import dataclasses
import io
import tempfile
import unittest
@@ -19,7 +18,6 @@ from bot_bottle.backend.docker.bottle_plan import DockerBottlePlan
from bot_bottle.backend.docker.consolidated_launch import LaunchContext
from bot_bottle.egress import EgressPlan
from bot_bottle.git_gate import GitGatePlan
from bot_bottle.log import Die
from bot_bottle.manifest import ManifestIndex
from tests.unit import use_bottle_root
@@ -94,7 +92,6 @@ class TestLaunchCommittedImage(unittest.TestCase):
mock.patch.object(launch_mod.docker_mod, "image_exists", return_value=image_present), \
mock.patch.object(launch_mod.docker_mod, "build_image", side_effect=_build), \
mock.patch.object(launch_mod.docker_mod, "verify_agent_image"), \
mock.patch.object(launch_mod.docker_mod, "image_created_at"), \
mock.patch.object(launch_mod, "launch_consolidated", return_value=_CTX), \
mock.patch.object(launch_mod, "teardown_consolidated"), \
mock.patch.object(launch_mod, "DockerGateway", return_value=gw), \
@@ -115,8 +112,7 @@ class TestLaunchCommittedImage(unittest.TestCase):
with self._patched(
committed_tag=committed_tag, image_present=image_present, compose=compose,
) as built:
images = launch_mod.build_or_load_images(plan)
with launch_mod.launch(plan, images, provision=mock.Mock(return_value=None)):
with launch_mod.launch(plan, provision=mock.Mock(return_value=None)):
pass
return built
@@ -131,34 +127,14 @@ class TestLaunchCommittedImage(unittest.TestCase):
captured.append(p)
return {"services": {"agent": {}}}
with self._patched(committed_tag=_COMMITTED_TAG, image_present=True, compose=compose) as _:
plan = _plan(self._tmp)
images = launch_mod.build_or_load_images(plan)
with launch_mod.launch(plan, images, provision=mock.Mock(return_value=None)):
with self._patched(committed_tag=_COMMITTED_TAG, image_present=True, compose=compose):
with launch_mod.launch(_plan(self._tmp), provision=mock.Mock(return_value=None)):
pass
self.assertEqual(_COMMITTED_TAG, captured[0].image)
def test_falls_back_to_build_when_no_committed_image(self) -> None:
self.assertEqual([_DEFAULT_IMAGE], self._run_launch(_plan(self._tmp), committed_tag=None))
def test_cached_images_skip_build_when_present(self) -> None:
base = _plan(self._tmp)
plan = dataclasses.replace(
base,
spec=dataclasses.replace(base.spec, image_policy="cached"),
)
built = self._run_launch(plan, committed_tag=None, image_present=True)
self.assertEqual([], built)
def test_cached_images_die_when_agent_missing(self) -> None:
base = _plan(self._tmp)
plan = dataclasses.replace(
base,
spec=dataclasses.replace(base.spec, image_policy="cached"),
)
with self.assertRaises(Die):
self._run_launch(plan, committed_tag=None, image_present=False)
def test_falls_back_to_build_when_committed_image_missing_from_daemon(self) -> None:
built = self._run_launch(_plan(self._tmp), committed_tag=_COMMITTED_TAG, image_present=False)
self.assertEqual([_DEFAULT_IMAGE], built)
+2 -4
View File
@@ -16,7 +16,7 @@ from pathlib import Path
from unittest import mock
from bot_bottle.agent_provider import AgentProvisionPlan
from bot_bottle.backend import BottleImages, BottleSpec
from bot_bottle.backend import BottleSpec
from bot_bottle.backend.docker import launch as launch_mod
from bot_bottle.backend.docker.bottle_plan import DockerBottlePlan
from bot_bottle.backend.docker.consolidated_launch import LaunchContext
@@ -93,8 +93,6 @@ class TestTeardownWarning(unittest.TestCase):
orchestrator_url="http://orch:8099",
)
images = BottleImages(agent="bot-bottle-claude:latest", sidecar="bot-bottle-sidecars:latest")
with mock.patch.object(launch_mod.docker_mod, "build_image"), \
mock.patch.object(launch_mod.docker_mod, "verify_agent_image"), \
mock.patch.object(launch_mod, "launch_consolidated", return_value=ctx), \
@@ -115,7 +113,7 @@ class TestTeardownWarning(unittest.TestCase):
), \
contextlib.redirect_stderr(buf):
provision = mock.Mock(return_value=None)
with launch_mod.launch(plan, images, provision=provision):
with launch_mod.launch(plan, provision=provision):
pass
output = buf.getvalue()
@@ -45,15 +45,12 @@ _PROVIDER = _Provider()
def _plan(*, git_user: dict | None = None, # type: ignore
git_repos: dict | None = None, # type: ignore
copy_cwd: bool = False,
user_cwd: str = "/tmp/x",
stage_dir: Path | None = None) -> DockerBottlePlan:
bottle_json: dict = {} # type: ignore
if git_user is not None:
bottle_json["git-gate"] = {"user": git_user}
if git_repos is not None:
bottle_json.setdefault("git-gate", {})["repos"] = git_repos
index = ManifestIndex.from_json_obj({
"bottles": {"dev": bottle_json},
"agents": {"demo": {"skills": [], "prompt": "", "bottle": "dev"}},
@@ -128,62 +125,6 @@ class TestProvisionGitUser(unittest.TestCase):
_PROVIDER.provision_git(bottle, _plan(stage_dir=self.stage))
self.assertEqual([], _git_config_exec_calls(bottle))
def test_repairs_git_xdg_directory_for_runtime_user(self):
bottle = _make_bottle()
_PROVIDER.provision_git(bottle, _plan(stage_dir=self.stage))
script, user = next(
(call.args[0], call.kwargs.get("user", "node"))
for call in bottle.exec.call_args_list
if "/home/node/.config/git" in call.args[0]
)
self.assertEqual("root", user)
self.assertIn("chown node:node /home/node", script)
self.assertIn("chmod 755 /home/node", script)
self.assertIn("mkdir -p /home/node/.config/git", script)
self.assertIn("chown -R node:node /home/node/.config", script)
self.assertIn("chmod -R u+rwX,go+rX /home/node/.config", script)
def test_fails_closed_when_home_permissions_cannot_be_repaired(self):
bottle = _make_bottle()
bottle.exec.return_value = ExecResult(1, "", "read-only filesystem")
with self.assertRaises(SystemExit):
_PROVIDER.provision_git(bottle, _plan(stage_dir=self.stage))
def _git_plan(self) -> DockerBottlePlan:
return _plan(
git_repos={
"repo": {
"url": "ssh://git@example.com/repo.git",
"key": {"provider": "static", "path": "/dev/null"},
"host_key": "ssh-ed25519 AAAA",
},
},
stage_dir=self.stage,
)
def test_fails_closed_when_gitconfig_permissions_cannot_be_set(self):
bottle = _make_bottle()
bottle.exec.side_effect = [
ExecResult(0, "", ""),
ExecResult(1, "", "chown failed"),
]
with self.assertRaises(SystemExit):
_PROVIDER.provision_git(bottle, self._git_plan())
def test_fails_closed_when_runtime_user_cannot_read_gitconfig(self):
bottle = _make_bottle()
bottle.exec.side_effect = [
ExecResult(0, "", ""),
ExecResult(0, "", ""),
ExecResult(1, "", "permission denied"),
]
with self.assertRaises(SystemExit):
_PROVIDER.provision_git(bottle, self._git_plan())
def test_sets_name_and_email(self):
plan = _plan(
git_user={"name": "Eric Bauerfeld", "email": "eric@dideric.is"},
-54
View File
@@ -9,7 +9,6 @@ from __future__ import annotations
import subprocess
import unittest
from datetime import timezone
from unittest.mock import patch
from bot_bottle.backend.docker import util as docker_mod
@@ -27,59 +26,6 @@ def _fail(stderr: str = "boom") -> subprocess.CompletedProcess: # type: ignore
)
class TestImageCreatedAt(unittest.TestCase):
def test_parses_docker_timestamp_with_nanoseconds(self):
with patch.object(
docker_mod.subprocess, "run",
return_value=_ok(stdout="2026-07-06T15:33:47.123456789Z\n"),
) as run:
created = docker_mod.image_created_at("bot-bottle-claude:latest")
self.assertIsNotNone(created)
assert created is not None
self.assertEqual(2026, created.year)
self.assertEqual(123456, created.microsecond)
self.assertEqual(timezone.utc, created.tzinfo)
self.assertEqual(
["docker", "image", "inspect", "--format", "{{.Created}}", "bot-bottle-claude:latest"],
run.call_args.args[0],
)
def test_dies_on_inspect_failure(self):
with patch.object(
docker_mod.subprocess, "run", return_value=_fail("No such image"),
), patch.object(
docker_mod, "die", side_effect=SystemExit("die"),
) as die:
with self.assertRaises(SystemExit):
docker_mod.image_created_at("missing:tag")
die.assert_called_once()
self.assertIn("missing:tag", die.call_args.args[0])
def test_returns_none_on_invalid_timestamp(self):
with patch.object(
docker_mod.subprocess, "run",
return_value=_ok(stdout="not-a-timestamp\n"),
):
result = docker_mod.image_created_at("some:tag")
self.assertIsNone(result)
def test_returns_none_on_empty_stdout(self):
with patch.object(
docker_mod.subprocess, "run",
return_value=_ok(stdout=""),
):
result = docker_mod.image_created_at("some:tag")
self.assertIsNone(result)
def test_parse_docker_timestamp_no_tzinfo_defaults_to_utc(self):
# A bare datetime with no tz offset should be treated as UTC.
dt = docker_mod._parse_docker_timestamp("2024-05-01T10:00:00.000000")
self.assertIsNotNone(dt.tzinfo)
self.assertEqual(timezone.utc, dt.tzinfo)
class TestCommitContainer(unittest.TestCase):
def test_runs_docker_commit(self):
with patch.object(
@@ -413,28 +413,6 @@ class TestAuthInjection(unittest.TestCase):
assert flow.response is not None
self.assertEqual(403, flow.response.status_code)
def test_preserve_auth_passes_agent_token_through(self) -> None:
route = Route(host="registry-1.docker.io", preserve_auth=True)
addon = _addon(Config(routes=(route,)))
flow = _Flow(_Request(
host="registry-1.docker.io",
headers={"authorization": "Bearer agent-registry-token"},
))
_run_request(addon, flow)
self.assertEqual("Bearer agent-registry-token", flow.request.headers.get("authorization"))
self.assertIsNone(flow.response)
def test_default_route_strips_agent_auth(self) -> None:
route = Route(host="registry-1.docker.io")
addon = _addon(Config(routes=(route,)))
flow = _Flow(_Request(
host="registry-1.docker.io",
headers={"authorization": "Bearer agent-registry-token"},
))
_run_request(addon, flow)
self.assertIsNone(flow.request.headers.get("authorization"))
self.assertIsNone(flow.response)
# ---------------------------------------------------------------------------
# git push / fetch over HTTPS
+1 -60
View File
@@ -4,14 +4,7 @@ from __future__ import annotations
import unittest
from bot_bottle.egress_addon_core import (
DENY_RESOLVER_ERROR,
DENY_UNATTRIBUTED,
DENY_UNPARSEABLE,
decide,
resolve_client_config,
resolve_client_context,
)
from bot_bottle.egress_addon_core import resolve_client_config, resolve_client_context
from bot_bottle.policy_resolver import PolicyResolveError
@@ -115,55 +108,3 @@ class TestResolveClientContext(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestDenyReasonNamesTheRealFault(unittest.TestCase):
"""A deny-all must not masquerade as a missing allowlist entry.
Regression: an unregistered bottle resolves no policy, so *every* host is
denied but the block message said `host X is not in the allowlist`,
which reads as a config problem and sends the operator hunting for a route
that was never missing. The structural reason wins over that wording.
"""
def _reason(self, resolver: object, host: str = "chatgpt.com") -> str:
cfg = resolve_client_config(resolver, "10.243.0.1") # type: ignore[arg-type]
return decide(cfg.routes, host, "/v1/x", {}, deny_reason=cfg.deny_reason).reason
def test_unattributed_says_unattributed_not_allowlist(self) -> None:
reason = self._reason(_FakeResolver(result=None))
self.assertEqual(DENY_UNATTRIBUTED, reason)
# The misleading claim is the one that must be gone: the host was
# never "not in the allowlist" — there was no allowlist at all.
self.assertNotIn("is not in the bottle's egress.routes allowlist", reason)
# Both causes must be named. `/resolve` fail-closes on a missing row
# *and* on a token mismatch, and the message pointing only at the row
# sent us hunting for a deregistered bottle that was registered fine.
self.assertIn("registry row", reason)
self.assertIn("identity token", reason)
def test_resolver_error_says_orchestrator_unreachable(self) -> None:
self.assertEqual(DENY_RESOLVER_ERROR, self._reason(_FakeResolver(raises=True)))
def test_unparseable_policy_says_so(self) -> None:
self.assertEqual(
DENY_UNPARSEABLE, self._reason(_FakeResolver(result="routes: notalist\n")))
def test_a_real_allowlist_miss_keeps_the_allowlist_wording(self) -> None:
"""The message only changes for structural deny-alls — a loaded policy
that genuinely lacks the host still points at the allowlist."""
reason = self._reason(_FakeResolver(result='routes:\n - host: "api.example.com"\n'))
self.assertIn("is not in the bottle's egress.routes allowlist", reason)
self.assertIn("chatgpt.com", reason)
def test_allowed_host_is_still_forwarded(self) -> None:
cfg = resolve_client_config(
_FakeResolver(result='routes:\n - host: "api.example.com"\n'), "10.243.0.1")
decision = decide(
cfg.routes, "api.example.com", "/v1/x", {}, deny_reason=cfg.deny_reason)
self.assertEqual("forward", decision.action)
def test_a_parsed_policy_carries_no_deny_reason(self) -> None:
cfg = resolve_client_config(
_FakeResolver(result='routes:\n - host: "api.example.com"\n'), "10.243.0.1")
self.assertEqual("", cfg.deny_reason)
+1 -18
View File
@@ -61,12 +61,6 @@ class TestNetpoolSlots(unittest.TestCase):
class TestNetpoolRenderers(unittest.TestCase):
def test_guest_init_restores_node_home_boundary(self):
from bot_bottle.backend.firecracker import util
self.assertIn("chown node:node /home/node", util._GUEST_INIT)
self.assertIn("chmod 755 /home/node", util._GUEST_INIT)
def test_nixos_module_is_non_invasive(self):
# The NixOS module must NOT flip the host firewall backend or
# hand interfaces to systemd-networkd; it brings the pool up via
@@ -428,15 +422,10 @@ class TestBottlePlanProperties(unittest.TestCase):
ap.command = "claude"
ap.prompt_mode = "append_file"
ap.template = "claude"
ap.guest_home = "/home/node"
ap.guest_env = {}
egress_plan = cast(Any, MagicMock())
egress_plan.canary = ""
egress_plan.canary_env = ""
fields = dict(
spec=cast(Any, MagicMock()), manifest=cast(Any, MagicMock()),
stage_dir=Path("/stage"), git_gate_plan=cast(Any, MagicMock()),
egress_plan=egress_plan, supervise_plan=None,
egress_plan=cast(Any, MagicMock()), supervise_plan=None,
agent_provision=ap, slug="demo-x", forwarded_env={},
)
fields.update(overrides)
@@ -462,12 +451,6 @@ class TestBottlePlanProperties(unittest.TestCase):
self.assertEqual("10.243.0.0:9420", p.git_gate_insteadof_host)
self.assertEqual("http", p.git_gate_insteadof_scheme)
def test_guest_env_pins_global_git_config(self):
from bot_bottle.backend.firecracker.launch import _agent_guest_env
env = _agent_guest_env(self._plan(), "10.243.0.0")
self.assertEqual("/home/node/.gitconfig", env["GIT_CONFIG_GLOBAL"])
if __name__ == "__main__":
unittest.main()
@@ -35,15 +35,6 @@ class TestBuildAgentRootfsDir(unittest.TestCase):
build.assert_not_called()
self.assertEqual(base, out)
def test_cached_lookup_requires_ready_marker(self):
digest = image_builder._rootfs_digest(self.dockerfile)
base = self.cache / "rootfs" / f"agent-{digest}"
base.mkdir(parents=True)
with patch.object(image_builder.util, "cache_dir", return_value=self.cache):
self.assertIsNone(image_builder.cached_agent_rootfs_dir(self.dockerfile))
(base / ".bb-ready").write_text("ok\n")
self.assertEqual(base, image_builder.cached_agent_rootfs_dir(self.dockerfile))
def test_cache_miss_builds_injects_and_marks_ready(self):
with patch.object(image_builder.util, "cache_dir", return_value=self.cache), \
patch.object(image_builder, "_build_in_infra") as build, \
-10
View File
@@ -71,16 +71,6 @@ class TestSshGatewayTransport(unittest.TestCase):
t.exec(["mkdir", "-p", "/git-gate"])
class TestGatewayCaPem(unittest.TestCase):
def test_dies_when_cert_never_appears(self) -> None:
from subprocess import CompletedProcess
infra = infra_vm.InfraVm(vm=None, guest_ip="10.0.0.1", private_key=Path("/k"))
with patch.object(infra_vm.subprocess, "run",
return_value=CompletedProcess([], 1, stdout="", stderr="")), \
self.assertRaises(SystemExit):
infra.gateway_ca_pem(timeout=0)
class TestRegistryVolume(unittest.TestCase):
def test_reuses_existing_volume(self):
import tempfile
-81
View File
@@ -1,81 +0,0 @@
"""Unit: image_cache.py — check_stale / check_stale_path."""
from __future__ import annotations
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from pathlib import Path
from unittest.mock import patch
from bot_bottle.image_cache import StaleImageError, check_stale, check_stale_path
class TestCheckStale(unittest.TestCase):
def _run(self, threshold: int, age_days: float) -> None:
created = datetime.now(tz=timezone.utc) - timedelta(days=age_days)
with patch("bot_bottle.image_cache.ConfigStore") as cs:
cs.return_value.cached_image_stale_warning_days.return_value = threshold
check_stale("test image", created)
def test_negative_threshold_never_raises(self):
# Threshold < 0 means the check is disabled — always passes.
self._run(threshold=-1, age_days=9999)
def test_zero_threshold_raises_immediately(self):
# threshold=0 means any image is stale the moment it exists.
with self.assertRaises(StaleImageError):
self._run(threshold=0, age_days=0.1)
def test_within_threshold_does_not_raise(self):
# Age well under threshold — should pass silently.
self._run(threshold=7, age_days=2)
def test_at_threshold_does_not_raise(self):
# Exactly at the boundary is fine (<=, not <).
self._run(threshold=1, age_days=0.9999)
def test_exceeds_threshold_raises(self):
created = datetime.now(tz=timezone.utc) - timedelta(days=3)
with patch("bot_bottle.image_cache.ConfigStore") as cs:
cs.return_value.cached_image_stale_warning_days.return_value = 1
with self.assertRaises(StaleImageError) as ctx:
check_stale("agent image 'bot-bottle:latest'", created)
self.assertIn("agent image", str(ctx.exception))
self.assertIn("day(s) old", str(ctx.exception))
def test_naive_datetime_treated_as_utc(self):
# check_stale calls .astimezone(utc) on the input; naive datetimes
# that would be interpreted as local time should still work.
# We can't control the local tz in a unit test, so just ensure
# no exception is thrown for a very recent naive datetime.
naive_now = datetime(2099, 1, 1) # far future, always "fresh"
with patch("bot_bottle.image_cache.ConfigStore") as cs:
cs.return_value.cached_image_stale_warning_days.return_value = 1
# Should not raise — the image is brand new.
check_stale("test image", naive_now)
class TestCheckStalePath(unittest.TestCase):
def test_delegates_to_check_stale_with_mtime(self):
with tempfile.NamedTemporaryFile() as f:
path = Path(f.name)
with patch("bot_bottle.image_cache.check_stale") as mock_check:
check_stale_path("some artifact", path)
mock_check.assert_called_once()
label, dt = mock_check.call_args.args
self.assertEqual("some artifact", label)
self.assertIsInstance(dt, datetime)
self.assertIsNotNone(dt.tzinfo)
def test_raises_stale_for_old_file(self):
with tempfile.NamedTemporaryFile() as f:
path = Path(f.name)
with patch("bot_bottle.image_cache.ConfigStore") as cs:
cs.return_value.cached_image_stale_warning_days.return_value = 0
with self.assertRaises(StaleImageError):
check_stale_path("cached artifact", path)
if __name__ == "__main__":
unittest.main()
+1 -1
View File
@@ -96,7 +96,7 @@ class TestVersionInputs(unittest.TestCase):
(pkg / "app.py").write_text("print('hi')\n")
(pkg / "egress_entrypoint.sh").write_text("#!/bin/sh\nexec mitmdump\n")
(pkg / "netpool.defaults.env").write_text("FOO=1\n")
for name in ("Dockerfile.orchestrator", "Dockerfile.gateway", "Dockerfile.infra", "Dockerfile.infra.fc"):
for name in ("Dockerfile.orchestrator", "Dockerfile.gateway", "Dockerfile.infra"):
(root / name).write_text(f"FROM scratch # {name}\n")
(root / "pyproject.toml").write_text("[project]\nname = 'bot-bottle'\n")
+3 -106
View File
@@ -17,7 +17,6 @@ from bot_bottle.git_gate import GitGatePlan
from bot_bottle.orchestrator.client import RegisteredBottle
_MOD = "bot_bottle.backend.macos_container.consolidated_launch"
_UTIL = "bot_bottle.backend.consolidated_util"
def _egress_plan() -> EgressPlan:
@@ -88,8 +87,7 @@ class TestRegisterAgent(unittest.TestCase):
*, source_ip: str = "192.168.128.9",
):
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_UTIL}.provision_git_gate", provision or Mock()), \
patch(f"{_MOD}.live_source_ips", return_value=[]):
patch(f"{_MOD}.provision_git_gate", provision or Mock()):
return register_agent(
_egress_plan(), _git_plan(),
source_ip=source_ip, endpoint=_endpoint(), image_ref="img:1",
@@ -127,8 +125,8 @@ class TestTeardown(unittest.TestCase):
def test_deregisters_and_deprovisions(self) -> None:
client = Mock()
deprovision = Mock()
with patch(f"{_UTIL}.OrchestratorClient", return_value=client), \
patch(f"{_UTIL}.deprovision_git_gate", deprovision):
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_MOD}.deprovision_git_gate", deprovision):
teardown_consolidated("b1", orchestrator_url="http://o:8099")
client.teardown_bottle.assert_called_once_with("b1")
self.assertEqual("b1", deprovision.call_args.args[1])
@@ -136,104 +134,3 @@ class TestTeardown(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestLiveSourceIps(unittest.TestCase):
"""The reconciliation input: the host enumerates its own bottles because
the orchestrator, inside the infra container, cannot see the backend."""
def _agents(self, *slugs: str) -> list[Mock]:
return [Mock(slug=s) for s in slugs]
def test_maps_slugs_to_container_addresses(self) -> None:
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \
patch(f"{_MOD}.container_mod.inspect_container_network_ip",
side_effect=["10.0.0.1", "10.0.0.2"]) as ip:
got = live_source_ips("net0")
self.assertEqual(["10.0.0.1", "10.0.0.2"], got)
self.assertEqual("bot-bottle-a", ip.call_args_list[0].args[0])
def test_containers_without_an_address_are_skipped(self) -> None:
"""A container that hasn't been given a DHCP address yet contributes
nothing the reap's grace window, not this list, protects it."""
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \
patch(f"{_MOD}.container_mod.inspect_container_network_ip",
side_effect=["", "10.0.0.2"]):
self.assertEqual(["10.0.0.2"], live_source_ips("net0"))
def test_container_list_failure_raises(self) -> None:
"""If container list fails, the live set is not authoritative and
reconciliation must be skipped."""
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with patch(f"{_MOD}.enumerate_active",
side_effect=EnumerationError("container list failed")):
with self.assertRaises(EnumerationError):
live_source_ips("net0")
def test_per_container_inspect_failure_raises(self) -> None:
"""If any individual inspect fails, the live set is not authoritative."""
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \
patch(f"{_MOD}.container_mod.inspect_container_network_ip",
side_effect=["10.0.0.1", None]):
with self.assertRaises(EnumerationError):
live_source_ips("net0")
class TestRegisterAgentReconciles(unittest.TestCase):
"""Registration self-heals the registry first: an orphan row at a recycled
address makes attribution ambiguous, which resolves no policy at all and
denies every host for the bottle being launched."""
def _register(self, client: Mock) -> None:
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_UTIL}.provision_git_gate"), \
patch(f"{_MOD}.live_source_ips", return_value=["10.0.0.7"]):
register_agent(
_egress_plan(), _git_plan(),
source_ip="10.0.0.7", endpoint=_endpoint(),
)
def test_reconciles_before_registering(self) -> None:
client = _client()
calls: list[str] = []
def _reconcile(*_args: object, **_kwargs: object) -> list[str]:
calls.append("reconcile")
return []
def _register_bottle(*_args: object, **_kwargs: object) -> RegisteredBottle:
calls.append("register")
return RegisteredBottle("b1", "tok")
client.reconcile.side_effect = _reconcile
client.register_bottle.side_effect = _register_bottle
self._register(client)
self.assertEqual(["reconcile", "register"], calls)
client.reconcile.assert_called_once_with(["10.0.0.7"])
def test_a_reconcile_failure_does_not_block_the_launch(self) -> None:
from bot_bottle.orchestrator.client import OrchestratorClientError
client = _client()
client.reconcile.side_effect = OrchestratorClientError("unreachable")
self._register(client)
client.register_bottle.assert_called_once()
def test_enumeration_error_does_not_block_the_launch(self) -> None:
"""A partial container listing must not abort the launch — skip
reconciliation and proceed, just as with an unreachable orchestrator."""
from bot_bottle.backend.macos_container.enumerate import EnumerationError
client = _client()
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_UTIL}.provision_git_gate"), \
patch(f"{_MOD}.live_source_ips",
side_effect=EnumerationError("container list failed")):
register_agent(
_egress_plan(), _git_plan(),
source_ip="10.0.0.7", endpoint=_endpoint(),
)
client.register_bottle.assert_called_once()
+2 -4
View File
@@ -66,10 +66,8 @@ class TestMacosContainerEnumerate(unittest.TestCase):
agents = self._enumerate("bot-bottle-mac-infra\nbot-bottle-dev-abc\n")
self.assertEqual(["dev-abc"], [a.slug for a in agents])
def test_raises_when_the_cli_fails(self):
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with self.assertRaises(EnumerationError):
self._enumerate("", returncode=1)
def test_empty_when_the_cli_fails(self):
self.assertEqual([], self._enumerate("", returncode=1))
if __name__ == "__main__":
@@ -12,7 +12,7 @@ import unittest
from pathlib import Path
from types import SimpleNamespace
from typing import cast
from unittest.mock import ANY, patch
from unittest.mock import patch
from bot_bottle.backend.macos_container.bottle import MacosContainerBottle
from bot_bottle.backend.macos_container.bottle_plan import MacosContainerBottlePlan
@@ -21,9 +21,7 @@ from bot_bottle.backend.macos_container.gateway_hosts import GATEWAY_HOSTNAME
from bot_bottle.backend.macos_container.launch import (
_agent_run_argv,
_identity_proxy_env,
build_or_load_images,
)
from bot_bottle.log import Die
from bot_bottle.manifest import ManifestIndex
_BOTTLE = "bot_bottle.backend.macos_container.bottle"
@@ -48,7 +46,6 @@ def _plan(
*,
agent_git_gate_url: str = "",
agent_supervise_url: str = "",
image_policy: str = "fresh",
) -> MacosContainerBottlePlan:
routes_path = stage_dir / "routes.yaml"
routes_path.write_text("routes: []\n", encoding="utf-8")
@@ -63,13 +60,12 @@ def _plan(
canary_env="",
)
return cast(MacosContainerBottlePlan, SimpleNamespace(
spec=SimpleNamespace(image_policy=image_policy),
spec=SimpleNamespace(),
manifest=_MANIFEST,
stage_dir=stage_dir,
slug="dev-abc",
container_name="bot-bottle-dev-abc",
image="bot-bottle-agent:latest",
dockerfile_path="/repo/Dockerfile",
forwarded_env={"OAUTH_TOKEN": "host-value"},
egress_plan=egress_plan,
git_gate_plan=SimpleNamespace(upstreams=()),
@@ -83,88 +79,6 @@ def _plan(
))
class TestBuildOrLoadImages(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.plan = _plan(Path(self._tmp.name))
def tearDown(self) -> None:
self._tmp.cleanup()
def test_reuses_present_committed_image(self) -> None:
with (
patch(
"bot_bottle.backend.macos_container.launch.read_committed_image",
return_value="committed:latest",
),
patch(
"bot_bottle.backend.macos_container.launch.container_mod.image_exists",
return_value=True,
),
patch(
"bot_bottle.backend.macos_container.launch.container_mod.build_image"
) as build,
):
images = build_or_load_images(self.plan)
self.assertEqual("committed:latest", images.agent)
build.assert_not_called()
def test_reuses_present_cached_image(self) -> None:
plan = _plan(Path(self._tmp.name), image_policy="cached")
with (
patch(
"bot_bottle.backend.macos_container.launch.read_committed_image",
return_value=None,
),
patch(
"bot_bottle.backend.macos_container.launch.container_mod.image_exists",
return_value=True,
),
patch(
"bot_bottle.backend.macos_container.launch.container_mod.build_image"
) as build,
):
images = build_or_load_images(plan)
self.assertEqual(plan.image, images.agent)
build.assert_not_called()
def test_cached_policy_rejects_missing_image(self) -> None:
plan = _plan(Path(self._tmp.name), image_policy="cached")
with (
patch(
"bot_bottle.backend.macos_container.launch.read_committed_image",
return_value=None,
),
patch(
"bot_bottle.backend.macos_container.launch.container_mod.image_exists",
return_value=False,
),
):
with self.assertRaises(Die):
build_or_load_images(plan)
def test_fresh_policy_builds_image(self) -> None:
with (
patch(
"bot_bottle.backend.macos_container.launch.read_committed_image",
return_value=None,
),
patch(
"bot_bottle.backend.macos_container.launch.container_mod.build_image"
) as build,
):
images = build_or_load_images(self.plan)
self.assertEqual(self.plan.image, images.agent)
build.assert_called_once_with(
self.plan.image,
ANY,
dockerfile=self.plan.dockerfile_path,
)
class TestAgentRunArgv(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
-111
View File
@@ -334,50 +334,6 @@ class TestInspectDigests(unittest.TestCase):
self.assertEqual({}, util.container_env("x"))
class TestInspectContainerNetworkIp(unittest.TestCase):
"""inspect_container_network_ip must distinguish inspect failure (None)
from 'no DHCP address yet' (""), which is the invariant live_source_ips
relies on to skip reconciliation on partial snapshots."""
_NETWORK = "bot-bottle-mac-gateway"
def _inspect(self, stdout: str, returncode: int = 0) -> str | None:
cp = util.subprocess.CompletedProcess(
args=[], returncode=returncode, stdout=stdout, stderr="",
)
with patch.object(util.subprocess, "run", return_value=cp):
return util.inspect_container_network_ip("bot-bottle-abc", self._NETWORK)
def _entry(self, ip: str = "192.168.128.5") -> str:
return (
f'[{{"status":{{"networks":['
f'{{"network":"{self._NETWORK}","ipv4Address":"{ip}"}}'
f']}}}}]'
)
def test_returns_ip_when_inspect_succeeds(self) -> None:
self.assertEqual("192.168.128.5", self._inspect(self._entry()))
def test_strips_cidr_prefix(self) -> None:
self.assertEqual("192.168.128.5", self._inspect(self._entry("192.168.128.5/24")))
def test_returns_empty_string_when_no_address_assigned_yet(self) -> None:
no_ip = f'[{{"status":{{"networks":[{{"network":"{self._NETWORK}","ipv4Address":""}}]}}}}]'
self.assertEqual("", self._inspect(no_ip))
def test_returns_empty_string_when_network_list_absent(self) -> None:
self.assertEqual("", self._inspect('[{"status":{}}]'))
def test_returns_none_on_nonzero_exit(self) -> None:
self.assertIsNone(self._inspect("", returncode=1))
def test_returns_none_on_malformed_json(self) -> None:
self.assertIsNone(self._inspect("not-json"))
def test_returns_none_on_unexpected_json_shape(self) -> None:
self.assertIsNone(self._inspect("null"))
class TestWaitContainerIpv4(unittest.TestCase):
def test_returns_address_once_dhcp_assigns_it(self):
with patch.object(util, "try_container_ipv4_on_network", side_effect=["", "", "192.168.128.4"]), \
@@ -391,72 +347,5 @@ class TestWaitContainerIpv4(unittest.TestCase):
self.assertEqual("", util.wait_container_ipv4_on_network("c", "net", timeout=-1))
class TestMacosContainerImageCreatedAt(unittest.TestCase):
def _ok(self, stdout: str) -> "util.subprocess.CompletedProcess": # type: ignore
return util.subprocess.CompletedProcess(
args=[], returncode=0, stdout=stdout, stderr="",
)
def _fail(self, stderr: str = "no such image") -> "util.subprocess.CompletedProcess": # type: ignore
return util.subprocess.CompletedProcess(
args=[], returncode=1, stdout="", stderr=stderr,
)
def test_parses_iso_timestamp_from_dict(self):
payload = '[{"created": "2025-06-01T12:00:00"}]'
with patch.object(util.subprocess, "run", return_value=self._ok(payload)):
dt = util.image_created_at("bot-bottle-agent:latest")
self.assertIsNotNone(dt)
assert dt is not None
self.assertEqual(2025, dt.year)
self.assertEqual(6, dt.month)
self.assertEqual(1, dt.day)
def test_accepts_list_or_dict_input(self):
# Container CLI may return a list; we take the first element.
payload = '[{"created": "2024-01-15T08:30:00"}]'
with patch.object(util.subprocess, "run", return_value=self._ok(payload)):
dt = util.image_created_at("some-image:latest")
self.assertIsNotNone(dt)
assert dt is not None
self.assertEqual(2024, dt.year)
def test_accepts_uppercase_Created_field(self):
payload = '[{"Created": "2024-03-20T10:00:00"}]'
with patch.object(util.subprocess, "run", return_value=self._ok(payload)):
dt = util.image_created_at("some-image:latest")
self.assertIsNotNone(dt)
assert dt is not None
self.assertEqual(2024, dt.year)
self.assertEqual(3, dt.month)
def test_dies_on_nonzero_returncode(self):
with patch.object(util.subprocess, "run", return_value=self._fail("not found")), \
patch.object(util, "die", side_effect=SystemExit("die")) as die:
with self.assertRaises(SystemExit):
util.image_created_at("missing:tag")
die.assert_called_once()
self.assertIn("missing:tag", die.call_args.args[0])
def test_dies_on_malformed_json(self):
with patch.object(util.subprocess, "run", return_value=self._ok("not-json {")), \
patch.object(util, "die", side_effect=SystemExit("die")) as die:
with self.assertRaises(SystemExit):
util.image_created_at("some:tag")
die.assert_called_once()
def test_returns_none_when_no_created_field(self):
payload = '[{"id": "sha256:abc123"}]'
with patch.object(util.subprocess, "run", return_value=self._ok(payload)):
result = util.image_created_at("some:tag")
self.assertIsNone(result)
def test_returns_none_on_invalid_timestamp_format(self):
payload = '[{"created": "not-a-date"}]'
with patch.object(util.subprocess, "run", return_value=self._ok(payload)):
result = util.image_created_at("some:tag")
self.assertIsNone(result)
if __name__ == "__main__":
unittest.main()
-22
View File
@@ -61,20 +61,6 @@ class TestInfraRun(unittest.TestCase):
mounts = [argv[i + 1] for i, a in enumerate(argv) if a == "--mount"]
self.assertTrue(all("bot-bottle.db" not in m for m in mounts))
def test_ca_is_persisted_on_the_host_not_the_container_volume(self) -> None:
"""The CA survives infra recreation and cannot be removed by Apple
Container's volume-prune command."""
argv = self._run_container(MacosInfraService(repo_root=Path("/r")))
mounts = [argv[i + 1] for i, a in enumerate(argv) if a == "--mount"]
ca_mounts = [
m for m in mounts
if "target=/home/mitmproxy/.mitmproxy" in m
]
self.assertEqual(1, len(ca_mounts))
self.assertIn("source=", ca_mounts[0])
self.assertIn("/gateway-ca", ca_mounts[0])
self.assertNotIn(",readonly", ca_mounts[0])
def test_nat_network_precedes_the_host_only_network(self) -> None:
argv = self._run_container(MacosInfraService(repo_root=Path("/r")))
nets = [argv[i + 1] for i, a in enumerate(argv) if a == "--network"]
@@ -167,14 +153,6 @@ class TestCaCertPem(unittest.TestCase):
argv = mod.run_container_argv.call_args.args[0]
self.assertEqual(["container", "exec", "bot-bottle-mac-infra", "cat"], argv[:4])
def test_raises_gateway_error_when_cert_never_appears(self) -> None:
from bot_bottle.backend.macos_container.gateway import GatewayError
svc = MacosInfraService(repo_root=Path("/r"))
with patch(f"{_INFRA}.container_mod") as mod:
mod.run_container_argv.return_value = _fail()
with self.assertRaises(GatewayError):
svc.ca_cert_pem(timeout=0)
class TestProbeControlPlane(unittest.TestCase):
def test_returns_url_when_running(self) -> None:
+1 -27
View File
@@ -80,19 +80,11 @@ class TestAgentProviderHostCredentials(unittest.TestCase):
"forward_host_credentials": "yes",
})
def test_forward_host_credentials_allowed_for_claude(self):
b = _provider_config_bottle({
"template": "claude",
"forward_host_credentials": True,
})
self.assertTrue(b.agent_provider.forward_host_credentials)
def test_forward_host_credentials_and_auth_token_rejected_together(self):
def test_forward_host_credentials_rejected_for_claude(self):
with self.assertRaises(ManifestError):
_provider_config_bottle({
"template": "claude",
"forward_host_credentials": True,
"auth_token": "SOME_TOKEN",
})
def test_auth_token_defaults_empty(self):
@@ -458,24 +450,6 @@ class TestRole(unittest.TestCase):
_bottle([{"host": "x.example", "role": ["x", 42]}])
class TestPreserveAuth(unittest.TestCase):
def test_omitted_defaults_false(self):
b = _bottle([{"host": "registry-1.docker.io"}])
self.assertFalse(b.egress.routes[0].PreserveAuth)
def test_true_accepted(self):
b = _bottle([{"host": "registry-1.docker.io", "preserve_auth": True}])
self.assertTrue(b.egress.routes[0].PreserveAuth)
def test_false_accepted(self):
b = _bottle([{"host": "registry-1.docker.io", "preserve_auth": False}])
self.assertFalse(b.egress.routes[0].PreserveAuth)
def test_non_bool_rejected(self):
with self.assertRaises(ManifestError):
_bottle([{"host": "registry-1.docker.io", "preserve_auth": "yes"}])
class TestPipelockKeyRejected(unittest.TestCase):
def test_pipelock_key_rejected_as_unknown(self):
with self.assertRaises(ManifestError):
+2 -14
View File
@@ -86,22 +86,10 @@ class TestAgentProviderValidation(unittest.TestCase):
"b", {"forward_host_credentials": True, "template": "weird"}
)
def test_forward_creds_pi_template_rejected(self) -> None:
def test_forward_creds_non_codex_template(self) -> None:
with self.assertRaises(ManifestError):
ManifestAgentProvider.from_dict(
"b", {"forward_host_credentials": True, "template": "pi"}
)
def test_forward_creds_claude_allowed(self) -> None:
p = ManifestAgentProvider.from_dict(
"b", {"forward_host_credentials": True, "template": "claude"}
)
self.assertTrue(p.forward_host_credentials)
def test_forward_creds_and_auth_token_rejected(self) -> None:
with self.assertRaises(ManifestError):
ManifestAgentProvider.from_dict(
"b", {"forward_host_credentials": True, "auth_token": "T", "template": "claude"}
"b", {"forward_host_credentials": True, "template": "claude"}
)
def test_valid_claude_auth_token(self) -> None:
-29
View File
@@ -104,32 +104,3 @@ class TestHealthAndPolicy(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestReconcile(unittest.TestCase):
def setUp(self) -> None:
self.c = OrchestratorClient("http://orch:8080")
def test_posts_live_ips_and_returns_reaped(self) -> None:
with patch(_URLOPEN, return_value=_resp(200, {"reaped": ["b1", "b2"]})) as m:
got = self.c.reconcile(["10.0.0.2", "10.0.0.3"])
self.assertEqual(["b1", "b2"], got)
sent = json.loads(m.call_args.args[0].data)
self.assertEqual(["10.0.0.2", "10.0.0.3"], sent["live_source_ips"])
self.assertNotIn("grace_seconds", sent) # omitted -> server default
def test_grace_seconds_is_forwarded_when_given(self) -> None:
with patch(_URLOPEN, return_value=_resp(200, {"reaped": []})) as m:
self.c.reconcile([], grace_seconds=30)
self.assertEqual(30, json.loads(m.call_args.args[0].data)["grace_seconds"])
def test_malformed_reaped_is_tolerated(self) -> None:
with patch(_URLOPEN, return_value=_resp(200, {"reaped": ["ok", 5, None]})):
self.assertEqual(["ok"], self.c.reconcile([]))
with patch(_URLOPEN, return_value=_resp(200, {})):
self.assertEqual([], self.c.reconcile([]))
def test_error_status_raises(self) -> None:
with patch(_URLOPEN, side_effect=_http_error(500)):
with self.assertRaises(OrchestratorClientError):
self.c.reconcile([])
@@ -1,147 +0,0 @@
"""Unit: OrchestratorConfigStore and resolve_teardown_timeout.
Also verifies the lifecycle ordering invariant: resolve_teardown_timeout()
must be called before launch_consolidated() / register_agent() so that a
resolver failure cannot leave an orphaned registration with no teardown
callback.
"""
from __future__ import annotations
import inspect
import os
import tempfile
import unittest
from pathlib import Path
from types import ModuleType
from bot_bottle.orchestrator.config_store import (
DEFAULT_TEARDOWN_TIMEOUT_SECONDS,
TEARDOWN_TIMEOUT_ENV,
OrchestratorConfigStore,
resolve_teardown_timeout,
)
class TestOrchestratorConfigStore(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db = Path(self._tmp.name) / "test.db"
self.store = OrchestratorConfigStore(self.db)
self.store.migrate()
def tearDown(self) -> None:
self._tmp.cleanup()
def test_get_returns_none_when_not_set(self) -> None:
self.assertIsNone(self.store.get_teardown_timeout_seconds())
def test_set_and_get_roundtrip(self) -> None:
self.store.set_teardown_timeout_seconds(42.5)
self.assertEqual(42.5, self.store.get_teardown_timeout_seconds())
def test_set_overwrites_existing_value(self) -> None:
self.store.set_teardown_timeout_seconds(10.0)
self.store.set_teardown_timeout_seconds(20.0)
self.assertEqual(20.0, self.store.get_teardown_timeout_seconds())
def test_delete_clears_value_and_returns_true(self) -> None:
self.store.set_teardown_timeout_seconds(30.0)
deleted = self.store.delete_teardown_timeout_seconds()
self.assertTrue(deleted)
self.assertIsNone(self.store.get_teardown_timeout_seconds())
def test_delete_absent_returns_false(self) -> None:
self.assertFalse(self.store.delete_teardown_timeout_seconds())
def test_is_migrated_true_after_migrate(self) -> None:
self.assertTrue(self.store.is_migrated())
def test_is_migrated_false_before_migrate(self) -> None:
store = OrchestratorConfigStore(Path(self._tmp.name) / "new.db")
self.assertFalse(store.is_migrated())
class TestResolveTeardownTimeout(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db = Path(self._tmp.name) / "cfg.db"
def tearDown(self) -> None:
self._tmp.cleanup()
os.environ.pop(TEARDOWN_TIMEOUT_ENV, None)
def test_returns_default_when_nothing_configured(self) -> None:
self.assertEqual(
DEFAULT_TEARDOWN_TIMEOUT_SECONDS,
resolve_teardown_timeout(self.db),
)
def test_env_var_overrides_default(self) -> None:
os.environ[TEARDOWN_TIMEOUT_ENV] = "99"
self.assertEqual(99.0, resolve_teardown_timeout(self.db))
def test_env_var_overrides_db_value(self) -> None:
store = OrchestratorConfigStore(self.db)
store.migrate()
store.set_teardown_timeout_seconds(55.0)
os.environ[TEARDOWN_TIMEOUT_ENV] = "77"
self.assertEqual(77.0, resolve_teardown_timeout(self.db))
def test_db_value_overrides_default(self) -> None:
store = OrchestratorConfigStore(self.db)
store.migrate()
store.set_teardown_timeout_seconds(42.0)
self.assertEqual(42.0, resolve_teardown_timeout(self.db))
def test_invalid_env_var_falls_through_to_default(self) -> None:
os.environ[TEARDOWN_TIMEOUT_ENV] = "not-a-number"
self.assertEqual(DEFAULT_TEARDOWN_TIMEOUT_SECONDS, resolve_teardown_timeout(self.db))
def test_non_positive_env_var_falls_through_to_default(self) -> None:
os.environ[TEARDOWN_TIMEOUT_ENV] = "0"
self.assertEqual(DEFAULT_TEARDOWN_TIMEOUT_SECONDS, resolve_teardown_timeout(self.db))
def test_non_positive_db_value_falls_through_to_default(self) -> None:
store = OrchestratorConfigStore(self.db)
store.migrate()
store.set_teardown_timeout_seconds(0.0)
self.assertEqual(DEFAULT_TEARDOWN_TIMEOUT_SECONDS, resolve_teardown_timeout(self.db))
def test_migrates_db_on_first_call(self) -> None:
result = resolve_teardown_timeout(self.db)
self.assertEqual(DEFAULT_TEARDOWN_TIMEOUT_SECONDS, result)
self.assertTrue(OrchestratorConfigStore(self.db).is_migrated())
class TestTeardownTimeoutResolvedBeforeRegistration(unittest.TestCase):
"""Ordering invariant: if resolve_teardown_timeout() raises, the bottle
must not yet be registered no orphaned state can result."""
def _src(self, module: ModuleType) -> str:
return inspect.getsource(module)
def test_docker_resolves_timeout_before_launch_consolidated(self) -> None:
from bot_bottle.backend.docker import launch
src = self._src(launch)
resolve_at = src.index("teardown_timeout = resolve_teardown_timeout()")
launch_at = src.index("ctx = launch_consolidated(")
self.assertLess(resolve_at, launch_at)
def test_firecracker_resolves_timeout_before_launch_consolidated(self) -> None:
from bot_bottle.backend.firecracker import launch
src = self._src(launch)
resolve_at = src.index("teardown_timeout = resolve_teardown_timeout()")
launch_at = src.index("ctx = launch_consolidated(")
self.assertLess(resolve_at, launch_at)
def test_macos_resolves_timeout_before_register_agent(self) -> None:
from bot_bottle.backend.macos_container import launch
src = self._src(launch)
resolve_at = src.index("teardown_timeout = resolve_teardown_timeout()")
register_at = src.index("ctx = register_agent(")
self.assertLess(resolve_at, register_at)
if __name__ == "__main__":
unittest.main()
@@ -8,13 +8,11 @@ from __future__ import annotations
import json
import secrets
import sqlite3
import tempfile
import threading
import unittest
import urllib.error
import urllib.request
from contextlib import closing
from pathlib import Path
from unittest.mock import patch
@@ -266,7 +264,6 @@ class TestControlPlaneAuth(unittest.TestCase):
("POST", "/bottles", _body({"source_ip": "10.0.0.1"})),
("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})),
("DELETE", "/bottles/x", b""),
("POST", "/reconcile", _body({"live_source_ips": []})),
("POST", "/resolve", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("POST", "/attribute", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("GET", "/supervise/proposals", b""),
@@ -390,56 +387,3 @@ class TestDispatchSupervise(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestReconcileRoute(unittest.TestCase):
"""`POST /reconcile` — the host tells the orchestrator which bottles are
actually up, since the orchestrator can't see the backend from inside the
infra container."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.orch = _orchestrator(Path(self._tmp.name) / "r.db")
def tearDown(self) -> None:
self._tmp.cleanup()
def _old(self, source_ip: str) -> str:
rec = self.orch.registry.register(source_ip)
with closing(sqlite3.connect(self.orch.registry.db_path)) as conn:
conn.execute(
"UPDATE orchestrator_bottles SET created_at = 0.0 WHERE bottle_id = ?",
(rec.bottle_id,))
conn.commit()
return rec.bottle_id
def test_reaps_absent_and_reports_ids(self) -> None:
dead = self._old("10.0.0.1")
alive = self._old("10.0.0.2")
status, payload = dispatch(
self.orch, "POST", "/reconcile", _body({"live_source_ips": ["10.0.0.2"]}))
self.assertEqual(200, status)
self.assertEqual([dead], payload["reaped"])
self.assertIsNone(self.orch.registry.get(dead))
self.assertIsNotNone(self.orch.registry.get(alive))
def test_missing_live_source_ips_is_400(self) -> None:
status, _ = dispatch(self.orch, "POST", "/reconcile", _body({}))
self.assertEqual(400, status)
def test_grace_seconds_is_honoured(self) -> None:
"""A grace window wide enough to cover the row protects it."""
self.orch.registry.register("10.0.0.3")
status, payload = dispatch(
self.orch, "POST", "/reconcile",
_body({"live_source_ips": [], "grace_seconds": 3600}))
self.assertEqual(200, status)
self.assertEqual([], payload["reaped"])
def test_non_string_entries_are_ignored(self) -> None:
dead = self._old("10.0.0.4")
status, payload = dispatch(
self.orch, "POST", "/reconcile",
_body({"live_source_ips": [None, 7, "10.0.0.9"]}))
self.assertEqual(200, status)
self.assertEqual([dead], payload["reaped"])
+2 -65
View File
@@ -2,9 +2,7 @@
from __future__ import annotations
import tempfile
import unittest
from pathlib import Path
from unittest.mock import Mock, patch
from bot_bottle.orchestrator.gateway import (
@@ -12,10 +10,7 @@ from bot_bottle.orchestrator.gateway import (
GATEWAY_NAME,
DockerGateway,
GatewayError,
rotate_gateway_ca,
)
from bot_bottle.paths import GATEWAY_CA_DIRNAME, host_gateway_ca_dir
from tests.unit import use_bottle_root
_CA_PEM = "-----BEGIN CERTIFICATE-----\nMII...\n-----END CERTIFICATE-----\n"
@@ -32,11 +27,6 @@ _ORCH_URL = "http://orchestrator:9000"
class TestDockerGateway(unittest.TestCase):
def setUp(self) -> None:
# Redirect the app-data root so ensure_running's host-dir mkdirs (CA +
# DB) land in a throwaway dir, not the real ~/.bot-bottle.
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.addCleanup(use_bottle_root(Path(self._tmp.name)))
# Resolver-only data plane (PRD 0070): running the gateway requires an
# orchestrator URL, so the fixture supplies one.
self.sc = DockerGateway("bot-bottle-gateway:latest", orchestrator_url=_ORCH_URL)
@@ -113,17 +103,8 @@ class TestDockerGateway(unittest.TestCase):
self.assertIn("bot-bottle-gateway:latest", runs[0])
# Runs on the shared gateway network so agents can reach it by IP.
self.assertEqual(self.sc.network, runs[0][runs[0].index("--network") + 1])
# Persists its CA on a HOST bind-mount (not a docker named volume, which
# `docker volume prune` would wipe — issue #450) so agents keep trusting
# it across restarts. The mount source is the host gateway-CA dir.
ca_mounts = [
a for a in runs[0]
if a.endswith(":/home/mitmproxy/.mitmproxy")
]
self.assertEqual(1, len(ca_mounts))
src = ca_mounts[0].rsplit(":", 1)[0]
self.assertTrue(src.endswith("/" + GATEWAY_CA_DIRNAME), src)
self.assertTrue(Path(src).is_absolute(), src)
# Persists its CA on a named volume so agents keep trusting it.
self.assertTrue(any("mitmproxy" in a for a in runs[0]))
# Shares the ONE host DB: the supervise daemon queues into the same
# file the orchestrator + operator (over HTTP) use.
self.assertTrue(any(
@@ -253,49 +234,5 @@ class TestDockerGatewayBuild(unittest.TestCase):
self.sc.ensure_built()
class TestRotateGatewayCa(unittest.TestCase):
"""rotate_gateway_ca clears the persisted CA so the next start remints it."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.addCleanup(use_bottle_root(Path(self._tmp.name)))
def _seed_ca(self) -> Path:
ca_dir = host_gateway_ca_dir()
# A representative mitmproxy confdir: the CA identity + derived encodings,
# plus one non-CA file that rotation must leave untouched.
for name in (
"mitmproxy-ca.pem",
"mitmproxy-ca-cert.pem",
"mitmproxy-ca-cert.cer",
"mitmproxy-ca-cert.p12",
):
(ca_dir / name).write_text("x")
(ca_dir / "combined-trust.pem").write_text("keep")
return ca_dir
def test_removes_ca_material_only(self) -> None:
ca_dir = self._seed_ca()
removed = rotate_gateway_ca(ca_dir)
self.assertEqual(4, len(removed))
self.assertTrue(all(p.name.startswith("mitmproxy-ca") for p in removed))
# The CA files are gone; the non-CA trust bundle survives.
self.assertEqual(
{"combined-trust.pem"}, {p.name for p in ca_dir.iterdir()}
)
def test_defaults_to_host_ca_dir(self) -> None:
self._seed_ca()
removed = rotate_gateway_ca() # no arg → host_gateway_ca_dir()
self.assertTrue(removed)
self.assertEqual(
[], list(host_gateway_ca_dir().glob("mitmproxy-ca*"))
)
def test_idempotent_when_no_ca(self) -> None:
self.assertEqual([], rotate_gateway_ca(host_gateway_ca_dir()))
if __name__ == "__main__":
unittest.main()
+55 -104
View File
@@ -1,4 +1,4 @@
"""Unit: infra container lifecycle — idempotent singleton (PRD 0070)."""
"""Unit: orchestrator+gateway container lifecycle — idempotent singleton (PRD 0070)."""
from __future__ import annotations
@@ -8,19 +8,19 @@ import urllib.error
from pathlib import Path
from unittest.mock import MagicMock, Mock, patch
from bot_bottle.orchestrator.gateway import GatewayError
from bot_bottle.orchestrator.lifecycle import (
INFRA_NAME,
INFRA_SOURCE_HASH_LABEL,
ORCHESTRATOR_IMAGE,
ORCHESTRATOR_NAME,
ORCHESTRATOR_SOURCE_HASH_LABEL,
OrchestratorService,
OrchestratorStartError,
source_hash,
)
from bot_bottle.paths import GATEWAY_CA_DIRNAME
from tests.unit import use_bottle_root
_URLOPEN = "bot_bottle.orchestrator.lifecycle.urllib.request.urlopen"
_RUN = "bot_bottle.orchestrator.lifecycle.run_docker"
_GATEWAY = "bot_bottle.orchestrator.lifecycle.DockerGateway"
_SLEEP = "bot_bottle.orchestrator.lifecycle.time.sleep"
_MONOTONIC = "bot_bottle.orchestrator.lifecycle.time.monotonic"
@@ -42,8 +42,10 @@ class TestOrchestratorService(unittest.TestCase):
self.addCleanup(use_bottle_root(Path(self._tmp.name)))
self.svc = OrchestratorService(port=8099)
def test_url(self) -> None:
def test_urls(self) -> None:
self.assertEqual("http://127.0.0.1:8099", self.svc.url)
# The gateway reaches the control plane by container name over docker DNS.
self.assertEqual(f"http://{ORCHESTRATOR_NAME}:8099", self.svc.internal_url)
def test_is_healthy(self) -> None:
with patch(_URLOPEN, return_value=_health(200)):
@@ -52,177 +54,126 @@ class TestOrchestratorService(unittest.TestCase):
self.assertFalse(self.svc.is_healthy())
def test_ensure_running_noop_when_healthy_and_source_unchanged(self) -> None:
# A healthy container on current source is left alone — recreating it
# on every launch drops in-memory egress tokens (#381).
# A healthy control plane already running the *current* bind-mounted
# source is left alone — recreating it on every launch would drop
# every other active bottle's in-memory egress tokens (#381).
current = source_hash(self.svc._repo_root)
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
if argv[:2] == ["docker", "ps"]:
return _proc(stdout=INFRA_NAME)
return _proc(stdout=ORCHESTRATOR_NAME)
if argv[:2] == ["docker", "inspect"]:
return _proc(stdout=current)
return _proc()
with patch(_URLOPEN, return_value=_health(200)), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
patch(_GATEWAY) as gw_cls, patch(_RUN, side_effect=fake), patch(_SLEEP):
self.assertEqual(self.svc.url, self.svc.ensure_running())
gw_cls.return_value.ensure_running.assert_called() # gateway kept up
runs = [c for c in calls if c[:2] == ["docker", "run"]]
rms = [c for c in calls if c[:3] == ["docker", "rm", "--force"] and INFRA_NAME in c]
self.assertEqual([], runs)
rms = [c for c in calls if c[:3] == ["docker", "rm", "--force"] and ORCHESTRATOR_NAME in c]
self.assertEqual([], runs) # not recreated
self.assertEqual([], rms)
def test_ensure_running_recreates_when_source_changed(self) -> None:
# Healthy, but the running container's label doesn't match the
# current source hash (a real code change) — recreate so it takes
# effect, same as the gateway's image-staleness check.
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
if argv[:2] == ["docker", "ps"]:
return _proc(stdout=INFRA_NAME)
return _proc(stdout=ORCHESTRATOR_NAME)
if argv[:2] == ["docker", "inspect"]:
return _proc(stdout="stale-hash")
return _proc()
with patch(_URLOPEN, return_value=_health(200)), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
patch(_GATEWAY), patch(_RUN, side_effect=fake), patch(_SLEEP):
self.assertEqual(self.svc.url, self.svc.ensure_running())
runs = [c for c in calls if c[:2] == ["docker", "run"]]
self.assertEqual(1, len(runs))
self.assertIn(INFRA_NAME, runs[0])
self.assertIn(ORCHESTRATOR_NAME, runs[0])
# the fresh container is labeled with the current hash, not the stale one
current = source_hash(self.svc._repo_root)
self.assertIn(f"{INFRA_SOURCE_HASH_LABEL}={current}", runs[0])
self.assertIn(f"{ORCHESTRATOR_SOURCE_HASH_LABEL}={current}", runs[0])
def test_ensure_running_starts_infra_container_when_absent(self) -> None:
def test_ensure_running_starts_orchestrator_container_when_absent(self) -> None:
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
if argv[:2] == ["docker", "ps"]:
return _proc(stdout="")
return _proc(stdout="") # not running
return _proc()
with patch(_URLOPEN, side_effect=[urllib.error.URLError("down"), _health(200)]), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
patch(_GATEWAY), patch(_RUN, side_effect=fake), patch(_SLEEP):
self.assertEqual(self.svc.url, self.svc.ensure_running())
runs = [c for c in calls if c[:2] == ["docker", "run"]]
self.assertEqual(1, len(runs))
argv = runs[0]
self.assertIn(INFRA_NAME, argv)
# Published on loopback — not exposed on external interfaces.
self.assertIn(ORCHESTRATOR_NAME, argv)
self.assertIn("--broker", argv)
self.assertEqual("stub", argv[argv.index("--broker") + 1]) # register-only, no socket
self.assertIn("bot_bottle.orchestrator", argv)
self.assertEqual("127.0.0.1:8099:8099", argv[argv.index("--publish") + 1])
# Both processes in one container — no separate entrypoint override.
self.assertNotIn("--entrypoint", argv)
# Gateway daemons + orchestrator explicitly opted in.
daemons_flag = "BOT_BOTTLE_GATEWAY_DAEMONS=egress,git-http,supervise,orchestrator"
self.assertIn("orchestrator", argv[argv.index(daemons_flag)])
# The mitmproxy CA persists on a HOST bind-mount under the app-data root
# (not a docker named volume `docker volume prune` would wipe — #450), so
# a restarted infra container keeps the CA every running bottle trusts.
ca_mounts = [a for a in argv if a.endswith(":/home/mitmproxy/.mitmproxy")]
self.assertEqual(1, len(ca_mounts))
src = ca_mounts[0].rsplit(":", 1)[0]
self.assertTrue(src.startswith(self._tmp.name), src)
self.assertTrue(src.endswith("/" + GATEWAY_CA_DIRNAME), src)
def test_ensure_running_builds_all_images(self) -> None:
def test_ensure_running_builds_lean_orchestrator_image_when_missing(self) -> None:
# The control plane runs its own lean image (#384), distinct from the
# gateway data plane — built from Dockerfile.orchestrator when absent.
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
if argv[:2] == ["docker", "ps"]:
return _proc(stdout="")
return _proc(stdout="") # orchestrator not running
if argv[:3] == ["docker", "image", "inspect"]:
return _proc(returncode=1) # image absent -> build
return _proc()
with patch(_URLOPEN, side_effect=[urllib.error.URLError("down"), _health(200)]), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
patch(_GATEWAY), patch(_RUN, side_effect=fake), patch(_SLEEP):
self.svc.ensure_running()
builds = [c for c in calls if c[:2] == ["docker", "build"]]
# Gateway base + orchestrator intermediate + infra image — all three built.
self.assertEqual(3, len(builds))
dockerfiles = [next(a for a in b if "Dockerfile" in a) for b in builds]
self.assertIn("Dockerfile.gateway", dockerfiles[0])
self.assertIn("Dockerfile.orchestrator", dockerfiles[1])
self.assertIn("Dockerfile.infra", dockerfiles[2])
# All three images are distinct.
tags = [b[b.index("-t") + 1] for b in builds]
self.assertEqual(3, len(set(tags)))
self.assertEqual(1, len(builds))
self.assertIn(ORCHESTRATOR_IMAGE, builds[0])
self.assertTrue(any(a.endswith("Dockerfile.orchestrator") for a in builds[0]))
# It is NOT the gateway image/dockerfile — the split is the point.
self.assertFalse(any("Dockerfile.gateway" in a for a in builds[0]))
def test_publish_maps_host_port_to_fixed_internal_port(self) -> None:
"""A non-default self.port is published to the fixed internal port 8099,
not to self.port:self.port the orchestrator always listens on 8099."""
def test_ensure_running_skips_orchestrator_image_build_when_present(self) -> None:
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
if argv[:2] == ["docker", "ps"]:
return _proc(stdout="")
if argv[:3] == ["docker", "image", "inspect"]:
return _proc(returncode=0) # image present -> no build
return _proc()
svc = OrchestratorService(port=20001)
with patch(_URLOPEN, side_effect=[urllib.error.URLError("down"), _health(200)]), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
svc.ensure_running()
runs = [c for c in calls if c[:2] == ["docker", "run"]]
argv = runs[0]
self.assertEqual("127.0.0.1:20001:8099", argv[argv.index("--publish") + 1])
orch_url = next(a for a in argv if "BOT_BOTTLE_ORCHESTRATOR_URL" in a)
self.assertIn(":8099", orch_url)
patch(_GATEWAY), patch(_RUN, side_effect=fake), patch(_SLEEP):
self.svc.ensure_running()
self.assertEqual([], [c for c in calls if c[:2] == ["docker", "build"]])
def test_ensure_running_raises_on_timeout(self) -> None:
with patch(_URLOPEN, side_effect=urllib.error.URLError("down")), \
patch(_RUN, return_value=Mock(returncode=0, stdout="", stderr="")), \
patch(_GATEWAY), patch(_RUN, return_value=Mock(returncode=0, stderr="")), \
patch(_SLEEP), patch(_MONOTONIC, side_effect=[0.0, 0.5, 2.0]):
with self.assertRaises(OrchestratorStartError):
self.svc.ensure_running(startup_timeout=1.0)
def test_noop_when_healthy_and_inspect_fails(self) -> None:
"""If docker inspect fails (e.g. docker daemon hiccup), leave the
working container alone rather than churning it."""
def fake(argv: list[str], **_kw: object) -> Mock:
if argv[:2] == ["docker", "ps"]:
return _proc(stdout=INFRA_NAME)
if argv[:2] == ["docker", "inspect"]:
return _proc(returncode=1, stderr="daemon error")
return _proc()
with patch(_URLOPEN, return_value=_health(200)), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
self.svc.ensure_running()
# no docker run — the working container was left alone
def test_build_failure_raises(self) -> None:
with patch(_URLOPEN, side_effect=urllib.error.URLError("down")), \
patch(_RUN, return_value=_proc(returncode=1, stderr="no space left on device")):
with self.assertRaises(GatewayError):
self.svc.ensure_running()
def test_ensure_network_creates_if_missing(self) -> None:
"""If the gateway network doesn't exist yet, create it."""
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
if argv[:3] == ["docker", "network", "inspect"]:
return _proc(returncode=1, stderr="not found")
if argv[:2] == ["docker", "ps"]:
return _proc(stdout="")
return _proc()
with patch(_URLOPEN, side_effect=[urllib.error.URLError("down"), _health(200)]), \
patch(_RUN, side_effect=fake), patch(_SLEEP):
self.svc.ensure_running()
creates = [c for c in calls if c[:3] == ["docker", "network", "create"]]
self.assertEqual(1, len(creates))
def test_stop_removes_infra_container(self) -> None:
with patch(_RUN) as run:
def test_stop_removes_orchestrator_and_gateway(self) -> None:
with patch(_RUN) as run, patch(_GATEWAY) as gw_cls:
self.svc.stop()
rms = [
c.args[0] for c in run.call_args_list
if c.args[0][:3] == ["docker", "rm", "--force"]
]
self.assertTrue(any(INFRA_NAME in a for a in rms))
rms = [c.args[0] for c in run.call_args_list if c.args[0][:3] == ["docker", "rm", "--force"]]
self.assertTrue(any(ORCHESTRATOR_NAME in a for a in rms))
gw_cls.return_value.stop.assert_called_once()
if __name__ == "__main__":
-80
View File
@@ -4,9 +4,7 @@ from __future__ import annotations
import sqlite3
import tempfile
import time
import unittest
from contextlib import closing
from pathlib import Path
from bot_bottle.orchestrator.registry import (
@@ -171,81 +169,3 @@ class TestRegistryStore(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestReapAbsent(unittest.TestCase):
"""`reap_absent` — the self-heal for rows whose bottle is gone.
An orphan is not merely untidy: source IPs get recycled, and
`by_source_ip` fail-closes on ambiguity, so a leftover row at a reused
address resolves *no* policy for the next bottle that lands there and
every host it asks for is denied.
"""
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 _aged(self, source_ip: str, *, age: float) -> BottleRecord:
"""Register a bottle and backdate it past the grace window."""
rec = self.store.register(source_ip)
with closing(sqlite3.connect(self.db)) as conn:
conn.execute(
"UPDATE orchestrator_bottles SET created_at = ? WHERE bottle_id = ?",
(time.time() - age, rec.bottle_id),
)
conn.commit()
return rec
def test_reaps_row_with_no_live_container(self) -> None:
gone = self._aged("10.243.0.9", age=600)
reaped = self.store.reap_absent([])
self.assertEqual([gone.bottle_id], [r.bottle_id for r in reaped])
self.assertIsNone(self.store.get(gone.bottle_id))
def test_keeps_row_whose_ip_is_live(self) -> None:
alive = self._aged("10.243.0.9", age=600)
self.assertEqual([], self.store.reap_absent(["10.243.0.9"]))
self.assertIsNotNone(self.store.get(alive.bottle_id))
def test_grace_window_protects_an_in_flight_launch(self) -> None:
"""A bottle registered moments ago is never reaped, even though the
caller's enumeration didn't see its address yet."""
fresh = self.store.register("10.243.0.10")
self.assertEqual([], self.store.reap_absent([]))
self.assertIsNotNone(self.store.get(fresh.bottle_id))
def test_reaping_the_orphan_unbricks_the_reused_address(self) -> None:
"""The regression this exists for: an orphan at an address that vmnet
later hands to a new bottle makes `by_source_ip` ambiguous, so the new
bottle resolves no policy at all."""
orphan = self._aged("10.243.0.11", age=600)
# A new bottle lands on the recycled address. Force the row in directly
# so `register`'s own supersede sweep doesn't mask the ambiguity.
with closing(sqlite3.connect(self.db)) as conn:
conn.execute(
"INSERT INTO orchestrator_bottles "
"(bottle_id, source_ip, identity_token, state, created_at, metadata, policy) "
"VALUES ('newbottle', '10.243.0.11', 'tok-new', 'active', ?, '', 'routes: []')",
(time.time(),),
)
conn.commit()
self.assertIsNone(self.store.by_source_ip("10.243.0.11")) # bricked
reaped = self.store.reap_absent(["10.243.0.11"], grace_seconds=60)
self.assertEqual([orphan.bottle_id], [r.bottle_id for r in reaped])
rec = self.store.by_source_ip("10.243.0.11")
assert rec is not None
self.assertEqual("newbottle", rec.bottle_id)
def test_ignores_empty_ips_in_the_live_set(self) -> None:
gone = self._aged("10.243.0.12", age=600)
self.assertEqual(
[gone.bottle_id],
[r.bottle_id for r in self.store.reap_absent(["", "10.243.0.99"])],
)
-59
View File
@@ -1,59 +0,0 @@
"""Unit: the `rotate_ca` one-shot CLI (issue #450). Docker mocked."""
from __future__ import annotations
import tempfile
import unittest
from pathlib import Path
from unittest.mock import Mock, patch
from bot_bottle.orchestrator import rotate_ca
from bot_bottle.orchestrator.gateway import GATEWAY_NAME
from bot_bottle.orchestrator.lifecycle import INFRA_NAME
from bot_bottle.paths import host_gateway_ca_dir
from tests.unit import use_bottle_root
_RUN = "bot_bottle.orchestrator.rotate_ca.run_docker"
def _proc(returncode: int = 0, stdout: str = "", stderr: str = "") -> Mock:
return Mock(returncode=returncode, stdout=stdout, stderr=stderr)
class TestRotateCaCli(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.addCleanup(use_bottle_root(Path(self._tmp.name)))
def test_clears_ca_and_drops_gateway_containers(self) -> None:
ca_dir = host_gateway_ca_dir()
(ca_dir / "mitmproxy-ca.pem").write_text("x")
(ca_dir / "mitmproxy-ca-cert.pem").write_text("x")
calls: list[list[str]] = []
def fake(argv: list[str], **_kw: object) -> Mock:
calls.append(argv)
# Report a removed container name so the CLI logs it.
return _proc(stdout=argv[-1])
with patch(_RUN, side_effect=fake):
self.assertEqual(0, rotate_ca.main([]))
# Persisted CA is gone → next start remints it.
self.assertEqual([], list(ca_dir.glob("mitmproxy-ca*")))
# Both the infra container and the standalone gateway are force-removed
# so no mitmproxy keeps serving the old CA from memory.
removed = {c[-1] for c in calls if c[:3] == ["docker", "rm", "--force"]}
self.assertEqual({INFRA_NAME, GATEWAY_NAME}, removed)
def test_succeeds_with_no_persisted_ca(self) -> None:
with patch(_RUN, return_value=_proc()) as m:
self.assertEqual(0, rotate_ca.main([]))
# Still tears down any running gateway even when there was no CA on disk.
self.assertTrue(m.called)
if __name__ == "__main__":
unittest.main()
-54
View File
@@ -4,10 +4,8 @@ from __future__ import annotations
import json
import secrets
import sqlite3
import tempfile
import unittest
from contextlib import closing
from pathlib import Path
from unittest.mock import patch
@@ -306,55 +304,3 @@ class TestOrchestratorSupervise(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestOrchestratorReconcile(unittest.TestCase):
"""`reconcile` — drop rows for bottles that are no longer running."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.secret = secrets.token_bytes(16)
self.db = Path(self._tmp.name) / "r.db"
self.store = RegistryStore(self.db)
self.store.migrate()
self.broker = StubBroker(self.secret)
self.orch = Orchestrator(self.store, self.broker, self.secret)
def tearDown(self) -> None:
self._tmp.cleanup()
def _age_all(self, seconds: float) -> None:
"""Backdate every row past the reap grace window."""
with closing(sqlite3.connect(self.db)) as conn:
conn.execute(
"UPDATE orchestrator_bottles SET created_at = created_at - ?", (seconds,))
conn.commit()
def test_reaps_dead_bottle_and_forgets_its_tokens(self) -> None:
dead = self.orch.launch_bottle("10.243.0.1", tokens={"EGRESS_TOKEN_0": "s3cret"})
live = self.orch.launch_bottle("10.243.0.2", tokens={"EGRESS_TOKEN_0": "keep"})
self._age_all(600)
self.assertEqual([dead.bottle_id], self.orch.reconcile(["10.243.0.2"]))
self.assertIsNone(self.store.get(dead.bottle_id))
self.assertIsNotNone(self.store.get(live.bottle_id))
# The in-memory egress credential goes with the row.
self.assertEqual({}, self.orch.tokens_for(dead.bottle_id))
self.assertEqual({"EGRESS_TOKEN_0": "keep"}, self.orch.tokens_for(live.bottle_id))
def test_reconcile_does_not_broker_a_teardown(self) -> None:
"""The container is already gone — there is nothing to stop, and a
broker error must not stop the sweep clearing the row."""
self.orch.launch_bottle("10.243.0.1")
self._age_all(600)
self.broker.launched.clear()
self.orch.reconcile([])
self.assertEqual([], self.broker.torn_down)
def test_reconcile_keeps_everything_when_all_are_live(self) -> None:
a = self.orch.launch_bottle("10.243.0.1")
b = self.orch.launch_bottle("10.243.0.2")
self._age_all(600)
self.assertEqual([], self.orch.reconcile(["10.243.0.1", "10.243.0.2"]))
self.assertIsNotNone(self.store.get(a.bottle_id))
self.assertIsNotNone(self.store.get(b.bottle_id))

Some files were not shown because too many files have changed in this diff Show More