Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fd295d4c14 | |||
| f33566941b | |||
| c7c3a79028 | |||
| bb1776a858 | |||
| a24fe0264d | |||
| 105538d3a6 | |||
| ffda40abae | |||
| 7dcce2ff12 | |||
| 31a7efc0ed | |||
| a25ea7c188 | |||
| 3dbf1780b4 | |||
| ff4da6f41e | |||
| 3bb90da11c | |||
| 0146450951 |
@@ -30,7 +30,7 @@ from pathlib import Path
|
||||
|
||||
from ...log import info
|
||||
from .. import EnumerationError
|
||||
from . import util
|
||||
from . import lifecycle_lock, util
|
||||
from .bottle_cleanup_plan import FirecrackerBottleCleanupPlan
|
||||
|
||||
|
||||
@@ -126,12 +126,42 @@ def prepare_cleanup() -> FirecrackerBottleCleanupPlan:
|
||||
|
||||
|
||||
def cleanup(plan: FirecrackerBottleCleanupPlan) -> None:
|
||||
for pid in plan.vm_pids:
|
||||
info(f"kill firecracker VM pid {pid}")
|
||||
"""Revalidate the preview under the launch lock, then remove its survivors."""
|
||||
with lifecycle_lock.hold():
|
||||
fresh = prepare_cleanup()
|
||||
approved_pids = set(plan.vm_pids).intersection(fresh.vm_pids)
|
||||
approved_dirs = set(plan.run_dirs).intersection(fresh.run_dirs)
|
||||
for pid in sorted(approved_pids):
|
||||
_terminate_orphan(pid, _run_root())
|
||||
for path in sorted(approved_dirs):
|
||||
info(f"rm -rf {path}")
|
||||
shutil.rmtree(path, ignore_errors=True)
|
||||
|
||||
|
||||
def _terminate_orphan(pid: int, run_root: Path) -> None:
|
||||
"""Signal exactly the process identity that still owns an orphan config."""
|
||||
try:
|
||||
pidfd = os.pidfd_open(pid)
|
||||
except ProcessLookupError:
|
||||
return
|
||||
except OSError as exc:
|
||||
raise EnumerationError(
|
||||
f"could not pin Firecracker pid {pid} for cleanup: {exc}"
|
||||
) from exc
|
||||
try:
|
||||
try:
|
||||
os.kill(pid, signal.SIGTERM)
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
for path in plan.run_dirs:
|
||||
info(f"rm -rf {path}")
|
||||
shutil.rmtree(path, ignore_errors=True)
|
||||
raw = Path(f"/proc/{pid}/cmdline").read_bytes()
|
||||
except FileNotFoundError:
|
||||
return
|
||||
except OSError as exc:
|
||||
raise EnumerationError(
|
||||
f"could not revalidate Firecracker pid {pid}: {exc}"
|
||||
) from exc
|
||||
command = raw.replace(b"\0", b" ").decode(errors="replace")
|
||||
run_dir = _run_dir_of(command, run_root)
|
||||
if run_dir is None or run_dir.is_dir():
|
||||
return
|
||||
info(f"kill firecracker VM pid {pid}")
|
||||
signal.pidfd_send_signal(pidfd, signal.SIGTERM)
|
||||
finally:
|
||||
os.close(pidfd)
|
||||
|
||||
@@ -46,7 +46,7 @@ from ...log import die, info, warn
|
||||
from ...supervisor.types import SUPERVISE_PORT
|
||||
from ..docker.egress import EGRESS_PORT
|
||||
from ..util import AGENT_CA_BUNDLE, AGENT_CA_PATH
|
||||
from . import firecracker_vm, image_builder, isolation_probe, netpool, util
|
||||
from . import firecracker_vm, image_builder, isolation_probe, lifecycle_lock, netpool, util
|
||||
from .bottle import FirecrackerBottle
|
||||
from .bottle_plan import FirecrackerBottlePlan
|
||||
from ...orchestrator.store.config_store import resolve_teardown_timeout
|
||||
@@ -164,25 +164,29 @@ def launch(
|
||||
)
|
||||
|
||||
# Step 6: build the per-bottle rootfs + SSH key, then boot.
|
||||
run_dir = util.cache_dir() / "run" / plan.slug
|
||||
run_dir.mkdir(parents=True, exist_ok=True)
|
||||
# Remove the run dir on teardown so the per-bottle rootfs.ext4 (~1G)
|
||||
# doesn't leak. Registered before vm.terminate below so it runs *after*
|
||||
# it (ExitStack is LIFO): the VM is gone before we rm its rootfs.
|
||||
stack.callback(lambda: shutil.rmtree(run_dir, ignore_errors=True))
|
||||
rootfs = run_dir / "rootfs.ext4"
|
||||
util.build_rootfs_ext4(agent_base, rootfs)
|
||||
private_key, pubkey = util.generate_keypair(run_dir)
|
||||
# Cleanup takes the same lock while refreshing its process snapshot.
|
||||
# Hold it until the VMM exists so a newly-created run dir can never be
|
||||
# mistaken for an orphan in the build-before-boot window.
|
||||
with lifecycle_lock.hold():
|
||||
run_dir = util.cache_dir() / "run" / plan.slug
|
||||
run_dir.mkdir(parents=True, exist_ok=True)
|
||||
# Remove the run dir on teardown so the per-bottle rootfs.ext4 (~1G)
|
||||
# doesn't leak. Registered before vm.terminate below so it runs
|
||||
# *after* it (ExitStack is LIFO).
|
||||
stack.callback(lambda: shutil.rmtree(run_dir, ignore_errors=True))
|
||||
rootfs = run_dir / "rootfs.ext4"
|
||||
util.build_rootfs_ext4(agent_base, rootfs)
|
||||
private_key, pubkey = util.generate_keypair(run_dir)
|
||||
|
||||
vm = firecracker_vm.boot(
|
||||
name=plan.container_name,
|
||||
rootfs=rootfs,
|
||||
tap=slot.iface,
|
||||
guest_ip=slot.guest_ip,
|
||||
host_ip=slot.host_ip,
|
||||
pubkey=pubkey,
|
||||
run_dir=run_dir,
|
||||
)
|
||||
vm = firecracker_vm.boot(
|
||||
name=plan.container_name,
|
||||
rootfs=rootfs,
|
||||
tap=slot.iface,
|
||||
guest_ip=slot.guest_ip,
|
||||
host_ip=slot.host_ip,
|
||||
pubkey=pubkey,
|
||||
run_dir=run_dir,
|
||||
)
|
||||
stack.callback(vm.terminate)
|
||||
firecracker_vm.wait_for_ssh(vm, private_key)
|
||||
persist_env_var_secret(private_key, slot.guest_ip, ctx.env_var_secret)
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
"""Serialize Firecracker run-directory creation with orphan cleanup."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import fcntl
|
||||
from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
from typing import Generator
|
||||
|
||||
from . import util
|
||||
|
||||
|
||||
def _lock_path() -> Path:
|
||||
return util.cache_dir() / "run.lifecycle.lock"
|
||||
|
||||
|
||||
@contextmanager
|
||||
def hold() -> Generator[None]:
|
||||
"""Exclude cleanup while a launch directory lacks a visible VMM."""
|
||||
path = _lock_path()
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
with path.open("a", encoding="utf-8") as handle:
|
||||
fcntl.flock(handle, fcntl.LOCK_EX)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
fcntl.flock(handle, fcntl.LOCK_UN)
|
||||
|
||||
|
||||
__all__ = ["hold"]
|
||||
@@ -4,7 +4,8 @@ from __future__ import annotations
|
||||
|
||||
import subprocess
|
||||
|
||||
from ...log import info, warn
|
||||
from .. import EnumerationError
|
||||
from ...log import info
|
||||
from . import util as container_mod
|
||||
from .bottle_cleanup_plan import MacosContainerBottleCleanupPlan
|
||||
|
||||
@@ -19,8 +20,8 @@ def _list_prefixed_containers() -> list[str]:
|
||||
check=False,
|
||||
)
|
||||
if result.returncode != 0:
|
||||
warn(f"container list failed: {result.stderr.strip()}")
|
||||
return []
|
||||
detail = result.stderr.strip() or f"exit {result.returncode}"
|
||||
raise EnumerationError(f"container list failed: {detail}")
|
||||
return sorted(
|
||||
name for name in (line.strip() for line in result.stdout.splitlines())
|
||||
if name.startswith(_PREFIX)
|
||||
@@ -35,7 +36,8 @@ def _list_prefixed_networks() -> list[str]:
|
||||
check=False,
|
||||
)
|
||||
if result.returncode != 0:
|
||||
return []
|
||||
detail = result.stderr.strip() or f"exit {result.returncode}"
|
||||
raise EnumerationError(f"container network list failed: {detail}")
|
||||
return sorted(
|
||||
name for name in (line.strip() for line in result.stdout.splitlines())
|
||||
if name.startswith(_PREFIX)
|
||||
|
||||
@@ -52,7 +52,11 @@ def cmd_cleanup(_argv: list[str]) -> int:
|
||||
info("cleanup: skipped")
|
||||
return 0
|
||||
|
||||
for name, backend, plan in prepared:
|
||||
# Confirmation authorizes a fresh authoritative snapshot, not blind use of
|
||||
# identities that may have changed while the operator reviewed the preview.
|
||||
refreshed = [(name, backend, backend.prepare_cleanup())
|
||||
for name, backend, _plan in prepared]
|
||||
for name, backend, plan in refreshed:
|
||||
if plan.empty:
|
||||
continue
|
||||
backend.cleanup(plan)
|
||||
|
||||
@@ -136,10 +136,18 @@ def _pump(name: str, stream: IO[bytes]) -> None:
|
||||
"""Read lines from `stream`, prefix with `[name]`, write to
|
||||
stdout. Runs in its own thread per child; daemon=True so a
|
||||
blocked read doesn't keep the process alive after main exits."""
|
||||
for raw in iter(stream.readline, b""):
|
||||
line = raw.decode("utf-8", errors="replace").rstrip("\n")
|
||||
sys.stdout.write(f"[{name}] {line}\n")
|
||||
sys.stdout.flush()
|
||||
try:
|
||||
for raw in iter(stream.readline, b""):
|
||||
line = raw.decode("utf-8", errors="replace").rstrip("\n")
|
||||
sys.stdout.write(f"[{name}] {line}\n")
|
||||
sys.stdout.flush()
|
||||
except (OSError, ValueError) as exc:
|
||||
# The manager closes a dead child's pipe after wait() and before a
|
||||
# restart. A pump can be between readline calls at that exact moment;
|
||||
# closed-stream errors are normal completion, not uncaught thread
|
||||
# failures. Preserve genuinely unexpected I/O diagnostics.
|
||||
if not stream.closed:
|
||||
_log(f"{name} output pump stopped: {type(exc).__name__}: {exc}")
|
||||
|
||||
|
||||
def _spawn(spec: _DaemonSpec) -> subprocess.Popen[bytes]:
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
"""Shared resource boundaries for gateway stdlib HTTP services."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import http.server
|
||||
import socket
|
||||
import threading
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Protocol
|
||||
|
||||
|
||||
class Readable(Protocol):
|
||||
def read(self, size: int = -1, /) -> bytes: ...
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class BodyReadError(Exception):
|
||||
status: int
|
||||
message: str
|
||||
|
||||
|
||||
def read_declared_body(
|
||||
stream: Readable,
|
||||
connection: socket.socket,
|
||||
raw_length: str | None,
|
||||
*,
|
||||
maximum: int,
|
||||
timeout_seconds: float,
|
||||
require_length: bool,
|
||||
) -> bytes:
|
||||
"""Validate and read exactly one declared body under a read deadline."""
|
||||
if raw_length is None:
|
||||
if require_length:
|
||||
raise BodyReadError(411, "Content-Length required")
|
||||
raw_length = "0"
|
||||
try:
|
||||
length = int(raw_length)
|
||||
except ValueError as exc:
|
||||
raise BodyReadError(400, "invalid Content-Length") from exc
|
||||
if length < 0:
|
||||
raise BodyReadError(400, "invalid Content-Length")
|
||||
if length > maximum:
|
||||
raise BodyReadError(413, "request body too large")
|
||||
previous_timeout = connection.gettimeout()
|
||||
deadline = time.monotonic() + timeout_seconds
|
||||
chunks: list[bytes] = []
|
||||
remaining = length
|
||||
try:
|
||||
while remaining:
|
||||
timeout = deadline - time.monotonic()
|
||||
if timeout <= 0:
|
||||
raise BodyReadError(408, "request body read timed out")
|
||||
connection.settimeout(timeout)
|
||||
chunk = stream.read(min(remaining, 64 * 1024))
|
||||
if not chunk:
|
||||
raise BodyReadError(400, "incomplete request body")
|
||||
chunks.append(chunk)
|
||||
remaining -= len(chunk)
|
||||
except TimeoutError as exc:
|
||||
raise BodyReadError(408, "request body read timed out") from exc
|
||||
finally:
|
||||
connection.settimeout(previous_timeout)
|
||||
return b"".join(chunks)
|
||||
|
||||
|
||||
class BoundedThreadingHTTPServer(http.server.ThreadingHTTPServer):
|
||||
"""ThreadingHTTPServer with a hard cap on in-flight request threads."""
|
||||
|
||||
daemon_threads = True
|
||||
|
||||
def __init__( # pylint: disable=consider-using-with
|
||||
self, *args, max_workers: int = 32, **kwargs, # type: ignore[no-untyped-def]
|
||||
):
|
||||
if max_workers < 1:
|
||||
raise ValueError("max_workers must be positive")
|
||||
self._request_slots = threading.BoundedSemaphore(max_workers)
|
||||
super().__init__(*args, **kwargs)
|
||||
|
||||
def process_request(
|
||||
self, request: Any, client_address: Any,
|
||||
) -> None:
|
||||
if not self._request_slots.acquire( # pylint: disable=consider-using-with
|
||||
blocking=False,
|
||||
):
|
||||
try:
|
||||
request.sendall(
|
||||
b"HTTP/1.1 503 Service Unavailable\r\n"
|
||||
b"Content-Length: 0\r\nConnection: close\r\n\r\n"
|
||||
)
|
||||
finally:
|
||||
self.shutdown_request(request)
|
||||
return
|
||||
try:
|
||||
super().process_request(request, client_address)
|
||||
except BaseException:
|
||||
self._request_slots.release()
|
||||
raise
|
||||
|
||||
def process_request_thread(
|
||||
self, request: Any, client_address: Any,
|
||||
) -> None:
|
||||
try:
|
||||
super().process_request_thread(request, client_address)
|
||||
finally:
|
||||
self._request_slots.release()
|
||||
|
||||
|
||||
__all__ = ["BodyReadError", "BoundedThreadingHTTPServer", "read_declared_body"]
|
||||
@@ -22,11 +22,16 @@ import os
|
||||
import subprocess
|
||||
import sys
|
||||
import typing
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
from http.server import BaseHTTPRequestHandler
|
||||
from pathlib import Path
|
||||
from urllib.parse import urlsplit
|
||||
|
||||
from bot_bottle.constants import GIT_GATE_TIMEOUT_SECS, IDENTITY_HEADER
|
||||
from bot_bottle.gateway.bounded_http import (
|
||||
BodyReadError,
|
||||
BoundedThreadingHTTPServer,
|
||||
read_declared_body,
|
||||
)
|
||||
from bot_bottle.gateway.policy_resolver import PolicyResolveError, PolicyResolver
|
||||
|
||||
|
||||
@@ -77,6 +82,8 @@ def resolve_sandbox_root(
|
||||
|
||||
# Bound memory use while still allowing ordinary git push packfiles.
|
||||
MAX_BODY_BYTES = 100 * 1024 * 1024
|
||||
REQUEST_BODY_TIMEOUT_SECONDS = 30.0
|
||||
MAX_REQUEST_WORKERS = 16
|
||||
|
||||
|
||||
class GitHttpHandler(BaseHTTPRequestHandler):
|
||||
@@ -184,19 +191,18 @@ class GitHttpHandler(BaseHTTPRequestHandler):
|
||||
value = self.headers.get(header)
|
||||
if value:
|
||||
env[variable] = value
|
||||
raw_length = self.headers.get("content-length", "0") or "0"
|
||||
try:
|
||||
length = int(raw_length)
|
||||
except ValueError:
|
||||
self.send_error(400, "Bad Content-Length")
|
||||
body = read_declared_body(
|
||||
self.rfile,
|
||||
self.connection,
|
||||
self.headers.get("content-length"),
|
||||
maximum=MAX_BODY_BYTES,
|
||||
timeout_seconds=REQUEST_BODY_TIMEOUT_SECONDS,
|
||||
require_length=False,
|
||||
)
|
||||
except BodyReadError as exc:
|
||||
self.send_error(exc.status, exc.message)
|
||||
return
|
||||
if length < 0:
|
||||
self.send_error(400, "Negative Content-Length")
|
||||
return
|
||||
if length > MAX_BODY_BYTES:
|
||||
self.send_error(413, "Request body too large")
|
||||
return
|
||||
body = self.rfile.read(length) if length else b""
|
||||
proc = subprocess.run(
|
||||
["git", "http-backend"],
|
||||
input=body,
|
||||
@@ -273,7 +279,9 @@ def main() -> int:
|
||||
"(no single-tenant flat-root fallback)\n"
|
||||
)
|
||||
return 1
|
||||
server = ThreadingHTTPServer(("0.0.0.0", port), GitHttpHandler)
|
||||
server = BoundedThreadingHTTPServer(
|
||||
("0.0.0.0", port), GitHttpHandler, max_workers=MAX_REQUEST_WORKERS,
|
||||
)
|
||||
# Resolve each request's sandbox namespace by source IP against the
|
||||
# orchestrator control plane.
|
||||
server.policy_resolver = PolicyResolver(orch_url) # type: ignore[attr-defined]
|
||||
|
||||
@@ -6,9 +6,8 @@ import json
|
||||
from dataclasses import dataclass
|
||||
from typing import Callable, Protocol
|
||||
|
||||
from bot_bottle.gateway.egress.context import resolve_client_context
|
||||
from bot_bottle.gateway.egress.schema import route_to_yaml_dict
|
||||
from bot_bottle.gateway.policy_resolver import PolicyResolver
|
||||
from bot_bottle.gateway.egress.schema import load_config, route_to_yaml_dict
|
||||
from bot_bottle.gateway.policy_resolver import PolicyResolveError, PolicyResolver
|
||||
from bot_bottle.supervisor import types as _sv
|
||||
|
||||
|
||||
@@ -24,6 +23,10 @@ class MethodNotFoundError(Exception):
|
||||
"""Raised when a JSON-RPC method has no MCP handler."""
|
||||
|
||||
|
||||
class RouteResolutionError(Exception):
|
||||
"""The caller's live route table could not be resolved authoritatively."""
|
||||
|
||||
|
||||
Handler = Callable[[dict[str, object]], object]
|
||||
|
||||
|
||||
@@ -60,10 +63,19 @@ def resolved_routes_payload(
|
||||
source_ip: str,
|
||||
identity_token: str,
|
||||
) -> dict[str, object]:
|
||||
"""Render the calling bottle's routes, failing closed to an empty list."""
|
||||
config, _slug, _tokens = resolve_client_context(
|
||||
resolver, source_ip, identity_token,
|
||||
)
|
||||
"""Render an authoritatively resolved route table for the calling bottle."""
|
||||
try:
|
||||
policy, bottle_id, _tokens = resolver.resolve_policy_and_bottle_id(
|
||||
source_ip, identity_token,
|
||||
)
|
||||
except PolicyResolveError as exc:
|
||||
raise RouteResolutionError("orchestrator unavailable") from exc
|
||||
if not bottle_id:
|
||||
raise RouteResolutionError("request source is not attributed to a bottle")
|
||||
try:
|
||||
config = load_config(policy or "")
|
||||
except ValueError as exc:
|
||||
raise RouteResolutionError("resolved policy is invalid") from exc
|
||||
body = json.dumps(
|
||||
{"routes": [route_to_yaml_dict(route) for route in config.routes]},
|
||||
indent=2,
|
||||
@@ -74,6 +86,7 @@ def resolved_routes_payload(
|
||||
__all__ = [
|
||||
"Handlers",
|
||||
"MethodNotFoundError",
|
||||
"RouteResolutionError",
|
||||
"dispatch",
|
||||
"resolved_routes_payload",
|
||||
]
|
||||
|
||||
@@ -51,19 +51,24 @@ from __future__ import annotations
|
||||
import http.server
|
||||
import json
|
||||
import os
|
||||
import socketserver
|
||||
import sys
|
||||
import time
|
||||
import typing
|
||||
from dataclasses import dataclass
|
||||
|
||||
from bot_bottle.constants import IDENTITY_HEADER
|
||||
from bot_bottle.gateway.bounded_http import (
|
||||
BodyReadError,
|
||||
BoundedThreadingHTTPServer,
|
||||
read_declared_body,
|
||||
)
|
||||
from bot_bottle.gateway.egress.schema import load_config
|
||||
from bot_bottle.gateway.egress.types import LOG_OFF
|
||||
from bot_bottle.gateway.policy_resolver import PolicyResolveError, PolicyResolver
|
||||
from bot_bottle.gateway.supervisor.mcp_dispatch import (
|
||||
Handlers as DispatchHandlers,
|
||||
MethodNotFoundError,
|
||||
RouteResolutionError,
|
||||
dispatch,
|
||||
resolved_routes_payload,
|
||||
)
|
||||
@@ -570,6 +575,8 @@ def format_unknown_proposal_text(proposal_id: str) -> str:
|
||||
# Max request body the server accepts. 1 MB is well above any realistic
|
||||
# routes.yaml proposal.
|
||||
MAX_BODY_BYTES = 1 * 1024 * 1024
|
||||
REQUEST_BODY_TIMEOUT_SECONDS = 10.0
|
||||
MAX_REQUEST_WORKERS = 32
|
||||
|
||||
|
||||
class MCPHandler(http.server.BaseHTTPRequestHandler):
|
||||
@@ -592,19 +599,18 @@ class MCPHandler(http.server.BaseHTTPRequestHandler):
|
||||
self._write_text(405, "use POST for MCP requests\n")
|
||||
|
||||
def do_POST(self) -> None:
|
||||
length_header = self.headers.get("Content-Length")
|
||||
if length_header is None:
|
||||
self._write_text(411, "Content-Length required\n")
|
||||
return
|
||||
try:
|
||||
length = int(length_header)
|
||||
except ValueError:
|
||||
self._write_text(400, "invalid Content-Length\n")
|
||||
body = read_declared_body(
|
||||
self.rfile,
|
||||
self.connection,
|
||||
self.headers.get("Content-Length"),
|
||||
maximum=MAX_BODY_BYTES,
|
||||
timeout_seconds=REQUEST_BODY_TIMEOUT_SECONDS,
|
||||
require_length=True,
|
||||
)
|
||||
except BodyReadError as exc:
|
||||
self._write_text(exc.status, exc.message + "\n")
|
||||
return
|
||||
if length < 0 or length > MAX_BODY_BYTES:
|
||||
self._write_text(413, "request body too large\n")
|
||||
return
|
||||
body = self.rfile.read(length)
|
||||
|
||||
try:
|
||||
req = parse_jsonrpc(body)
|
||||
@@ -660,20 +666,25 @@ class MCPHandler(http.server.BaseHTTPRequestHandler):
|
||||
identity_token=self._identity_token(),
|
||||
)
|
||||
|
||||
return dispatch(
|
||||
req,
|
||||
DispatchHandlers(
|
||||
initialize=handle_initialize,
|
||||
tools_list=handle_tools_list,
|
||||
list_routes=lambda _params: resolved_routes_payload(
|
||||
def list_routes(_params: dict[str, object]) -> object:
|
||||
try:
|
||||
return resolved_routes_payload(
|
||||
self._resolver_or_fail(),
|
||||
self.client_address[0],
|
||||
self._identity_token(),
|
||||
),
|
||||
check_proposal=check,
|
||||
propose=propose,
|
||||
),
|
||||
)
|
||||
)
|
||||
except RouteResolutionError as exc:
|
||||
raise _RpcInternalError(
|
||||
f"could not resolve live egress routes: {exc}"
|
||||
) from exc
|
||||
|
||||
return dispatch(req, DispatchHandlers(
|
||||
initialize=handle_initialize,
|
||||
tools_list=handle_tools_list,
|
||||
list_routes=list_routes,
|
||||
check_proposal=check,
|
||||
propose=propose,
|
||||
))
|
||||
|
||||
def _identity_token(self) -> str:
|
||||
"""The agent's per-bottle identity token from the request header (the
|
||||
@@ -711,7 +722,7 @@ class MCPHandler(http.server.BaseHTTPRequestHandler):
|
||||
self.wfile.write(encoded)
|
||||
|
||||
|
||||
class MCPServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
|
||||
class MCPServer(BoundedThreadingHTTPServer):
|
||||
allow_reuse_address = True
|
||||
daemon_threads = True
|
||||
config: ServerConfig = ServerConfig()
|
||||
@@ -720,6 +731,9 @@ class MCPServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
|
||||
# closed per request (see `_resolver_or_fail`).
|
||||
policy_resolver: "PolicyResolver | None" = None
|
||||
|
||||
def __init__(self, *args, **kwargs): # type: ignore[no-untyped-def]
|
||||
super().__init__(*args, max_workers=MAX_REQUEST_WORKERS, **kwargs)
|
||||
|
||||
|
||||
# --- Entry point -----------------------------------------------------------
|
||||
|
||||
|
||||
@@ -13,7 +13,9 @@ resource-consuming boundary revalidate the assumptions it acts on. This
|
||||
finishes the focused quality work begun under #444 without broad rewrites:
|
||||
cleanup cannot act on stale identities, policy introspection cannot publish a
|
||||
fabricated empty policy, gateway servers bound untrusted work, and daemon
|
||||
shutdown does not emit uncaught background-thread failures.
|
||||
shutdown does not emit uncaught background-thread failures. Shared
|
||||
control-plane storage and gateway credential provisioning also enforce their
|
||||
filesystem security contract before sensitive data is written.
|
||||
|
||||
## Problem
|
||||
|
||||
@@ -36,6 +38,27 @@ misleading behavior:
|
||||
destructive plan.
|
||||
5. Gateway log-pump threads race stream closure during shutdown and emit
|
||||
uncaught exceptions even when shutdown otherwise succeeds.
|
||||
6. Firecracker discovers VMs through whitespace-split `pgrep -a` output.
|
||||
A configured cache path containing spaces can hide a live VM from the
|
||||
snapshot and make its run directory appear orphaned.
|
||||
7. Docker cleanup asks compose for its project snapshot in best-effort mode.
|
||||
A transient query failure can therefore become an empty stopped-project
|
||||
set and authorize deletion of associated state directories.
|
||||
8. Firecracker artifact downloads and registry publication have no network
|
||||
deadline, so an unresponsive registry can hold setup or release work
|
||||
indefinitely.
|
||||
9. Authenticated secret blobs select the unauthenticated legacy decoder when
|
||||
their in-band version prefix is changed, allowing storage tampering to
|
||||
bypass tag verification.
|
||||
10. Cleanup executes the entire post-confirmation snapshot rather than the
|
||||
intersection with what the operator saw, and mutation failures are not
|
||||
reflected in the command result.
|
||||
11. Git smart-HTTP can retain sixteen 100 MiB request bodies concurrently,
|
||||
cleanup mutations have no subprocess deadline, and Firecracker signalling
|
||||
failures bypass shared mutation accounting.
|
||||
12. SQLite creates the shared control-plane database before its mode is
|
||||
restricted, then suppresses permission-repair failures. Gateway transports
|
||||
also differ in whether copied deploy-key modes are preserved.
|
||||
|
||||
These are one design problem: state used to authorize deletion, replacement,
|
||||
or resource allocation must be authoritative at the point of use.
|
||||
@@ -46,6 +69,8 @@ or resource allocation must be authoritative at the point of use.
|
||||
appeared in a pre-confirmation snapshot.
|
||||
- Firecracker cleanup proves immediately before action that a PID is still the
|
||||
same Firecracker process and that a run directory is still orphaned.
|
||||
- Firecracker process discovery reads NUL-delimited argv from `/proc`; paths
|
||||
are never reconstructed from whitespace-delimited process listings.
|
||||
- All backend cleanup discovery primitives raise a typed enumeration error on
|
||||
operational failure. No backend may independently continue from a partial
|
||||
snapshot.
|
||||
@@ -60,6 +85,22 @@ or resource allocation must be authoritative at the point of use.
|
||||
callers because bottles themselves are untrusted.
|
||||
- Gateway child-output pumping treats expected stream closure during shutdown
|
||||
as completion while preserving diagnostics for unexpected failures.
|
||||
- Artifact pull, existence-check, and publication requests use explicit
|
||||
network deadlines.
|
||||
- Persisted secrets accept only the authenticated format. The schema migration
|
||||
intentionally clears legacy rows; local agents are reprovisioned rather
|
||||
than retaining a ciphertext-controlled downgrade path.
|
||||
- Cleanup executes only resources present in both the displayed and current
|
||||
authoritative plans, attempts every approved mutation, and returns failure
|
||||
when any mutation does not complete.
|
||||
- Git request bodies spool to disk behind a separate heavy-work semaphore;
|
||||
cleanup commands have configurable deadlines; Firecracker signalling
|
||||
failures aggregate while identity-verification uncertainty still aborts.
|
||||
- The shared database directory and file are private before SQLite writes any
|
||||
control-plane state; an inability to enforce those modes aborts startup.
|
||||
- Gateway credential directories and files receive explicit private modes
|
||||
inside the gateway, independent of Docker, Apple Container, or SSH copy
|
||||
semantics.
|
||||
- Unit tests cover PID/path reuse, partial backend enumeration, transient
|
||||
policy resolution failure, slow bodies, concurrency saturation, and stream
|
||||
closure races.
|
||||
@@ -94,6 +135,13 @@ Backend-specific primitives define how to identify a resource. Firecracker
|
||||
uses process start identity plus canonical config/run paths; container
|
||||
backends use authoritative CLI queries and stable resource names/labels.
|
||||
|
||||
Container engines expose destructive name-based commands without a portable
|
||||
compare-and-delete operation. Cleanup therefore refreshes after confirmation
|
||||
and requires every discovery query to succeed, minimizing but not claiming to
|
||||
eliminate the final name-reuse race. A future engine-specific stable-ID
|
||||
primitive may close that residual window without moving control flow back
|
||||
into each backend.
|
||||
|
||||
### Enforcement state versus introspection state
|
||||
|
||||
Egress enforcement retains its deny-all fallback because uncertainty must not
|
||||
@@ -115,6 +163,17 @@ The gateway output pump catches only stream-closure exceptions expected after
|
||||
the supervisor closes child pipes. Other I/O failures remain visible and are
|
||||
reported through the supervisor's normal diagnostic channel.
|
||||
|
||||
### Shared filesystem security
|
||||
|
||||
The common SQLite store owns database creation for every backend. It creates
|
||||
the parent directory and an empty database with private modes before opening
|
||||
SQLite, repairs existing modes, verifies the resulting state, and propagates
|
||||
every enforcement failure. Backend launchers do not duplicate this policy.
|
||||
|
||||
The backend-neutral gateway provisioner likewise applies directory and file
|
||||
modes after transport copies complete. This avoids relying on copy behavior
|
||||
that differs among Docker, Apple Container, and Firecracker's SSH transport.
|
||||
|
||||
## Implementation chunks
|
||||
|
||||
1. Existing fail-closed security and backend enumeration fixes.
|
||||
@@ -124,6 +183,14 @@ reported through the supervisor's normal diagnostic channel.
|
||||
5. Shared cleanup refresh/revalidation plus authoritative macOS discovery.
|
||||
6. Strict supervisor introspection and bounded supervisor/Git HTTP work.
|
||||
7. Gateway shutdown log-pump closure handling.
|
||||
8. Lossless Firecracker process identities, authoritative Docker cleanup
|
||||
queries, and bounded Firecracker artifact transfers.
|
||||
9. Mandatory authenticated secret storage, shared cleanup-plan intersection
|
||||
and mutation accounting, and contained Git backend process failures.
|
||||
10. Disk-spooled and separately bounded Git bodies, cleanup command deadlines,
|
||||
and classified Firecracker signalling failures.
|
||||
11. Fail-closed shared database creation and backend-neutral gateway credential
|
||||
permissions.
|
||||
|
||||
## Open questions
|
||||
|
||||
|
||||
@@ -41,8 +41,8 @@ class TestCmdCleanup(unittest.TestCase):
|
||||
):
|
||||
self.assertEqual(0, cmd.cmd_cleanup([]))
|
||||
|
||||
docker.prepare_cleanup.assert_called_once()
|
||||
fc.prepare_cleanup.assert_called_once()
|
||||
self.assertEqual(2, docker.prepare_cleanup.call_count)
|
||||
self.assertEqual(2, fc.prepare_cleanup.call_count)
|
||||
docker.cleanup.assert_called_once_with(docker_plan)
|
||||
fc.cleanup.assert_called_once_with(fc_plan)
|
||||
|
||||
@@ -68,7 +68,7 @@ class TestCmdCleanup(unittest.TestCase):
|
||||
):
|
||||
self.assertEqual(0, cmd.cmd_cleanup([]))
|
||||
|
||||
docker.prepare_cleanup.assert_called_once()
|
||||
self.assertEqual(2, docker.prepare_cleanup.call_count)
|
||||
docker.cleanup.assert_called_once_with(docker_plan)
|
||||
macos.prepare_cleanup.assert_not_called()
|
||||
|
||||
@@ -135,6 +135,25 @@ class TestCmdCleanup(unittest.TestCase):
|
||||
docker.cleanup.assert_called_once_with(docker_plan)
|
||||
fc.cleanup.assert_not_called()
|
||||
|
||||
def test_executes_refreshed_plan_after_confirmation(self):
|
||||
backend = MagicMock()
|
||||
preview = MagicMock(empty=False)
|
||||
refreshed = MagicMock(empty=False)
|
||||
backend.prepare_cleanup.side_effect = [preview, refreshed]
|
||||
|
||||
with patch.object(
|
||||
cmd, "known_backend_names", return_value=("firecracker",),
|
||||
), patch.object(
|
||||
cmd, "get_bottle_backend", return_value=backend,
|
||||
), patch.object(
|
||||
cmd, "has_backend", return_value=True,
|
||||
), patch.object(
|
||||
cmd, "_prompt_yes", return_value=True,
|
||||
):
|
||||
self.assertEqual(0, cmd.cmd_cleanup([]))
|
||||
|
||||
backend.cleanup.assert_called_once_with(refreshed)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -132,18 +132,46 @@ class TestCleanupRemoval(unittest.TestCase):
|
||||
vm_pids=(101,),
|
||||
run_dirs=("/run/dev-x",),
|
||||
)
|
||||
with patch.object(fc_cleanup.os, "kill") as kill, \
|
||||
with patch.object(fc_cleanup, "prepare_cleanup", return_value=plan), \
|
||||
patch.object(fc_cleanup, "_run_root", return_value=Path("/run")), \
|
||||
patch.object(fc_cleanup, "_terminate_orphan") as terminate, \
|
||||
patch.object(fc_cleanup.shutil, "rmtree") as rmtree, \
|
||||
patch.object(fc_cleanup, "info"):
|
||||
fc_cleanup.cleanup(plan)
|
||||
kill.assert_called_once()
|
||||
terminate.assert_called_once_with(101, Path("/run"))
|
||||
rmtree.assert_called_once_with("/run/dev-x", ignore_errors=True)
|
||||
|
||||
def test_cleanup_tolerates_dead_pid(self):
|
||||
plan = FirecrackerBottleCleanupPlan(vm_pids=(999,))
|
||||
with patch.object(fc_cleanup.os, "kill", side_effect=ProcessLookupError), \
|
||||
patch.object(fc_cleanup, "info"):
|
||||
fc_cleanup.cleanup(plan) # must not raise
|
||||
def test_cleanup_skips_resources_no_longer_in_refreshed_plan(self):
|
||||
preview = FirecrackerBottleCleanupPlan(
|
||||
vm_pids=(999,), run_dirs=("/run/reused",),
|
||||
)
|
||||
with patch.object(
|
||||
fc_cleanup, "prepare_cleanup",
|
||||
return_value=FirecrackerBottleCleanupPlan(),
|
||||
), patch.object(fc_cleanup, "_terminate_orphan") as terminate, \
|
||||
patch.object(fc_cleanup.shutil, "rmtree") as rmtree:
|
||||
fc_cleanup.cleanup(preview)
|
||||
terminate.assert_not_called()
|
||||
rmtree.assert_not_called()
|
||||
|
||||
def test_pidfd_prevents_pid_reuse_from_signalling_unrelated_process(self):
|
||||
with patch.object(fc_cleanup.os, "pidfd_open", return_value=7), \
|
||||
patch.object(
|
||||
fc_cleanup.Path, "read_bytes",
|
||||
return_value=b"/usr/bin/python\0worker.py\0",
|
||||
), patch.object(fc_cleanup.signal, "pidfd_send_signal") as send, \
|
||||
patch.object(fc_cleanup.os, "close"):
|
||||
fc_cleanup._terminate_orphan(101, Path("/run"))
|
||||
send.assert_not_called()
|
||||
|
||||
def test_pidfd_signals_revalidated_orphan(self):
|
||||
command = b"firecracker\0--config-file\0/run/gone/config.json\0"
|
||||
with patch.object(fc_cleanup.os, "pidfd_open", return_value=7), \
|
||||
patch.object(fc_cleanup.Path, "read_bytes", return_value=command), \
|
||||
patch.object(fc_cleanup.signal, "pidfd_send_signal") as send, \
|
||||
patch.object(fc_cleanup.os, "close"), patch.object(fc_cleanup, "info"):
|
||||
fc_cleanup._terminate_orphan(101, Path("/run"))
|
||||
send.assert_called_once_with(7, fc_cleanup.signal.SIGTERM)
|
||||
|
||||
|
||||
class TestCleanupPlan(unittest.TestCase):
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
"""Unit tests for shared gateway stdlib HTTP resource boundaries."""
|
||||
# pylint: disable=protected-access
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import io
|
||||
import socket
|
||||
import unittest
|
||||
|
||||
from bot_bottle.gateway.bounded_http import (
|
||||
BodyReadError,
|
||||
BoundedThreadingHTTPServer,
|
||||
read_declared_body,
|
||||
)
|
||||
|
||||
|
||||
class _Handler:
|
||||
pass
|
||||
|
||||
|
||||
class _TimeoutStream:
|
||||
def read(self, _size: int = -1, /) -> bytes:
|
||||
raise TimeoutError
|
||||
|
||||
|
||||
class TestDeclaredBody(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.left, self.right = socket.socketpair()
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self.left.close()
|
||||
self.right.close()
|
||||
|
||||
def test_rejects_incomplete_body(self) -> None:
|
||||
with self.assertRaisesRegex(BodyReadError, "incomplete"):
|
||||
read_declared_body(
|
||||
io.BytesIO(b"short"),
|
||||
self.left,
|
||||
"10",
|
||||
maximum=100,
|
||||
timeout_seconds=1,
|
||||
require_length=True,
|
||||
)
|
||||
|
||||
def test_maps_read_timeout(self) -> None:
|
||||
with self.assertRaises(BodyReadError) as caught:
|
||||
read_declared_body(
|
||||
_TimeoutStream(),
|
||||
self.left,
|
||||
"1",
|
||||
maximum=100,
|
||||
timeout_seconds=1,
|
||||
require_length=True,
|
||||
)
|
||||
self.assertEqual(408, caught.exception.status)
|
||||
|
||||
|
||||
class TestBoundedServer(unittest.TestCase):
|
||||
def test_saturated_server_rejects_without_spawning_thread(self) -> None:
|
||||
client, peer = socket.socketpair()
|
||||
with BoundedThreadingHTTPServer(
|
||||
("127.0.0.1", 0), _Handler, max_workers=1, # type: ignore[arg-type]
|
||||
) as server:
|
||||
with client, peer:
|
||||
# Directly reserve the only slot to model an in-flight handler.
|
||||
self.assertTrue(
|
||||
server._request_slots.acquire( # pylint: disable=consider-using-with
|
||||
blocking=False,
|
||||
),
|
||||
)
|
||||
try:
|
||||
server.process_request(client, ("127.0.0.1", 1))
|
||||
self.assertIn(b"503 Service Unavailable", peer.recv(1024))
|
||||
finally:
|
||||
server._request_slots.release()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -23,6 +23,7 @@ from bot_bottle.gateway.bootstrap import (
|
||||
_DaemonManager,
|
||||
_argv_for_daemon,
|
||||
_env_for_daemon,
|
||||
_pump,
|
||||
_selected_daemons,
|
||||
)
|
||||
from tests._bin import SLEEP
|
||||
@@ -565,5 +566,26 @@ class TestMainEndToEnd(unittest.TestCase):
|
||||
self.assertIn("no daemons selected", out)
|
||||
|
||||
|
||||
class _FailingStream:
|
||||
def __init__(self, *, closed: bool) -> None:
|
||||
self.closed = closed
|
||||
|
||||
def readline(self) -> bytes:
|
||||
raise ValueError("I/O operation on closed file")
|
||||
|
||||
|
||||
class TestOutputPump(unittest.TestCase):
|
||||
def test_closed_stream_race_is_normal_completion(self) -> None:
|
||||
with patch("bot_bottle.gateway.bootstrap._log") as log:
|
||||
_pump("egress", _FailingStream(closed=True)) # type: ignore[arg-type]
|
||||
log.assert_not_called()
|
||||
|
||||
def test_unexpected_io_failure_is_logged(self) -> None:
|
||||
with patch("bot_bottle.gateway.bootstrap._log") as log:
|
||||
_pump("egress", _FailingStream(closed=False)) # type: ignore[arg-type]
|
||||
log.assert_called_once()
|
||||
self.assertIn("output pump stopped", log.call_args.args[0])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -42,6 +42,22 @@ class TestMacosContainerCleanup(unittest.TestCase):
|
||||
run.call_args_list[1].args[0],
|
||||
)
|
||||
|
||||
def test_container_enumeration_failure_aborts(self):
|
||||
completed = cleanup.subprocess.CompletedProcess(
|
||||
args=[], returncode=1, stdout="", stderr="service unavailable",
|
||||
)
|
||||
with patch.object(cleanup.subprocess, "run", return_value=completed), \
|
||||
self.assertRaisesRegex(EnumerationError, "service unavailable"):
|
||||
cleanup._list_prefixed_containers()
|
||||
|
||||
def test_network_enumeration_failure_aborts(self):
|
||||
completed = cleanup.subprocess.CompletedProcess(
|
||||
args=[], returncode=1, stdout="", stderr="service unavailable",
|
||||
)
|
||||
with patch.object(cleanup.subprocess, "run", return_value=completed), \
|
||||
self.assertRaisesRegex(EnumerationError, "service unavailable"):
|
||||
cleanup._list_prefixed_networks()
|
||||
|
||||
|
||||
class TestMacosContainerEnumerate(unittest.TestCase):
|
||||
"""The backend launches bottles again (PRD 0070), so enumeration is real
|
||||
|
||||
@@ -717,18 +717,26 @@ class TestResolvedRoutesPayload(unittest.TestCase):
|
||||
hosts = {r["host"] for r in data["routes"]}
|
||||
self.assertEqual({"api.anthropic.com", "www.google.com"}, hosts)
|
||||
|
||||
def test_orchestrator_error_fails_closed_to_empty(self) -> None:
|
||||
# resolve_client_context swallows resolver errors → deny-all (empty),
|
||||
# never another bottle's routes.
|
||||
def test_orchestrator_error_is_not_reported_as_empty_policy(self) -> None:
|
||||
with self.assertRaises(supervise_server.RouteResolutionError):
|
||||
resolved_routes_payload(
|
||||
typing.cast(
|
||||
supervise_server.PolicyResolver,
|
||||
_FakeSuperviseResolver(raises=True),
|
||||
),
|
||||
_SRC,
|
||||
_TOK,
|
||||
)
|
||||
|
||||
def test_authoritative_empty_policy_remains_successful(self) -> None:
|
||||
payload = resolved_routes_payload(
|
||||
typing.cast(
|
||||
supervise_server.PolicyResolver,
|
||||
_FakeSuperviseResolver(raises=True),
|
||||
_FakeSuperviseResolver(bottle_id="b1", policy="routes: []\n"),
|
||||
),
|
||||
_SRC,
|
||||
_TOK,
|
||||
)
|
||||
assert payload is not None
|
||||
data = json.loads(payload["content"][0]["text"]) # type: ignore[index]
|
||||
self.assertEqual([], data["routes"])
|
||||
|
||||
|
||||
Reference in New Issue
Block a user