Compare commits

..

2 Commits

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

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

Closes #446
2026-07-21 04:04:30 +00:00
24 changed files with 255 additions and 1253 deletions
+4
View File
@@ -1,6 +1,10 @@
[run]
branch = True
source = .
# Store paths relative to the project root so .coverage.* files produced on
# different runners (ubuntu-latest vs self-hosted KVM) can be combined by the
# coverage job without a [paths] remapping section.
relative_files = True
[report]
# Coverage policy: see docs/decisions/0004-coverage-policy.md.
+99 -112
View File
@@ -9,10 +9,12 @@
# tests/canaries/ — upstream regression canaries; run on a separate
# schedule (see canaries.yml), not here
#
# Integration tests run once per backend in separate jobs. Each job sets
# BOT_BOTTLE_BACKEND explicitly so the test suite uses the right backend.
# Backends that aren't available on the runner fail the preflight step
# rather than silently skipping inside the test output.
# Each test job runs once under coverage and uploads a small .coverage.*
# artifact. The `coverage` job combines them — no test reruns, no KVM
# dependency on that job. For main-branch pushes only, the tested rootfs
# and matching dropbear are uploaded so `publish-infra` can publish the
# byte-identical artifact that was tested. PRs avoid the ~194 MB rootfs
# transfer entirely.
name: test
@@ -40,53 +42,6 @@ on:
workflow_dispatch:
jobs:
stage-firecracker-inputs:
runs-on: [self-hosted, kvm]
# Same guard as the other KVM-runner jobs: don't spin the privileged
# runner for fork PRs (this only copies a non-secret static binary, but
# keep the posture consistent — build-infra/integration/coverage all
# depend on it, so gating here gates the whole Firecracker chain).
if: >-
github.event_name == 'push' ||
github.event_name == 'workflow_dispatch' ||
(github.event_name == 'pull_request' &&
github.event.pull_request.head.repo.full_name == github.repository)
steps:
- name: Stage the provisioned static dropbear
run: |
mkdir -p firecracker-inputs
cp /var/cache/bot-bottle-fc/dropbear firecracker-inputs/dropbear
- name: Upload Firecracker build inputs
uses: actions/upload-artifact@v3
with:
name: firecracker-inputs
path: firecracker-inputs/
build-infra:
needs: stage-firecracker-inputs
runs-on: ubuntu-latest
steps:
- name: Checkout
uses: actions/checkout@v4
- name: Download Firecracker build inputs
uses: actions/download-artifact@v3
with:
name: firecracker-inputs
path: firecracker-inputs
- name: Build infra candidate from this checkout
env:
BOT_BOTTLE_FC_DROPBEAR: ${{ github.workspace }}/firecracker-inputs/dropbear
run: python3 -m bot_bottle.backend.firecracker.publish_infra --output infra-candidate
- name: Upload infra candidate
uses: actions/upload-artifact@v3
with:
name: infra-candidate
path: infra-candidate/
unit:
runs-on: ubuntu-latest
steps:
@@ -101,11 +56,17 @@ jobs:
- name: Install dev requirements
run: python3 -m pip install --break-system-packages -r requirements-dev.txt
- name: Run unit tests
run: python3 -m coverage run -m unittest discover -t . -s tests/unit -v
- name: Run unit tests with coverage
run: python3 -m coverage run --data-file=.coverage.unit -m unittest discover -t . -s tests/unit -v
- name: Report unit coverage
run: python3 -m coverage report -m
run: python3 -m coverage report --data-file=.coverage.unit -m
- name: Upload unit coverage artifact
uses: actions/upload-artifact@v3
with:
name: coverage-unit
path: ${{ github.workspace }}/.coverage.unit
integration-docker:
runs-on: ubuntu-latest
@@ -115,6 +76,9 @@ jobs:
# No actions/setup-python (see the note in the `unit` job); the
# container's system Python 3.12 runs the stdlib test suite directly.
- name: Install coverage
run: python3 -m pip install --break-system-packages coverage
- name: Show environment
run: |
python3 --version
@@ -124,10 +88,16 @@ jobs:
echo "docker not on PATH — integration tests will skip"
fi
- name: Run integration tests (docker)
- name: Run integration tests (docker) with coverage
env:
BOT_BOTTLE_BACKEND: docker
run: python3 -m unittest discover -t . -s tests/integration -v
run: python3 -m coverage run --data-file=.coverage.docker -m unittest discover -t . -s tests/integration -v
- name: Upload docker coverage artifact
uses: actions/upload-artifact@v3
with:
name: coverage-docker
path: ${{ github.workspace }}/.coverage.docker
# Integration tests against the Firecracker backend. Runs on a self-hosted
# KVM runner (label `kvm`) where /dev/kvm and the TAP/nft pool are available.
@@ -137,9 +107,16 @@ jobs:
#
# Runner prerequisites (provision once; see README "Firecracker on Linux"):
# `firecracker` on PATH, `/dev/kvm` accessible, cached kernel +
# static dropbear, and the pool as a persistent systemd unit.
# static dropbear at /var/cache/bot-bottle-fc/dropbear, and the pool as a
# persistent systemd unit.
#
# The infra candidate is built here directly (no artifact download) to
# eliminate the ~70 s ubuntu-latest upload + ~83 s combined download that
# the old build-infra → integration-firecracker + coverage chain incurred.
# For main-branch pushes the tested rootfs and matching dropbear are
# uploaded so publish-infra can publish the byte-identical artifact; PRs
# skip those uploads entirely.
integration-firecracker:
needs: build-infra
runs-on: [self-hosted, kvm]
if: >-
github.event_name == 'push' ||
@@ -159,49 +136,58 @@ jobs:
# range overlap; it prints the exact `backend setup` fix.
python3 cli.py backend status --backend=firecracker
- name: Download the candidate built from this checkout
uses: actions/download-artifact@v3
with:
name: infra-candidate
path: infra-candidate
- name: Build infra candidate from this checkout
env:
BOT_BOTTLE_FC_DROPBEAR: /var/cache/bot-bottle-fc/dropbear
run: python3 -m bot_bottle.backend.firecracker.publish_infra --output infra-candidate
- name: Replace the persistent infra VM with the candidate
run: python3 -c 'from bot_bottle.backend.firecracker import infra_vm; infra_vm.stop()'
# No dev-requirements install: the integration suite runs on stdlib
# `unittest` (pylint/pyright are lint.yml's concern, not this job's),
# and the self-hosted runner's Nix python env has no `pip` module
# (`python3 -m pip` → "No module named pip"). Nothing to install.
- name: Run integration tests (firecracker)
# No dev-requirements install: `coverage` is already provided by the
# self-hosted runner's Nix python env, and that env has no `pip`
# module to install into anyway.
- name: Run integration tests (firecracker) with coverage
env:
BOT_BOTTLE_BACKEND: firecracker
BOT_BOTTLE_INFRA_ARTIFACT_DIR: ${{ github.workspace }}/infra-candidate
run: python3 -m unittest discover -t . -s tests/integration -v
run: python3 -m coverage run --data-file=.coverage.firecracker -m unittest discover -t . -s tests/integration -v
# Combined unit+integration coverage + the diff-coverage gate (the hard
# gate: new/changed lines >= 90%). See docs/decisions/0004-coverage-policy.md.
- name: Upload firecracker coverage artifact
uses: actions/upload-artifact@v3
with:
name: coverage-firecracker
path: ${{ github.workspace }}/.coverage.firecracker
# Only upload the large rootfs artifact on main-branch pushes;
# PRs avoid the ~194 MB transfer. publish-infra only runs on main
# and downloads these to publish the byte-identical tested rootfs.
- name: Upload tested rootfs (main branch only)
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
uses: actions/upload-artifact@v3
with:
name: infra-candidate
path: infra-candidate/
- name: Upload dropbear for publish verification (main branch only)
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
uses: actions/upload-artifact@v3
with:
name: firecracker-inputs
path: /var/cache/bot-bottle-fc/dropbear
# Combined coverage gate: aggregates .coverage.* artifacts uploaded by each
# test job, then runs the diff-coverage gate (new/changed lines >= 90%).
#
# This runs on a self-hosted KVM runner (label `kvm`), NOT ubuntu-latest,
# because the Firecracker backend's subprocess/VM orchestration
# (launch/boot/SSH/isolation-probe) is covered by the integration suite,
# and that suite needs `/dev/kvm` + the provisioned TAP/nft pool — which a
# container-based runner doesn't have. On such a runner the firecracker
# integration test skips and its ~230 orchestration lines read as
# uncovered, so the gate can't pass there.
# Runs on ubuntu-latest — no KVM needed, no test reruns. Coverage files use
# relative_files = True (.coveragerc) so they combine cleanly across runners.
#
# Restricted to the same events as integration-firecracker (same-repo PRs,
# push, workflow_dispatch) for the same security reason.
#
# See #414 for the planned follow-up: artifact-based coverage combination
# (run tests once in their respective jobs, combine .coverage files here).
#
# build-infra creates one candidate from the checkout. This job boots that
# same candidate after integration-firecracker has exercised it; the main
# push path publishes the identical bytes only after every required job.
# Restricted to the same events as integration-firecracker: it depends on
# that job's coverage artifact and skips for fork PRs alongside it.
coverage:
needs: [build-infra, integration-firecracker]
needs: [unit, integration-docker, integration-firecracker]
timeout-minutes: 15
runs-on: [self-hosted, kvm]
runs-on: ubuntu-latest
if: >-
github.event_name == 'push' ||
github.event_name == 'workflow_dispatch' ||
@@ -213,29 +199,29 @@ jobs:
with:
fetch-depth: 0
- name: Preflight — Firecracker host is ready
run: |
command -v firecracker >/dev/null || {
echo "firecracker not on PATH — provision the runner (README: Firecracker on Linux)"; exit 1; }
test -e /dev/kvm || { echo "/dev/kvm missing — KVM not available on this runner"; exit 1; }
# `backend status` exits non-zero unless the TAP pool is up + no
# range overlap; it prints the exact `backend setup` fix.
python3 cli.py backend status --backend=firecracker
- name: Install coverage
run: python3 -m pip install --break-system-packages coverage
- name: Download the candidate already exercised by integration
- name: Download unit coverage artifact
uses: actions/download-artifact@v3
with:
name: infra-candidate
path: infra-candidate
name: coverage-unit
path: ${{ github.workspace }}
- name: Download docker coverage artifact
uses: actions/download-artifact@v3
with:
name: coverage-docker
path: ${{ github.workspace }}
- name: Download firecracker coverage artifact
uses: actions/download-artifact@v3
with:
name: coverage-firecracker
path: ${{ github.workspace }}
# No dev-requirements install: `coverage` is already provided by the
# self-hosted runner's Nix python env, and that env has no `pip`
# module to install into anyway. `scripts/coverage.sh` +
# `diff_coverage.py` need only `coverage` (not pylint/pyright).
- name: Combined coverage (unit + integration, incl. firecracker)
env:
BOT_BOTTLE_CI_INFRA_ARTIFACT_DIR: ${{ github.workspace }}/infra-candidate
run: PYTHON=python3 bash scripts/coverage.sh critical
run: PYTHON=python3 bash scripts/coverage.sh aggregate critical
- name: Diff-coverage gate (changed lines >= 90%)
run: |
@@ -243,14 +229,14 @@ jobs:
python3 scripts/diff_coverage.py --base origin/main --min 90
publish-infra:
needs: [stage-firecracker-inputs, build-infra, unit, integration-docker, integration-firecracker, coverage]
needs: [unit, integration-docker, integration-firecracker, coverage]
runs-on: ubuntu-latest
if: github.event_name == 'push' && github.ref == 'refs/heads/main'
steps:
- name: Checkout the tested revision
uses: actions/checkout@v4
- name: Download the tested candidate
- name: Download the tested rootfs
uses: actions/download-artifact@v3
with:
name: infra-candidate
@@ -258,9 +244,10 @@ jobs:
# publish_infra re-derives the version from the checkout to confirm the
# bundle matches before uploading, and the version hashes the dropbear
# bytes. Stage the SAME dropbear build-infra used, or the recheck
# computes a "<missing>"-dropbear version and rejects the candidate.
- name: Download the staged dropbear (matches build-infra's version)
# bytes. Download the SAME dropbear integration-firecracker used, or
# the recheck computes a "<missing>"-dropbear version and rejects the
# candidate.
- name: Download the staged dropbear (matches build's version)
uses: actions/download-artifact@v3
with:
name: firecracker-inputs
@@ -36,12 +36,9 @@ from dataclasses import dataclass
from ...egress import EgressPlan
from ...git_gate import GitGatePlan
from ...log import info
from ...orchestrator.client import OrchestratorClient, OrchestratorClientError
from ...orchestrator.client import OrchestratorClient
from ...orchestrator.registration import registration_inputs
from ..docker.gateway_provision import deprovision_git_gate, provision_git_gate
from . import util as container_mod
from .enumerate import CONTAINER_NAME_PREFIX, EnumerationError, enumerate_active
from .gateway import GATEWAY_NETWORK
from .gateway_provision import AppleGatewayTransport
from .infra import MacosInfraService, OrchestratorStartError
@@ -92,32 +89,6 @@ def ensure_gateway(
)
def live_source_ips(network: str) -> list[str]:
"""Every running agent container's address on `network`.
The reconciliation input: the orchestrator lives inside the infra
container and cannot enumerate the host's containers, so the host has to
tell it which bottles are actually up. Containers that have not been
assigned an address yet contribute nothing — the reap's grace window, not
this list, is what protects an in-flight launch.
Raises `EnumerationError` when the live set cannot be determined
authoritatively: either the container listing fails or any individual
inspect fails. Callers must skip reconciliation in that case to avoid
unregistering healthy bottles."""
ips: list[str] = []
for agent in enumerate_active():
name = f"{CONTAINER_NAME_PREFIX}{agent.slug}"
ip = container_mod.inspect_container_network_ip(name, network)
if ip is None:
raise EnumerationError(
f"container inspect {name!r} failed; live set is not authoritative"
)
if ip:
ips.append(ip)
return ips
def register_agent(
egress_plan: EgressPlan,
git_gate_plan: GitGatePlan,
@@ -132,16 +103,6 @@ def register_agent(
container — it is the attribution key the gateway resolves policy by.
Raises on failure; the caller tears down."""
client = OrchestratorClient(endpoint.orchestrator_url)
# Self-heal before registering: a launcher that died hard (SIGKILL, closed
# terminal, host sleep) never ran its teardown callback, leaving an active
# row with no container. vmnet recycles addresses, so such a row can
# collide with this bottle's — and `by_source_ip` fail-closes on ambiguity,
# which would resolve no policy at all and deny every host. Best-effort: a
# reconciliation failure must not block an otherwise-fine launch.
try:
client.reconcile(live_source_ips(endpoint.network))
except (OrchestratorClientError, EnumerationError) as e:
info(f"registry reconciliation skipped: {e}")
inputs = registration_inputs(egress_plan)
reg = client.register_bottle(
source_ip, image_ref=image_ref, policy=inputs.policy,
@@ -175,7 +136,6 @@ __all__ = [
"GatewayEndpoint",
"LaunchContext",
"ensure_gateway",
"live_source_ips",
"register_agent",
"teardown_consolidated",
"ConsolidatedLaunchError",
@@ -19,10 +19,6 @@ CONTAINER_NAME_PREFIX = "bot-bottle-"
_INFRA_NAMES = frozenset({INFRA_NAME})
class EnumerationError(RuntimeError):
"""container list failed; the resulting live set is not authoritative."""
def enumerate_active() -> list[ActiveAgent]:
result = subprocess.run(
["container", "list", "--quiet"],
@@ -31,10 +27,7 @@ def enumerate_active() -> list[ActiveAgent]:
check=False,
)
if result.returncode != 0:
raise EnumerationError(
f"container list failed: "
f"{(result.stderr or '').strip() or '<no stderr>'}"
)
return []
out: list[ActiveAgent] = []
for name in sorted(line.strip() for line in result.stdout.splitlines()):
if not name.startswith(CONTAINER_NAME_PREFIX) or name in _INFRA_NAMES:
@@ -572,41 +572,6 @@ def try_container_ipv4_on_network(name: str, network: str) -> str:
return ""
def inspect_container_network_ip(name: str, network: str) -> str | None:
"""IP of `name` on `network`, distinguishing inspect failure from "not yet".
Returns:
- the IP string when the container has one on `network`
- "" when inspect succeeds but no address is assigned yet (in-flight DHCP)
- None when the inspect command itself fails (authoritative list impossible)
"""
result = subprocess.run(
[_CONTAINER, "inspect", name],
capture_output=True, text=True, check=False,
)
if result.returncode != 0:
return None
try:
data = json.loads(result.stdout or "[]")
except json.JSONDecodeError:
return None
if isinstance(data, list):
data = data[0] if data else {}
if not isinstance(data, dict):
return None
status = data.get("status")
networks = status.get("networks") if isinstance(status, dict) else None
if not isinstance(networks, list):
return ""
for entry in networks:
if not isinstance(entry, dict) or entry.get("network") != network:
continue
raw = entry.get("ipv4Address")
if isinstance(raw, str) and raw:
return raw.split("/", 1)[0]
return ""
def wait_container_ipv4_on_network(
name: str, network: str, *, timeout: float = 15.0, poll: float = 0.25,
) -> str:
+1 -4
View File
@@ -19,7 +19,6 @@ from .commit import cmd_commit
from .edit import cmd_edit
from .info import cmd_info
from .init import cmd_init
from .login import cmd_login
from .resume import cmd_resume
from .start import cmd_start
from .supervise import cmd_supervise
@@ -34,7 +33,6 @@ COMMANDS = {
"info": cmd_info,
"init": cmd_init,
"list": cmd_list,
"login": cmd_login,
"resume": cmd_resume,
"start": cmd_start,
"supervise": cmd_supervise,
@@ -45,7 +43,7 @@ COMMANDS = {
# the host (TAP pool, /dev/kvm, firecracker) and never opens the store, so
# gating it on the schema breaks preflight on a fresh CI runner where stdin
# isn't a TTY and the migration prompt can't be answered.
NO_MIGRATION_COMMANDS = frozenset({"backend", "login"})
NO_MIGRATION_COMMANDS = frozenset({"backend"})
def usage() -> None:
@@ -58,7 +56,6 @@ def usage() -> None:
sys.stderr.write(" info print env, skills, and prompt details for a named agent\n")
sys.stderr.write(" init interactively create a new agent and add it to bot-bottle.json\n")
sys.stderr.write(" list list available agents or active containers\n")
sys.stderr.write(" login register this host with a bot-bottle console\n")
sys.stderr.write(
" resume re-launch a bottle by its identity "
"(continues state from PRD 0016)\n"
-168
View File
@@ -1,168 +0,0 @@
"""bb login — register this host with a bot-bottle console.
Opens a device-authorization flow against the target console, waits for the
operator to approve, then writes access and refresh tokens to
~/.bot-bottle/console.json (or $BOT_BOTTLE_ROOT/console.json).
Usage:
bb login [--console-url URL] [--label LABEL]
Flags:
--console-url URL Target console URL (overrides BB_CONSOLE_URL env var)
--label LABEL Host label shown in the console (default: hostname)
"""
from __future__ import annotations
import json
import os
import socket
import sys
import tempfile
import time
import urllib.error
import urllib.request
from pathlib import Path
from typing import Any
from ..paths import bot_bottle_root
_CONSOLE_URL_ENV = "BB_CONSOLE_URL"
_POLL_SLEEP = 2 # seconds between polls; matches console's poll_interval default
def _usage() -> None:
sys.stderr.write(
"usage: bb login [--console-url URL] [--label LABEL]\n"
"\n"
"Options:\n"
" --console-url URL Console base URL (or BB_CONSOLE_URL env var)\n"
" --label LABEL Host label shown in the console (default: hostname)\n"
)
def _flag(argv: list[str], name: str) -> str | None:
for i, arg in enumerate(argv):
if arg == name and i + 1 < len(argv):
return argv[i + 1]
if arg.startswith(f"{name}="):
return arg[len(name) + 1:]
return None
def _post(url: str, payload: dict[str, Any]) -> dict[str, Any]:
data = json.dumps(payload).encode()
req = urllib.request.Request(
url, data=data, headers={"Content-Type": "application/json"}
)
with urllib.request.urlopen(req, timeout=10) as resp:
return json.loads(resp.read())
def _get(url: str) -> tuple[int, dict[str, Any]]:
req = urllib.request.Request(url)
try:
with urllib.request.urlopen(req, timeout=10) as resp:
return resp.status, json.loads(resp.read())
except urllib.error.HTTPError as e:
return e.code, {}
def _save_credentials(
console_url: str, host_id: str, access_token: str, refresh_token: str
) -> Path:
path = bot_bottle_root() / "console.json"
path.parent.mkdir(parents=True, exist_ok=True)
content = (
json.dumps(
{
"url": console_url,
"host_id": host_id,
"access_token": access_token,
"refresh_token": refresh_token,
},
indent=2,
)
+ "\n"
)
fd, tmp_path_str = tempfile.mkstemp(dir=path.parent, prefix=".console-")
tmp = Path(tmp_path_str)
try:
tmp.chmod(0o600)
with os.fdopen(fd, "w") as f:
f.write(content)
os.replace(tmp, path)
except Exception:
try:
tmp.unlink()
except OSError:
pass
raise
return path
def cmd_login(argv: list[str]) -> int:
if "--help" in argv or "-h" in argv:
_usage()
return 0
console_url = _flag(argv, "--console-url") or os.environ.get(_CONSOLE_URL_ENV)
if not console_url:
sys.stderr.write(
"bb login: --console-url or BB_CONSOLE_URL is required\n"
)
return 1
console_url = console_url.rstrip("/")
label = _flag(argv, "--label") or socket.gethostname()
try:
resp = _post(f"{console_url}/api/v1/hosts/authorize", {"label": label})
except Exception as exc:
sys.stderr.write(f"bb login: failed to start authorization: {exc}\n")
return 1
device_code = resp["device_code"]
user_code = resp["user_code"]
expires_in = resp.get("expires_in", 300)
poll_sleep = max(1, min(int(resp.get("poll_interval", _POLL_SLEEP)), 60))
sys.stderr.write(
f"\nOpen this URL in your browser to authorize this host:\n\n"
f" {console_url}/authorize?code={user_code}\n\n"
f"Waiting for approval"
)
deadline = time.monotonic() + expires_in
while time.monotonic() < deadline:
sys.stderr.write(".")
sys.stderr.flush()
time.sleep(poll_sleep)
try:
code, result = _get(
f"{console_url}/api/v1/hosts/authorize/{device_code}"
)
except Exception:
continue
if code == 410:
break
st = result.get("status")
if st == "approved":
sys.stderr.write("\n\nApproved.\n")
path = _save_credentials(
console_url,
result["host_id"],
result["access_token"],
result["refresh_token"],
)
sys.stderr.write(f"Credentials saved to {path}\n")
return 0
if st == "denied":
sys.stderr.write("\n\nDenied by operator.\n")
return 1
sys.stderr.write("\n\nAuthorization timed out.\n")
return 1
-1
View File
@@ -379,7 +379,6 @@ class EgressAddon:
env,
request_method=flow.request.method,
request_headers=req_headers,
deny_reason=config.deny_reason,
)
if decision.action == "block":
+6 -45
View File
@@ -89,14 +89,6 @@ LOG_FULL = 2 # log block/warn events + full request and response bodies
class Config:
routes: tuple[Route, ...]
log: int = LOG_OFF
# Why this Config is a deny-all, when it is one for a reason *other* than
# the bottle's own policy genuinely not listing the host. A deny-all is
# indistinguishable from "policy loaded, host not allowed" at the decision
# point — both are simply "no matching route" — so without this the
# operator sees `host X is not in the allowlist` and goes hunting for a
# missing route that was never the problem. Empty for a normally-parsed
# policy; `decide` prefers it over the allowlist wording when set.
deny_reason: str = ""
@dataclass(frozen=True)
@@ -413,40 +405,16 @@ class PolicyResolverLike(typing.Protocol):
...
# Deny-all explanations. Each names the *actual* failure so an operator isn't
# sent looking for a missing egress route when the bottle never had a policy
# to begin with — the failure mode that made a bricked registration read like
# a misconfigured allowlist.
DENY_UNATTRIBUTED = (
"egress: this request was not attributed to any bottle, so no egress "
"policy applies and every host is denied. Either the bottle's registry "
"row is missing/ambiguous (torn down, or another bottle claimed its "
"source IP), or the request carried no matching identity token — check "
"that the caller's proxy URL includes it. This is not an allowlist problem."
)
DENY_UNPARSEABLE = (
"egress: this bottle's egress policy could not be parsed, so it is being "
"treated as deny-all. Fix the bottle's egress.routes; every host is denied "
"until it loads."
)
DENY_RESOLVER_ERROR = (
"egress: the orchestrator could not be reached to resolve this bottle's "
"egress policy, so every host is denied (fail-closed). Check that the "
"control plane is up; this is not an allowlist problem."
)
def _config_from_policy(policy: "str | None") -> "Config":
"""Parse a resolved policy blob into a Config, fail-closed: None / empty /
unparseable all become a deny-all Config (no routes every request
blocked). Each deny-all carries the reason it is one, so the block message
names the real fault instead of blaming the allowlist."""
blocked)."""
if not policy:
return Config(routes=(), deny_reason=DENY_UNATTRIBUTED)
return Config(routes=()) # unattributed or empty → deny-all
try:
return load_config(policy)
except ValueError:
return Config(routes=(), deny_reason=DENY_UNPARSEABLE)
return Config(routes=()) # unparseable policy → deny
def resolve_client_config(
@@ -460,7 +428,7 @@ def resolve_client_config(
try:
policy = resolver.resolve(client_ip, identity_token)
except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught
return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR)
return Config(routes=()) # orchestrator unreachable/errored → deny
return _config_from_policy(policy)
@@ -489,7 +457,7 @@ def resolve_client_context(
client_ip, identity_token,
)
except Exception: # noqa: BLE001 # pylint: disable=broad-exception-caught
return Config(routes=(), deny_reason=DENY_RESOLVER_ERROR), "", {}
return Config(routes=()), "", {} # orchestrator unreachable/errored → deny
return _config_from_policy(policy), (bottle_id or ""), tokens
@@ -604,16 +572,12 @@ def decide(
*,
request_method: str = "GET",
request_headers: typing.Mapping[str, str] | None = None,
deny_reason: str = "",
) -> Decision:
"""`deny_reason` is `Config.deny_reason`: when the deny-all came from a
missing/unparseable policy rather than the bottle's own allowlist, report
that instead of implying a route is merely absent."""
route = match_route(routes, request_host)
if route is None:
return Decision(
action="block",
reason=deny_reason or (
reason=(
f"egress: host {request_host!r} is not in the "
f"bottle's egress.routes allowlist. Declare a "
f"route for it or remove the request."
@@ -888,9 +852,6 @@ __all__ = [
"is_git_push_request",
"is_git_fetch_request",
"load_config",
"DENY_UNATTRIBUTED",
"DENY_UNPARSEABLE",
"DENY_RESOLVER_ERROR",
"resolve_client_config",
"resolve_client_context",
"PolicyResolverLike",
-15
View File
@@ -15,7 +15,6 @@ from __future__ import annotations
import json
import urllib.error
import urllib.request
from collections.abc import Iterable
from dataclasses import dataclass
from ..paths import host_control_plane_token
@@ -148,20 +147,6 @@ class OrchestratorClient:
raise OrchestratorClientError(f"teardown {bottle_id}: HTTP {status}")
return True
def reconcile(
self, live_source_ips: Iterable[str], *, grace_seconds: float | None = None,
) -> list[str]:
"""Drop registry rows for bottles that are no longer running
(`POST /reconcile`), returning the reaped bottle ids. `live_source_ips`
is the caller's enumeration of its live bottles — the orchestrator
can't see the backend from inside the infra container."""
body: dict[str, object] = {"live_source_ips": list(live_source_ips)}
if grace_seconds is not None:
body["grace_seconds"] = grace_seconds
payload = self._ok("POST", "/reconcile", body)
reaped = payload.get("reaped")
return [r for r in reaped if isinstance(r, str)] if isinstance(reaped, list) else []
def set_policy(self, bottle_id: str, policy: str) -> bool:
"""Live-reload a bottle's policy (`PUT /bottles/<id>/policy`). False on
404 (unknown bottle)."""
-24
View File
@@ -13,9 +13,6 @@ vsock / unix-socket portability caveats):
PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload)
body: {"policy"}
DELETE /bottles/<bottle_id> -> 200 {"torn_down": true} | 404 (teardown)
POST /reconcile -> 200 {"reaped": [bottle_id, ...]}
body: {"live_source_ips": [...],
["grace_seconds"]}
POST /attribute -> 200 {"bottle_id"} | 403
POST /resolve -> 200 {"bottle_id","policy"} | 403
body: {"source_ip","identity_token"}
@@ -144,27 +141,6 @@ def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
return 200, {"torn_down": True}
return 404, {"error": "no such bottle"}
if method == "POST" and route == "/reconcile":
# Host-driven self-heal: the caller enumerates its live bottles (only
# the host can see the backend) and the orchestrator drops rows for
# every other active bottle. Trusted-caller only — an agent that could
# reach this would be able to unregister its neighbours.
try:
data = _parse_json_object(body)
except ValueError as e:
return 400, {"error": f"invalid JSON: {e}"}
raw_ips = data.get("live_source_ips")
if not isinstance(raw_ips, list):
return 400, {"error": "live_source_ips (list of strings) is required"}
live = [ip for ip in raw_ips if isinstance(ip, str) and ip]
grace = data.get("grace_seconds")
kwargs = (
{"grace_seconds": float(grace)}
if isinstance(grace, (int, float)) and not isinstance(grace, bool)
else {}
)
return 200, {"reaped": orch.reconcile(live, **kwargs)}
if method == "POST" and route == "/attribute":
try:
data = _parse_json_object(body)
-72
View File
@@ -32,7 +32,6 @@ import hmac
import secrets
import sqlite3
import time
from collections.abc import Iterable
from dataclasses import dataclass
from pathlib import Path
@@ -43,12 +42,6 @@ from ..paths import host_db_path
# 256 bits of urandom, URL-safe — unguessable per-bottle identity token.
IDENTITY_TOKEN_BYTES = 32
# How recently a row must have been registered to be exempt from
# `reap_absent`. Covers the window between `container run` and the address
# becoming visible to another launch's enumeration, so reconciliation never
# reaps a bottle that is still coming up.
DEFAULT_REAP_GRACE_SECONDS = 120.0
def new_identity_token() -> str:
"""A fresh per-bottle identity token (PRD 0070 attribution defence)."""
@@ -232,70 +225,6 @@ class RegistryStore(DbStore):
).fetchall()
return [_row_to_record(r) for r in rows]
def reap_absent(
self,
live_source_ips: Iterable[str],
*,
grace_seconds: float = DEFAULT_REAP_GRACE_SECONDS,
now: float | None = None,
) -> list[BottleRecord]:
"""Delete active rows whose source IP is not held by a live bottle.
A row only ever leaves the registry two ways: an explicit
`teardown_bottle` (the launcher's cleanup callback) or the supersede
sweep in `register`. Neither runs when the launching CLI dies hard
SIGKILL, a closed terminal, a host sleep/crash so the row outlives
its container. That orphan is not inert: source IPs are recycled by
the backend's DHCP, and `by_source_ip` fail-closes on ambiguity, so a
leftover row at a reused address can brick the *next* bottle that
lands on it (no policy resolved -> every host denied, reported to the
agent as "not in the allowlist"). Reconciling against the live set at
launch keeps the registry from accumulating those landmines.
Restores the invariant the data plane needs: **at most one active row
per live address, and none at all for a dead one.** Two cases, because
a dead bottle's address may already have been handed to a live one:
* no live bottle holds the address every row there is an orphan;
* a live bottle holds it but several rows claim it the newest
registration is authoritative and the rest are orphans, the same
rule `register`'s same-IP supersede sweep applies. Without this
second case a recycled address stays ambiguous, which is exactly
the state that resolves no policy.
`grace_seconds` protects an in-flight launch: registration happens
moments after `container run`, and a concurrent launch's address may
not be visible to the caller's enumeration yet. Rows younger than the
grace window are never reaped, so reconciliation can't race a bottle
that is still coming up. Returns the deleted records."""
live = {ip for ip in live_source_ips if ip}
cutoff = (time.time() if now is None else now) - grace_seconds
with self._connection() as conn:
rows = conn.execute(
"SELECT * FROM orchestrator_bottles WHERE state = 'active'",
).fetchall()
by_ip: dict[str, list[BottleRecord]] = {}
for row in rows:
rec = _row_to_record(row)
by_ip.setdefault(rec.source_ip, []).append(rec)
candidates: list[BottleRecord] = []
for ip, recs in by_ip.items():
if ip not in live:
candidates.extend(recs)
continue
# Keep the newest claim on a live address; supersede the rest.
recs.sort(key=lambda r: r.created_at)
candidates.extend(recs[:-1])
doomed = [r for r in candidates if r.created_at <= cutoff]
for rec in doomed:
conn.execute(
"DELETE FROM orchestrator_bottles WHERE bottle_id = ?",
(rec.bottle_id,),
)
if doomed:
self._chmod()
return doomed
def by_source_ip(self, source_ip: str) -> BottleRecord | None:
"""Network-layer attribution: the single active bottle at this source
IP, or None if unknown or ambiguous (more than one a
@@ -333,5 +262,4 @@ __all__ = [
"new_identity_token",
"default_db_path",
"IDENTITY_TOKEN_BYTES",
"DEFAULT_REAP_GRACE_SECONDS",
]
+1 -30
View File
@@ -13,20 +13,15 @@ Launch lifecycle:
and returns the record. If the broker rejects/fails, the registry entry
is rolled back so a failed launch leaves no orphan.
* `teardown_bottle` sends a signed teardown request, then deregisters.
* `reconcile` sweeps rows whose bottle is no longer running the
self-heal for the teardown paths that never got to run (a hard-killed
launcher), since an orphan row at a recycled source IP bricks the next
bottle that lands on it.
"""
from __future__ import annotations
import json
from collections.abc import Iterable
from datetime import datetime, timezone
from .broker import LaunchBroker, LaunchRequest, sign_request
from .registry import DEFAULT_REAP_GRACE_SECONDS, BottleRecord, RegistryStore
from .registry import BottleRecord, RegistryStore
from .gateway import Gateway
from ..supervise import (
AuditEntry,
@@ -122,30 +117,6 @@ class Orchestrator:
self._tokens.pop(bottle_id, None)
return True
def reconcile(
self,
live_source_ips: Iterable[str],
*,
grace_seconds: float = DEFAULT_REAP_GRACE_SECONDS,
) -> list[str]:
"""Drop registry rows for bottles that are no longer running, and
forget their in-memory egress tokens. Returns the reaped bottle ids.
The caller supplies the live set because only the host can enumerate
its own containers the orchestrator runs *inside* the infra
container and has no view of the backend. Deliberately does not
broker a teardown: the container is already gone, so there is nothing
to stop, and a broker error must not stop the sweep from clearing
the row that would otherwise brick the next bottle at that address.
See `RegistryStore.reap_absent` for why orphans accumulate and why
they are harmful rather than merely untidy."""
reaped = self.registry.reap_absent(
live_source_ips, grace_seconds=grace_seconds)
for rec in reaped:
self._tokens.pop(rec.bottle_id, None)
return [rec.bottle_id for rec in reaped]
def tokens_for(self, bottle_id: str) -> dict[str, str]:
"""The bottle's in-memory egress auth tokens (env_name -> value), or
empty. The gateway injects these per request; they are never
+110
View File
@@ -0,0 +1,110 @@
# PRD prd-new: CI artifact-based coverage and local Firecracker candidate flow
- **Status:** Active
- **Author:** Claude
- **Created:** 2026-07-21
- **Issue:** #446
## Summary
Restructure the CI test pipeline to run each test suite exactly once, upload
small `.coverage.*` artifacts, and combine them in a lightweight aggregation
job. Move the infra build onto the KVM runner so the ~194 MB rootfs never
crosses the network for PRs. On main-branch pushes, publish the byte-identical
rootfs that was tested.
## Motivation
The prior pipeline had two redundant costs:
1. **Duplicate artifact transfers.** `build-infra` (ubuntu-latest) built and
uploaded the ~194 MB rootfs; `integration-firecracker` downloaded it; the
`coverage` job downloaded it a second time. Combined download overhead: ~83
seconds per run, plus the ~70-second upload.
2. **Duplicate test execution.** `integration-firecracker` ran the Firecracker
integration suite; `coverage` ran the entire unit + integration suite again
on the same KVM runner to collect coverage data. Every line of Firecracker
code was tested twice per CI run.
## Goals
- Each test suite (unit, integration-docker, integration-firecracker) executes
exactly once per workflow run.
- PRs incur no large artifact transfers — the rootfs stays on the KVM runner.
- Main-branch pushes publish a byte-for-byte identical rootfs to the one that
passed the integration tests.
- Concurrent workflow runs cannot cross-publish candidates (naturally enforced
by Gitea Actions' per-run artifact scoping).
- Failed or cancelled runs block publication (enforced by the `needs:` chain on
`publish-infra`).
## Non-goals
- Changing test semantics or the coverage policy (ADR 0004).
- Removing the KVM runner guard on `integration-firecracker` and `coverage`.
- Changing how `publish_infra.py` builds or uploads the rootfs.
## Design
### Job graph
```
unit ──────────────────────────────────┐
integration-docker ────────────────────┤──► coverage ──► publish-infra (main only)
integration-firecracker (KVM) ─────────┘
```
### `unit`
Unchanged except: `coverage run` writes `--data-file=.coverage.unit`; the file
is uploaded as the `coverage-unit` artifact.
### `integration-docker`
Adds a `coverage` install step. `coverage run` writes `--data-file=.coverage.docker`;
the file is uploaded as `coverage-docker`.
### `integration-firecracker` (KVM runner)
Replaces the old `stage-firecracker-inputs``build-infra` → download chain:
1. Builds the infra candidate locally with
`BOT_BOTTLE_FC_DROPBEAR=/var/cache/bot-bottle-fc/dropbear`.
2. Boots the candidate and runs integration tests with coverage, writing
`.coverage.firecracker`.
3. Uploads the small `coverage-firecracker` artifact unconditionally.
4. On main-branch pushes only, uploads the rootfs as `infra-candidate` and the
dropbear as `firecracker-inputs` so `publish-infra` can verify and publish
the byte-identical artifact.
### `coverage`
Moves from a KVM runner to `ubuntu-latest`. No tests are re-executed:
1. Downloads `coverage-unit`, `coverage-docker`, and `coverage-firecracker`.
2. Runs `scripts/coverage.sh aggregate critical`, which calls
`coverage combine` then `coverage report`.
3. Runs the diff-coverage gate (`scripts/diff_coverage.py`).
Coverage files use `relative_files = True` (`.coveragerc`) so they combine
cleanly across runners with different absolute workspace paths.
### `publish-infra`
Depends on all four predecessor jobs (unchanged gate). Downloads `infra-candidate`
and `firecracker-inputs` that were uploaded by `integration-firecracker` on
main — the same byte sequence that passed the integration tests.
### Eliminated jobs
- `stage-firecracker-inputs`: existed only to copy the dropbear to ubuntu-latest
for `build-infra`. No longer needed.
- `build-infra`: the infra candidate is now built on the KVM runner in
`integration-firecracker`.
### Script changes
`scripts/coverage.sh` gains an `aggregate` mode (`coverage.sh aggregate [critical]`)
that combines pre-existing `.coverage.*` files instead of re-running tests.
The existing run mode (`coverage.sh [critical]`) is preserved for local dev.
+28 -8
View File
@@ -1,15 +1,19 @@
#!/usr/bin/env bash
# Combined unit + integration coverage (see docs/decisions/0004-coverage-policy.md).
#
# Runs the unit suite, then appends the integration suite (which skips
# cleanly when Docker / the backend CLIs are unavailable), and prints one
# combined report. The integration suite is what scores the subprocess /
# backend orchestration modules, so the number here is the policy's
# yardstick — not the unit-only badge.
# Two modes:
#
# Usage:
# scripts/coverage.sh # combined report
# scripts/coverage.sh critical # also report just the critical modules
# scripts/coverage.sh [critical]
# Run mode (default, for local dev): executes the unit suite then the
# integration suite under coverage and prints a combined report.
#
# scripts/coverage.sh aggregate [critical]
# Aggregate mode (used by CI): combines pre-existing .coverage.* files
# produced by individual test jobs and prints a combined report. No tests
# are re-executed; no KVM or Docker dependency.
#
# Pass "critical" as the last argument in either mode to also report just the
# critical modules (ADR 0004 target: 90%).
set -euo pipefail
cd "$(dirname "$0")/.."
@@ -21,6 +25,22 @@ PY="${PYTHON:-python3}"
# README "core coverage" badge can't drift; comma-join it for --include.
CRITICAL=$(grep -vE '^[[:space:]]*(#|$)' scripts/critical-modules.txt | paste -sd, -)
if [ "${1:-}" = "aggregate" ]; then
# Aggregate mode: combine .coverage.* artifacts already in the workspace.
echo "== combining coverage artifacts ==" >&2
"$PY" -m coverage combine
echo "== combined report ==" >&2
"$PY" -m coverage report -m
if [ "${2:-}" = "critical" ]; then
echo "== critical modules (ADR 0004 target: 90%) ==" >&2
"$PY" -m coverage report --include="$CRITICAL"
fi
exit 0
fi
# Run mode (default): execute both suites under coverage in this process.
rm -f .coverage
echo "== unit ==" >&2
-260
View File
@@ -1,260 +0,0 @@
"""Unit tests for bb login command."""
from __future__ import annotations
import json
import os
import tempfile
import unittest
import urllib.error
from email.message import Message
from typing import Any
from unittest.mock import MagicMock, patch
class TestFlagParsing(unittest.TestCase):
def test_console_url_flag(self) -> None:
from bot_bottle.cli.login import _flag
self.assertEqual(_flag(["--console-url", "http://x"], "--console-url"), "http://x")
def test_console_url_equals_form(self) -> None:
from bot_bottle.cli.login import _flag
self.assertEqual(
_flag(["--console-url=http://x"], "--console-url"), "http://x"
)
def test_label_flag(self) -> None:
from bot_bottle.cli.login import _flag
self.assertEqual(_flag(["--label", "my-mac"], "--label"), "my-mac")
def test_missing_flag_returns_none(self) -> None:
from bot_bottle.cli.login import _flag
self.assertIsNone(_flag([], "--console-url"))
class TestHttpHelpers(unittest.TestCase):
def test_post_sends_json_and_decodes_response(self) -> None:
from bot_bottle.cli.login import _post
response = MagicMock()
response.__enter__.return_value.read.return_value = b'{"ok": true}'
with patch("urllib.request.urlopen", return_value=response) as urlopen:
self.assertEqual(
_post("http://console/start", {"label": "host"}), {"ok": True}
)
request = urlopen.call_args.args[0]
self.assertEqual(request.data, b'{"label": "host"}')
self.assertEqual(request.get_header("Content-type"), "application/json")
def test_get_decodes_success_response(self) -> None:
from bot_bottle.cli.login import _get
response = MagicMock()
response.__enter__.return_value.status = 200
response.__enter__.return_value.read.return_value = b'{"status": "pending"}'
with patch("urllib.request.urlopen", return_value=response):
self.assertEqual(
_get("http://console/status"), (200, {"status": "pending"})
)
def test_get_returns_http_error_status(self) -> None:
from bot_bottle.cli.login import _get
error = urllib.error.HTTPError(
"http://console/status", 410, "gone", Message(), None
)
with patch("urllib.request.urlopen", side_effect=error):
self.assertEqual(_get("http://console/status"), (410, {}))
class TestSaveCredentials(unittest.TestCase):
def test_writes_json_and_sets_perms(self) -> None:
from bot_bottle.cli.login import _save_credentials
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
path = _save_credentials("http://c", "hid", "at", "rt")
self.assertTrue(path.exists())
data = json.loads(path.read_text())
self.assertEqual(data["url"], "http://c")
self.assertEqual(data["host_id"], "hid")
self.assertEqual(data["access_token"], "at")
self.assertEqual(data["refresh_token"], "rt")
self.assertEqual(oct(path.stat().st_mode & 0o777), oct(0o600))
def test_temp_file_is_private_before_replace(self) -> None:
"""Temp file must be 0600 at the moment os.replace is called."""
from bot_bottle.cli.login import _save_credentials
from pathlib import Path as _Path
tmp_perms_at_replace: list[int] = []
real_replace = os.replace
def _spy_replace(
src: str | os.PathLike[str], dst: str | os.PathLike[str]
) -> None:
tmp_perms_at_replace.append(_Path(src).stat().st_mode & 0o777)
real_replace(src, dst)
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("os.replace", side_effect=_spy_replace):
path = _save_credentials("http://c", "hid", "at", "rt")
self.assertEqual(len(tmp_perms_at_replace), 1)
self.assertEqual(oct(tmp_perms_at_replace[0]), oct(0o600))
self.assertEqual(oct(path.stat().st_mode & 0o777), oct(0o600))
class TestCmdLoginMissingUrl(unittest.TestCase):
def test_help_returns_0(self) -> None:
from bot_bottle.cli.login import cmd_login
self.assertEqual(cmd_login(["--help"]), 0)
def test_returns_1_without_url(self) -> None:
from bot_bottle.cli.login import cmd_login
with patch.dict(os.environ, {}, clear=True):
os.environ.pop("BB_CONSOLE_URL", None)
result = cmd_login([])
self.assertEqual(result, 1)
def test_reads_env_var(self) -> None:
"""Exits 1 (network error) not because of missing URL when env var is set."""
from bot_bottle.cli.login import cmd_login
def _fail_post(_url: str, _payload: dict[str, Any]) -> dict[str, Any]:
raise OSError("connection refused")
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(
os.environ,
{
"BB_CONSOLE_URL": "http://localhost:9999",
"BOT_BOTTLE_ROOT": tmp,
},
):
with patch("bot_bottle.cli.login._post", side_effect=_fail_post):
result = cmd_login([])
self.assertEqual(result, 1)
class TestCmdLoginFlow(unittest.TestCase):
def _run_with_mocks(
self, poll_responses: list[dict[str, Any]], tmp: str
) -> int:
from bot_bottle.cli.login import cmd_login
start_resp = {
"device_code": "dc123",
"user_code": "ABC-DEF",
"expires_in": 300,
"poll_interval": 0,
}
poll_iter = iter(poll_responses)
def _fake_post(
_url: str, _payload: dict[str, Any]
) -> dict[str, Any]:
return start_resp
def _fake_get(_url: str) -> tuple[int, dict[str, Any]]:
try:
resp = next(poll_iter)
except StopIteration:
return 200, {"status": "pending"}
return 200, resp
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("bot_bottle.cli.login._post", side_effect=_fake_post):
with patch("bot_bottle.cli.login._get", side_effect=_fake_get):
with patch("time.sleep"):
return cmd_login(["--console-url", "http://console"])
def test_approved_flow_returns_0(self) -> None:
approved = {
"status": "approved",
"host_id": "hid",
"access_token": "at",
"refresh_token": "rt",
}
with tempfile.TemporaryDirectory() as tmp:
result = self._run_with_mocks(
[{"status": "pending"}, approved], tmp
)
self.assertEqual(result, 0)
with open(os.path.join(tmp, "console.json")) as f:
creds = json.loads(f.read())
self.assertEqual(creds["host_id"], "hid")
def test_denied_flow_returns_1(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
result = self._run_with_mocks([{"status": "denied"}], tmp)
self.assertEqual(result, 1)
def test_timeout_returns_1(self) -> None:
from bot_bottle.cli.login import cmd_login
start_resp = {
"device_code": "dc",
"user_code": "ZZZ-ZZZ",
"expires_in": 0, # already expired; loop never runs
"poll_interval": 2,
}
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("bot_bottle.cli.login._post", return_value=start_resp):
result = cmd_login(["--console-url", "http://console"])
self.assertEqual(result, 1)
def test_poll_interval_from_server_is_used(self) -> None:
"""time.sleep must be called with the server-provided poll_interval."""
from bot_bottle.cli.login import cmd_login
server_interval = 7
start_resp = {
"device_code": "dc",
"user_code": "ABC-DEF",
"expires_in": 300,
"poll_interval": server_interval,
}
approved = {
"status": "approved",
"host_id": "hid",
"access_token": "at",
"refresh_token": "rt",
}
poll_iter = iter([{"status": "pending"}, approved])
def _fake_get(_url: str) -> tuple[int, dict[str, str]]:
return 200, next(poll_iter)
with tempfile.TemporaryDirectory() as tmp:
with patch.dict(os.environ, {"BOT_BOTTLE_ROOT": tmp}):
with patch("bot_bottle.cli.login._post", return_value=start_resp):
with patch(
"bot_bottle.cli.login._get", side_effect=_fake_get
):
with patch("time.sleep") as mock_sleep:
result = cmd_login(["--console-url", "http://console"])
self.assertEqual(result, 0)
self.assertTrue(mock_sleep.called)
for call in mock_sleep.call_args_list:
self.assertEqual(call.args[0], server_interval)
class TestDispatcherRegistration(unittest.TestCase):
def test_login_in_commands(self) -> None:
from bot_bottle.cli import COMMANDS
self.assertIn("login", COMMANDS)
def test_login_in_no_migration(self) -> None:
from bot_bottle.cli import NO_MIGRATION_COMMANDS
self.assertIn("login", NO_MIGRATION_COMMANDS)
if __name__ == "__main__":
unittest.main()
+1 -60
View File
@@ -4,14 +4,7 @@ from __future__ import annotations
import unittest
from bot_bottle.egress_addon_core import (
DENY_RESOLVER_ERROR,
DENY_UNATTRIBUTED,
DENY_UNPARSEABLE,
decide,
resolve_client_config,
resolve_client_context,
)
from bot_bottle.egress_addon_core import resolve_client_config, resolve_client_context
from bot_bottle.policy_resolver import PolicyResolveError
@@ -115,55 +108,3 @@ class TestResolveClientContext(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestDenyReasonNamesTheRealFault(unittest.TestCase):
"""A deny-all must not masquerade as a missing allowlist entry.
Regression: an unregistered bottle resolves no policy, so *every* host is
denied but the block message said `host X is not in the allowlist`,
which reads as a config problem and sends the operator hunting for a route
that was never missing. The structural reason wins over that wording.
"""
def _reason(self, resolver: object, host: str = "chatgpt.com") -> str:
cfg = resolve_client_config(resolver, "10.243.0.1") # type: ignore[arg-type]
return decide(cfg.routes, host, "/v1/x", {}, deny_reason=cfg.deny_reason).reason
def test_unattributed_says_unattributed_not_allowlist(self) -> None:
reason = self._reason(_FakeResolver(result=None))
self.assertEqual(DENY_UNATTRIBUTED, reason)
# The misleading claim is the one that must be gone: the host was
# never "not in the allowlist" — there was no allowlist at all.
self.assertNotIn("is not in the bottle's egress.routes allowlist", reason)
# Both causes must be named. `/resolve` fail-closes on a missing row
# *and* on a token mismatch, and the message pointing only at the row
# sent us hunting for a deregistered bottle that was registered fine.
self.assertIn("registry row", reason)
self.assertIn("identity token", reason)
def test_resolver_error_says_orchestrator_unreachable(self) -> None:
self.assertEqual(DENY_RESOLVER_ERROR, self._reason(_FakeResolver(raises=True)))
def test_unparseable_policy_says_so(self) -> None:
self.assertEqual(
DENY_UNPARSEABLE, self._reason(_FakeResolver(result="routes: notalist\n")))
def test_a_real_allowlist_miss_keeps_the_allowlist_wording(self) -> None:
"""The message only changes for structural deny-alls — a loaded policy
that genuinely lacks the host still points at the allowlist."""
reason = self._reason(_FakeResolver(result='routes:\n - host: "api.example.com"\n'))
self.assertIn("is not in the bottle's egress.routes allowlist", reason)
self.assertIn("chatgpt.com", reason)
def test_allowed_host_is_still_forwarded(self) -> None:
cfg = resolve_client_config(
_FakeResolver(result='routes:\n - host: "api.example.com"\n'), "10.243.0.1")
decision = decide(
cfg.routes, "api.example.com", "/v1/x", {}, deny_reason=cfg.deny_reason)
self.assertEqual("forward", decision.action)
def test_a_parsed_policy_carries_no_deny_reason(self) -> None:
cfg = resolve_client_config(
_FakeResolver(result='routes:\n - host: "api.example.com"\n'), "10.243.0.1")
self.assertEqual("", cfg.deny_reason)
+1 -103
View File
@@ -87,8 +87,7 @@ class TestRegisterAgent(unittest.TestCase):
*, source_ip: str = "192.168.128.9",
):
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_MOD}.provision_git_gate", provision or Mock()), \
patch(f"{_MOD}.live_source_ips", return_value=[]):
patch(f"{_MOD}.provision_git_gate", provision or Mock()):
return register_agent(
_egress_plan(), _git_plan(),
source_ip=source_ip, endpoint=_endpoint(), image_ref="img:1",
@@ -135,104 +134,3 @@ class TestTeardown(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestLiveSourceIps(unittest.TestCase):
"""The reconciliation input: the host enumerates its own bottles because
the orchestrator, inside the infra container, cannot see the backend."""
def _agents(self, *slugs: str) -> list[Mock]:
return [Mock(slug=s) for s in slugs]
def test_maps_slugs_to_container_addresses(self) -> None:
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \
patch(f"{_MOD}.container_mod.inspect_container_network_ip",
side_effect=["10.0.0.1", "10.0.0.2"]) as ip:
got = live_source_ips("net0")
self.assertEqual(["10.0.0.1", "10.0.0.2"], got)
self.assertEqual("bot-bottle-a", ip.call_args_list[0].args[0])
def test_containers_without_an_address_are_skipped(self) -> None:
"""A container that hasn't been given a DHCP address yet contributes
nothing the reap's grace window, not this list, protects it."""
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \
patch(f"{_MOD}.container_mod.inspect_container_network_ip",
side_effect=["", "10.0.0.2"]):
self.assertEqual(["10.0.0.2"], live_source_ips("net0"))
def test_container_list_failure_raises(self) -> None:
"""If container list fails, the live set is not authoritative and
reconciliation must be skipped."""
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with patch(f"{_MOD}.enumerate_active",
side_effect=EnumerationError("container list failed")):
with self.assertRaises(EnumerationError):
live_source_ips("net0")
def test_per_container_inspect_failure_raises(self) -> None:
"""If any individual inspect fails, the live set is not authoritative."""
from bot_bottle.backend.macos_container.consolidated_launch import live_source_ips
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with patch(f"{_MOD}.enumerate_active", return_value=self._agents("a", "b")), \
patch(f"{_MOD}.container_mod.inspect_container_network_ip",
side_effect=["10.0.0.1", None]):
with self.assertRaises(EnumerationError):
live_source_ips("net0")
class TestRegisterAgentReconciles(unittest.TestCase):
"""Registration self-heals the registry first: an orphan row at a recycled
address makes attribution ambiguous, which resolves no policy at all and
denies every host for the bottle being launched."""
def _register(self, client: Mock) -> None:
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_MOD}.provision_git_gate"), \
patch(f"{_MOD}.live_source_ips", return_value=["10.0.0.7"]):
register_agent(
_egress_plan(), _git_plan(),
source_ip="10.0.0.7", endpoint=_endpoint(),
)
def test_reconciles_before_registering(self) -> None:
client = _client()
calls: list[str] = []
def _reconcile(*_args: object, **_kwargs: object) -> list[str]:
calls.append("reconcile")
return []
def _register_bottle(*_args: object, **_kwargs: object) -> RegisteredBottle:
calls.append("register")
return RegisteredBottle("b1", "tok")
client.reconcile.side_effect = _reconcile
client.register_bottle.side_effect = _register_bottle
self._register(client)
self.assertEqual(["reconcile", "register"], calls)
client.reconcile.assert_called_once_with(["10.0.0.7"])
def test_a_reconcile_failure_does_not_block_the_launch(self) -> None:
from bot_bottle.orchestrator.client import OrchestratorClientError
client = _client()
client.reconcile.side_effect = OrchestratorClientError("unreachable")
self._register(client)
client.register_bottle.assert_called_once()
def test_enumeration_error_does_not_block_the_launch(self) -> None:
"""A partial container listing must not abort the launch — skip
reconciliation and proceed, just as with an unreachable orchestrator."""
from bot_bottle.backend.macos_container.enumerate import EnumerationError
client = _client()
with patch(f"{_MOD}.OrchestratorClient", return_value=client), \
patch(f"{_MOD}.provision_git_gate"), \
patch(f"{_MOD}.live_source_ips",
side_effect=EnumerationError("container list failed")):
register_agent(
_egress_plan(), _git_plan(),
source_ip="10.0.0.7", endpoint=_endpoint(),
)
client.register_bottle.assert_called_once()
+2 -4
View File
@@ -66,10 +66,8 @@ class TestMacosContainerEnumerate(unittest.TestCase):
agents = self._enumerate("bot-bottle-mac-infra\nbot-bottle-dev-abc\n")
self.assertEqual(["dev-abc"], [a.slug for a in agents])
def test_raises_when_the_cli_fails(self):
from bot_bottle.backend.macos_container.enumerate import EnumerationError
with self.assertRaises(EnumerationError):
self._enumerate("", returncode=1)
def test_empty_when_the_cli_fails(self):
self.assertEqual([], self._enumerate("", returncode=1))
if __name__ == "__main__":
-44
View File
@@ -334,50 +334,6 @@ class TestInspectDigests(unittest.TestCase):
self.assertEqual({}, util.container_env("x"))
class TestInspectContainerNetworkIp(unittest.TestCase):
"""inspect_container_network_ip must distinguish inspect failure (None)
from 'no DHCP address yet' (""), which is the invariant live_source_ips
relies on to skip reconciliation on partial snapshots."""
_NETWORK = "bot-bottle-mac-gateway"
def _inspect(self, stdout: str, returncode: int = 0) -> str | None:
cp = util.subprocess.CompletedProcess(
args=[], returncode=returncode, stdout=stdout, stderr="",
)
with patch.object(util.subprocess, "run", return_value=cp):
return util.inspect_container_network_ip("bot-bottle-abc", self._NETWORK)
def _entry(self, ip: str = "192.168.128.5") -> str:
return (
f'[{{"status":{{"networks":['
f'{{"network":"{self._NETWORK}","ipv4Address":"{ip}"}}'
f']}}}}]'
)
def test_returns_ip_when_inspect_succeeds(self) -> None:
self.assertEqual("192.168.128.5", self._inspect(self._entry()))
def test_strips_cidr_prefix(self) -> None:
self.assertEqual("192.168.128.5", self._inspect(self._entry("192.168.128.5/24")))
def test_returns_empty_string_when_no_address_assigned_yet(self) -> None:
no_ip = f'[{{"status":{{"networks":[{{"network":"{self._NETWORK}","ipv4Address":""}}]}}}}]'
self.assertEqual("", self._inspect(no_ip))
def test_returns_empty_string_when_network_list_absent(self) -> None:
self.assertEqual("", self._inspect('[{"status":{}}]'))
def test_returns_none_on_nonzero_exit(self) -> None:
self.assertIsNone(self._inspect("", returncode=1))
def test_returns_none_on_malformed_json(self) -> None:
self.assertIsNone(self._inspect("not-json"))
def test_returns_none_on_unexpected_json_shape(self) -> None:
self.assertIsNone(self._inspect("null"))
class TestWaitContainerIpv4(unittest.TestCase):
def test_returns_address_once_dhcp_assigns_it(self):
with patch.object(util, "try_container_ipv4_on_network", side_effect=["", "", "192.168.128.4"]), \
-29
View File
@@ -104,32 +104,3 @@ class TestHealthAndPolicy(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestReconcile(unittest.TestCase):
def setUp(self) -> None:
self.c = OrchestratorClient("http://orch:8080")
def test_posts_live_ips_and_returns_reaped(self) -> None:
with patch(_URLOPEN, return_value=_resp(200, {"reaped": ["b1", "b2"]})) as m:
got = self.c.reconcile(["10.0.0.2", "10.0.0.3"])
self.assertEqual(["b1", "b2"], got)
sent = json.loads(m.call_args.args[0].data)
self.assertEqual(["10.0.0.2", "10.0.0.3"], sent["live_source_ips"])
self.assertNotIn("grace_seconds", sent) # omitted -> server default
def test_grace_seconds_is_forwarded_when_given(self) -> None:
with patch(_URLOPEN, return_value=_resp(200, {"reaped": []})) as m:
self.c.reconcile([], grace_seconds=30)
self.assertEqual(30, json.loads(m.call_args.args[0].data)["grace_seconds"])
def test_malformed_reaped_is_tolerated(self) -> None:
with patch(_URLOPEN, return_value=_resp(200, {"reaped": ["ok", 5, None]})):
self.assertEqual(["ok"], self.c.reconcile([]))
with patch(_URLOPEN, return_value=_resp(200, {})):
self.assertEqual([], self.c.reconcile([]))
def test_error_status_raises(self) -> None:
with patch(_URLOPEN, side_effect=_http_error(500)):
with self.assertRaises(OrchestratorClientError):
self.c.reconcile([])
@@ -8,13 +8,11 @@ from __future__ import annotations
import json
import secrets
import sqlite3
import tempfile
import threading
import unittest
import urllib.error
import urllib.request
from contextlib import closing
from pathlib import Path
from unittest.mock import patch
@@ -266,7 +264,6 @@ class TestControlPlaneAuth(unittest.TestCase):
("POST", "/bottles", _body({"source_ip": "10.0.0.1"})),
("PUT", "/bottles/x/policy", _body({"policy": "routes: []"})),
("DELETE", "/bottles/x", b""),
("POST", "/reconcile", _body({"live_source_ips": []})),
("POST", "/resolve", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("POST", "/attribute", _body({"source_ip": "10.0.0.1", "identity_token": "t"})),
("GET", "/supervise/proposals", b""),
@@ -390,56 +387,3 @@ class TestDispatchSupervise(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestReconcileRoute(unittest.TestCase):
"""`POST /reconcile` — the host tells the orchestrator which bottles are
actually up, since the orchestrator can't see the backend from inside the
infra container."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.orch = _orchestrator(Path(self._tmp.name) / "r.db")
def tearDown(self) -> None:
self._tmp.cleanup()
def _old(self, source_ip: str) -> str:
rec = self.orch.registry.register(source_ip)
with closing(sqlite3.connect(self.orch.registry.db_path)) as conn:
conn.execute(
"UPDATE orchestrator_bottles SET created_at = 0.0 WHERE bottle_id = ?",
(rec.bottle_id,))
conn.commit()
return rec.bottle_id
def test_reaps_absent_and_reports_ids(self) -> None:
dead = self._old("10.0.0.1")
alive = self._old("10.0.0.2")
status, payload = dispatch(
self.orch, "POST", "/reconcile", _body({"live_source_ips": ["10.0.0.2"]}))
self.assertEqual(200, status)
self.assertEqual([dead], payload["reaped"])
self.assertIsNone(self.orch.registry.get(dead))
self.assertIsNotNone(self.orch.registry.get(alive))
def test_missing_live_source_ips_is_400(self) -> None:
status, _ = dispatch(self.orch, "POST", "/reconcile", _body({}))
self.assertEqual(400, status)
def test_grace_seconds_is_honoured(self) -> None:
"""A grace window wide enough to cover the row protects it."""
self.orch.registry.register("10.0.0.3")
status, payload = dispatch(
self.orch, "POST", "/reconcile",
_body({"live_source_ips": [], "grace_seconds": 3600}))
self.assertEqual(200, status)
self.assertEqual([], payload["reaped"])
def test_non_string_entries_are_ignored(self) -> None:
dead = self._old("10.0.0.4")
status, payload = dispatch(
self.orch, "POST", "/reconcile",
_body({"live_source_ips": [None, 7, "10.0.0.9"]}))
self.assertEqual(200, status)
self.assertEqual([dead], payload["reaped"])
-80
View File
@@ -4,9 +4,7 @@ from __future__ import annotations
import sqlite3
import tempfile
import time
import unittest
from contextlib import closing
from pathlib import Path
from bot_bottle.orchestrator.registry import (
@@ -171,81 +169,3 @@ class TestRegistryStore(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestReapAbsent(unittest.TestCase):
"""`reap_absent` — the self-heal for rows whose bottle is gone.
An orphan is not merely untidy: source IPs get recycled, and
`by_source_ip` fail-closes on ambiguity, so a leftover row at a reused
address resolves *no* policy for the next bottle that lands there and
every host it asks for is denied.
"""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db = Path(self._tmp.name) / "registry.db"
self.store = RegistryStore(self.db)
self.store.migrate()
def tearDown(self) -> None:
self._tmp.cleanup()
def _aged(self, source_ip: str, *, age: float) -> BottleRecord:
"""Register a bottle and backdate it past the grace window."""
rec = self.store.register(source_ip)
with closing(sqlite3.connect(self.db)) as conn:
conn.execute(
"UPDATE orchestrator_bottles SET created_at = ? WHERE bottle_id = ?",
(time.time() - age, rec.bottle_id),
)
conn.commit()
return rec
def test_reaps_row_with_no_live_container(self) -> None:
gone = self._aged("10.243.0.9", age=600)
reaped = self.store.reap_absent([])
self.assertEqual([gone.bottle_id], [r.bottle_id for r in reaped])
self.assertIsNone(self.store.get(gone.bottle_id))
def test_keeps_row_whose_ip_is_live(self) -> None:
alive = self._aged("10.243.0.9", age=600)
self.assertEqual([], self.store.reap_absent(["10.243.0.9"]))
self.assertIsNotNone(self.store.get(alive.bottle_id))
def test_grace_window_protects_an_in_flight_launch(self) -> None:
"""A bottle registered moments ago is never reaped, even though the
caller's enumeration didn't see its address yet."""
fresh = self.store.register("10.243.0.10")
self.assertEqual([], self.store.reap_absent([]))
self.assertIsNotNone(self.store.get(fresh.bottle_id))
def test_reaping_the_orphan_unbricks_the_reused_address(self) -> None:
"""The regression this exists for: an orphan at an address that vmnet
later hands to a new bottle makes `by_source_ip` ambiguous, so the new
bottle resolves no policy at all."""
orphan = self._aged("10.243.0.11", age=600)
# A new bottle lands on the recycled address. Force the row in directly
# so `register`'s own supersede sweep doesn't mask the ambiguity.
with closing(sqlite3.connect(self.db)) as conn:
conn.execute(
"INSERT INTO orchestrator_bottles "
"(bottle_id, source_ip, identity_token, state, created_at, metadata, policy) "
"VALUES ('newbottle', '10.243.0.11', 'tok-new', 'active', ?, '', 'routes: []')",
(time.time(),),
)
conn.commit()
self.assertIsNone(self.store.by_source_ip("10.243.0.11")) # bricked
reaped = self.store.reap_absent(["10.243.0.11"], grace_seconds=60)
self.assertEqual([orphan.bottle_id], [r.bottle_id for r in reaped])
rec = self.store.by_source_ip("10.243.0.11")
assert rec is not None
self.assertEqual("newbottle", rec.bottle_id)
def test_ignores_empty_ips_in_the_live_set(self) -> None:
gone = self._aged("10.243.0.12", age=600)
self.assertEqual(
[gone.bottle_id],
[r.bottle_id for r in self.store.reap_absent(["", "10.243.0.99"])],
)
-54
View File
@@ -4,10 +4,8 @@ from __future__ import annotations
import json
import secrets
import sqlite3
import tempfile
import unittest
from contextlib import closing
from pathlib import Path
from unittest.mock import patch
@@ -306,55 +304,3 @@ class TestOrchestratorSupervise(unittest.TestCase):
if __name__ == "__main__":
unittest.main()
class TestOrchestratorReconcile(unittest.TestCase):
"""`reconcile` — drop rows for bottles that are no longer running."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.secret = secrets.token_bytes(16)
self.db = Path(self._tmp.name) / "r.db"
self.store = RegistryStore(self.db)
self.store.migrate()
self.broker = StubBroker(self.secret)
self.orch = Orchestrator(self.store, self.broker, self.secret)
def tearDown(self) -> None:
self._tmp.cleanup()
def _age_all(self, seconds: float) -> None:
"""Backdate every row past the reap grace window."""
with closing(sqlite3.connect(self.db)) as conn:
conn.execute(
"UPDATE orchestrator_bottles SET created_at = created_at - ?", (seconds,))
conn.commit()
def test_reaps_dead_bottle_and_forgets_its_tokens(self) -> None:
dead = self.orch.launch_bottle("10.243.0.1", tokens={"EGRESS_TOKEN_0": "s3cret"})
live = self.orch.launch_bottle("10.243.0.2", tokens={"EGRESS_TOKEN_0": "keep"})
self._age_all(600)
self.assertEqual([dead.bottle_id], self.orch.reconcile(["10.243.0.2"]))
self.assertIsNone(self.store.get(dead.bottle_id))
self.assertIsNotNone(self.store.get(live.bottle_id))
# The in-memory egress credential goes with the row.
self.assertEqual({}, self.orch.tokens_for(dead.bottle_id))
self.assertEqual({"EGRESS_TOKEN_0": "keep"}, self.orch.tokens_for(live.bottle_id))
def test_reconcile_does_not_broker_a_teardown(self) -> None:
"""The container is already gone — there is nothing to stop, and a
broker error must not stop the sweep clearing the row."""
self.orch.launch_bottle("10.243.0.1")
self._age_all(600)
self.broker.launched.clear()
self.orch.reconcile([])
self.assertEqual([], self.broker.torn_down)
def test_reconcile_keeps_everything_when_all_are_live(self) -> None:
a = self.orch.launch_bottle("10.243.0.1")
b = self.orch.launch_bottle("10.243.0.2")
self._age_all(600)
self.assertEqual([], self.orch.reconcile(["10.243.0.1", "10.243.0.2"]))
self.assertIsNotNone(self.store.get(a.bottle_id))
self.assertIsNotNone(self.store.get(b.bottle_id))