Files
bot-bottle/bot_bottle/orchestrator/service.py
T
didericis-claude 7ab85e9ea6
prd-number-check / require-numbered-prds (pull_request) Failing after 12s
test / integration-docker (pull_request) Successful in 19s
tracker-policy-pr / check-pr (pull_request) Successful in 9s
test / unit (pull_request) Failing after 52s
lint / lint (push) Successful in 1m1s
test / coverage (pull_request) Has been skipped
fix(orchestrator): address review on host control server transport (#468)
Codex review on #496:

- **High — ambiguous delivery no longer orphans a launched bottle.** A
  timeout / dropped response from the host controller is now the ambiguous
  BrokerUnavailableError (distinct from the definite BrokerAuthError /
  BrokerClientError). OrchestratorCore.launch_bottle keeps the registry
  row on the ambiguous case instead of deregistering — deregistering would
  orphan a running container with no record (reconcile reaps rows, never
  containers). The row is left for reconcile to reap iff the bottle is not
  actually live. Definite failures still roll back, so a real failure
  leaves no orphan row.
- **Medium — the privileged endpoint bounds request bodies.** The host
  server rejects an oversized Content-Length with 413 before reading it,
  and sets a per-request socket timeout, so a caller that can merely reach
  the socket (no signed token) can't exhaust memory or a handler thread.

Tests: ambiguous-keep vs definite-rollback in the launch path; the
BrokerUnavailableError/BrokerClientError split in BrokerClient; the 413
body cap + handler error paths (driven in-thread, since daemon request
threads lose coverage); and the __main__ entrypoint broker selection.
Diff-coverage 98%; pyright clean; pylint 9.88.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-26 08:43:00 +00:00

403 lines
18 KiB
Python

"""The per-host orchestrator core (PRD 0070).
`OrchestratorCore` is the single backend-neutral object the control plane talks
to: it owns the registry (runtime state) and brokers agent launches. It
never branches on backend — the `LaunchBroker` abstracts the backend-native
launch, so this same object drives docker / firecracker / apple once a real
broker is wired in.
Launch lifecycle:
* `launch_bottle` mints the bottle (registry: source IP + identity
token), sends a *signed, structured* launch request through the broker,
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 BrokerUnavailableError, LaunchRequest, SubmitBroker, sign_request
from .store.registry_store import DEFAULT_REAP_GRACE_SECONDS, BottleRecord, RegistryStore
from .supervisor import (
AuditEntry,
COMPONENT_FOR_TOOL,
POLL_STATUS_PENDING,
POLL_STATUS_UNKNOWN,
Proposal,
Response,
STATUS_APPROVED,
STATUS_MODIFIED,
STATUS_REJECTED,
TOOL_EGRESS_ALLOW,
TOOL_EGRESS_BLOCK,
Supervisor,
)
from ..util import render_diff
# Operator decision → Response.status. The apply half (egress tools) runs
# for approve/modify only.
_RESPOND_STATUS = {
"approve": STATUS_APPROVED,
"modify": STATUS_MODIFIED,
"reject": STATUS_REJECTED,
}
_APPLY_TOOLS = (TOOL_EGRESS_ALLOW, TOOL_EGRESS_BLOCK)
class OrchestratorCore:
"""Owns the registry + brokers launches, and manages the single
consolidated per-host gateway. Backend-neutral (broker and gateway
abstract the backend-native pieces)."""
def __init__(
self,
registry: RegistryStore,
broker: SubmitBroker,
sign_secret: bytes,
supervisor: Supervisor | None = None,
) -> None:
self.registry = registry
self._broker = broker
self._secret = sign_secret
# The supervise service (queue + audit I/O). Injectable so tests can
# point it at a temp DB; defaults to the host DB.
self._supervisor = supervisor or Supervisor()
# Per-bottle egress auth tokens (env_name -> value), keyed by bottle_id.
# Held **in memory only** — never written to the registry DB — so the
# gateway can inject each bottle's upstream credential without secrets
# at rest. Lost on restart (re-launch re-registers them); the future
# SecretProvider (#355) replaces this with per-request minting.
self._tokens: dict[str, dict[str, str]] = {}
def launch_bottle(
self,
source_ip: str,
*,
image_ref: str = "",
slot: int | None = None,
metadata: str = "",
policy: str = "",
tokens: dict[str, str] | None = None,
env_var_secret: str = "",
) -> BottleRecord:
"""Register a bottle (with its gateway policy + in-memory egress auth
tokens) and broker its launch. Rolls the registry entry back if the
launch doesn't take, so a failure leaves no orphan.
When *env_var_secret* is provided alongside *tokens*, the token values
are also encrypted and written to ``bottled_agent_secrets`` so they can
survive an orchestrator restart (see ``reprovision_from_secret``)."""
rec = self.registry.register(source_ip, metadata=metadata, policy=policy)
if tokens:
self._tokens[rec.bottle_id] = dict(tokens)
if env_var_secret:
from .store.secret_store import encrypt_value
encrypted = {k: encrypt_value(env_var_secret, v) for k, v in tokens.items()}
self.registry.store_agent_secrets(rec.bottle_id, encrypted)
req = LaunchRequest(
op="launch",
bottle_id=rec.bottle_id,
source_ip=source_ip,
image_ref=image_ref,
slot=slot,
)
try:
self._broker.submit(sign_request(req, self._secret))
except BrokerUnavailableError:
# Ambiguous delivery failure (timeout / dropped response): the broker
# may already have launched the bottle before the response was lost.
# Do NOT deregister — that would orphan a running container with no
# registry row (reconcile reaps rows, never containers). Keep the row
# so reconcile reaps it iff the bottle is not actually live; surface
# the error so the caller knows the launch is unconfirmed.
raise
except Exception:
# A definite failure — a fail-closed rejection, a backend launch
# error, or the host reporting it did not launch: nothing is running,
# so roll the registry entry back to leave no orphan.
self.registry.deregister(rec.bottle_id)
self._tokens.pop(rec.bottle_id, None)
raise
return rec
def teardown_bottle(self, bottle_id: str) -> bool:
"""Broker teardown then deregister. False if the bottle is unknown."""
rec = self.registry.get(bottle_id)
if rec is None:
return False
req = LaunchRequest(op="teardown", bottle_id=bottle_id, source_ip=rec.source_ip)
self._broker.submit(sign_request(req, self._secret))
self.registry.deregister(bottle_id)
self._tokens.pop(bottle_id, None)
# Reap any supervise proposals the gone bottle never acknowledged.
self._supervisor.archive_all_proposals(bottle_id)
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)
# Reap any supervise proposals the gone bottle never acknowledged.
self._supervisor.archive_all_proposals(rec.bottle_id)
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
persisted."""
return dict(self._tokens.get(bottle_id, {}))
def attribute(self, source_ip: str, identity_token: str) -> BottleRecord | None:
"""Fail-closed attribution (delegates to the registry)."""
return self.registry.attribute(source_ip, identity_token)
def resolve(self, source_ip: str, identity_token: str) -> BottleRecord | None:
"""Resolve the bottle behind a request — the per-request lookup the
multi-tenant gateway makes; the returned record carries its `policy`.
**Mandatory pair**: requires a matching `(source_ip, identity_token)`
(constant-time). There is no source-IP-only fallback — the app-layer
token is delivered on every attributed data plane (egress proxy
credentials, git-gate/supervise headers), so a missing or mismatched
token fail-closes. This keeps a spoofed source IP (which the /31 TAP
alone does not prevent) from selecting another bottle's policy/tokens
without also holding that bottle's unguessable token."""
return self.registry.attribute(source_ip, identity_token)
def set_policy(self, bottle_id: str, policy: str) -> bool:
"""Update a bottle's gateway policy in place (live reload). False if
the bottle is unknown."""
return self.registry.set_policy(bottle_id, policy)
# --- supervise queue (operator approvals) ------------------------------
#
# The orchestrator owns the single DB *and* the live policy, so operator
# decisions are applied here, server-side, and reached over HTTP by the
# host TUI (no direct-DB access, one path for every backend).
#
# The *agent* half of the queue (queue a proposal, poll for its response)
# is served here too (PRD 0070): the data plane no longer opens the DB, so
# supervise_server / egress / git-gate reach the queue over the control
# plane. Both are keyed on the caller's orchestrator-resolved bottle id —
# never a caller-supplied slug — so a bottle can only queue or read its own
# proposals even if the data plane is compromised.
def supervise_queue_proposal(
self, bottle_id: str, *, tool: str, proposed_file: str, justification: str,
) -> str:
"""Queue an agent-side proposal under `bottle_id` and return its id.
`bottle_id` is the caller resolved from `(source_ip, identity_token)` by
the control plane, so the proposal is attributed to the calling bottle,
not to anything the (possibly hostile) data plane asserts. The
proposed-file self-hash is computed here so the agent never supplies it."""
proposal = Proposal.new(
bottle_slug=bottle_id,
tool=tool,
proposed_file=proposed_file,
justification=justification,
)
self._supervisor.write_proposal(proposal)
return proposal.id
def supervise_poll_response(self, bottle_id: str, proposal_id: str) -> dict[str, object]:
"""Non-blocking, **idempotent** poll of one of `bottle_id`'s proposals.
Returns `{"status": ...}`: a terminal decision (`approved`/`modified`/
`rejected`, with `notes` + `final_file`) once the operator has responded,
`pending` while it's still queued undecided, or `unknown` when no such
queued proposal exists for this bottle (wrong id, or already archived).
Poll does **not** archive — a decided proposal keeps returning the same
decision so a client whose connection drops after the decision but before
it consumes the response still gets it on retry (issue #469 review). A
decided proposal is already excluded from the operator's pending list (a
response row exists), so it doesn't linger there; the row itself is reaped
when the bottle is torn down / reconciled.
Reads are scoped to `bottle_id` (the queue key), so a caller can never
read another bottle's proposal even with a guessed id."""
try:
response = self._supervisor.read_response(bottle_id, proposal_id)
except FileNotFoundError:
try:
self._supervisor.read_proposal(bottle_id, proposal_id)
except FileNotFoundError:
return {"status": POLL_STATUS_UNKNOWN}
return {"status": POLL_STATUS_PENDING}
return {
"status": response.status,
"notes": response.notes,
"final_file": response.final_file,
}
def supervise_pending(self) -> list[dict[str, object]]:
"""All pending proposals across bottles, FIFO, as JSON dicts
(`Proposal.to_dict`, round-trippable via `Proposal.from_dict`).
Each dict carries an extra `bottle_label`: the bottle's human slug
resolved from the registry (the proposal itself is keyed by the
orchestrator-assigned bottle_id, which is opaque to an operator). The
CLI renders the label but still responds against `bottle_slug`."""
out: list[dict[str, object]] = []
for p in self._supervisor.list_all_pending_proposals():
d = p.to_dict()
d["bottle_label"] = self._label_for(p.bottle_slug)
out.append(d)
return out
def _label_for(self, bottle_slug: str) -> str:
"""The human slug recorded in registry metadata for a proposal's
bottle, or the bottle_slug unchanged when the bottle is gone or has no
recorded slug — so the label is always non-empty."""
rec = self.registry.get(bottle_slug)
if rec is None:
return bottle_slug
try:
meta = json.loads(rec.metadata) if rec.metadata else {}
except ValueError:
meta = {}
slug = meta.get("slug") if isinstance(meta, dict) else None
return slug if isinstance(slug, str) and slug else bottle_slug
def _record_for_slug(self, slug: str) -> BottleRecord | None:
"""The live registry record for a proposal's bottle, or None (e.g. the
bottle was torn down before the operator responded).
In consolidated mode the supervise server attributes each proposal to
the orchestrator-assigned bottle_id and stores that as the proposal's
`bottle_slug` (see supervise_server `_attributed_config`), so the fast
path is a direct bottle_id lookup. The metadata-slug scan is the
fallback for legacy single-tenant proposals keyed by the human slug."""
rec = self.registry.get(slug)
if rec is not None:
return rec
for rec in self.registry.all():
try:
meta = json.loads(rec.metadata) if rec.metadata else {}
except ValueError:
meta = {}
if isinstance(meta, dict) and meta.get("slug") == slug:
return rec
return None
def supervise_respond(
self,
proposal_id: str,
*,
bottle_slug: str,
decision: str,
notes: str = "",
final_file: str | None = None,
) -> tuple[bool, str]:
"""Record an operator decision on a queued proposal, applying it
server-side. `decision` is approve/modify/reject.
Approve/modify on an egress tool rewrites the bottle's policy so the
gateway serves the new routes on its next `/resolve` (the live apply);
then the queued Response is written (unblocking the agent's MCP call)
and an audit entry recorded — all against the one DB. Returns
(ok, error): ok=False with a message when the proposal or decision is
unknown, or the bottle is gone so an approval can't be applied."""
status = _RESPOND_STATUS.get(decision)
if status is None:
return False, f"unknown decision {decision!r}"
try:
proposal = self._supervisor.read_proposal(bottle_slug, proposal_id)
except FileNotFoundError:
return False, "no such proposal"
diff_before, diff_after = "", ""
if status in (STATUS_APPROVED, STATUS_MODIFIED) and proposal.tool in _APPLY_TOOLS:
new_policy = final_file if final_file is not None else proposal.proposed_file
rec = self._record_for_slug(bottle_slug)
if rec is None:
return False, (
f"bottle {bottle_slug!r} is no longer registered; "
"cannot apply the route change"
)
diff_before, diff_after = rec.policy, new_policy
self.set_policy(rec.bottle_id, new_policy)
self._supervisor.write_response(bottle_slug, Response(
proposal_id=proposal_id, status=status, notes=notes, final_file=final_file,
))
component = COMPONENT_FOR_TOOL.get(proposal.tool)
if component is not None:
self._supervisor.write_audit_entry(AuditEntry(
timestamp=datetime.now(timezone.utc).isoformat(),
bottle_slug=bottle_slug,
component=component,
operator_action=status,
operator_notes=notes,
justification=proposal.justification,
diff=render_diff(
diff_before, diff_after, label=component,
before_title="current", after_label="proposed",
),
))
return True, ""
# --- secret reprovision -----------------------------------------------
def reprovision_from_secret(self, bottle_id: str, env_var_secret: str) -> bool:
"""Re-inject a bottle's egress tokens from its ENV_VAR_SECRET.
Reads the encrypted rows from ``bottled_agent_secrets``, decrypts each
value with *env_var_secret*, and restores ``_tokens[bottle_id]``.
Returns True on success, False when no stored secrets exist for this
bottle or decryption fails (wrong key / corrupt data)."""
from .store.secret_store import decrypt_value
encrypted = self.registry.get_agent_secrets(bottle_id)
if not encrypted:
return False
try:
self._tokens[bottle_id] = {k: decrypt_value(env_var_secret, v)
for k, v in encrypted.items()}
except ValueError:
return False
return True
# --- consolidated gateway ----------------------------------------------
def gateway_status(self) -> dict[str, object]:
"""Report the shared gateway for the control plane / console.
The orchestrator no longer owns a standalone gateway lifecycle — the
consolidated flow runs the gateway data plane inside the per-host infra
container/VM (see `backend/*/gateway`), so this reports `configured:
false`. Retained for the documented `GET /gateway` control-plane
contract."""
return {"configured": False}
__all__ = ["OrchestratorCore"]