Compare commits

...

23 Commits

Author SHA1 Message Date
didericis-claude 892299aac5 docs: defer broker replay protection to its own issue (#494)
prd-number-check / require-numbered-prds (pull_request) Failing after 6s
tracker-policy-pr / check-pr (pull_request) Successful in 13s
Per PR review, replay protection is too heavy for the host control
server MVP. Drop it from the four-gap framing (now three gaps), remove
the enforcement design section and implementation chunk, and track the
iat-window + jti-cache work in #494 as an independent in-process change.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-26 13:32:41 -04:00
didericis-claude 4860eead0c docs: PRD for host-side control server (#468)
Promote the in-process launch broker into a standalone host control
server: the single privileged host component that brokers launches, owns
orchestrator lifecycle, and is the sole writer of host-durable state.

Closes the four broker gaps (transport, durable provisioned secret,
replay protection, disciplined op vocabulary) and splits host state by
owner and lifetime (orchestrator SQLite / host JSONL audit / gateway
none). The payoff is dropping the Docker socket from the CLI.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-26 13:32:41 -04:00
didericis-claude 0edd46d56d chore(ci): only run PRD number check for PRs into main
prd-number-check / require-numbered-prds (pull_request) Successful in 5s
tracker-policy-pr / check-pr (pull_request) Successful in 12s
The require-numbered-prds gate previously ran on every pull request.
Scope its trigger to PRs whose base branch is main, so numbering is
only enforced at the point of merging into main.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-26 17:19:36 +00:00
Quality Badge Bot 4682dd441f chore: update quality badges
- Coverage: 84%
- Core coverage: 94%

[skip ci]
2026-07-26 17:06:23 +00:00
didericis-codex 1827593b89 test(docker): authenticate multitenant egress probes
prd-number-check / require-numbered-prds (pull_request) Successful in 9s
test / unit (pull_request) Successful in 53s
test / integration-docker (pull_request) Successful in 1m7s
test / coverage (pull_request) Successful in 17s
tracker-policy-pr / check-pr (pull_request) Successful in 32s
test / unit (push) Successful in 57s
lint / lint (push) Successful in 1m4s
test / integration-docker (push) Successful in 1m11s
test / coverage (push) Successful in 21s
Update Quality Badges / update-badges (push) Successful in 56s
2026-07-26 08:41:00 +00:00
didericis-codex 73e70e326c test(docker): wait for multitenant probe readiness
prd-number-check / require-numbered-prds (pull_request) Successful in 6s
tracker-policy-pr / check-pr (pull_request) Successful in 26s
test / unit (pull_request) Successful in 49s
lint / lint (push) Successful in 59s
test / integration-docker (pull_request) Failing after 1m1s
test / coverage (pull_request) Has been skipped
2026-07-26 08:38:16 +00:00
didericis-codex d744bec7b1 fix(docker): configure the attributed gateway subnet
tracker-policy-pr / check-pr (pull_request) Successful in 12s
prd-number-check / require-numbered-prds (pull_request) Successful in 42s
lint / lint (push) Successful in 1m4s
test / unit (pull_request) Successful in 57s
test / integration-docker (pull_request) Failing after 1m5s
test / coverage (pull_request) Has been skipped
2026-07-26 08:35:37 +00:00
didericis-codex f3fbfb3cc3 fix(ci): join control planes to the runner network
prd-number-check / require-numbered-prds (pull_request) Successful in 13s
tracker-policy-pr / check-pr (pull_request) Successful in 8s
test / integration-docker (pull_request) Failing after 35s
test / unit (pull_request) Successful in 52s
test / coverage (pull_request) Has been skipped
lint / lint (push) Successful in 1m7s
2026-07-26 08:30:57 +00:00
didericis-codex b2245ae1f3 test(orchestrator): make wrong-key coverage deterministic 2026-07-26 08:30:57 +00:00
didericis-codex f0fe33b1d0 fix(ci): bypass proxies for Docker host probes
prd-number-check / require-numbered-prds (pull_request) Successful in 6s
tracker-policy-pr / check-pr (pull_request) Successful in 7s
test / unit (pull_request) Failing after 45s
test / integration-docker (pull_request) Failing after 7m1s
test / coverage (pull_request) Has been skipped
2026-07-26 08:24:21 +00:00
didericis-codex d4889663d1 fix(ci): enable the scheduled canary suite
prd-number-check / require-numbered-prds (pull_request) Successful in 9s
tracker-policy-pr / check-pr (pull_request) Successful in 9s
test / unit (pull_request) Successful in 50s
test / integration-docker (pull_request) Failing after 7m45s
test / coverage (pull_request) Has been skipped
2026-07-26 08:15:21 +00:00
didericis-codex 1022247ce5 test(ci): cover assurance gate entry points
prd-number-check / require-numbered-prds (pull_request) Successful in 7s
tracker-policy-pr / check-pr (pull_request) Successful in 8s
lint / lint (push) Successful in 58s
test / unit (pull_request) Successful in 48s
test / integration-docker (pull_request) Failing after 19m28s
test / coverage (pull_request) Has been skipped
2026-07-26 08:14:00 +00:00
didericis-codex 671d91070e docs(ci): document the assurance topology 2026-07-26 08:02:12 +00:00
didericis-codex 1c6d30ffd8 test(canary): verify the pinned gitleaks release 2026-07-26 08:01:07 +00:00
didericis-codex 583ff98b27 ci(docker): require the full integration suite 2026-07-26 08:00:15 +00:00
didericis-codex fafb828bb7 ci(coverage): enforce the critical core contract 2026-07-26 07:50:44 +00:00
Quality Badge Bot f9ad6c85aa chore: update quality badges
- Coverage: 83%
- Core coverage: 93%

[skip ci]
2026-07-26 07:38:37 +00:00
didericis-codex 7488110e71 chore: tighten upkeep boundaries and static checks
prd-number-check / require-numbered-prds (pull_request) Successful in 8s
test / integration-docker (pull_request) Successful in 12s
test / unit (pull_request) Successful in 49s
test / coverage (pull_request) Successful in 15s
tracker-policy-pr / check-pr (pull_request) Successful in 5s
test / integration-docker (push) Successful in 13s
test / unit (push) Successful in 48s
lint / lint (push) Successful in 56s
Update Quality Badges / update-badges (push) Successful in 50s
test / coverage (push) Successful in 13s
2026-07-26 06:55:49 +00:00
didericis-codex 22dde95561 fix(diagnostics): make optional failures observable and safe 2026-07-26 06:55:49 +00:00
didericis-codex e29b79d517 refactor(backend): extract shared bottle preparation planner 2026-07-26 06:55:49 +00:00
didericis-codex 1d85acfd99 refactor(egress): make addon core a compatibility facade 2026-07-26 06:55:49 +00:00
didericis-codex 90defdc9cd refactor(egress): separate matching, DLP, and context concerns 2026-07-26 06:55:49 +00:00
didericis-codex 2039ef635f refactor: validate reconciliation inputs and neutralize CLI helpers 2026-07-26 06:55:49 +00:00
59 changed files with 2505 additions and 1336 deletions
+6 -3
View File
@@ -2,7 +2,7 @@
# digest, etc.) without coupling every dev push to upstream registry
# availability.
#
# Opt-in via CLAUDE_BOTTLE_RUN_CANARIES=1 so the same files can be run
# Opt-in via BOT_BOTTLE_RUN_CANARIES=1 so the same files can be run
# locally with the same gating.
name: canaries
@@ -17,7 +17,7 @@ jobs:
canaries:
runs-on: ubuntu-latest
env:
CLAUDE_BOTTLE_RUN_CANARIES: "1"
BOT_BOTTLE_RUN_CANARIES: "1"
steps:
- name: Checkout
uses: actions/checkout@v4
@@ -25,4 +25,7 @@ jobs:
# No actions/setup-python: canaries are stdlib unittest on the image's
# system Python 3.12 (older act_runner mishandles setup-python's PATH).
- name: Run canaries
run: python3 -m unittest discover -t . -s tests/canaries -v
run: |
python3 -m scripts.unittest_gate \
-t . -s tests/canaries -v \
--minimum-executed 1 --fail-on-skip
+1
View File
@@ -3,6 +3,7 @@ name: prd-number-check
on:
pull_request:
types: [opened, reopened, synchronize]
branches: [main]
jobs:
require-numbered-prds:
+27 -1
View File
@@ -60,6 +60,9 @@ jobs:
integration-docker:
runs-on: ubuntu-latest
concurrency:
group: integration-docker-infra
cancel-in-progress: false
steps:
- name: Checkout
uses: actions/checkout@v4
@@ -84,7 +87,30 @@ jobs:
env:
BOT_BOTTLE_BACKEND: docker
COVERAGE_FILE: ${{ github.workspace }}/.coverage.docker
run: python3 -m coverage run -m unittest discover -t . -s tests/integration -v
run: |
set -euo pipefail
DOCKER_CLIENT_NETWORK=$(
docker inspect "$(hostname)" |
python3 -c 'import json,sys; n=json.load(sys.stdin)[0]["NetworkSettings"]["Networks"]; print(next(iter(n)))'
)
test -n "$DOCKER_CLIENT_NETWORK"
RUN_KEY="${GITHUB_RUN_ID:-${GITHUB_RUN_NUMBER:-0}}"
export NO_PROXY="*"
export no_proxy="*"
export BOT_BOTTLE_DOCKER_CLIENT_NETWORK="$DOCKER_CLIENT_NETWORK"
export BOT_BOTTLE_DOCKER_ROOT_MOUNT="bot-bottle-ci-root-$RUN_KEY"
export BOT_BOTTLE_DOCKER_CA_MOUNT="bot-bottle-ci-ca-$RUN_KEY"
python3 -m coverage run -m scripts.unittest_gate \
-t . -s tests/integration -v \
--minimum-executed 22 --fail-on-skip
- name: Clean Docker integration volumes
if: always()
run: |
RUN_KEY="${GITHUB_RUN_ID:-${GITHUB_RUN_NUMBER:-0}}"
docker volume rm --force \
"bot-bottle-ci-root-$RUN_KEY" \
"bot-bottle-ci-ca-$RUN_KEY" 2>/dev/null || true
# Non-dot name so upload-artifact's dotfile-skipping glob picks it up.
- name: Stage docker coverage for upload
+44 -5
View File
@@ -12,10 +12,14 @@ on:
- 'bot_bottle/**'
- 'tests/**/*.py'
- 'cli.py'
- 'install.sh'
- 'setup.py'
- 'MANIFEST.in'
- 'flake.nix'
- 'nix/firecracker-netpool.nix'
- 'scripts/coverage.sh'
- 'scripts/critical-modules.txt'
- 'scripts/diff_coverage.py'
- 'scripts/tracker_policy.py'
- 'scripts/**/*.py'
- 'scripts/firecracker-netpool.sh'
- 'Dockerfile*'
- 'pyproject.toml'
@@ -23,15 +27,20 @@ on:
- '.coveragerc'
- '.dockerignore'
- '.gitea/workflows/test.yml'
- '.gitea/workflows/pre-release-test.yml'
pull_request:
paths:
- 'bot_bottle/**'
- 'tests/**/*.py'
- 'cli.py'
- 'install.sh'
- 'setup.py'
- 'MANIFEST.in'
- 'flake.nix'
- 'nix/firecracker-netpool.nix'
- 'scripts/coverage.sh'
- 'scripts/critical-modules.txt'
- 'scripts/diff_coverage.py'
- 'scripts/tracker_policy.py'
- 'scripts/**/*.py'
- 'scripts/firecracker-netpool.sh'
- 'Dockerfile*'
- 'pyproject.toml'
@@ -39,6 +48,7 @@ on:
- '.coveragerc'
- '.dockerignore'
- '.gitea/workflows/test.yml'
- '.gitea/workflows/pre-release-test.yml'
jobs:
unit:
@@ -71,6 +81,9 @@ jobs:
integration-docker:
runs-on: ubuntu-latest
concurrency:
group: integration-docker-infra
cancel-in-progress: false
steps:
- name: Checkout
uses: actions/checkout@v4
@@ -87,7 +100,33 @@ jobs:
env:
BOT_BOTTLE_BACKEND: docker
COVERAGE_FILE: ${{ github.workspace }}/.coverage.docker
run: python3 -m coverage run -m unittest discover -t . -s tests/integration -v
run: |
set -euo pipefail
# act_runner executes this job in a container while sharing the host
# Docker socket. Attach control-plane siblings to the job's network,
# and use named volumes for state the host daemon must mount.
DOCKER_CLIENT_NETWORK=$(
docker inspect "$(hostname)" |
python3 -c 'import json,sys; n=json.load(sys.stdin)[0]["NetworkSettings"]["Networks"]; print(next(iter(n)))'
)
test -n "$DOCKER_CLIENT_NETWORK"
RUN_KEY="${GITHUB_RUN_ID:-${GITHUB_RUN_NUMBER:-0}}"
export NO_PROXY="*"
export no_proxy="*"
export BOT_BOTTLE_DOCKER_CLIENT_NETWORK="$DOCKER_CLIENT_NETWORK"
export BOT_BOTTLE_DOCKER_ROOT_MOUNT="bot-bottle-ci-root-$RUN_KEY"
export BOT_BOTTLE_DOCKER_CA_MOUNT="bot-bottle-ci-ca-$RUN_KEY"
python3 -m coverage run -m scripts.unittest_gate \
-t . -s tests/integration -v \
--minimum-executed 22 --fail-on-skip
- name: Clean Docker integration volumes
if: always()
run: |
RUN_KEY="${GITHUB_RUN_ID:-${GITHUB_RUN_NUMBER:-0}}"
docker volume rm --force \
"bot-bottle-ci-root-$RUN_KEY" \
"bot-bottle-ci-ca-$RUN_KEY" 2>/dev/null || true
- name: Stage docker coverage for upload
run: cp .coverage.docker coverage-docker.dat
+15 -6
View File
@@ -33,19 +33,28 @@ jobs:
- name: Run coverage and extract percentage
id: coverage
run: |
python3 -m coverage run -m unittest discover -t . -s tests/unit > /dev/null 2>&1 || true
PERCENT=$(python3 -m coverage report 2>/dev/null | grep '^TOTAL' | grep -oP '\d+(?=%)' | tail -1)
set -euo pipefail
# Never publish a badge from a failed or partial test run.
python3 -m coverage run -m unittest discover -t . -s tests/unit
REPORT=$(python3 -m coverage report)
printf '%s\n' "$REPORT"
PERCENT=$(printf '%s\n' "$REPORT" | awk '$1 == "TOTAL" {gsub("%", "", $NF); print $NF}')
test -n "$PERCENT"
echo "percent=$PERCENT" >> $GITHUB_OUTPUT
echo "Coverage: $PERCENT%"
- name: Extract core (critical-module) coverage percentage
id: core_coverage
run: |
set -euo pipefail
# Reuses the .coverage data from the previous step. The core list is
# the single source of truth in scripts/critical-modules.txt; every
# core module is unit-tested, so the unit-only run is accurate for it.
INCLUDE=$(grep -vE '^[[:space:]]*(#|$)' scripts/critical-modules.txt | paste -sd, -)
PERCENT=$(python3 -m coverage report --include="$INCLUDE" 2>/dev/null | grep '^TOTAL' | grep -oP '\d+(?=%)' | tail -1)
# validated single source of truth. Fail if a listed path disappeared
# or if the measured core falls below ADR 0004's 90% minimum.
INCLUDE=$(python3 scripts/critical_modules.py)
REPORT=$(python3 -m coverage report --include="$INCLUDE" --fail-under=90)
printf '%s\n' "$REPORT"
PERCENT=$(printf '%s\n' "$REPORT" | awk '$1 == "TOTAL" {gsub("%", "", $NF); print $NF}')
test -n "$PERCENT"
echo "percent=$PERCENT" >> $GITHUB_OUTPUT
echo "Core coverage: $PERCENT%"
+3 -3
View File
@@ -5,7 +5,7 @@
# bot-bottle
[![test](https://gitea.dideric.is/didericis/bot-bottle/actions/workflows/test.yml/badge.svg?branch=main)](https://gitea.dideric.is/didericis/bot-bottle/actions?workflow=test.yml)
[![coverage](https://img.shields.io/badge/coverage-83%25-brightgreen)](https://coverage.readthedocs.io/)
[![coverage](https://img.shields.io/badge/coverage-84%25-brightgreen)](https://coverage.readthedocs.io/)
[![core coverage](https://img.shields.io/badge/core%20coverage-94%25-brightgreen)](https://gitea.dideric.is/didericis/bot-bottle/src/branch/main/docs/decisions/0004-coverage-policy.md)
**Problem:** Developer wants to run a coding agent without supervision, but they don't want a prompt injected or misbehaving agent wrecking their environment or exfiltrating sensitive data.
@@ -75,7 +75,7 @@ On compatible macOS hosts, the default backend requires Apple's `container` CLI
Use `BOT_BOTTLE_BACKEND=docker ./cli.py start <agent>` on hosts where neither Apple Container nor KVM is available and Docker is the desired backend.
> **CI (macOS Apple Container):** the `integration-macos` job (`.gitea/workflows/test.yml`) runs the integration suite against `BOT_BOTTLE_BACKEND=macos-container` on a self-hosted macOS runner labelled `macos`, because Apple Container needs the host virtualization framework and cannot run in a Linux container (so it can't reuse the `kvm` runner). Provision an Apple Silicon host with the `container` CLI on `PATH` and `container system status` running, then register the runner in **host mode** (not docker mode) with the `macos` label — `brew install gitea-runner` (the `act_runner` rename). Give it a Python ≥ 3.11 with `coverage` importable on the launchd service's `PATH` (a launchd service doesn't inherit your shell profile, so pin `node` and the Python env explicitly). The job is **advisory** — `workflow_dispatch` (manual) only, never triggered by push or PR — since a single laptop that sleeps/roams must not block merges or churn on every push to main; its coverage doesn't feed the gate. The infra container is a singleton (`bot-bottle-mac-infra`), so keep runner concurrency at 1.
> **CI (macOS Apple Container):** the advisory `integration-macos` job in `.gitea/workflows/pre-release-test.yml` runs only on manual dispatch. It targets a self-hosted host-mode runner labelled `macos`; Apple Container cannot run inside the Linux pull-request runner. Provision an Apple Silicon host with the `container` CLI running and Python ≥ 3.11 plus `coverage` on the launchd service's explicit `PATH`. The infra container is a singleton (`bot-bottle-mac-infra`), so keep runner concurrency at 1. Its coverage is reported separately and never feeds the required pull-request gate.
### Containers inside a bottle
@@ -174,7 +174,7 @@ BOT_BOTTLE_BACKEND=firecracker ./cli.py start <agent>
> **NixOS:** enable `virtualisation.docker`, ensure the KVM module is loaded (`boot.kernelModules = [ "kvm-intel" ];` or `kvm-amd`), and add your user to the `kvm` and `docker` groups. For the network pool, consume the flake module — `imports = [ inputs.bot-bottle.nixosModules.firecracker-netpool ]; services.bot-bottle-firecracker = { enable = true; owner = "you"; };` — then `nixos-rebuild switch` (imperative nft/TAP rules don't survive a rebuild; channel users can `imports = [ <bot-bottle>/nix/firecracker-netpool.nix ]`). `firecracker` isn't in nixpkgs by default as a user binary — install the release binary (pin the version) and put it on `PATH`.
> **CI:** the coverage gate (`.gitea/workflows/test.yml` → `coverage` job) runs on a self-hosted runner labelled `kvm`, because the Firecracker backend's VM/SSH orchestration is exercised only by the integration suite, which needs `/dev/kvm` + the provisioned pool (a container runner would skip it and read as uncovered). Provision that runner exactly like a normal Firecracker host `firecracker` on `PATH`, `/dev/kvm`, the cached guest kernel + static dropbear, and the pool installed as the persistent systemd unit — then register it with the `kvm` label. A Docker-capable hosted job builds the candidate once; KVM tests boot those exact bytes, and a successful main run publishes them. The unit/lint jobs still run on `ubuntu-latest`.
> **CI:** Firecracker integration runs in the manually dispatched `.gitea/workflows/pre-release-test.yml` on a self-hosted runner labelled `kvm`; privileged KVM hosts never execute unreviewed PR code automatically. Provision it like a normal Firecracker host: `firecracker` on `PATH`, `/dev/kvm`, the cached guest kernel and static dropbear, and the persistent TAP/nft pool. The required pull-request workflow runs unit plus the complete Docker integration suite on `ubuntu-latest`; see `docs/ci.md`.
```sh
./cli.py start <agent> # builds the image on first run, drops you into claude
+11 -75
View File
@@ -23,14 +23,14 @@ from dataclasses import dataclass
from pathlib import Path
from typing import Generator, Generic, Sequence, TypeVar
from ..agent_provider import AgentProvisionPlan, get_provider, build_agent_provision_plan
from ..agent_provider import AgentProvisionPlan, get_provider
from ..egress import EgressPlan
from ..git_gate import GitGatePlan
from ..log import die, info
from ..util import expand_tilde
from ..manifest import Manifest, ManifestIndex
from ..supervisor.plan import SupervisePlan
from ..env import resolve_env, ResolvedEnv
from ..env import ResolvedEnv
from ..workspace import WorkspacePlan, workspace_plan
from .print_util import print_multi, visible_agent_env_names
from .util import host_skill_dir
@@ -296,82 +296,18 @@ class BottleBackend(ABC, Generic[PlanT, CleanupT]):
backend-specific resolution (names, scratch files, etc.). The
validation step is enforced here so a future backend cannot
accidentally skip it. No remote/runtime resources are created."""
from .resolve_common import (
merge_provision_env_vars,
mint_slug,
prepare_agent_state_dir,
prepare_egress,
prepare_git_gate,
prepare_supervise,
reject_nested_containers,
resolve_manifest_dockerfile,
write_launch_metadata,
)
manifest = self._validate(spec)
if not self.supports_nested_containers:
reject_nested_containers(self.name, manifest)
self._preflight()
from ..git_gate import GitGate
manifest = GitGate().preflight_host_keys(
manifest,
headless=spec.headless,
home_md=spec.manifest.home_md,
)
manifest_bottle = manifest.bottle
manifest_agent_provider = manifest_bottle.agent_provider
agent_provider = get_provider(manifest_agent_provider.template)
resolved_env = resolve_env(manifest)
workspace = workspace_plan(spec, guest_home=agent_provider.guest_home)
slug = mint_slug(spec)
write_launch_metadata(slug, spec, compose_project="", backend=self.name)
# Manifest may override the Dockerfile per-bottle; otherwise fall
# back to the provider plugin's bundled Dockerfile (next to its
# agent_provider.py module).
if manifest_agent_provider.dockerfile:
agent_dockerfile_path = resolve_manifest_dockerfile(
manifest_agent_provider.dockerfile, spec,
)
else:
agent_dockerfile_path = str(agent_provider.dockerfile)
agent_dir, prompt_file = prepare_agent_state_dir(slug, manifest)
agent_provision_plan = build_agent_provision_plan(
template=manifest_agent_provider.template,
dockerfile=agent_dockerfile_path,
state_dir=agent_dir,
instance_name=f"bot-bottle-{slug}",
prompt_file=prompt_file,
guest_env=self._build_guest_env(resolved_env),
forward_host_credentials=manifest_agent_provider.forward_host_credentials,
auth_token=manifest_agent_provider.auth_token,
host_env=dict(os.environ),
trusted_project_path=workspace.workdir,
label=spec.label,
color=spec.color,
provider_settings=manifest_agent_provider.settings,
)
agent_provision_plan = merge_provision_env_vars(agent_provision_plan)
egress_plan = prepare_egress(manifest_bottle, slug, agent_provision_plan)
supervise_plan = prepare_supervise(manifest_bottle, slug)
git_gate_plan = prepare_git_gate(manifest_bottle, slug)
from .preparation import BottlePreparationPlanner
prepared = BottlePreparationPlanner(self).prepare(spec)
return self._resolve_plan(
spec,
manifest=manifest,
slug=slug,
resolved_env=resolved_env,
agent_provision_plan=agent_provision_plan,
egress_plan=egress_plan,
supervise_plan=supervise_plan,
git_gate_plan=git_gate_plan,
manifest=prepared.manifest,
slug=prepared.slug,
resolved_env=prepared.resolved_env,
agent_provision_plan=prepared.agent_provision_plan,
egress_plan=prepared.egress_plan,
supervise_plan=prepared.supervise_plan,
git_gate_plan=prepared.git_gate_plan,
stage_dir=stage_dir,
)
+46 -7
View File
@@ -17,6 +17,10 @@ from ...gateway import (
DEFAULT_CA_TIMEOUT_SECONDS, CA_POLL_SECONDS, GATEWAY_CA_CERT, GatewayError
)
DEFAULT_GATEWAY_SUBNET = "10.242.255.0/24"
_GATEWAY_SUBNET_LABEL = "bot-bottle.gateway-subnet"
class DockerGateway(Gateway):
"""The consolidated gateway as a single, fixed-name Docker container.
@@ -35,6 +39,8 @@ class DockerGateway(Gateway):
build_context: Path | None = None,
dockerfile: str | None = GATEWAY_DOCKERFILE,
host_port_bindings: tuple[int, ...] = (),
ca_mount_source: str | Path | None = None,
subnet: str | None = None,
) -> None:
self.image_ref = image_ref
self.name = name
@@ -59,6 +65,15 @@ class DockerGateway(Gateway):
# backend's dev-harness gateway so VMs can reach it via their TAP link;
# Docker's DNAT + the nft `ct status dnat accept` rule handle the rest.
self._host_port_bindings = host_port_bindings
self._subnet = (
subnet
or os.environ.get("BOT_BOTTLE_DOCKER_GATEWAY_SUBNET", "").strip()
or DEFAULT_GATEWAY_SUBNET
)
configured_ca = os.environ.get("BOT_BOTTLE_DOCKER_CA_MOUNT", "").strip()
self._ca_mount_source = str(
ca_mount_source or configured_ca or host_gateway_ca_dir()
)
def image_exists(self) -> bool:
return run_docker(["docker", "image", "inspect", self.image_ref]).returncode == 0
@@ -109,10 +124,34 @@ class DockerGateway(Gateway):
def _ensure_network(self) -> None:
"""Create the shared gateway network if it doesn't exist. Idempotent —
a concurrent create loses harmlessly (the loser sees 'already exists').
Docker picks the subnet; the launcher reads it back to allocate IPs."""
if run_docker(["docker", "network", "inspect", self.network]).returncode == 0:
return
proc = run_docker(["docker", "network", "create", self.network])
The explicit subnet is required because bottle attribution pins source
IPs; Docker rejects static endpoint addresses on an auto-IPAM network."""
inspected = run_docker([
"docker", "network", "inspect",
"--format", f'{{{{index .Labels "{_GATEWAY_SUBNET_LABEL}"}}}}',
self.network,
])
if inspected.returncode == 0:
marker = inspected.stdout.strip()
if marker in {"", self._subnet}:
return
if inspected.returncode == 0:
# Migrate the stale auto-IPAM network created by older releases.
# Removing the fixed gateway is safe here: this launch recreates it.
run_docker(["docker", "rm", "--force", self.name])
removed = run_docker(["docker", "network", "rm", self.network])
if removed.returncode != 0:
raise GatewayError(
f"gateway network {self.network} needs explicit subnet "
f"{self._subnet} but could not be replaced: "
f"{removed.stderr.strip()}"
)
proc = run_docker([
"docker", "network", "create",
"--subnet", self._subnet,
"--label", f"{_GATEWAY_SUBNET_LABEL}={self._subnet}",
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()}"
@@ -143,9 +182,9 @@ class DockerGateway(Gateway):
# Recreate when the running container's image is stale (a rebuild),
# so source changes to the gateway's flat daemons take effect — not
# just when the container is absent.
self._ensure_network()
if self.is_running() and self._running_image_is_current():
return
self._ensure_network()
# Clear any stale (stopped OR outdated-image) container holding the
# fixed name, then start fresh. `rm --force` on an absent name is a
# tolerated no-op.
@@ -158,7 +197,7 @@ class DockerGateway(Gateway):
# 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}",
"--volume", f"{self._ca_mount_source}:{MITMPROXY_HOME}",
# No DB mount: the data plane (egress / supervise / git-gate) reaches
# the supervise queue over the control-plane RPC and never opens
# bot-bottle.db, so the gateway container gets no file handle on it
@@ -253,4 +292,4 @@ class DockerGateway(Gateway):
def provisioning_transport(self) -> GatewayTransport:
"""The exec/cp transport git-gate provisioning stages per-bottle repos +
deploy keys through (over the docker socket)."""
return DockerGatewayTransport(self.name)
return DockerGatewayTransport(self.name)
+11 -4
View File
@@ -33,7 +33,6 @@ from .orchestrator import (
ORCHESTRATOR_NAME,
ORCHESTRATOR_NETWORK,
)
from ...paths import bot_bottle_root
from ... import resources
from ...gateway import (
GATEWAY_IMAGE,
@@ -69,6 +68,8 @@ class DockerInfraService(InfraService):
gateway_image: str = GATEWAY_IMAGE,
repo_root: Path | None = None,
host_root: Path | None = None,
root_mount_source: str | Path | None = None,
gateway_ca_mount_source: str | Path | None = None,
orchestrator_name: str = ORCHESTRATOR_NAME,
orchestrator_label: str = ORCHESTRATOR_LABEL,
gateway_name: str = GATEWAY_NAME,
@@ -78,10 +79,14 @@ class DockerInfraService(InfraService):
self.control_network = control_network
self.orchestrator_image = orchestrator_image
self.gateway_image = gateway_image
# Build context / bind-mount source: the repo root in a checkout, a
# staged copy from the installed wheel otherwise (bot_bottle.resources).
# Build context: the repo root in a checkout, a staged copy from the
# installed wheel otherwise (bot_bottle.resources).
self._repo_root = repo_root if repo_root is not None else resources.build_root()
self._host_root = host_root or bot_bottle_root()
if host_root is not None and root_mount_source is not None:
raise ValueError("pass host_root or root_mount_source, not both")
self._host_root = host_root
self._root_mount_source = root_mount_source
self._gateway_ca_mount_source = gateway_ca_mount_source
self._orchestrator_name = orchestrator_name
self._orchestrator_label = orchestrator_label
self._gateway_name = gateway_name
@@ -98,6 +103,7 @@ class DockerInfraService(InfraService):
control_network=self.control_network,
repo_root=self._repo_root,
host_root=self._host_root,
root_mount_source=self._root_mount_source,
)
def gateway(self) -> DockerGateway:
@@ -112,6 +118,7 @@ class DockerInfraService(InfraService):
network=self.network,
control_network=self.control_network,
build_context=self._repo_root,
ca_mount_source=self._gateway_ca_mount_source,
)
def ensure_running(
+55 -19
View File
@@ -42,13 +42,9 @@ ORCHESTRATOR_IMAGE = os.environ.get(
)
ORCHESTRATOR_DOCKERFILE = "Dockerfile.orchestrator"
# Baked as a container label so `ensure_running` can detect whether the running
# orchestrator is executing the current bind-mounted source.
# orchestrator image was built from the current source.
ORCHESTRATOR_SOURCE_HASH_LABEL = "bot-bottle-orchestrator-source-hash"
# The bind-mount path for the live control-plane source inside the container.
# PYTHONPATH points here so a code change takes effect on the next launch
# without an image rebuild.
_SRC_IN_CONTAINER = "/bot-bottle-src"
# Bot-bottle host-root bind-mount (DB + state) inside the orchestrator. The
# control plane opens bot-bottle.db under here (via BOT_BOTTLE_ROOT ->
# host_db_path()); it is the ONLY container with a handle on it (issue #469).
@@ -72,23 +68,51 @@ class DockerOrchestrator(Orchestrator):
control_network: str = ORCHESTRATOR_NETWORK,
repo_root: Path | None = None,
host_root: Path | None = None,
root_mount_source: str | Path | None = None,
client_host: str | None = None,
client_network: str | None = None,
bind_host: str | None = None,
dockerfile: str | None = ORCHESTRATOR_DOCKERFILE,
) -> None:
if host_root is not None and root_mount_source is not None:
raise ValueError("pass host_root or root_mount_source, not both")
self.image_ref = image_ref
self.name = name
self.label = label
self.port = port
self.control_network = control_network
# Build context / bind-mount source: the repo root in a checkout, a
# staged copy from the installed wheel otherwise (bot_bottle.resources).
# Build context: the repo root in a checkout, a staged copy from the
# installed wheel otherwise (bot_bottle.resources).
self._repo_root = repo_root if repo_root is not None else resources.build_root()
self._host_root = host_root or bot_bottle_root()
configured_root = os.environ.get("BOT_BOTTLE_DOCKER_ROOT_MOUNT", "").strip()
self._root_mount_source = str(
root_mount_source or configured_root or host_root or bot_bottle_root()
)
configured_network = os.environ.get(
"BOT_BOTTLE_DOCKER_CLIENT_NETWORK", ""
).strip()
self._client_network = client_network or configured_network or None
configured_host = os.environ.get(
"BOT_BOTTLE_DOCKER_HOST_ADDRESS", ""
).strip()
self._client_host = (
client_host or configured_host
or (self.name if self._client_network else "127.0.0.1")
)
# A socket-shared CI runner reaches published ports through its Docker
# network rather than its own loopback. Production stays bound to host
# loopback unless a caller explicitly selects another client.
self._bind_host = bind_host or (
"0.0.0.0"
if not self._client_network and self._client_host != "127.0.0.1"
else "127.0.0.1"
)
self._dockerfile = dockerfile
def url(self) -> str:
"""Host-side control-plane URL — the orchestrator's published loopback,
which the CLI reaches."""
return f"http://127.0.0.1:{self.port}"
"""Control-plane URL reachable by this Docker client."""
port = DEFAULT_PORT if self._client_network else self.port
return f"http://{self._client_host}:{port}"
def gateway_url(self) -> str:
"""The URL the gateway's data plane resolves policy against — the
@@ -119,8 +143,7 @@ class DockerOrchestrator(Orchestrator):
return self.name in proc.stdout.split()
def _source_current(self, current_hash: str) -> bool:
"""True iff the running orchestrator was started from the current
bind-mounted source."""
"""True iff the running orchestrator image matches current source."""
if not self.is_running():
return False
proc = run_docker([
@@ -182,15 +205,19 @@ class DockerOrchestrator(Orchestrator):
# Control network only — agents are never on it, so they have no
# route to the control plane (the L3 block, not just the JWT).
"--network", self.control_network,
# Host CLI reaches the control plane here (loopback only). The
# Host CLI reaches the control plane here (loopback by default).
# Socket-shared CI joins the container directly to the job network;
# the host-side mapping remains loopback-only in that topology. The
# orchestrator listens on the fixed DEFAULT_PORT inside the
# container; self.port is the host-side published port.
"--publish", f"127.0.0.1:{self.port}:{DEFAULT_PORT}",
# Live control-plane source (code changes without an image rebuild).
"--volume", f"{self._repo_root}:{_SRC_IN_CONTAINER}:ro",
"--env", f"PYTHONPATH={_SRC_IN_CONTAINER}",
"--publish", f"{self._bind_host}:{self.port}:{DEFAULT_PORT}",
# The image was rebuilt from `_repo_root` immediately before this
# launch. Running its baked package avoids a host-path bind mount,
# which is both more production-like and works with socket-shared
# CI where the daemon cannot see the job container's workspace.
# Orchestrator registry DB on the host (sole writer: control plane).
"--volume", f"{self._host_root}:{_ROOT_IN_CONTAINER}",
# `root_mount_source` may be a host path or a named Docker volume.
"--volume", f"{self._root_mount_source}:{_ROOT_IN_CONTAINER}",
"--env", f"BOT_BOTTLE_ROOT={_ROOT_IN_CONTAINER}",
# The signing key — held ONLY by the orchestrator (it verifies
# tokens); the gateway gets the pre-minted `gateway` JWT, never the
@@ -205,6 +232,15 @@ class DockerOrchestrator(Orchestrator):
raise OrchestratorStartError(
f"orchestrator container failed to start: {proc.stderr.strip()}"
)
if self._client_network:
proc = run_docker([
"docker", "network", "connect", self._client_network, self.name,
])
if proc.returncode != 0:
raise OrchestratorStartError(
f"orchestrator container failed to join client network "
f"{self._client_network}: {proc.stderr.strip()}"
)
def stop(self) -> None:
"""Remove the control-plane container (idempotent)."""
+6 -14
View File
@@ -6,12 +6,17 @@ from __future__ import annotations
import os
from datetime import datetime, timezone
import re
import shutil
import subprocess
from typing import Iterator
from ...log import die, info
from ...util import slugify as _slugify
def slugify(name: str) -> str:
"""Compatibility wrapper; new generic callers import ``bot_bottle.util``."""
return _slugify(name)
def run_docker(
@@ -114,19 +119,6 @@ def docker_cp(src: str, dest: str) -> None:
f"{(result.stderr or '').strip() or '<no stderr>'}")
_SLUG_RE = re.compile(r"[^a-z0-9]+")
def slugify(name: str) -> str:
"""Lowercase, non-alnum runs → '-', trimmed. Dies on empty result."""
if not name:
die("slugify: missing name")
slug = _SLUG_RE.sub("-", name.lower()).strip("-")
if not slug:
die(f"name '{name}' produced an empty slug; use alphanumeric characters")
return slug
def build_image(ref: str, context: str, *, dockerfile: str = "") -> None:
"""Invokes `docker build` every call. Layer cache makes no-change
rebuilds cheap; running every time means Dockerfile edits land
+2 -1
View File
@@ -11,7 +11,8 @@ from pathlib import Path
from ..bottle_state import egress_state_dir
from ..egress import EGRESS_ROUTES_FILENAME
from ..gateway.egress.addon_core import LOG_OFF, load_config
from ..gateway.egress.schema import load_config
from ..gateway.egress.types import LOG_OFF
class EgressApplyError(RuntimeError):
+124
View File
@@ -0,0 +1,124 @@
"""Backend-neutral preparation planner.
This module owns the shared transformation from a CLI ``BottleSpec`` to the
typed inputs consumed by a concrete backend's ``_resolve_plan``. Backend
classes retain only their validation/preflight/env hooks and their
backend-specific final resolution.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import TYPE_CHECKING, Protocol
from ..agent_provider import AgentProvisionPlan, build_agent_provision_plan, get_provider
from ..egress import EgressPlan
from ..env import ResolvedEnv, resolve_env
from ..git_gate import GitGate, GitGatePlan
from ..manifest import Manifest
from ..supervisor.plan import SupervisePlan
from ..workspace import workspace_plan
from .resolve_common import (
merge_provision_env_vars,
mint_slug,
prepare_agent_state_dir,
prepare_egress,
prepare_git_gate,
prepare_supervise,
reject_nested_containers,
resolve_manifest_dockerfile,
write_launch_metadata,
)
if TYPE_CHECKING:
from .base import BottleSpec
class PreparationBackend(Protocol):
"""Backend hooks needed by the shared planner."""
name: str
supports_nested_containers: bool
def _validate(self, spec: BottleSpec) -> Manifest: ...
def _preflight(self) -> None: ...
def _build_guest_env(self, resolved_env: ResolvedEnv) -> dict[str, str]: ...
@dataclass(frozen=True)
class PreparedBottle:
"""Typed, backend-neutral result of shared launch preparation."""
manifest: Manifest
slug: str
resolved_env: ResolvedEnv
agent_provision_plan: AgentProvisionPlan
egress_plan: EgressPlan
git_gate_plan: GitGatePlan
supervise_plan: SupervisePlan | None
class BottlePreparationPlanner:
"""Run the common, side-effect-limited part of bottle preparation."""
def __init__(self, backend: PreparationBackend) -> None:
self._backend = backend
def prepare(self, spec: BottleSpec) -> PreparedBottle:
backend = self._backend
# These are deliberately protected backend hooks: only this shared
# planner orchestrates them, while concrete backends provide the
# implementation.
manifest = backend._validate(spec) # pylint: disable=protected-access
if not backend.supports_nested_containers:
reject_nested_containers(backend.name, manifest)
backend._preflight() # pylint: disable=protected-access
manifest = GitGate().preflight_host_keys(
manifest,
headless=spec.headless,
home_md=spec.manifest.home_md,
)
bottle = manifest.bottle
provider_config = bottle.agent_provider
provider = get_provider(provider_config.template)
resolved_env = resolve_env(manifest)
workspace = workspace_plan(spec, guest_home=provider.guest_home)
slug = mint_slug(spec)
write_launch_metadata(slug, spec, compose_project="", backend=backend.name)
dockerfile = (
resolve_manifest_dockerfile(provider_config.dockerfile, spec)
if provider_config.dockerfile
else str(provider.dockerfile)
)
agent_dir, prompt_file = prepare_agent_state_dir(slug, manifest)
provision = build_agent_provision_plan(
template=provider_config.template,
dockerfile=dockerfile,
state_dir=agent_dir,
instance_name=f"bot-bottle-{slug}",
prompt_file=prompt_file,
guest_env=backend._build_guest_env( # pylint: disable=protected-access
resolved_env
),
forward_host_credentials=provider_config.forward_host_credentials,
auth_token=provider_config.auth_token,
host_env=dict(os.environ),
trusted_project_path=workspace.workdir,
label=spec.label,
color=spec.color,
provider_settings=provider_config.settings,
)
provision = merge_provision_env_vars(provision)
return PreparedBottle(
manifest=manifest,
slug=slug,
resolved_env=resolved_env,
agent_provision_plan=provision,
egress_plan=prepare_egress(bottle, slug, provision),
git_gate_plan=prepare_git_gate(bottle, slug),
supervise_plan=prepare_supervise(bottle, slug),
)
+2 -2
View File
@@ -30,6 +30,7 @@ from ..log import die
from ..manifest import Manifest, ManifestBottle
from ..supervisor.plan import SupervisePlan
from ..orchestrator.supervisor import Supervisor
from ..util import slugify
from . import BottleSpec
@@ -44,8 +45,7 @@ def mint_slug(spec: BottleSpec) -> str:
if spec.identity:
return spec.identity
if spec.label:
from .docker import util as docker_mod
return docker_mod.slugify(spec.label)
return slugify(spec.label)
return bottle_identity(spec.agent_name)
+8 -9
View File
@@ -25,12 +25,11 @@ from typing import Callable
from ...agent_provider import get_provider, runtime_for
from ...backend import (
Bottle,
BottlePlan,
BottleSpec,
enumerate_active_agents,
get_bottle_backend,
)
from ...backend.docker import util as docker_mod
from ...backend.docker.bottle_plan import DockerBottlePlan
from ...bottle_state import (
cleanup_state,
is_preserved,
@@ -40,7 +39,7 @@ from ...image_cache import StaleImageError
from ...log import info, die
from ...manifest import Manifest, ManifestIndex
from ..constants import PROG
from ...util import read_tty_line
from ...util import read_tty_line, slugify
from .. import tui
@@ -257,10 +256,10 @@ def _uniquify_label_headless(label: str) -> str:
logging the chosen label. Orchestrators fire-and-forget many bottles,
so silently picking a free name beats erroring on every collision."""
active_slugs = {a.slug for a in enumerate_active_agents()}
if docker_mod.slugify(label) not in active_slugs:
if slugify(label) not in active_slugs:
return label
n = 2
while docker_mod.slugify(f"{label}-{n}") in active_slugs:
while slugify(f"{label}-{n}") in active_slugs:
n += 1
chosen = f"{label}-{n}"
info(f"label '{label}' already in use; using '{chosen}'")
@@ -274,11 +273,11 @@ def prepare_with_preflight(
spec: BottleSpec,
*,
stage_dir: Path,
render_preflight: Callable[[DockerBottlePlan, str], None],
render_preflight: Callable[[BottlePlan, str], None],
prompt_yes: Callable[[], bool],
dry_run: bool = False,
backend_name: str | None = None,
) -> tuple[DockerBottlePlan | None, str]:
) -> tuple[BottlePlan | None, str]:
"""Run `backend.prepare`, render the preflight summary via the
injected callable, prompt y/N via the injected callable.
@@ -405,7 +404,7 @@ def _resolve_unique_label(label: str, color: str) -> tuple[str, str]:
in use among running bottles. Passes through unchanged when no
collision is found on the first check."""
while True:
slug_candidate = docker_mod.slugify(label)
slug_candidate = slugify(label)
active_slugs = {a.slug for a in enumerate_active_agents()}
if slug_candidate not in active_slugs:
return label, color
@@ -432,7 +431,7 @@ def _select_image_policy() -> str | None:
def _text_render_preflight():
def _render(plan: DockerBottlePlan, backend_name: str) -> None:
def _render(plan: BottlePlan, backend_name: str) -> None:
print(file=sys.stderr)
print(f"backend: {backend_name}", file=sys.stderr)
print(_manifest_to_yaml(plan.manifest), file=sys.stderr)
+22 -19
View File
@@ -368,7 +368,7 @@ def _main_loop(stdscr: "curses._CursesWindow") -> None: # type: ignore # pragm
elif key in (curses.KEY_UP, ord("k")):
selected = max(selected - 1, 0)
elif key in (curses.KEY_ENTER, 10, 13):
_detail_view(stdscr, qp, green_attr=green_attr)
status_line = _detail_view(stdscr, qp, green_attr=green_attr)
elif key == ord("a"):
try:
status_line = _approve_from_tui(stdscr, qp)
@@ -456,7 +456,7 @@ def _detail_view(
qp: QueuedProposal,
*,
green_attr: int = 0,
) -> None: # pragma: no cover
) -> str: # pragma: no cover
"""Render the full proposal. Scrollable. Press q to return."""
lines = _detail_lines(qp, green_attr=green_attr)
offset = 0
@@ -473,7 +473,7 @@ def _detail_view(
stdscr.refresh()
key = stdscr.getch()
if key in (ord("q"), 27):
return
return ""
if key in (curses.KEY_DOWN, ord("j")):
offset = min(offset + 1, max(0, len(lines) - 1))
elif key in (curses.KEY_UP, ord("k")):
@@ -484,31 +484,34 @@ def _detail_view(
offset = max(0, len(lines) - 1)
elif key == ord("a"):
try:
_approve_from_tui(stdscr, qp)
except ApplyError:
pass
return
return _approve_from_tui(stdscr, qp)
except ApplyError as exc:
return f"apply failed: {exc}"
elif key == ord("m"):
if qp.proposal.tool in _REPORT_ONLY_TOOLS:
return
return f"modify unavailable for {qp.proposal.tool}"
edited = _modify(stdscr, qp)
if edited is not None:
try:
_approve_from_tui(
stdscr, qp, final_file=edited,
notes="operator modified before approving",
)
except ApplyError:
pass
return
if edited is None:
return "modify aborted (no change)"
try:
return _approve_from_tui(
stdscr, qp, final_file=edited,
notes="operator modified before approving",
)
except ApplyError as exc:
return f"apply failed: {exc}"
elif key == ord("r"):
reason = _prompt(stdscr, "reject reason: ")
if reason:
reject(qp, reason=reason)
return
return f"rejected {qp.proposal.tool} for [{qp.label}]"
return "reject aborted (empty reason)"
def _modify(stdscr: "curses._CursesWindow", qp: QueuedProposal) -> str | None: # type: ignore # pragma: no cover
def _modify(
stdscr: "curses._CursesWindow", # type: ignore
qp: QueuedProposal,
) -> str | None: # pragma: no cover
"""Suspend curses, open $EDITOR on the proposed file, return edited content."""
suffix = _suffix_for_tool(qp.proposal.tool)
curses.endwin()
+32 -6
View File
@@ -16,6 +16,8 @@ import os
import sys
from typing import Any, Optional
from ..log import debug
def filter_multiselect(
items: list[str],
@@ -42,7 +44,11 @@ def filter_multiselect(
try:
tty_fd = open(tty_path, "r+b", buffering=0)
except OSError:
except OSError as exc:
debug(
"multi-select unavailable; treating it as cancellation",
context={"error_type": type(exc).__name__, "tty": tty_path},
)
return None
try:
@@ -73,7 +79,11 @@ def filter_select(
try:
tty_fd = open(tty_path, "r+b", buffering=0)
except OSError:
except OSError as exc:
debug(
"filter-select unavailable; treating it as cancellation",
context={"error_type": type(exc).__name__, "tty": tty_path},
)
return None
try:
@@ -129,7 +139,11 @@ def _run_picker(items: list[str], *, title: str, tty_fd: int) -> Optional[str]:
curses.nocbreak()
curses.echo()
curses.endwin()
except Exception: # noqa: W0718 — curses can raise many error types
except Exception as exc: # noqa: W0718 — curses can raise many error types
debug(
"filter-select display failed; treating it as cancellation",
context={"error_type": type(exc).__name__},
)
return None
finally:
sys.__stdin__ = orig_stdin # type: ignore[assignment]
@@ -292,7 +306,11 @@ def _run_multiselect(
curses.nocbreak()
curses.echo()
curses.endwin()
except Exception: # noqa: W0718
except Exception as exc: # noqa: W0718
debug(
"multi-select display failed; treating it as cancellation",
context={"error_type": type(exc).__name__},
)
return None
finally:
sys.__stdin__ = orig_stdin # type: ignore[assignment]
@@ -558,13 +576,21 @@ def name_color_modal(
"""
try:
tty_fd = open(tty_path, "r+b", buffering=0) # pylint: disable=consider-using-with
except OSError:
except OSError as exc:
debug(
"name/color picker unavailable; using defaults",
context={"error_type": type(exc).__name__, "tty": tty_path},
)
return default_label, ""
try:
fd_dup = os.dup(tty_fd.fileno())
return _run_name_color(default_label, tty_fd=fd_dup, disclaimer=disclaimer)
except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught
except Exception as exc: # noqa: BLE001 # pylint: disable=broad-exception-caught
debug(
"name/color picker failed; using defaults",
context={"error_type": type(exc).__name__},
)
return default_label, ""
finally:
tty_fd.close()
+2 -2
View File
@@ -11,7 +11,7 @@ from __future__ import annotations
from dataclasses import dataclass
from pathlib import Path
from ..gateway.egress.addon_core import Route
from ..gateway.egress.types import Route
@dataclass(frozen=True)
@@ -19,7 +19,7 @@ class EgressRoute(Route):
"""Host-side extension of the addon's `Route`.
Inherits `host`, `matches`, `auth_scheme`, and `token_env`
from `egress_addon_core.Route` those are the fields that cross the
from the gateway's wire `Route` — those are the fields that cross the
YAML wire into the gateway. The fields below are host-only and
are never serialised to the addon.
+2 -2
View File
@@ -14,8 +14,8 @@ import secrets
from pathlib import Path
from typing import TYPE_CHECKING
from ..gateway.egress.addon_core import (
ON_MATCH_REDACT,
from ..gateway.egress.dlp_config import ON_MATCH_REDACT
from ..gateway.egress.types import (
HeaderMatch as CoreHeaderMatch,
MatchEntry as CoreMatchEntry,
PathMatch as CorePathMatch,
+17 -11
View File
@@ -17,28 +17,34 @@ from mitmproxy import http # type: ignore[import-not-found] # pylint: disable=
from bot_bottle.constants import IDENTITY_HEADER
from bot_bottle.gateway.egress.dlp_detectors import redact_tokens, strip_crlf
from bot_bottle.gateway.egress.addon_core import (
LOG_BLOCKS,
LOG_FULL,
from bot_bottle.gateway.egress.dlp_config import (
DEFAULT_OUTBOUND_ON_MATCH,
ON_MATCH_BLOCK,
ON_MATCH_REDACT,
Config,
Route,
ScanResult,
)
from bot_bottle.gateway.egress.context import resolve_client_context
from bot_bottle.gateway.egress.dlp import (
build_inbound_scan_text,
build_outbound_scan_text,
build_token_allow_payload,
outbound_scan_headers,
scan_inbound,
scan_outbound,
)
from bot_bottle.gateway.egress.matching import (
decide,
decide_git_fetch,
is_git_fetch_request,
is_git_push_request,
match_route,
resolve_client_context,
outbound_scan_headers,
route_to_yaml_dict,
scan_inbound,
scan_outbound,
)
from bot_bottle.gateway.egress.schema import route_to_yaml_dict
from bot_bottle.gateway.egress.types import (
LOG_BLOCKS,
LOG_FULL,
Config,
Route,
ScanResult,
)
from bot_bottle.gateway.policy_resolver import PolicyResolveError, PolicyResolver
from bot_bottle.supervisor.types import (
File diff suppressed because it is too large Load Diff
+77
View File
@@ -0,0 +1,77 @@
"""Fail-closed resolution of a client's policy and egress credentials."""
from __future__ import annotations
import typing
from ...log import debug
from .types import Config
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."
)
class PolicyResolverLike(typing.Protocol):
def resolve(self, source_ip: str, identity_token: str = ...) -> str | None: ...
class ContextResolverLike(typing.Protocol):
def resolve_policy_and_bottle_id(
self, source_ip: str, identity_token: str = ...,
) -> tuple[str | None, str | None, dict[str, str]]: ...
def _config_from_policy(policy: str | None) -> Config:
# Local import keeps schema parsing independent of resolver protocols.
from .schema import load_config
if not policy:
return Config(routes=(), deny_reason=DENY_UNATTRIBUTED)
try:
return load_config(policy)
except ValueError:
return Config(routes=(), deny_reason=DENY_UNPARSEABLE)
def resolve_client_config(
resolver: PolicyResolverLike, client_ip: str, identity_token: str = "",
) -> Config:
try:
policy = resolver.resolve(client_ip, identity_token)
except Exception as exc: # noqa: BLE001 - a policy lookup failure must deny
debug(
"egress policy resolution failed; applying deny-all",
context={"error_type": type(exc).__name__},
)
return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR)
return _config_from_policy(policy)
def resolve_client_context(
resolver: ContextResolverLike, client_ip: str, identity_token: str = "",
) -> tuple[Config, str, dict[str, str]]:
try:
policy, bottle_id, tokens = resolver.resolve_policy_and_bottle_id(
client_ip, identity_token)
except Exception as exc: # noqa: BLE001 - a policy lookup failure must deny
debug(
"egress context resolution failed; applying deny-all",
context={"error_type": type(exc).__name__},
)
return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR), "", {}
return _config_from_policy(policy), (bottle_id or ""), tokens
+99
View File
@@ -0,0 +1,99 @@
"""DLP scan dispatch and safe proposal rendering for egress requests."""
from __future__ import annotations
import typing
from .types import Route, ScanResult
def build_outbound_scan_text(host: str, path: str, query: str,
headers: typing.Mapping[str, str], body: str) -> str:
parts = [host, path]
if query:
parts.append(query)
parts.extend(f"{name}: {value}" for name, value in headers.items())
if body:
parts.append(body)
return "\n".join(parts)
def outbound_scan_headers(route: Route, headers: typing.Mapping[str, str]) -> dict[str, str]:
"""Drop agent Authorization when the route injects gateway-owned auth."""
skip_auth = bool(route.auth_scheme and route.token_env)
return {name: value for name, value in headers.items()
if not (skip_auth and name.lower() == "authorization")}
def build_inbound_scan_text(headers: typing.Mapping[str, str], body: str) -> str:
parts = [f"{name}: {value}" for name, value in headers.items()]
if body:
parts.append(body)
return "\n".join(parts)
def _enabled(configured: tuple[str, ...] | None, name: str) -> bool:
return configured is None or name in configured
def scan_outbound(route: Route, body: str | bytes, environ: typing.Mapping[str, str], *,
safe_tokens: typing.AbstractSet[str] | None = None,
crlf_text: str | None = None) -> ScanResult | None:
if not route.inspect:
return None
try:
from dlp_detectors import ( # type: ignore[import-not-found]
scan_crlf_injection, scan_entropy, scan_known_secrets, scan_token_patterns)
except ImportError: # pragma: no cover - gateway's flat module path
from .dlp_detectors import (
scan_crlf_injection, scan_entropy, scan_known_secrets, scan_token_patterns)
if isinstance(body, bytes):
try:
text = body.decode("utf-8")
except UnicodeDecodeError:
text = body.decode("latin-1")
else:
text = body
result = scan_crlf_injection(text if crlf_text is None else crlf_text)
if result is not None:
return result
if _enabled(route.outbound_detectors, "token_patterns"):
result = scan_token_patterns(text, location="body", safe_tokens=safe_tokens)
if result is not None:
return result
if _enabled(route.outbound_detectors, "known_secrets"):
extra = tuple(prefix for prefix in environ.get(
"BOT_BOTTLE_SENSITIVE_PREFIXES", "").split(",") if prefix)
result = scan_known_secrets(text, location="body", env=environ,
sensitive_prefixes=("EGRESS_TOKEN_",) + extra,
safe_tokens=safe_tokens)
if result is not None:
return result
if route.outbound_detectors is not None and "entropy" in route.outbound_detectors:
return scan_entropy(text, location="body")
return None
def build_token_allow_payload(host: str, method: str, path: str, result: ScanResult) -> str:
"""Render redacted operator context; the raw matched secret is excluded."""
lines = [
"egress blocked an outbound request carrying a detected token",
f"host: {host}", f"method: {method}", f"path: {path}",
f"detector: {result.reason}",
]
if result.context:
lines.append(f"context: {result.context}")
return "\n".join(lines) + "\n"
def scan_inbound(route: Route, body: str | bytes) -> ScanResult | None:
if not route.inspect:
return None
try:
from dlp_detectors import scan_naive_injection # type: ignore[import-not-found]
except ImportError: # pragma: no cover - gateway's flat module path
from .dlp_detectors import scan_naive_injection
text = body if isinstance(body, str) else body.decode("utf-8", errors="replace")
if _enabled(route.inbound_detectors, "naive_injection_detection"):
return scan_naive_injection(text)
return None
+1 -1
View File
@@ -19,7 +19,7 @@ from math import log2
from collections import Counter
from urllib.parse import quote as url_quote
from .addon_core import ScanResult
from .types import ScanResult
# ---------------------------------------------------------------------------
+112
View File
@@ -0,0 +1,112 @@
"""Route matching and request-policy decisions for the egress gateway."""
from __future__ import annotations
import typing
from .types import Decision, MatchEntry, PathMatch, Route
def _path_matches(pm: PathMatch, request_path: str) -> bool:
if pm.type == "exact":
return request_path == pm.value
if pm.type == "prefix":
if request_path == pm.value:
return True
if not pm.value.endswith("/"):
return request_path.startswith(pm.value + "/")
return request_path.startswith(pm.value)
return (
pm.type == "regex"
and pm.compiled is not None
and pm.compiled.search(request_path) is not None
)
def _entry_matches(
entry: MatchEntry, request_path: str, request_method: str,
request_headers: typing.Mapping[str, str],
) -> bool:
if entry.paths and not any(_path_matches(pm, request_path) for pm in entry.paths):
return False
if entry.methods and request_method.upper() not in entry.methods:
return False
for match in entry.headers:
value = request_headers.get(match.name.lower())
if value is None:
return False
if match.type == "exact" and value != match.value:
return False
if match.type == "regex" and (
match.compiled is None or match.compiled.search(value) is None
):
return False
return True
def evaluate_matches(
route: Route, request_path: str, request_method: str = "GET",
request_headers: typing.Mapping[str, str] | None = None,
) -> bool:
"""Return whether a request satisfies a route's optional match entries."""
if not route.matches:
return True
return any(_entry_matches(entry, request_path, request_method, request_headers or {})
for entry in route.matches)
def is_git_push_request(path: str, query: str) -> bool:
return path.endswith("/git-receive-pack") or (
path.endswith("/info/refs") and any(
pair.partition("=") == ("service", "=", "git-receive-pack")
for pair in query.split("&")
)
)
def is_git_fetch_request(path: str, query: str) -> bool:
return path.endswith("/git-upload-pack") or (
path.endswith("/info/refs") and any(
pair.partition("=") == ("service", "=", "git-upload-pack")
for pair in query.split("&")
)
)
def match_route(routes: typing.Sequence[Route], request_host: str) -> Route | None:
target = request_host.lower()
return next((route for route in routes if route.host.lower() == target), None)
def decide(
routes: typing.Sequence[Route], request_host: str, request_path: str,
environ: typing.Mapping[str, str], *, request_method: str = "GET",
request_headers: typing.Mapping[str, str] | None = None, deny_reason: str = "",
) -> Decision:
route = match_route(routes, request_host)
if route is None:
return Decision("block", deny_reason or (
f"egress: host {request_host!r} is not in the bottle's egress.routes "
"allowlist. Declare a route for it or remove the request."))
if not evaluate_matches(route, request_path, request_method, request_headers):
return Decision("block", (
f"egress: request {request_method} {request_path!r} does not match any "
f"entry in matches for {route.host!r}"))
if route.auth_scheme and route.token_env:
token = environ.get(route.token_env, "")
if not token:
return Decision("block", (
f"egress: route for {route.host!r} declared auth but env var "
f"{route.token_env!r} is unset"))
return Decision("forward", inject_authorization=f"{route.auth_scheme} {token}")
return Decision("forward")
def decide_git_fetch(routes: typing.Sequence[Route], request_host: str) -> Decision:
route = match_route(routes, request_host)
if route is not None and route.git_fetch:
return Decision("forward")
return Decision("block", (
"egress: git fetch/clone over HTTPS is not allowed by default; use git-gate "
"for declared repos or set egress.routes[].git.fetch=true for explicit "
"read-only HTTPS Git access."))
+349
View File
@@ -0,0 +1,349 @@
"""Egress policy schema parsing and serialization (PRD 0017 / 0053)."""
from __future__ import annotations
import re
import typing
from ...yaml_subset import YamlSubsetError, parse_yaml_subset
from .dlp_config import parse_inspect_block
from .types import (
HEADER_MATCH_TYPES,
LOG_BLOCKS,
LOG_FULL,
LOG_OFF,
PATH_MATCH_TYPES,
VALID_METHODS,
Config,
HeaderMatch,
MatchEntry,
PathMatch,
Route,
)
# Parsing
# ---------------------------------------------------------------------------
def _parse_path_match(idx: int, j: int, raw: object) -> PathMatch:
label = f"route[{idx}] matches paths[{j}]"
if not isinstance(raw, dict):
raise ValueError(f"{label}: must be an object")
raw_dict: dict[str, object] = typing.cast(dict[str, object], raw)
ptype = raw_dict.get("type", "prefix")
if not isinstance(ptype, str) or ptype not in PATH_MATCH_TYPES:
raise ValueError(
f"{label}: 'type' must be one of {', '.join(PATH_MATCH_TYPES)} "
f"(got {ptype!r})"
)
value = raw_dict.get("value")
if not isinstance(value, str) or not value:
raise ValueError(f"{label}: 'value' must be a non-empty string")
if ptype in ("exact", "prefix") and not value.startswith("/"):
raise ValueError(
f"{label}: value {value!r} must start with '/' for "
f"type {ptype!r}"
)
compiled: re.Pattern[str] | None = None
if ptype == "regex":
try:
compiled = re.compile(value)
except re.error as e:
raise ValueError(
f"{label}: regex {value!r} failed to compile: {e}"
) from e
for k in raw_dict:
if k not in ("type", "value"):
raise ValueError(f"{label}: unknown key {k!r}")
return PathMatch(type=ptype, value=value, compiled=compiled)
def _parse_header_match(idx: int, j: int, raw: object) -> HeaderMatch:
label = f"route[{idx}] matches headers[{j}]"
if not isinstance(raw, dict):
raise ValueError(f"{label}: must be an object")
raw_dict: dict[str, object] = typing.cast(dict[str, object], raw)
name = raw_dict.get("name")
if not isinstance(name, str) or not name:
raise ValueError(f"{label}: 'name' must be a non-empty string")
value = raw_dict.get("value")
if not isinstance(value, str):
raise ValueError(f"{label}: 'value' must be a string")
htype = raw_dict.get("type", "exact")
if not isinstance(htype, str) or htype not in HEADER_MATCH_TYPES:
raise ValueError(
f"{label}: 'type' must be one of {', '.join(HEADER_MATCH_TYPES)} "
f"(got {htype!r})"
)
compiled: re.Pattern[str] | None = None
if htype == "regex":
try:
compiled = re.compile(value)
except re.error as e:
raise ValueError(
f"{label}: regex {value!r} failed to compile: {e}"
) from e
for k in raw_dict:
if k not in ("name", "value", "type"):
raise ValueError(f"{label}: unknown key {k!r}")
return HeaderMatch(name=name, value=value, type=htype, compiled=compiled)
def _parse_match_entry(idx: int, k: int, raw: object) -> MatchEntry:
label = f"route[{idx}] matches[{k}]"
if not isinstance(raw, dict):
raise ValueError(f"{label}: must be an object")
raw_dict: dict[str, object] = typing.cast(dict[str, object], raw)
paths: tuple[PathMatch, ...] = ()
paths_raw = raw_dict.get("paths")
if paths_raw is not None:
if not isinstance(paths_raw, list):
raise ValueError(f"{label}: 'paths' must be a list")
paths_list = typing.cast(list[object], paths_raw)
paths = tuple(_parse_path_match(idx, j, p) for j, p in enumerate(paths_list))
methods: tuple[str, ...] = ()
methods_raw = raw_dict.get("methods")
if methods_raw is not None:
if not isinstance(methods_raw, list):
raise ValueError(f"{label}: 'methods' must be a list")
methods_list = typing.cast(list[object], methods_raw)
normalised: list[str] = []
for j, m in enumerate(methods_list):
if not isinstance(m, str):
raise ValueError(f"{label}: methods[{j}] must be a string")
upper = m.upper()
if upper not in VALID_METHODS:
raise ValueError(
f"{label}: methods[{j}] {m!r} is not a valid HTTP method"
)
normalised.append(upper)
methods = tuple(normalised)
headers: tuple[HeaderMatch, ...] = ()
headers_raw = raw_dict.get("headers")
if headers_raw is not None:
if not isinstance(headers_raw, list):
raise ValueError(f"{label}: 'headers' must be a list")
headers_list = typing.cast(list[object], headers_raw)
headers = tuple(
_parse_header_match(idx, j, h) for j, h in enumerate(headers_list)
)
for key in raw_dict:
if key not in ("paths", "methods", "headers"):
raise ValueError(f"{label}: unknown key {key!r}")
return MatchEntry(paths=paths, methods=methods, headers=headers)
def parse_routes(payload: object) -> tuple[Route, ...]:
if not isinstance(payload, dict):
raise ValueError("routes payload: top-level must be an object")
payload_dict: dict[str, object] = typing.cast(dict[str, object], payload)
raw: object = payload_dict.get("routes")
if not isinstance(raw, list):
raise ValueError("routes payload: 'routes' must be a list")
raw_list: list[object] = typing.cast(list[object], raw)
out: list[Route] = []
for i, r in enumerate(raw_list):
out.append(_parse_one(i, r))
return tuple(out)
def _parse_one(idx: int, raw: object) -> Route:
label = f"route[{idx}]"
if not isinstance(raw, dict):
raise ValueError(f"{label}: must be an object (got {type(raw).__name__})")
raw_dict: dict[str, object] = typing.cast(dict[str, object], raw)
host: object = raw_dict.get("host")
if not isinstance(host, str) or not host:
raise ValueError(f"{label}: 'host' must be a non-empty string")
legacy_flat = "inspect" not in raw_dict
inspect_raw = raw_dict.get("inspect", {})
if inspect_raw is False:
inspect = False
settings: dict[str, object] = {}
elif isinstance(inspect_raw, dict):
inspect = True
settings = (
{k: v for k, v in raw_dict.items() if k != "host"}
if legacy_flat
else typing.cast(dict[str, object], inspect_raw)
)
legacy_dlp = settings.pop("dlp", None)
if isinstance(legacy_dlp, dict):
settings.update(typing.cast(dict[str, object], legacy_dlp))
elif legacy_dlp is not None:
raise ValueError(
f"{label} ({host}): legacy 'dlp' must be an object"
)
else:
raise ValueError(f"{label} ({host}): 'inspect' must be false or an object")
# matches
matches: tuple[MatchEntry, ...] = ()
matches_raw = settings.get("matches")
if matches_raw is not None:
if not isinstance(matches_raw, list):
raise ValueError(f"{label} ({host}): 'matches' must be a list")
matches_list = typing.cast(list[object], matches_raw)
matches = tuple(
_parse_match_entry(idx, k, m) for k, m in enumerate(matches_list)
)
# auth (unchanged wire format)
auth_scheme: object = settings.get("auth_scheme", "")
token_env: object = settings.get("token_env", "")
if not isinstance(auth_scheme, str):
raise ValueError(f"{label} ({host}): 'auth_scheme' must be a string")
if not isinstance(token_env, str):
raise ValueError(f"{label} ({host}): 'token_env' must be a string")
if bool(auth_scheme) != bool(token_env):
raise ValueError(
f"{label} ({host}): 'auth_scheme' and 'token_env' must be both "
f"set or both empty (got auth_scheme={auth_scheme!r}, "
f"token_env={token_env!r})"
)
# git-over-HTTPS policy
git_fetch = False
git_raw = settings.get("git")
if git_raw is not None:
if not isinstance(git_raw, dict):
raise ValueError(f"{label} ({host}): 'git' must be an object")
git_dict: dict[str, object] = typing.cast(dict[str, object], git_raw)
fetch_raw = git_dict.get("fetch", False)
if fetch_raw is True or fetch_raw is False:
git_fetch = fetch_raw
else:
raise ValueError(f"{label} ({host}): 'git.fetch' must be a boolean")
for k in git_dict:
if k != "fetch":
raise ValueError(
f"{label} ({host}): git has unknown key {k!r}; "
"accepted key is 'fetch'"
)
# dlp detectors
outbound_detectors, inbound_detectors, outbound_on_match = parse_inspect_block(
idx, host, settings,
)
preserve_auth_raw = settings.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 settings:
if k not in (
"matches", "auth_scheme", "token_env", "git", "preserve_auth",
"outbound_detectors", "inbound_detectors", "outbound_on_match",
):
raise ValueError(
f"{label} ({host}): inspect has unknown key {k!r}"
)
for k in raw_dict:
if not legacy_flat and k not in ("host", "inspect"):
raise ValueError(
f"{label} ({host}): unknown key {k!r}; accepted keys "
f"are 'host' and 'inspect'"
)
return Route(
host=host,
matches=matches,
auth_scheme=auth_scheme,
token_env=token_env,
git_fetch=git_fetch,
outbound_detectors=outbound_detectors,
inbound_detectors=inbound_detectors,
outbound_on_match=outbound_on_match,
preserve_auth=preserve_auth,
inspect=inspect,
)
def _path_match_to_dict(pm: PathMatch) -> dict[str, object]:
d: dict[str, object] = {"value": pm.value}
if pm.type != "prefix":
d["type"] = pm.type
return d
def _header_match_to_dict(hm: HeaderMatch) -> dict[str, object]:
d: dict[str, object] = {"name": hm.name, "value": hm.value}
if hm.type != "exact":
d["type"] = hm.type
return d
def _match_entry_to_dict(me: MatchEntry) -> dict[str, object]:
d: dict[str, object] = {}
if me.paths:
d["paths"] = [_path_match_to_dict(p) for p in me.paths]
if me.methods:
d["methods"] = list(me.methods)
if me.headers:
d["headers"] = [_header_match_to_dict(h) for h in me.headers]
return d
def route_to_yaml_dict(r: Route) -> dict[str, object]:
"""Serialize a Route to YAML-schema-compatible dict.
Uses the same field names the YAML parser accepts, so the output
can be round-tripped directly into an `allow` or `egress-block`
proposal without translation. Fields that are empty/default are
omitted so the agent doesn't copy irrelevant keys."""
d: dict[str, object] = {"host": r.host}
if not r.inspect:
d["inspect"] = False
return d
inspected: dict[str, object] = {}
if r.auth_scheme:
inspected["auth_scheme"] = r.auth_scheme
inspected["token_env"] = r.token_env
if r.matches:
inspected["matches"] = [_match_entry_to_dict(m) for m in r.matches]
if r.git_fetch:
inspected["git"] = {"fetch": True}
if r.outbound_detectors is not None:
inspected["outbound_detectors"] = list(r.outbound_detectors)
if r.inbound_detectors is not None:
inspected["inbound_detectors"] = list(r.inbound_detectors)
if r.outbound_on_match:
inspected["outbound_on_match"] = r.outbound_on_match
if r.preserve_auth:
inspected["preserve_auth"] = True
if inspected:
d["inspect"] = inspected
return d
def parse_config(payload: object) -> "Config":
"""Parse a full egress config payload (top-level log level + routes)."""
if not isinstance(payload, dict):
raise ValueError("routes payload: top-level must be an object")
payload_dict: dict[str, object] = typing.cast(dict[str, object], payload)
log_raw: object = payload_dict.get("log", LOG_OFF)
if log_raw is True or log_raw is False or not isinstance(log_raw, int) \
or log_raw not in (LOG_OFF, LOG_BLOCKS, LOG_FULL):
raise ValueError(
f"routes payload: 'log' must be {LOG_OFF}, {LOG_BLOCKS}, or {LOG_FULL}"
)
routes = parse_routes(payload)
return Config(routes=routes, log=log_raw)
def load_config(text: str) -> "Config":
"""Parse YAML text → Config (routes + log flag)."""
try:
payload = parse_yaml_subset(text)
except YamlSubsetError as e:
raise ValueError(f"routes payload: invalid YAML: {e}") from e
return parse_config(payload)
+81
View File
@@ -0,0 +1,81 @@
"""Shared egress policy value objects.
Kept dependency-free so the schema parser, matcher, DLP scanner, and addon
adapter can use the same immutable public shapes without importing each other.
"""
from __future__ import annotations
import re
from dataclasses import dataclass
PATH_MATCH_TYPES = ("exact", "prefix", "regex")
HEADER_MATCH_TYPES = ("exact", "regex")
VALID_METHODS = frozenset({
"GET", "HEAD", "POST", "PUT", "DELETE", "PATCH", "OPTIONS", "TRACE",
"CONNECT",
})
LOG_OFF = 0
LOG_BLOCKS = 1
LOG_FULL = 2
@dataclass(frozen=True)
class PathMatch:
type: str
value: str
compiled: re.Pattern[str] | None = None
@dataclass(frozen=True)
class HeaderMatch:
name: str
value: str
type: str = "exact"
compiled: re.Pattern[str] | None = None
@dataclass(frozen=True)
class MatchEntry:
paths: tuple[PathMatch, ...] = ()
methods: tuple[str, ...] = ()
headers: tuple[HeaderMatch, ...] = ()
@dataclass(frozen=True)
class Route:
host: str
matches: tuple[MatchEntry, ...] = ()
auth_scheme: str = ""
token_env: str = ""
git_fetch: bool = False
outbound_detectors: tuple[str, ...] | None = None
inbound_detectors: tuple[str, ...] | None = None
outbound_on_match: str = ""
preserve_auth: bool = False
inspect: bool = True
@dataclass(frozen=True)
class Config:
routes: tuple[Route, ...]
log: int = LOG_OFF
deny_reason: str = ""
@dataclass(frozen=True)
class Decision:
action: str
reason: str = ""
inject_authorization: str | None = None
@dataclass(frozen=True)
class ScanResult:
severity: str
reason: str
location: str = ""
context: str = ""
matched: str = ""
+3 -3
View File
@@ -58,9 +58,9 @@ import typing
from dataclasses import dataclass
from bot_bottle.constants import IDENTITY_HEADER
from bot_bottle.gateway.egress.addon_core import (
LOG_OFF, load_config, resolve_client_context, route_to_yaml_dict,
)
from bot_bottle.gateway.egress.context import resolve_client_context
from bot_bottle.gateway.egress.schema import load_config, route_to_yaml_dict
from bot_bottle.gateway.egress.types import LOG_OFF
from bot_bottle.gateway.policy_resolver import PolicyResolveError, PolicyResolver
from bot_bottle.supervisor import types as _sv
+33 -6
View File
@@ -18,6 +18,7 @@ import urllib.request
from collections.abc import Iterable
from dataclasses import dataclass
from ..log import debug
from ..orchestrator_auth import ROLE_CLI
from ..trust_domain import CONTROL_PLANE
from .server import ORCHESTRATOR_AUTH_HEADER
@@ -53,6 +54,23 @@ class RegisteredBottle:
env_var_secret: str = ""
@dataclass(frozen=True)
class BackendProbeFailure:
"""Safe diagnostic for an optional backend discovery probe."""
backend: str
error_type: str
def _probe_failure(backend: str, exc: BaseException) -> BackendProbeFailure:
failure = BackendProbeFailure(backend, type(exc).__name__)
debug(
"orchestrator discovery probe unavailable",
context={"backend": failure.backend, "error_type": failure.error_type},
)
return failure
class OrchestratorClient:
"""Trusted host-side client for the orchestrator control plane.
@@ -245,32 +263,41 @@ def discover_orchestrator_url(*, timeout: float = 2.0) -> str:
orchestrator TAP. Returns the first that answers `/health`; raises if none
do (no orchestrator up launch a bottle first)."""
candidates: list[str] = []
failures: list[BackendProbeFailure] = []
try: # docker: loopback-published control plane
from .lifecycle import DEFAULT_PORT as _DOCKER_PORT
candidates.append(f"http://127.0.0.1:{_DOCKER_PORT}")
except Exception: # noqa: BLE001 — backend optional
except Exception as exc: # noqa: BLE001 — backend optional
failures.append(_probe_failure("docker", exc))
candidates.append("http://127.0.0.1:8099")
try: # firecracker: infra VM control plane on the orchestrator TAP
from ..backend.firecracker import netpool
from ..backend.firecracker.infra_vm import ORCHESTRATOR_PORT
candidates.append(
f"http://{netpool.orch_slot().guest_ip}:{ORCHESTRATOR_PORT}")
except Exception: # noqa: BLE001 — backend optional / not firecracker
pass
except Exception as exc: # noqa: BLE001 — backend optional / not firecracker
failures.append(_probe_failure("firecracker", exc))
try: # macOS: orchestrator container on its host-only address
from ..backend.macos_container.infra import probe_orchestrator_url
url = probe_orchestrator_url()
if url:
candidates.append(url)
except Exception: # noqa: BLE001 — backend optional / not macOS
pass
except Exception as exc: # noqa: BLE001 — backend optional / not macOS
failures.append(_probe_failure("macos-container", exc))
for url in candidates:
if OrchestratorClient(url, timeout=timeout).health():
return url
detail = ""
if failures:
detail = "; optional probes unavailable: " + ", ".join(
f"{failure.backend} ({failure.error_type})" for failure in failures
)
raise OrchestratorClientError(
"no running orchestrator control plane found (tried "
+ ", ".join(candidates)
+ "); launch a bottle first"
+ ")"
+ detail
+ "; launch a bottle first"
)
+9 -1
View File
@@ -2,6 +2,7 @@
from __future__ import annotations
from ..log import debug
from .client import OrchestratorClient, OrchestratorClientError
@@ -27,7 +28,14 @@ def reprovision_bottles(
try:
if client.reprovision_gateway(bottle_id, secret):
restored += 1
except OrchestratorClientError:
except OrchestratorClientError as exc:
debug(
"gateway secret reprovision failed; continuing with other bottles",
context={
"bottle_id": bottle_id,
"error_type": type(exc).__name__,
},
)
continue
return restored
+20 -8
View File
@@ -57,6 +57,7 @@ from __future__ import annotations
import http.server
import json
import math
import os
import socketserver
import sys
@@ -217,13 +218,18 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
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]
if any(not isinstance(ip, str) or not ip for ip in raw_ips):
return 400, {"error": "live_source_ips must contain non-empty strings"}
live = raw_ips
grace = data.get("grace_seconds")
kwargs = (
{"grace_seconds": float(grace)}
if isinstance(grace, (int, float)) and not isinstance(grace, bool)
else {}
)
kwargs: dict[str, float] = {}
if grace is not None:
if isinstance(grace, bool) or not isinstance(grace, (int, float)):
return 400, {"error": "grace_seconds must be a non-negative finite number"}
parsed_grace = float(grace)
if not math.isfinite(parsed_grace) or parsed_grace < 0:
return 400, {"error": "grace_seconds must be a non-negative finite number"}
kwargs["grace_seconds"] = parsed_grace
return 200, {"reaped": orch.reconcile(live, **kwargs)}
if method == "POST" and route == "/attribute":
@@ -373,9 +379,15 @@ class Handler(http.server.BaseHTTPRequestHandler):
status, payload = dispatch(
server.orchestrator, method, self.path, body, role=role)
except Exception as e: # noqa: BLE001 — the control plane must stay up
sys.stderr.write(f"orchestrator: {method} {self.path} failed: {e!r}\n")
# Do not echo exception messages to the caller or logs: broker and
# persistence exceptions can contain request data. The operation,
# route, and exception type are enough to correlate a traceback.
sys.stderr.write(
f"orchestrator: {method} {self.path} failed "
f"[error_type={type(e).__name__}]\n"
)
sys.stderr.flush()
status, payload = 500, {"error": f"internal error: {e}"}
status, payload = 500, {"error": "internal error"}
data = json.dumps(payload).encode()
self.send_response(status)
self.send_header("Content-Type", "application/json")
+20
View File
@@ -9,8 +9,11 @@ import difflib
import hashlib
import ipaddress
import os
import re
import sys
from .log import die
def sha256_hex(content: str) -> str:
"""Hex SHA-256 of a UTF-8 string."""
@@ -67,3 +70,20 @@ def expand_tilde(path: str) -> str:
home = os.environ.get("HOME", "")
return home + path[1:]
return path
_SLUG_RE = re.compile(r"[^a-z0-9]+")
def slugify(name: str) -> str:
"""Return a portable bottle identifier from a human-readable name.
This is deliberately a root utility: names are part of the generic CLI
and state model, not a Docker container concern.
"""
if not name:
die("slugify: missing name")
slug = _SLUG_RE.sub("-", name.lower()).strip("-")
if not slug:
die(f"name '{name}' produced an empty slug; use alphanumeric characters")
return slug
+45 -42
View File
@@ -1,50 +1,53 @@
# CI
The test workflow lives at [`.gitea/workflows/test.yml`](../.gitea/workflows/test.yml).
It runs the unit suite plus one integration job per backend
(`integration-docker`, `integration-firecracker`, `integration-macos`) on:
## Required pull-request gate
- every push to a branch with an open pull request, and
- every push to `main`.
[`.gitea/workflows/test.yml`](../.gitea/workflows/test.yml) runs the unit
suite, Docker integration suite, combined coverage report, and diff-coverage
gate when tested package/build inputs change on a pull request or on `main`.
`integration-macos` is the exception: it is **advisory**, running only on
`workflow_dispatch` (manual dispatch), never on push or pull requests. It targets the
Apple Container backend on a self-hosted macOS runner (label `macos`,
registered in host mode — Apple Container can't run in a Linux container, so it
can't reuse the `kvm` runner). A single non-redundant laptop must not be able
to block a PR merge, so the job stays out of the `coverage` job's `needs` and
its coverage never feeds the diff-coverage gate. Because the infra container is
a singleton (`bot-bottle-mac-infra`), the job declares a `concurrency` group
and tears the container down on exit; keep runner concurrency at 1. See the
README "macOS Apple Container" CI note for runner provisioning.
The Docker job preflights the backend before discovery. Gitea's `act_runner`
runs the job in a container with the host Docker socket, so the test process
reaches control-plane siblings through the job's Docker network and uses named
Docker volumes for orchestrator/CA state the host daemon must mount. The
orchestrator runs the package baked into the image built from the checkout; it
does not bind the job container's invisible workspace into a sibling container.
Docker integration jobs share fixed singleton names, so required and manual
runs use one non-cancelling concurrency group. The shared agent/gateway network
has an explicit subnet, which Docker requires for the pinned source IPs used as
the isolation/attribution key.
Each integration job selects its backend via `BOT_BOTTLE_BACKEND` and
runs a **preflight** (`./cli.py backend status --backend=<name>`) that
prints a clear per-check readiness summary and fails the job when the
backend is missing — so absent infrastructure is visible at the job level
rather than hidden among per-test `unittest.skip` lines. The skip guards in
[`tests/_backend.py`](../tests/_backend.py) gate on the same readiness
check (`bot_bottle.backend.has_backend`): backend-agnostic tests use
`skip_unless_selected_backend_available()` and run through whichever
backend is selected (checking, e.g., Linux + `/dev/kvm` for Firecracker
rather than unrelated Docker availability); Docker-implementation tests use
`skip_unless_backend("docker")` and no-op under a non-Docker run.
`scripts.unittest_gate` enforces the Docker job's contract: all 22 integration
tests must execute and none may skip. This includes the real gateway-image,
control-plane authentication, multitenant policy/token isolation,
sandbox-escape, and orphan-network tests. Backend skip decorators remain useful
for local runs, but the CI preflight plus execution-count gate prevents a
missing backend or runner-topology regression from becoming a green job.
A small subset of integration tests skip when running specifically
under Gitea Actions (`GITEA_ACTIONS=true`), because `act_runner` runs
the job inside a container with the host's `/var/run/docker.sock`
mounted in. That topology breaks two assumptions those tests make:
Combined unit + Docker coverage is informational globally. Two focused gates
are enforced:
- networks created via the host daemon aren't always visible to a
same-process `docker network ls` call from inside the job container,
and
- ports published by sibling containers land on the host's loopback,
not on the job container's `127.0.0.1` — so HTTP probes against
`http://127.0.0.1:<host_port>` from inside the job time out.
- changed executable Python lines must be at least 90% covered; and
- the validated critical security/logic core must remain at least 90% covered.
The affected tests (`test_orphan_cleanup.test_create_and_remove`,
`test_gateway_image.TestGatewayImage`) still run
locally where the test process and Docker daemon share a host.
Making them work in CI is a follow-up: either re-write them to
discover container IPs via `docker inspect`, or reconfigure the
runner with host networking.
## Privileged pre-release matrix
[`.gitea/workflows/pre-release-test.yml`](../.gitea/workflows/pre-release-test.yml)
is manually dispatched before a release. It repeats unit and Docker integration
coverage, then runs:
- Firecracker integration on the self-hosted `kvm` runner; and
- advisory Apple Container integration on the self-hosted `macos` runner.
These privileged host-mode runners never execute unreviewed pull-request code
automatically. Firecracker coverage is combined in the manual pre-release
report; macOS reports advisory coverage in its own job. The macOS infra
container is a singleton, so its job uses a concurrency group and always tears
the service down.
## Scheduled canary
[`.gitea/workflows/canaries.yml`](../.gitea/workflows/canaries.yml) runs weekly
and on manual dispatch. It verifies the pinned gitleaks release URL, checksum,
archive shape, and executable. The same unittest execution gate requires at
least one executed canary and rejects skips.
+10 -6
View File
@@ -34,12 +34,13 @@ a regression (Goodhart's law).
Coverage is **risk-weighted**, measured over the **combined unit +
integration** suites, with three rules:
1. **Critical modules target ≥ 90%.** The security/logic core
`egress_addon{,_core}.py`, `dlp_detectors.py`, `egress.py`,
`manifest*.py`, `git_gate.py`, `git_http_backend.py`, `supervise.py`,
`yaml_subset.py`, `bottle_state.py` — is Docker-independent and
unit-testable, so it carries the high bar. We ratchet toward 90% as
these modules are touched; new gaps in them are not acceptable.
1. **Critical modules must remain ≥ 90%.** The curated security/logic core
covers the host and gateway egress policy, manifest trust boundary,
git-gate enforcement, supervise protocol/server, YAML parser, and bottle
state. The concrete module list lives in `scripts/critical-modules.txt`;
`scripts/critical_modules.py` rejects stale or ambiguous entries before
Coverage.py can silently ignore them. These modules are unit-testable, so
CI enforces the aggregate minimum independently of diff coverage.
2. **Subprocess/backend orchestration is covered by the integration
suite, not omitted.** `scripts/coverage.sh` runs unit + integration
@@ -82,6 +83,9 @@ omit list.
(critical-module standard + diff coverage) are Docker-independent.
- "We're at N%" is now a curated figure; outsiders should read the
policy, not just the badge.
- A rename or removal in the curated list fails CI. Updating the list is an
explicit review of where the security-critical behavior moved, not a way to
improve the percentage by omission.
## Links
+273
View File
@@ -0,0 +1,273 @@
# PRD prd-new: Host control server
- **Status:** Draft
- **Author:** Claude
- **Created:** 2026-07-26
- **Issue:** #468
## Summary
Promote the in-process launch broker into a standalone **host control
server**: the single privileged component on the host. Both the CLI and the
orchestrator drive it over HTTP; it brokers agent launches, owns the
orchestrator's own lifecycle, and is the sole writer of host-durable state (the
tamper-evident audit record). This closes the three gaps between today's
well-formed broker *contract* ([`orchestrator/broker.py`](../../bot_bottle/orchestrator/broker.py))
and a real out-of-process service — transport, durable provisioned secret,
and a disciplined op vocabulary — and splits host state by
owner and lifetime. The prize: **the CLI no longer needs the Docker socket**,
which is what finally lets a dedicated Gitea runner user drop the
root-equivalent `docker` group (PRD 0070, "Relationship to other work").
## Problem
Container launches run directly from a short-lived CLI process against the
Docker socket. That socket is root-equivalent, so every host that launches
bottles hands root to whoever invokes the CLI — including a CI runner user we
want to keep unprivileged. PRD 0070 already argues for replacing the fat socket
with a **thin, structured, auditable** launch broker, and the contract for that
broker exists and is tested in-process. But it is *only* in-process:
`LaunchBroker.submit(token)` is a method call from
`OrchestratorCore.launch_bottle` ([`service.py:116`](../../bot_bottle/orchestrator/service.py)),
and `DockerBroker` is on no production path — every backend starts the
orchestrator with `--broker stub` ([`__main__.py:54`](../../bot_bottle/orchestrator/__main__.py)).
Three gaps stand between that scaffold and a host service:
1. **No transport.** `submit` is an in-process call. A real service needs a
`BrokerClient` that POSTs the signed token and a host-side HTTP server that
verifies and acts.
2. **The signing secret is ephemeral and self-generated.**
[`__main__.py:53`](../../bot_bottle/orchestrator/__main__.py) does
`secrets.token_bytes(32)` and hands the *same value* to signer and verifier —
viable only because they share a process. A separate daemon needs the secret
provisioned out of band and durable across orchestrator restarts.
3. **The op vocabulary is `launch` / `teardown` only.** Everything else
host-privileged still lives in the CLI, so the schema has to grow — carefully,
since PRD 0070's security argument rests on "structured requests only, static
flags + ids."
Separately, host state has no clear owner. `OrchestratorCore.reconcile` takes
`live_source_ips` as a parameter *only because the orchestrator cannot see the
backend* ([`service.py:137`](../../bot_bottle/orchestrator/service.py)); the
egress traffic log is written to the container's stderr; and there is no durable,
tamper-evident home for the audit record that survives orchestrator destruction.
## Goals / Success Criteria
- A standalone host control server that the CLI and orchestrator reach over
**HTTP**, with three entry paths working end to end:
- `web console -(iroh)-> orchestrator -(http)-> host controller -> launch`
- `cli -(http)-> orchestrator -(http)-> host controller -> launch`
- `cli -(http)-> host controller` — start / restart / status of the
orchestrator **itself** (the bootstrap/recovery path #391 targets).
- The launch op is expressed as a **signed JWT of static flags + ids only**,
verified against a closed schema.
- The signing secret is **provisioned out of band and durable** across
orchestrator restarts (a `TrustDomain` per #476, with a key the orchestrator
never holds for the host controller's *own* endpoints).
- Host-privileged operations move off the CLI to the control server; **the CLI
no longer opens the Docker socket** for bottle operations.
- `Orchestrator.reconcile` no longer takes `live_source_ips` — live-bottle
enumeration becomes an internal control-server call.
- Host-durable state lands as an **append-only, hash-chained JSONL** audit log
owned solely by the host controller; operational state stays SQLite owned
solely by the orchestrator.
## Non-goals
- **Removing standing privilege.** This converts on-demand privilege (a CLI the
user invokes) into standing privilege (a daemon under launchd/systemd). The
win is that the privilege is *narrower* (structured requests vs. a raw socket),
not that it disappears. "Always running" is an accepted new property.
- **Asymmetric signing.** We stay HS256 — see Design / "Signing stays
symmetric."
- **Integrity against a live compromised orchestrator.** Host-location of the
audit log does not buy this: the orchestrator makes the decisions being audited
and can forge or omit entries wherever the file lives. An off-box copy is the
answer, tracked separately.
- **A single unified DB for all state.** Impossible over a guest-kernel share
(SQLite locking is not coherent); state is split by owner and lifetime instead.
- **The generic `SecretProvider` (#355)** and **remote terminal design (#478)**
both ride the same door but are their own work.
## Design
### Topology
The host controller is the sole privileged component. The orchestrator becomes a
client of it for launches, and the CLI becomes a client of it for *both* bottle
operations (indirectly, through the orchestrator) and orchestrator lifecycle
(directly, for bootstrap/recovery — startup can't route through the thing being
started).
```
web console ─(iroh)─▶ orchestrator ─┐
├─(http, signed JWT)─▶ host controller ─▶ launch
cli ────────(http)──▶ orchestrator ─┘
cli ────────(http, bearer)──────────────────────────────▶ host controller (orchestrator lifecycle)
```
### Transport: `BrokerClient` + host server
`LaunchBroker.submit(token)` keeps its exact signature and semantics; only the
*wire* changes. A new `BrokerClient` implements the same submit contract by
POSTing the signed token to the host controller (stdlib `urllib`, like the
existing [`orchestrator/client.py`](../../bot_bottle/orchestrator/client.py)),
and the host controller's launch handler is the existing `verify_request` +
`_launch`/`_teardown` path, now reached over HTTP instead of a method call. The
in-process `StubBroker` stays for the dev-harness and tests; `DockerBroker`'s
`_launch`/`_teardown` bodies move behind the server unchanged. Because the client
satisfies the same interface `OrchestratorCore` already depends on, the core does
not change to gain a real backend.
### Signing stays symmetric (HS256)
PRD 0070 nominally specifies asymmetric; the code is HS256 and we keep it.
Asymmetric matters when the verifier is *less* privileged than the signer — here
it is the reverse: the host controller (verifier) is strictly more privileged
than the orchestrator (signer), and a controller that could forge orchestrator
requests gains nothing, since it is already the component that launches. Staying
symmetric also honors the no-runtime-deps policy (stdlib has no Ed25519). This
matches the reasoning already inlined in `broker.py`'s module docstring.
### Replay protection is out of scope (tracked in #494)
Once the launch token travels over a wire, a captured token could be replayed —
`sign_request` already emits `jti`/`iat` but `verify_request` reads neither, so
there is no expiry window or `jti` cache today. Enforcing that (an `iat` window +
a self-trimming `jti` cache) is a pure in-process change that lands independently
of this work, and it is deferred to **#494** rather than gating the MVP of the
host control server. Nothing here depends on it; it can merge before or after.
### Op vocabulary and the "ids + static flags" rule (gap 3)
Each op moved off the CLI widens the privileged surface, so growth is governed by
one explicit rule, enforced in `verify_request`'s schema check:
> A broker op carries **only ids and enumerated static flags** — a bottle id, a
> pool slot, a **content-addressed** image ref chosen from a fixed set, an op
> name from a closed vocabulary. Never a free-form path, argv, command, or
> caller-supplied filesystem location. If an operation cannot be expressed that
> way, it does not become a broker op.
Operations that fit and move off the CLI (all today in
`backend/*/consolidated_launch.py`, driven by a short-lived CLI process):
| Op | What it does | Fits the rule because |
|---|---|---|
| `launch` / `teardown` | existing | ids + slot + image ref |
| `orchestrator.ensure_running` | start the infra container | no arguments |
| `orchestrator.{start,restart,status}` | lifecycle (the #391 path) | no arguments |
| `list_live` | enumerate running bottles for reconcile | no arguments; returns ids/IPs |
| `allocate_ip` | `next_free_ip` over `_network_container_ips` | no arguments; returns an IP |
| `provision_git_gate` | `cp`/`exec` a per-bottle deploy key into the gateway | bottle id + key handle, no path |
| `reprovision` | `docker exec printenv <ENV_VAR_SECRET>` on a live agent | bottle id + secret *name* |
Image **builds** stay with the orchestrator for v1 (PRD 0070 §Memory: builds run
control-plane-side; a dedicated slim build unit is later, #468-adjacent), so no
`build` broker op is added here.
With `list_live` as an internal control-server call, `Orchestrator.reconcile`'s
`live_source_ips` parameter goes away — the tell PRD 0070 called out that the
orchestrator couldn't see the backend disappears with it.
### Secret provisioning (gap 2)
The shared HS256 secret becomes a durable, out-of-band artifact via the
**`TrustDomain`** seam (#476,
[`trust_domain.py`](../../bot_bottle/trust_domain.py)):
- The **launch-broker secret** is a `TrustDomain` whose key
(`host_signing_key(<file>)`, minted 0600 on first use, durable under
`bot_bottle_root()`) is provisioned to the orchestrator (signer) and the host
controller (verifier). Durability across orchestrator restarts is what makes
re-adoption work — a restart re-verifies against the same key.
- The **host controller's own lifecycle endpoints** (the direct `cli -> host
controller` path) get a **separate** `TrustDomain` key the orchestrator never
holds — exactly the second domain #476's PRD reserves. The orchestrator must
not be able to mint the credentials used to start and stop it.
This reuses the seam #476 landed rather than re-deriving provisioning per
backend (the PR #471 bug class).
### One daemon, structurally separate handlers (open decision 1)
The audit writer and the broker live in **one daemon** for install simplicity,
but with **no shared parsing** and **different credentials per handler**:
- the **launch** handler requires the signed launch **JWT** (provenance +
un-coercible schema);
- the **audit-append** handler takes a plain **bearer token** and writes to the
JSONL log.
This does not defend against orchestrator compromise (it holds both creds) — it
stops a bug in the boring audit path from reaching the privileged launch path.
The launcher stays small enough to audit line-by-line, per PRD 0070.
### State ownership: split by owner and lifetime
A single mounted DB is impossible — SQLite locking is not coherent across guest
kernels over a share, which is why the macOS backend already uses a container-only
volume (`INFRA_DB_VOLUME`). So state splits three ways (depends on #469, which
gets `bot-bottle.db` off the data plane first):
| Owner | State | Home | Shape |
|---|---|---|---|
| **Orchestrator** | `orchestrator_bottles` registry; `bottled_agent_secrets` (encrypted egress tokens); `supervise_proposals` / `supervise_responses` | volume nothing else mounts (generalizing the macOS design) | **SQLite** — mutable, transactional, queried |
| **Host controller** | supervise audit entries; egress traffic log (today → container stderr); host-side config | host filesystem, survives orchestrator/volume destruction | **JSONL** — append-only |
| **Gateway** | none | — | after #469 the data plane holds no DB state |
The historical record is **JSONL, not SQLite**, because it is append-only, never
updated, never transactionally queried: `O_APPEND` writes are atomic, there is no
locking protocol to get wrong, hash-chaining for tamper-evidence is cheap, and it
survives container-runtime volume pruning (the #450 lesson) and stays readable
without the orchestrator running. Both halves of "the audit record" — supervise
decisions and the egress traffic log — land in the one place.
The orchestrator is **sole mounter and sole writer** of its SQLite volume; the
host controller is **sole writer** of the JSONL log, over the authenticated
audit-append channel.
## Implementation chunks
Ordered, each independently mergeable:
1. **`BrokerClient` + host launch server** over HTTP, reusing `verify_request`
and the existing `DockerBroker` bodies. Wire `OrchestratorCore` to a
`BrokerClient` behind a flag; keep `StubBroker` for the dev-harness. Closes
gap 1.
2. **Durable secret via `TrustDomain`** — provision the launch-broker key to
signer + verifier; add the host controller's own lifecycle `TrustDomain`.
Closes gap 2.
3. **Grow the op vocabulary** one op at a time (`list_live` first — it also
removes `reconcile`'s `live_source_ips`), each behind the ids + static-flags
rule. Closes gap 3.
4. **JSONL audit log** — the host-controller-owned, hash-chained historical
record with the plain-bearer audit-append handler; redirect the egress traffic
log into it.
5. **Drop the Docker socket from the CLI** once every host-privileged op it used
is a broker op — the payoff that unblocks the unprivileged Gitea runner user.
## Open questions
1. **Schema-width rule enforcement.** The "ids + static flags" rule is stated;
should `verify_request` reject unknown claim keys outright (strict schema) to
keep the surface from drifting? Leaning yes.
2. **Audit-append back-pressure.** What the audit handler does if the JSONL sink
is unavailable (fail-closed vs. buffer) — resolve before shipping chunk 5.
## References
- **PRD 0070** — the contract, the launch broker, and the state tiers this
implements.
- **#469** — get `bot-bottle.db` off the data plane (lands underneath this).
- **#476** ([`prd-new-control-plane-auth-provisioning`](prd-new-control-plane-auth-provisioning.md))
— the `TrustDomain` seam this plugs the host controller's key into.
- **#391** — backend-agnostic orchestrator restart (the bootstrap path).
- **#494** — enforce broker replay protection (`iat` window + `jti` cache); split
out of this PRD as an independent in-process change.
- **#386** — prebuilt images from the Gitea OCI registry (the fixed image set the
broker validates against).
- **#355** — generic `SecretProvider`.
- **#478** — remote terminal design.
+8 -8
View File
@@ -20,10 +20,10 @@ cd "$(dirname "$0")/.."
PY="${PYTHON:-python3}"
# Critical security/logic core held to the high bar by ADR 0004. The list
# lives in one place (scripts/critical-modules.txt) so this report and the
# README "core coverage" badge can't drift; comma-join it for --include.
CRITICAL=$(grep -vE '^[[:space:]]*(#|$)' scripts/critical-modules.txt | paste -sd, -)
# Critical security/logic core held to the high bar by ADR 0004. The helper
# fails before coverage when a curated path was renamed or removed; Coverage.py
# itself would silently ignore that stale include and inflate the score.
CRITICAL=$("$PY" scripts/critical_modules.py)
if [ "${1:-}" = "aggregate" ]; then
# Aggregate mode: combine .coverage.* artifacts already in the workspace.
@@ -34,8 +34,8 @@ if [ "${1:-}" = "aggregate" ]; then
"$PY" -m coverage report -m
if [ "${2:-}" = "critical" ]; then
echo "== critical modules (ADR 0004 target: 90%) ==" >&2
"$PY" -m coverage report --include="$CRITICAL"
echo "== critical modules (ADR 0004 minimum: 90%) ==" >&2
"$PY" -m coverage report --include="$CRITICAL" --fail-under=90
fi
exit 0
fi
@@ -55,6 +55,6 @@ echo "== combined report ==" >&2
"$PY" -m coverage report -m
if [ "${1:-}" = "critical" ]; then
echo "== critical modules (ADR 0004 target: 90%) ==" >&2
"$PY" -m coverage report --include="$CRITICAL"
echo "== critical modules (ADR 0004 minimum: 90%) ==" >&2
"$PY" -m coverage report --include="$CRITICAL" --fail-under=90
fi
+38 -9
View File
@@ -7,19 +7,48 @@
# number that silently stops measuring a module is worse than no badge.
#
# One module path per line, relative to the repo root. Blank lines and
# `#` comments are ignored.
# `#` comments are ignored. scripts/critical_modules.py rejects missing,
# duplicate, non-Python, and out-of-repository entries before coverage runs.
# Host-side egress planning and secret preparation.
bot_bottle/egress/plan.py
bot_bottle/egress/service.py
# Gateway egress policy, matching, and DLP enforcement.
bot_bottle/gateway/egress/addon.py
bot_bottle/gateway/egress/addon_core.py
bot_bottle/gateway/egress/context.py
bot_bottle/gateway/egress/dlp.py
bot_bottle/gateway/egress/dlp_config.py
bot_bottle/gateway/egress/dlp_detectors.py
bot_bottle/egress.py
bot_bottle/manifest.py
bot_bottle/manifest_egress.py
bot_bottle/manifest_agent.py
bot_bottle/manifest_schema.py
bot_bottle/git_gate.py
bot_bottle/gateway/egress/matching.py
bot_bottle/gateway/egress/schema.py
bot_bottle/gateway/egress/types.py
# Manifest trust boundary and schema.
bot_bottle/manifest/agent.py
bot_bottle/manifest/bottle.py
bot_bottle/manifest/egress.py
bot_bottle/manifest/extends.py
bot_bottle/manifest/git.py
bot_bottle/manifest/index.py
bot_bottle/manifest/loader.py
bot_bottle/manifest/schema.py
bot_bottle/manifest/util.py
# Host-side and gateway-side git policy enforcement.
bot_bottle/git_gate/host_key.py
bot_bottle/git_gate/plan.py
bot_bottle/git_gate/provision.py
bot_bottle/git_gate/service.py
bot_bottle/gateway/git_gate/render.py
bot_bottle/git_gate_provision.py
bot_bottle/gateway/git_gate/http_backend.py
bot_bottle/supervise.py
# Supervise proposal protocol and data plane.
bot_bottle/supervisor/plan.py
bot_bottle/supervisor/types.py
bot_bottle/gateway/supervisor/server.py
# Shared parsers and state validation.
bot_bottle/yaml_subset.py
bot_bottle/bottle_state.py
+101
View File
@@ -0,0 +1,101 @@
#!/usr/bin/env python3
"""Validate and render the critical-module coverage manifest.
Coverage.py silently ignores an ``--include`` path that does not exist. That
is useful for broad globs, but dangerous for bot-bottle's curated security
core: a rename could otherwise improve the reported percentage by removing a
module from the measurement. Keep the validation in one small stdlib helper
and make every coverage consumer call it.
"""
from __future__ import annotations
import argparse
import sys
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parents[1]
DEFAULT_MANIFEST = REPO_ROOT / "scripts" / "critical-modules.txt"
class CriticalModulesError(ValueError):
"""The critical-module manifest is empty, ambiguous, or stale."""
def load_critical_modules(manifest: Path, *, root: Path) -> list[str]:
"""Return validated module paths relative to *root*.
Entries must be unique, concrete Python files inside the repository.
Globs are deliberately rejected by the file check: each rename must update
this explicit security review surface.
"""
root = root.resolve()
try:
lines = manifest.read_text(encoding="utf-8").splitlines()
except OSError as exc:
raise CriticalModulesError(
f"cannot read critical-module manifest {manifest}: {exc}"
) from exc
modules: list[str] = []
seen: set[str] = set()
errors: list[str] = []
for line_number, raw in enumerate(lines, start=1):
entry = raw.strip()
if not entry or entry.startswith("#"):
continue
path = Path(entry)
prefix = f"{manifest}:{line_number}: {entry!r}"
if path.is_absolute():
errors.append(f"{prefix} must be relative to the repository root")
continue
try:
resolved = (root / path).resolve()
resolved.relative_to(root)
except ValueError:
errors.append(f"{prefix} escapes the repository root")
continue
if entry in seen:
errors.append(f"{prefix} is duplicated")
continue
seen.add(entry)
if path.suffix != ".py":
errors.append(f"{prefix} is not a Python module")
continue
if not resolved.is_file():
errors.append(f"{prefix} does not exist")
continue
modules.append(path.as_posix())
if not modules and not errors:
errors.append(f"{manifest}: contains no critical modules")
if errors:
raise CriticalModulesError("\n".join(errors))
return modules
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
description="validate and print the critical coverage include list"
)
parser.add_argument("--manifest", type=Path, default=DEFAULT_MANIFEST)
parser.add_argument("--root", type=Path, default=REPO_ROOT)
parser.add_argument(
"--check", action="store_true",
help="validate only; do not print the comma-separated include list",
)
args = parser.parse_args(argv)
try:
modules = load_critical_modules(args.manifest, root=args.root)
except CriticalModulesError as exc:
print(f"critical-modules: {exc}", file=sys.stderr)
return 1
if not args.check:
print(",".join(modules))
return 0
if __name__ == "__main__":
raise SystemExit(main())
+67
View File
@@ -0,0 +1,67 @@
#!/usr/bin/env python3
"""Run unittest discovery with explicit execution-count assurances.
The standard unittest CLI exits successfully when a suite contains skipped
tests. That is normally useful, but it let the Docker integration job stay
green while its security-boundary classes were all skipped under act_runner.
This wrapper keeps normal unittest output and adds opt-in minimum-executed and
no-skip gates for jobs that promise a concrete integration surface.
"""
from __future__ import annotations
import argparse
import sys
import unittest
def assurance_errors(
*, tests_run: int, skipped: int, minimum_executed: int, fail_on_skip: bool
) -> list[str]:
"""Return human-readable assurance failures for a completed suite."""
executed = tests_run - skipped
errors: list[str] = []
if executed < minimum_executed:
errors.append(
f"executed {executed} test(s), below required minimum "
f"{minimum_executed} (discovered {tests_run}, skipped {skipped})"
)
if fail_on_skip and skipped:
errors.append(f"{skipped} test(s) skipped in a no-skip suite")
return errors
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
description="unittest discovery with execution-count assurance"
)
parser.add_argument("-s", "--start-directory", default=".")
parser.add_argument("-t", "--top-level-directory", default=None)
parser.add_argument("-p", "--pattern", default="test*.py")
parser.add_argument("--minimum-executed", type=int, default=0)
parser.add_argument("--fail-on-skip", action="store_true")
parser.add_argument("-v", "--verbose", action="store_true")
args = parser.parse_args(argv)
suite = unittest.defaultTestLoader.discover(
args.start_directory,
pattern=args.pattern,
top_level_dir=args.top_level_directory,
)
result = unittest.TextTestRunner(
verbosity=2 if args.verbose else 1,
).run(suite)
failures = assurance_errors(
tests_run=result.testsRun,
skipped=len(result.skipped),
minimum_executed=args.minimum_executed,
fail_on_skip=args.fail_on_skip,
)
for failure in failures:
print(f"unittest-gate: {failure}", file=sys.stderr)
return 0 if result.wasSuccessful() and not failures else 1
if __name__ == "__main__":
raise SystemExit(main())
+12 -8
View File
@@ -20,10 +20,11 @@ tests/
... # many others; see unit/ directory
integration/
test_gateway_image.py
test_dry_run_plan.py
test_sandbox_escape.py
test_orphan_cleanup.py
...
canaries/ # opt-in; see below (currently empty)
canaries/
test_gitleaks_release.py # opt-in upstream artifact check
```
Classification falls out of the directory — no hand-maintained list to
@@ -43,24 +44,27 @@ Discovery is invoked with `-t .` (top-level dir = repo root) so the
## What the integration tests cover
- `test_dry_run_plan.py``cli.py start --dry-run --format=json` emits
a structured plan that contains the resolved egress allowlist and
the bottle's runtime, and creates zero Docker resources.
- `test_orphan_cleanup.py``network_remove` is idempotent against
missing resources, so the EXIT trap can call it unconditionally.
- `test_gateway_image.py` — builds Dockerfile.gateway and
probes that gitleaks / mitmdump / supervise are all reachable
inside the gateway image.
- `test_orchestrator_docker_auth.py` — drives the real control-plane
container and verifies role-scoped authentication.
- `test_multitenant_isolation.py` and `test_sandbox_escape.py` — exercise
token/allowlist separation and end-to-end escape attempts.
## Canaries
`tests/canaries/` holds upstream-regression checks gated on
`BOT_BOTTLE_RUN_CANARIES=1` and not part of the per-push suite.
They're invoked by the scheduled `canaries` workflow. Currently
no canaries are defined.
They're invoked by the scheduled `canaries` workflow. The gitleaks canary
downloads the exact release archive pinned by `Dockerfile.gateway`, verifies
its architecture-specific checksum, and executes the binary.
```bash
BOT_BOTTLE_RUN_CANARIES=1 python -m unittest discover -t . -s tests/canaries -v
BOT_BOTTLE_RUN_CANARIES=1 python -m scripts.unittest_gate \
-t . -s tests/canaries -v --minimum-executed 1 --fail-on-skip
```
## What's NOT covered
+85
View File
@@ -0,0 +1,85 @@
"""Canary: the pinned gitleaks release remains downloadable and executable.
The gateway Dockerfile verifies this archive during an image build. Repeating
the upstream check weekly keeps registry/release drift out of normal pull
requests while proving that the pinned URL, architecture checksum, archive
shape, and binary still agree.
"""
from __future__ import annotations
import hashlib
import os
import platform
import re
import subprocess
import tarfile
import tempfile
import unittest
import urllib.request
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
DOCKERFILE = ROOT / "Dockerfile.gateway"
def _docker_arg(text: str, name: str) -> str:
match = re.search(rf"^ARG {re.escape(name)}=(\S+)$", text, re.MULTILINE)
if match is None:
raise AssertionError(f"Dockerfile.gateway has no concrete ARG {name}")
return match.group(1)
@unittest.skipUnless(
os.environ.get("BOT_BOTTLE_RUN_CANARIES") == "1",
"canary suite is opt-in; set BOT_BOTTLE_RUN_CANARIES=1 to run",
)
class TestGitleaksRelease(unittest.TestCase):
def test_pinned_archive_checksum_and_binary(self) -> None:
dockerfile = DOCKERFILE.read_text(encoding="utf-8")
version = _docker_arg(dockerfile, "GITLEAKS_VERSION")
machine = platform.machine().lower()
architectures = {
"x86_64": ("linux_x64", "GITLEAKS_SHA256_AMD64"),
"amd64": ("linux_x64", "GITLEAKS_SHA256_AMD64"),
"aarch64": ("linux_arm64", "GITLEAKS_SHA256_ARM64"),
"arm64": ("linux_arm64", "GITLEAKS_SHA256_ARM64"),
}
if machine not in architectures:
self.fail(f"unsupported canary runner architecture: {machine}")
asset, checksum_arg = architectures[machine]
expected_checksum = _docker_arg(dockerfile, checksum_arg)
url = (
"https://github.com/gitleaks/gitleaks/releases/download/"
f"v{version}/gitleaks_{version}_{asset}.tar.gz"
)
with tempfile.TemporaryDirectory(prefix="bot-bottle-gitleaks-canary.") as tmp:
archive = Path(tmp) / "gitleaks.tar.gz"
urllib.request.urlretrieve(url, archive)
self.assertEqual(
expected_checksum,
hashlib.sha256(archive.read_bytes()).hexdigest(),
"the pinned upstream archive no longer matches Dockerfile.gateway",
)
with tarfile.open(archive, "r:gz") as bundle:
member = bundle.getmember("gitleaks")
source = bundle.extractfile(member)
if source is None:
self.fail("gitleaks archive member is not a regular file")
binary = Path(tmp) / "gitleaks"
binary.write_bytes(source.read())
binary.chmod(0o755)
result = subprocess.run(
[str(binary), "version"],
capture_output=True,
text=True,
check=False,
)
self.assertEqual(0, result.returncode, result.stderr)
self.assertIn(version, result.stdout + result.stderr)
if __name__ == "__main__":
unittest.main()
+11 -13
View File
@@ -14,9 +14,7 @@ the chunk-1 contract:
expected "no daemons selected" line when the supervisor is
pointed at an empty daemon set.
Skips cleanly when docker is unavailable, or under act_runner
where the host bind-mount topology breaks multi-stage builds
that pull large bases.
Skips cleanly only when the selected Docker backend is unavailable.
"""
from __future__ import annotations
@@ -33,12 +31,6 @@ _DOCKERFILE = "Dockerfile.gateway"
@skip_unless_backend("docker")
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: multi-stage build pulls a 200+MB "
"mitmproxy base + two upstream gateway images; runner storage "
"+ time budget make this an interactive-only test",
)
class TestGatewayImage(unittest.TestCase):
"""Builds the image once for the class, then runs a few
`docker run` probes against it."""
@@ -51,10 +43,11 @@ class TestGatewayImage(unittest.TestCase):
"-f", _DOCKERFILE, "."],
cwd=repo_root,
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
check=False,
)
if proc.returncode != 0:
raise unittest.SkipTest(
f"docker build failed; skipping image probes.\n"
raise AssertionError(
f"docker build failed; image probes cannot run.\n"
f"{proc.stdout.decode('utf-8', errors='replace')[-2000:]}"
)
@@ -63,14 +56,16 @@ class TestGatewayImage(unittest.TestCase):
subprocess.run(
["docker", "image", "rm", "-f", _IMAGE],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
check=False,
)
def _run_in_image(self, *cmd: str, timeout: float = 30.0) -> tuple[int, str]:
proc = subprocess.run(
["docker", "run", "--rm", "--entrypoint", cmd[0], _IMAGE,
*cmd[1:]],
*cmd[1:]],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
timeout=timeout,
check=False,
)
return proc.returncode, proc.stdout.decode("utf-8", errors="replace")
@@ -91,7 +86,9 @@ class TestGatewayImage(unittest.TestCase):
# Probe that the package imports resolve inside the image.
rc, out = self._run_in_image(
"python3", "-c",
"from bot_bottle.supervisor import types; from bot_bottle.gateway.supervisor import server as supervise_server; print('ok')",
"from bot_bottle.supervisor import types; "
"from bot_bottle.gateway.supervisor import server as supervise_server; "
"print('ok')",
)
self.assertEqual(0, rc, msg=out)
self.assertIn("ok", out)
@@ -106,6 +103,7 @@ class TestGatewayImage(unittest.TestCase):
_IMAGE],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
timeout=10.0,
check=False,
)
out = proc.stdout.decode("utf-8", errors="replace")
self.assertEqual(0, proc.returncode, msg=out)
+53 -34
View File
@@ -16,11 +16,10 @@ throwaway BOT_BOTTLE_ROOT for a clean registry and tears everything down.
from __future__ import annotations
import os
import secrets
import subprocess
import tempfile
import time
import unittest
from pathlib import Path
from bot_bottle.backend.docker.consolidated_launch import (
_network_cidr,
@@ -73,19 +72,12 @@ _PROBE_SRC = (
@skip_unless_backend("docker")
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: the orchestrator container bind-mounts the repo "
"path into a container on the socket-shared host daemon, which can't see the "
"runner's /workspace — same host-bind-mount constraint as the other "
"bottle-bringup integration tests",
)
class TestMultitenantIsolation(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
# Throwaway root → a clean registry DB, independent of the host's.
self.svc = DockerInfraService(host_root=Path(self._tmp.name))
# Named volume → a clean registry DB that is also visible to a
# socket-shared host daemon when the test process runs in act_runner.
self._root_volume = "bot-bottle-mtitest-root-" + secrets.token_hex(4)
self.svc = DockerInfraService(root_mount_source=self._root_volume)
self.addCleanup(self._teardown_docker)
# ensure_running builds the bundle image (slow on a cold cache) and
# brings up the shared network + gateway + orchestrator.
@@ -100,13 +92,8 @@ class TestMultitenantIsolation(unittest.TestCase):
self.svc.stop()
subprocess.run(["docker", "network", "rm", GATEWAY_NETWORK],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False)
# The orchestrator container wrote the registry DB as root into the
# throwaway root; chown it back so the (non-root) tempdir cleanup can
# remove it.
subprocess.run(
["docker", "run", "--rm", "-v", f"{self._tmp.name}:/r",
"--entrypoint", "chown", GATEWAY_IMAGE, "-R",
f"{os.getuid()}:{os.getgid()}", "/r"],
["docker", "volume", "rm", "--force", self._root_volume],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False)
@staticmethod
@@ -132,29 +119,61 @@ class TestMultitenantIsolation(unittest.TestCase):
taken = _network_container_ips(GATEWAY_NETWORK) + extra_taken
return next_free_ip(_network_cidr(GATEWAY_NETWORK), taken)
def _probe(self, source_ip: str, host: str) -> str:
proc = subprocess.run(
["docker", "run", "--rm", "--network", GATEWAY_NETWORK, "--ip", source_ip,
"--entrypoint", "python3", GATEWAY_IMAGE, "-c", _PROBE_SRC,
f"http://{self.gw_ip}:{EGRESS_PORT}", host],
stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, check=False, timeout=90,
def _probe(self, source_ip: str, identity_token: str, host: str) -> str:
deadline = time.monotonic() + 30
last = subprocess.CompletedProcess([], 1, "", "probe not attempted")
while time.monotonic() < deadline:
last = subprocess.run(
[
"docker", "run", "--rm",
"--network", GATEWAY_NETWORK, "--ip", source_ip,
"--entrypoint", "python3", GATEWAY_IMAGE, "-c", _PROBE_SRC,
f"http://bottle:{identity_token}@{self.gw_ip}:{EGRESS_PORT}",
host,
],
stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True,
check=False, timeout=90,
)
output = last.stdout.strip()
if last.returncode == 0 and output:
return output
time.sleep(0.25)
self.fail(
f"gateway probe did not become ready: "
f"exit={last.returncode}, stderr={last.stderr.strip()!r}"
)
return proc.stdout.strip()
def test_two_bottles_share_gateway_with_isolated_tokens_and_allowlists(self) -> None:
ip_a = self._free_ip([])
ip_b = self._free_ip([ip_a])
self.client.register_bottle(ip_a, policy=_POLICY_A, tokens={"EGRESS_TOKEN_0": _TOKEN_A})
self.client.register_bottle(ip_b, policy=_POLICY_B, tokens={"EGRESS_TOKEN_0": _TOKEN_B})
bottle_a = self.client.register_bottle(
ip_a, policy=_POLICY_A, tokens={"EGRESS_TOKEN_0": _TOKEN_A}
)
bottle_b = self.client.register_bottle(
ip_b, policy=_POLICY_B, tokens={"EGRESS_TOKEN_0": _TOKEN_B}
)
# Each bottle gets its OWN token injected on the shared route — no bleed.
self.assertEqual(f"200 AUTH=Bearer {_TOKEN_A}", self._probe(ip_a, "echo-shared"))
self.assertEqual(f"200 AUTH=Bearer {_TOKEN_B}", self._probe(ip_b, "echo-shared"))
self.assertEqual(
f"200 AUTH=Bearer {_TOKEN_A}",
self._probe(ip_a, bottle_a.identity_token, "echo-shared"),
)
self.assertEqual(
f"200 AUTH=Bearer {_TOKEN_B}",
self._probe(ip_b, bottle_b.identity_token, "echo-shared"),
)
# Allowlist is per-bottle: echo-bonly is only in B's policy.
self.assertTrue(self._probe(ip_a, "echo-bonly").startswith("403"), # fail-closed for A
"A reached a host outside its allowlist")
self.assertEqual("200 AUTH=NONE", self._probe(ip_b, "echo-bonly")) # allowed, unauthed for B
self.assertTrue(
self._probe(
ip_a, bottle_a.identity_token, "echo-bonly"
).startswith("403"), # fail-closed for A
"A reached a host outside its allowlist",
)
self.assertEqual(
"200 AUTH=NONE",
self._probe(ip_b, bottle_b.identity_token, "echo-bonly"),
) # allowed, unauthed for B
if __name__ == "__main__":
@@ -23,7 +23,6 @@ import secrets
import subprocess
import tempfile
import unittest
from pathlib import Path
from bot_bottle.orchestrator_auth import ROLE_CLI, ROLE_GATEWAY, mint
from bot_bottle.orchestrator.client import OrchestratorClient
@@ -38,13 +37,6 @@ _TEST_GATEWAY_IMAGE = "bot-bottle-gateway:itest"
@skip_unless_backend("docker")
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: the orchestrator container bind-mounts the repo "
"path into a container on the socket-shared host daemon, which can't see the "
"runner's /workspace — same host-bind-mount constraint as the other "
"bottle-bringup integration tests",
)
class TestDockerOrchestratorAuthIntegration(unittest.TestCase):
@classmethod
def setUpClass(cls) -> None:
@@ -74,10 +66,10 @@ class TestDockerOrchestratorAuthIntegration(unittest.TestCase):
gateway_name = f"bot-bottle-gateway-itest-{suffix}"
network = f"bot-bottle-net-itest-{suffix}"
control_network = f"bot-bottle-ctrl-itest-{suffix}"
host_root = Path(cls._tmp.name)
root_volume = f"bot-bottle-root-itest-{suffix}"
cls.addClassCleanup(
cls._teardown_docker,
orchestrator_name, gateway_name, network, control_network, host_root,
orchestrator_name, gateway_name, network, control_network, root_volume,
)
cls.svc = DockerInfraService(
@@ -88,7 +80,7 @@ class TestDockerOrchestratorAuthIntegration(unittest.TestCase):
orchestrator_image=_TEST_ORCHESTRATOR_IMAGE,
gateway_image=_TEST_GATEWAY_IMAGE,
port=20000 + secrets.randbelow(10000),
host_root=host_root,
root_mount_source=root_volume,
)
cls.svc.ensure_running()
# The control plane now verifies role-scoped signed tokens, not the raw
@@ -100,7 +92,7 @@ class TestDockerOrchestratorAuthIntegration(unittest.TestCase):
@staticmethod
def _teardown_docker(
orchestrator_name: str, gateway_name: str,
network: str, control_network: str, host_root: Path,
network: str, control_network: str, root_volume: str,
) -> None:
for name in (gateway_name, orchestrator_name):
subprocess.run(
@@ -112,14 +104,8 @@ class TestDockerOrchestratorAuthIntegration(unittest.TestCase):
["docker", "network", "rm", net],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False,
)
# 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_ORCHESTRATOR_IMAGE, "-R",
f"{os.getuid()}:{os.getgid()}", "/r"],
["docker", "volume", "rm", "--force", root_volume],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False,
)
-5
View File
@@ -42,11 +42,6 @@ class TestOrphanCleanup(unittest.TestCase):
# Returning True == idempotent success.
self.assertTrue(network_remove(f"bot-bottle-net-{self.slug}-does-not-exist"))
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true",
"skipped under act_runner: docker socket mount topology breaks "
"in-process visibility of networks created on the host daemon",
)
def test_create_and_remove(self):
self.internal_name = network_create_internal(self.slug)
self.egress_name = network_create_egress(self.slug)
+1 -20
View File
@@ -67,26 +67,7 @@ _DUMMY_HOST_KEY = (
)
# Backends whose CI runner is HOST-mode (self-hosted), so the test process
# and the backend share a host. The containerized act_runner (docker on
# ubuntu-latest) is the one that can't see the host bind mount egress_tls_init
# uses and hides sibling-gateway network topology; host-mode runners
# (firecracker/KVM, macos-container) don't have those constraints, so the test
# runs there. Keep this in sync with the `runs-on` labels in
# .gitea/workflows/test.yml.
_HOST_MODE_CI_BACKENDS = frozenset({"firecracker", "macos-container"})
@skip_unless_selected_backend_available()
@unittest.skipIf(
os.environ.get("GITEA_ACTIONS") == "true"
and os.environ.get("BOT_BOTTLE_BACKEND") not in _HOST_MODE_CI_BACKENDS,
"skipped under the containerized act_runner (docker on ubuntu-latest): "
"egress_tls_init uses a host bind mount the runner container can't "
"see, and the network topology hides sibling-gateway visibility — "
"these constraints don't apply on the self-hosted host-mode runners "
"(firecracker/KVM, macos-container)",
)
class TestSandboxEscape(unittest.TestCase):
"""End-to-end attacks against a real bottle. The bottle stays
up for the whole class bringup is ~10-30s, so per-test
@@ -189,7 +170,7 @@ class TestSandboxEscape(unittest.TestCase):
missing.append(tool)
if missing:
cls._teardown_resources()
raise unittest.SkipTest(
raise AssertionError(
f"agent missing required tools: {', '.join(missing)}"
f"add them to the backend's base image"
)
@@ -0,0 +1,96 @@
"""Architecture rules that should fail before coupling becomes entrenched."""
from __future__ import annotations
import ast
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
class TestCliBackendBoundaries(unittest.TestCase):
def test_cli_does_not_import_a_concrete_backend(self) -> None:
forbidden = (
"backend.docker", "backend.firecracker", "backend.macos_container",
"bot_bottle.backend.docker", "bot_bottle.backend.firecracker",
"bot_bottle.backend.macos_container",
)
violations: list[str] = []
for path in (ROOT / "bot_bottle" / "cli").rglob("*.py"):
tree = ast.parse(path.read_text(), filename=str(path))
for node in ast.walk(tree):
if isinstance(node, ast.ImportFrom):
module = node.module
if module and module.startswith(forbidden):
violations.append(
f"{path.relative_to(ROOT)}:{node.lineno}: {module}"
)
if isinstance(node, ast.Import):
violations.extend(
f"{path.relative_to(ROOT)}:{node.lineno}: {alias.name}"
for alias in node.names if alias.name.startswith(forbidden)
)
self.assertEqual([], violations, "generic CLI imports concrete backend internals:\n" +
"\n".join(violations))
class TestRuntimeModuleSizes(unittest.TestCase):
def test_no_runtime_module_grows_beyond_global_ceiling(self) -> None:
"""A coarse ceiling catches new monoliths; focused caps stay tighter."""
ceiling = 850
oversized = [
f"{path.relative_to(ROOT)} ({len(path.read_text().splitlines())})"
for path in (ROOT / "bot_bottle").rglob("*.py")
if len(path.read_text().splitlines()) > ceiling
]
self.assertEqual(
[], oversized,
f"runtime modules must stay at or below {ceiling} lines: "
+ ", ".join(oversized),
)
def test_egress_modules_stay_focused(self) -> None:
caps = {
"addon_core.py": 100,
"schema.py": 400,
"types.py": 180,
"matching.py": 180,
"dlp.py": 180,
"context.py": 140,
}
directory = ROOT / "bot_bottle" / "gateway" / "egress"
oversized = [f"{name} ({len((directory / name).read_text().splitlines())}>{cap})"
for name, cap in caps.items()
if len((directory / name).read_text().splitlines()) > cap]
self.assertEqual([], oversized, "split a module rather than raising its cap: " +
", ".join(oversized))
def test_runtime_code_uses_focused_egress_modules(self) -> None:
"""addon_core is compatibility-only, never an internal dependency."""
violations: list[str] = []
package = ROOT / "bot_bottle"
facade = package / "gateway" / "egress" / "addon_core.py"
package_init = package / "gateway" / "egress" / "__init__.py"
for path in package.rglob("*.py"):
if path in (facade, package_init):
continue
text = path.read_text()
if "gateway.egress.addon_core import" in text or \
".addon_core import" in text:
violations.append(str(path.relative_to(ROOT)))
self.assertEqual([], violations)
def test_backend_contract_does_not_absorb_preparation_logic(self) -> None:
caps = {
ROOT / "bot_bottle" / "backend" / "base.py": 580,
ROOT / "bot_bottle" / "backend" / "preparation.py": 160,
}
oversized = [
f"{path.relative_to(ROOT)} "
f"({len(path.read_text().splitlines())}>{cap})"
for path, cap in caps.items()
if len(path.read_text().splitlines()) > cap
]
self.assertEqual([], oversized)
@@ -48,12 +48,13 @@ class TestSharedReprovision(unittest.TestCase):
client.reprovision_gateway.side_effect = [
OrchestratorClientError("bad key"), True,
]
self.assertEqual(
1,
reprovision_bottles(
with patch("bot_bottle.orchestrator.reprovision.debug") as debug:
count = reprovision_bottles(
client, {"10.0.0.1": "key-1", "10.0.0.2": "key-2"},
),
)
)
self.assertEqual(1, count)
self.assertEqual("b1", debug.call_args.kwargs["context"]["bottle_id"])
self.assertNotIn("bad key", repr(debug.call_args))
class TestMacosReprovision(unittest.TestCase):
+7 -2
View File
@@ -9,6 +9,7 @@ from __future__ import annotations
import unittest
from typing import Any, Optional
from unittest.mock import patch
from bot_bottle.cli.tui import _filter_items, _multiselect_loop, filter_multiselect, filter_select
@@ -49,8 +50,10 @@ class TestFilterSelectEmptyItems(unittest.TestCase):
def test_returns_none_when_tty_unavailable(self):
# /nonexistent is guaranteed to not open.
result = filter_select(["a", "b"], tty_path="/nonexistent/tty")
with patch("bot_bottle.cli.tui.debug") as debug:
result = filter_select(["a", "b"], tty_path="/nonexistent/tty")
self.assertIsNone(result)
self.assertEqual("FileNotFoundError", debug.call_args.kwargs["context"]["error_type"])
class TestFilterMultiselectEmptyItems(unittest.TestCase):
@@ -60,8 +63,10 @@ class TestFilterMultiselectEmptyItems(unittest.TestCase):
self.assertEqual([], result)
def test_returns_none_when_tty_unavailable(self):
result = filter_multiselect(["a", "b"], tty_path="/nonexistent/tty")
with patch("bot_bottle.cli.tui.debug") as debug:
result = filter_multiselect(["a", "b"], tty_path="/nonexistent/tty")
self.assertIsNone(result)
self.assertEqual("FileNotFoundError", debug.call_args.kwargs["context"]["error_type"])
class TestMultiselectLoopReordering(unittest.TestCase):
+101
View File
@@ -0,0 +1,101 @@
"""Tests for the fail-closed critical coverage manifest."""
from __future__ import annotations
import tempfile
import unittest
from contextlib import redirect_stderr, redirect_stdout
from io import StringIO
from pathlib import Path
from scripts.critical_modules import (
DEFAULT_MANIFEST,
REPO_ROOT,
CriticalModulesError,
load_critical_modules,
main,
)
class TestCriticalModules(unittest.TestCase):
def test_repository_manifest_is_valid(self) -> None:
modules = load_critical_modules(DEFAULT_MANIFEST, root=REPO_ROOT)
self.assertGreater(len(modules), 20)
self.assertEqual(len(modules), len(set(modules)))
def test_missing_module_fails(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
manifest = root / "critical-modules.txt"
manifest.write_text("bot_bottle/renamed.py\n", encoding="utf-8")
with self.assertRaisesRegex(CriticalModulesError, "does not exist"):
load_critical_modules(manifest, root=root)
def test_duplicate_module_fails(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
module = root / "bot_bottle" / "core.py"
module.parent.mkdir()
module.write_text("", encoding="utf-8")
manifest = root / "critical-modules.txt"
manifest.write_text(
"bot_bottle/core.py\nbot_bottle/core.py\n", encoding="utf-8"
)
with self.assertRaisesRegex(CriticalModulesError, "duplicated"):
load_critical_modules(manifest, root=root)
def test_entry_cannot_escape_repository(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
manifest = root / "critical-modules.txt"
manifest.write_text("../outside.py\n", encoding="utf-8")
with self.assertRaisesRegex(CriticalModulesError, "escapes"):
load_critical_modules(manifest, root=root)
def test_invalid_entry_forms_are_reported_together(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
manifest = root / "critical-modules.txt"
manifest.write_text(
f"{root / 'absolute.py'}\nREADME.md\n",
encoding="utf-8",
)
with self.assertRaises(CriticalModulesError) as raised:
load_critical_modules(manifest, root=root)
self.assertIn("must be relative", str(raised.exception))
self.assertIn("is not a Python module", str(raised.exception))
def test_empty_manifest_fails(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
manifest = root / "critical-modules.txt"
manifest.write_text("# comments do not define modules\n", encoding="utf-8")
with self.assertRaisesRegex(CriticalModulesError, "contains no"):
load_critical_modules(manifest, root=root)
def test_unreadable_manifest_fails(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
missing = Path(tmp) / "missing.txt"
with self.assertRaisesRegex(CriticalModulesError, "cannot read"):
load_critical_modules(missing, root=Path(tmp))
def test_main_prints_include_list_or_checks_silently(self) -> None:
output = StringIO()
with redirect_stdout(output):
self.assertEqual(0, main([]))
self.assertIn("bot_bottle/manifest/egress.py", output.getvalue())
output = StringIO()
with redirect_stdout(output):
self.assertEqual(0, main(["--check"]))
self.assertEqual("", output.getvalue())
def test_main_reports_manifest_error(self) -> None:
error = StringIO()
with redirect_stderr(error):
self.assertEqual(1, main(["--manifest", "/definitely/missing"]))
self.assertIn("critical-modules:", error.getvalue())
if __name__ == "__main__":
unittest.main()
+16
View File
@@ -9,6 +9,7 @@ a freshly minted token."""
from __future__ import annotations
import unittest
from pathlib import Path
from unittest.mock import MagicMock, patch
from bot_bottle.backend.docker.infra import (
@@ -67,6 +68,21 @@ class TestDockerInfraService(unittest.TestCase):
self.assertTrue(any(GATEWAY_NAME in a for a in rms))
self.assertTrue(any(ORCHESTRATOR_NAME in a for a in rms))
def test_named_mounts_propagate_to_both_planes(self) -> None:
svc = DockerInfraService(
root_mount_source="registry-volume",
gateway_ca_mount_source="ca-volume",
)
self.assertEqual("registry-volume", svc.orchestrator()._root_mount_source)
self.assertEqual("ca-volume", svc.gateway()._ca_mount_source)
def test_host_path_and_named_root_mount_are_mutually_exclusive(self) -> None:
with self.assertRaisesRegex(ValueError, "host_root or root_mount_source"):
DockerInfraService(
host_root=Path("/host/path"),
root_mount_source="registry-volume",
)
if __name__ == "__main__":
unittest.main()
+52
View File
@@ -4,6 +4,7 @@ from __future__ import annotations
import unittest
import urllib.error
from pathlib import Path
from unittest.mock import MagicMock, Mock, patch
from bot_bottle.backend.docker.orchestrator import (
@@ -47,6 +48,55 @@ class TestDockerOrchestrator(unittest.TestCase):
def test_url_is_host_loopback(self) -> None:
self.assertEqual("http://127.0.0.1:8099", self.orch.url())
def test_socket_shared_client_uses_explicit_host_and_open_bind(self) -> None:
orch = DockerOrchestrator(
port=8099, client_host="172.17.0.1", root_mount_source="state-volume"
)
with patch(_TOKEN, return_value="k"), \
patch(_URLOPEN, side_effect=[urllib.error.URLError("down"), _health(200)]), \
patch(_RUN, return_value=_proc()) as run, patch(_SLEEP):
orch.ensure_running()
self.assertEqual("http://172.17.0.1:8099", orch.url())
argv = next(c.args[0] for c in run.call_args_list
if c.args[0][:2] == ["docker", "run"])
self.assertEqual("0.0.0.0:8099:8099", argv[argv.index("--publish") + 1])
self.assertIn("state-volume:/bot-bottle-root", argv)
def test_host_path_and_named_root_mount_are_mutually_exclusive(self) -> None:
with self.assertRaisesRegex(ValueError, "host_root or root_mount_source"):
DockerOrchestrator(
host_root=Path("/host/path"),
root_mount_source="state-volume",
)
def test_socket_shared_job_network_uses_container_dns(self) -> None:
orch = DockerOrchestrator(
name="orchestrator-itest",
port=22001,
client_network="runner-job-network",
root_mount_source="state-volume",
)
with patch(_TOKEN, return_value="k"), \
patch(
_URLOPEN,
side_effect=[urllib.error.URLError("down"), _health(200)],
), patch(_RUN, return_value=_proc()) as run, patch(_SLEEP):
orch.ensure_running()
self.assertEqual("http://orchestrator-itest:8099", orch.url())
calls = [call.args[0] for call in run.call_args_list]
self.assertIn(
[
"docker", "network", "connect",
"runner-job-network", "orchestrator-itest",
],
calls,
)
argv = next(call for call in calls if call[:2] == ["docker", "run"])
self.assertEqual(
"127.0.0.1:22001:8099",
argv[argv.index("--publish") + 1],
)
def test_gateway_url_is_the_container_dns_name(self) -> None:
# The gateway reaches the orchestrator by name on the control network.
self.assertEqual(f"http://{ORCHESTRATOR_NAME}:8099", self.orch.gateway_url())
@@ -110,6 +160,8 @@ class TestDockerOrchestrator(unittest.TestCase):
# The lean control plane: no mitmproxy CA mount, no gateway daemons.
self.assertFalse([a for a in argv if a.endswith(":/home/mitmproxy/.mitmproxy")])
self.assertNotIn("BOT_BOTTLE_GATEWAY_DAEMONS", " ".join(argv))
self.assertNotIn("/bot-bottle-src", " ".join(argv))
self.assertNotIn("PYTHONPATH", " ".join(argv))
# Orchestrator entrypoint args (image ENTRYPOINT is `-m bot_bottle.orchestrator`).
self.assertIn("--broker", argv)
self.assertIn("stub", argv)
+8 -6
View File
@@ -8,17 +8,19 @@ from __future__ import annotations
import unittest
from bot_bottle.gateway.egress.addon_core import (
HeaderMatch,
MatchEntry,
PathMatch,
Route,
evaluate_matches,
from bot_bottle.gateway.egress.matching import evaluate_matches
from bot_bottle.gateway.egress.schema import (
load_config,
parse_config,
parse_routes,
route_to_yaml_dict,
)
from bot_bottle.gateway.egress.types import (
HeaderMatch,
MatchEntry,
PathMatch,
Route,
)
def _route(d: dict[str, object]) -> Route:
+10 -3
View File
@@ -3,15 +3,16 @@
from __future__ import annotations
import unittest
from unittest.mock import patch
from bot_bottle.gateway.egress.addon_core import (
from bot_bottle.gateway.egress.context import (
DENY_RESOLVER_ERROR,
DENY_UNATTRIBUTED,
DENY_UNPARSEABLE,
decide,
resolve_client_config,
resolve_client_context,
)
from bot_bottle.gateway.egress.matching import decide
from bot_bottle.gateway.policy_resolver import PolicyResolveError
@@ -44,7 +45,13 @@ class TestResolveClientConfig(unittest.TestCase):
def test_resolver_error_denies_all(self) -> None:
# Orchestrator unreachable/errored must never widen egress.
self.assertEqual((), resolve_client_config(_FakeResolver(raises=True), "10.243.0.1").routes)
with patch("bot_bottle.gateway.egress.context.debug") as debug:
config = resolve_client_config(_FakeResolver(raises=True), "10.243.0.1")
self.assertEqual((), config.routes)
self.assertEqual(
"PolicyResolveError", debug.call_args.kwargs["context"]["error_type"],
)
self.assertNotIn("orchestrator down", repr(debug.call_args))
def test_unparseable_policy_denies_all(self) -> None:
cfg = resolve_client_config(_FakeResolver(result="routes: notalist\n"), "10.243.0.1")
+11
View File
@@ -12,7 +12,9 @@ from bot_bottle.orchestrator.client import (
OrchestratorClient,
OrchestratorClientError,
RegisteredBottle,
BackendProbeFailure,
_host_auth_token,
_probe_failure,
)
_URLOPEN = "bot_bottle.orchestrator.client.urllib.request.urlopen"
@@ -33,6 +35,15 @@ class TestHostAuthToken(unittest.TestCase):
self.assertEqual("", _host_auth_token())
class TestBackendProbeFailure(unittest.TestCase):
def test_records_safe_typed_diagnostic(self) -> None:
with patch("bot_bottle.orchestrator.client.debug") as debug:
result = _probe_failure("firecracker", RuntimeError("secret detail"))
self.assertEqual(BackendProbeFailure("firecracker", "RuntimeError"), result)
rendered = repr(debug.call_args)
self.assertNotIn("secret detail", rendered)
def _resp(status: int, payload: object) -> MagicMock:
m = MagicMock()
inner = m.__enter__.return_value
+52 -2
View File
@@ -7,7 +7,10 @@ import unittest
from pathlib import Path
from unittest.mock import Mock, patch
from bot_bottle.backend.docker.gateway import DockerGateway
from bot_bottle.backend.docker.gateway import (
DEFAULT_GATEWAY_SUBNET,
DockerGateway,
)
from bot_bottle.gateway import (
GATEWAY_CA_CERT,
GATEWAY_NAME,
@@ -142,6 +145,22 @@ class TestDockerGateway(unittest.TestCase):
# Data plane resolves policy against the orchestrator control plane.
self.assertIn(f"BOT_BOTTLE_ORCHESTRATOR_URL={_ORCH_URL}", runs[0])
def test_named_ca_volume_supports_socket_shared_runner(self) -> None:
sc = DockerGateway(
"bot-bottle-gateway:latest", ca_mount_source="ci-ca-volume"
)
def fake(argv: list[str], **_kw: object) -> Mock:
return _proc(stdout="") if argv[:2] == ["docker", "ps"] else _proc()
with patch(_RUN_DOCKER, side_effect=fake) as run:
sc.connect_to_orchestrator(_ORCH_URL, _TOKEN)
argv = next(c.args[0] for c in run.call_args_list
if c.args[0][:2] == ["docker", "run"])
self.assertIn(
"ci-ca-volume:/home/mitmproxy/.mitmproxy", argv
)
def test_connect_injects_the_pre_minted_gateway_token(self) -> None:
# The gateway presents the token the orchestrator handed it — it never
# mints (holds no signing key). The value rides the env (bare `--env
@@ -193,7 +212,38 @@ class TestDockerGateway(unittest.TestCase):
with patch(_RUN_DOCKER, side_effect=fake):
self.sc.connect_to_orchestrator(_ORCH_URL, _TOKEN)
creates = [c for c in calls if c[:3] == ["docker", "network", "create"]]
self.assertEqual([["docker", "network", "create", self.sc.network]], creates)
self.assertEqual(
[[
"docker", "network", "create",
"--subnet", DEFAULT_GATEWAY_SUBNET,
"--label",
f"bot-bottle.gateway-subnet={DEFAULT_GATEWAY_SUBNET}",
self.sc.network,
]],
creates,
)
def test_ensure_running_replaces_stale_auto_ipam_network(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", "network", "inspect"]:
return _proc(stdout="<no value>\n")
return _proc()
with patch(_RUN_DOCKER, side_effect=fake):
self.sc.connect_to_orchestrator(_ORCH_URL, _TOKEN)
self.assertIn(
["docker", "rm", "--force", self.sc.name],
calls,
)
self.assertIn(
["docker", "network", "rm", self.sc.network],
calls,
)
def test_ca_cert_pem_reads_from_container(self) -> None:
with patch(_RUN_DOCKER, return_value=_proc(stdout=_CA_PEM)) as m:
+39 -5
View File
@@ -7,6 +7,7 @@ server tests), plus one real-socket round-trip to prove the handler wiring.
from __future__ import annotations
import base64
import io
import json
import secrets
import sqlite3
@@ -17,7 +18,7 @@ import urllib.error
import urllib.request
from contextlib import closing
from pathlib import Path
from unittest.mock import patch
from unittest.mock import MagicMock, patch
from bot_bottle.orchestrator_auth import ROLE_CLI, ROLE_GATEWAY, mint
from bot_bottle.orchestrator.broker import StubBroker
@@ -283,6 +284,25 @@ class TestServerRoundTrip(unittest.TestCase):
))
self.assertEqual(reg["bottle_id"], attr["bottle_id"])
def test_internal_failure_is_contextual_but_redacted(self) -> None:
orch = MagicMock()
orch.registry.all.side_effect = RuntimeError("SENSITIVE request value")
with patch("sys.stderr", io.StringIO()) as stderr:
server = make_server(orch, "127.0.0.1", 0)
self.addCleanup(server.server_close)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
self.addCleanup(server.shutdown)
host, port = server.server_address[0], server.server_address[1]
with self.assertRaises(urllib.error.HTTPError) as raised:
urllib.request.urlopen(f"http://{host}:{port}/bottles", timeout=5)
payload = json.loads(raised.exception.read())
output = stderr.getvalue()
self.assertEqual({"error": "internal error"}, payload)
self.assertIn("GET /bottles", output)
self.assertIn("RuntimeError", output)
self.assertNotIn("SENSITIVE", output)
class TestOrchestratorAuth(unittest.TestCase):
"""Role-scoped control-plane tokens (issue #400 / #469 review): every route
@@ -647,10 +667,24 @@ class TestReconcileRoute(unittest.TestCase):
self.assertEqual(200, status)
self.assertEqual([], payload["reaped"])
def test_non_string_entries_are_ignored(self) -> None:
dead = self._old("10.0.0.4")
def test_non_string_entries_are_rejected(self) -> None:
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"])
self.assertEqual(400, status)
self.assertIn("live_source_ips", str(payload["error"]))
def test_empty_live_source_ip_is_rejected(self) -> None:
status, payload = dispatch(
self.orch, "POST", "/reconcile", _body({"live_source_ips": [""]}))
self.assertEqual(400, status)
self.assertIn("live_source_ips", str(payload["error"]))
def test_invalid_grace_seconds_is_rejected(self) -> None:
for value in (True, "30", -1, float("inf"), float("nan")):
with self.subTest(value=value):
status, payload = dispatch(
self.orch, "POST", "/reconcile",
_body({"live_source_ips": [], "grace_seconds": value}))
self.assertEqual(400, status)
self.assertIn("grace_seconds", str(payload["error"]))
+13 -5
View File
@@ -105,11 +105,19 @@ class TestOrchestrator(unittest.TestCase):
def test_reprovision_rejects_missing_rows_and_wrong_key(self) -> None:
self.assertFalse(self.orch.reprovision_from_secret("missing", new_env_var_secret()))
rec = self.orch.launch_bottle(
"10.243.0.13", tokens={"K": "value"},
env_var_secret=new_env_var_secret(),
)
self.assertFalse(self.orch.reprovision_from_secret(rec.bottle_id, new_env_var_secret()))
key = "AQEBAQEBAQEBAQEBAQEBAQEBAQEBAQEBAQEBAQEBAQE"
wrong_key = "FBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQUFBQ"
# Pin the nonce so this is a deterministic wrong-key/decryption vector
# instead of a probabilistic assertion over random bytes.
with patch(
"bot_bottle.orchestrator.store.secret_store.secrets.token_bytes",
return_value=b"\0" * 16,
):
rec = self.orch.launch_bottle(
"10.243.0.13", tokens={"K": "value"},
env_var_secret=key,
)
self.assertFalse(self.orch.reprovision_from_secret(rec.bottle_id, wrong_key))
def test_set_policy_live_reload(self) -> None:
rec = self.orch.launch_bottle("10.243.0.3")
+91
View File
@@ -0,0 +1,91 @@
"""Unit tests for CI's unittest execution-count gate."""
from __future__ import annotations
import unittest
from contextlib import redirect_stderr
from io import StringIO
from unittest.mock import Mock, patch
from scripts.unittest_gate import assurance_errors, main
class TestAssuranceErrors(unittest.TestCase):
def test_accepts_suite_that_meets_minimum_without_skips(self) -> None:
self.assertEqual(
[],
assurance_errors(
tests_run=22, skipped=0, minimum_executed=22, fail_on_skip=True
),
)
def test_rejects_green_suite_below_execution_minimum(self) -> None:
errors = assurance_errors(
tests_run=22, skipped=18, minimum_executed=22, fail_on_skip=False
)
self.assertEqual(1, len(errors))
self.assertIn("executed 4", errors[0])
def test_rejects_any_skip_when_required(self) -> None:
errors = assurance_errors(
tests_run=23, skipped=1, minimum_executed=22, fail_on_skip=True
)
self.assertEqual(["1 test(s) skipped in a no-skip suite"], errors)
def test_main_accepts_successful_assured_suite(self) -> None:
result = Mock(
testsRun=22,
skipped=[],
wasSuccessful=Mock(return_value=True),
)
runner = Mock()
runner.run.return_value = result
with patch(
"scripts.unittest_gate.unittest.defaultTestLoader.discover",
return_value=Mock(),
) as discover, patch(
"scripts.unittest_gate.unittest.TextTestRunner",
return_value=runner,
) as runner_type:
self.assertEqual(
0,
main([
"-s", "tests/integration",
"-t", ".",
"-p", "test_*.py",
"--minimum-executed", "22",
"--fail-on-skip",
"-v",
]),
)
discover.assert_called_once_with(
"tests/integration", pattern="test_*.py", top_level_dir="."
)
runner_type.assert_called_once_with(verbosity=2)
def test_main_rejects_unsuccessful_underfilled_suite(self) -> None:
result = Mock(
testsRun=1,
skipped=[(Mock(), "not available")],
wasSuccessful=Mock(return_value=False),
)
runner = Mock()
runner.run.return_value = result
error = StringIO()
with patch(
"scripts.unittest_gate.unittest.defaultTestLoader.discover",
return_value=Mock(),
), patch(
"scripts.unittest_gate.unittest.TextTestRunner",
return_value=runner,
), redirect_stderr(error):
self.assertEqual(
1,
main(["--minimum-executed", "2", "--fail-on-skip"]),
)
self.assertIn("below required minimum", error.getvalue())
self.assertIn("skipped in a no-skip suite", error.getvalue())
if __name__ == "__main__":
unittest.main()