refactor(supervise): split the supervise plane by tier; per-service store managers
test / integration-docker (pull_request) Successful in 13s
tracker-policy-pr / check-pr (pull_request) Successful in 14s
test / unit (pull_request) Successful in 42s
lint / lint (push) Failing after 54s
test / integration-firecracker (pull_request) Successful in 3m17s
test / coverage (pull_request) Successful in 17s
test / publish-infra (pull_request) Has been skipped

Give each service its own store package + manager, and cut the supervise module
along the control/data-plane boundary so nothing in the shared layer reaches up
into the orchestrator.

Stores, by owner:
  - bot_bottle/store/ keeps only the shared base (DbStore, migrations) and the
    concrete stores that aren't service-owned (audit_store, config_store).
  - bot_bottle/orchestrator/store/ now houses the orchestrator-owned stores —
    queue_store (supervise queue), secret_store, config_store — plus a new
    orchestrator store_manager that migrates them (composing audit/config
    downward from the base). The old shared store_manager is gone.

Supervise plane, by tier:
  - bot_bottle/supervisor/ (NEUTRAL, importable by every tier including the
    gateway): types.py (the Proposal/Response/AuditEntry wire types + the tool/
    status/poll constants + the shared daemon constants moved out of
    supervise.py) and plan.py (SupervisePlan, a pure DTO).
  - bot_bottle/orchestrator/supervisor/ (orchestrator-only): queue.py (the
    queue/audit I/O wrappers + render_diff + sha256_hex) and supervise.py (the
    Supervise lifecycle that stages the DB via the store manager). Its __init__
    re-exports the neutral vocabulary so orchestrator-side callers import from
    one place.

The gateway now imports only bot_bottle.supervisor.types (never
bot_bottle.supervise), so the data plane holds no code dependency on the
orchestrator — it reaches the queue over the control-plane RPC. This removes
the circular import that moving queue_store under orchestrator introduced
(supervise -> orchestrator -> service -> supervise).

supervise_types.py -> supervisor/types.py; supervise.py deleted (split). Full
unit suite green (2251).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-24 14:15:04 -04:00
parent f77023db1d
commit 44e2b5a897
52 changed files with 380 additions and 385 deletions
+10
View File
@@ -0,0 +1,10 @@
"""Orchestrator-owned SQLite stores.
The stores whose tables only the orchestrator (control plane) opens: the
supervise proposal/response `queue_store`, the agent-secret `secret_store`, and
the orchestrator `config_store`. They build on the shared `DbStore` /
`migrations` base in `bot_bottle.store`.
Callers import the concrete module directly, e.g.
`from bot_bottle.orchestrator.store.queue_store import QueueStore`.
"""
@@ -0,0 +1,107 @@
"""Per-host orchestrator configuration store (settings in bot-bottle.db).
Co-tenants the shared `bot-bottle.db` via the `DbStore` framework. Settings
are readable by the host launch path directly (no HTTP round-trip to the
orchestrator), so they take effect even before the orchestrator is reachable.
"""
from __future__ import annotations
import os
import sqlite3
from pathlib import Path
from ...store.db_store import DbStore
from ...store.migrations import TableMigrations
from ...paths import host_db_path
TEARDOWN_TIMEOUT_ENV = "BOT_BOTTLE_ORCHESTRATOR_TEARDOWN_TIMEOUT_SECONDS"
DEFAULT_TEARDOWN_TIMEOUT_SECONDS = 30.0
_MIGRATIONS = TableMigrations(
"orchestrator_config",
[
"""
CREATE TABLE IF NOT EXISTS orchestrator_config (
id INTEGER PRIMARY KEY CHECK (id = 1),
teardown_timeout_seconds REAL
)
""",
],
)
class OrchestratorConfigStore(DbStore):
"""Orchestrator settings in the shared host DB."""
def __init__(self, db_path: Path | None = None) -> None:
super().__init__(db_path or host_db_path(), _MIGRATIONS)
def _connect(self) -> sqlite3.Connection:
conn = super()._connect()
conn.execute("PRAGMA busy_timeout=5000")
return conn
def get_teardown_timeout_seconds(self) -> float | None:
"""Return the configured teardown timeout, or None if not set."""
try:
with self._connection() as conn:
row = conn.execute(
"SELECT teardown_timeout_seconds FROM orchestrator_config WHERE id = 1"
).fetchone()
except sqlite3.OperationalError:
return None
return row["teardown_timeout_seconds"] if row else None
def set_teardown_timeout_seconds(self, value: float) -> None:
"""Persist the teardown timeout."""
with self._connection() as conn:
conn.execute(
"INSERT OR REPLACE INTO orchestrator_config"
" (id, teardown_timeout_seconds) VALUES (1, ?)",
(value,),
)
self._chmod()
def delete_teardown_timeout_seconds(self) -> bool:
"""Clear the stored teardown timeout. Returns True if a value existed."""
with self._connection() as conn:
cur = conn.execute(
"UPDATE orchestrator_config SET teardown_timeout_seconds = NULL"
" WHERE id = 1 AND teardown_timeout_seconds IS NOT NULL"
)
return cur.rowcount > 0
def resolve_teardown_timeout(db_path: Path | None = None) -> float:
"""Return the teardown timeout to use, in priority order:
1. ``BOT_BOTTLE_ORCHESTRATOR_TEARDOWN_TIMEOUT_SECONDS`` env var
2. ``teardown_timeout_seconds`` in the orchestrator config DB
3. ``DEFAULT_TEARDOWN_TIMEOUT_SECONDS`` (30 s)
"""
raw = os.environ.get(TEARDOWN_TIMEOUT_ENV, "").strip()
if raw:
try:
value = float(raw)
if value > 0:
return value
except ValueError:
pass
store = OrchestratorConfigStore(db_path)
if not store.is_migrated():
store.migrate()
db_value = store.get_teardown_timeout_seconds()
if db_value is not None and db_value > 0:
return db_value
return DEFAULT_TEARDOWN_TIMEOUT_SECONDS
__all__ = [
"OrchestratorConfigStore",
"resolve_teardown_timeout",
"TEARDOWN_TIMEOUT_ENV",
"DEFAULT_TEARDOWN_TIMEOUT_SECONDS",
]
@@ -0,0 +1,234 @@
"""SQLite-backed queue store for supervise proposals and responses (PRD 0013)."""
from __future__ import annotations
import os
import sqlite3
from pathlib import Path
try:
from ...supervisor.types import Proposal, Response
from ...paths import host_db_path
from ...store.db_store import DbStore
from ...store.migrations import TableMigrations
except ImportError:
from supervisor.types import Proposal, Response # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from paths import host_db_path # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from db_store import DbStore # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
from migrations import TableMigrations # type: ignore[import-not-found] # pylint: disable=import-error,no-name-in-module
class QueueStore(DbStore):
"""SQLite-backed persistent store for supervise proposals and responses."""
def __init__(self, queue_key: str, db_path: Path | None = None) -> None:
self.queue_key = queue_key
if db_path is not None:
resolved = db_path
else:
# In the gateway container SUPERVISE_DB_PATH points at the
# bind-mounted host DB. On the host this env var is never set,
# so we always fall through to host_db_path().
env_path = os.environ.get("SUPERVISE_DB_PATH", "").strip()
resolved = Path(env_path) if env_path else host_db_path()
# One entry per schema version: migrations[0] brings a fresh DB to
# version 1, [1] to version 2, etc. Add new entries at the end; never
# edit existing ones.
migrations = TableMigrations("queue_store", [
# v1 — proposals table
"""
CREATE TABLE IF NOT EXISTS supervise_proposals (
queue_key TEXT NOT NULL,
id TEXT NOT NULL,
bottle_slug TEXT NOT NULL,
tool TEXT NOT NULL,
proposed_file TEXT NOT NULL,
justification TEXT NOT NULL,
arrival_timestamp TEXT NOT NULL,
current_file_hash TEXT NOT NULL,
archived INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (queue_key, id)
)
""",
# v2 — responses table
"""
CREATE TABLE IF NOT EXISTS supervise_responses (
queue_key TEXT NOT NULL,
proposal_id TEXT NOT NULL,
status TEXT NOT NULL,
notes TEXT NOT NULL,
final_file TEXT,
archived INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (queue_key, proposal_id)
)
""",
])
super().__init__(resolved, migrations)
def write_proposal(self, proposal: Proposal) -> Path:
with self._connection() as conn:
conn.execute(
"""
INSERT OR REPLACE INTO supervise_proposals (
queue_key, id, bottle_slug, tool, proposed_file, justification,
arrival_timestamp, current_file_hash, archived
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, 0)
""",
(
self.queue_key,
proposal.id,
proposal.bottle_slug,
proposal.tool,
proposal.proposed_file,
proposal.justification,
proposal.arrival_timestamp,
proposal.current_file_hash,
),
)
self._chmod()
return self.db_path
def read_proposal(self, proposal_id: str) -> Proposal:
with self._connection() as conn:
row = conn.execute(
"""
SELECT * FROM supervise_proposals
WHERE queue_key = ? AND id = ? AND archived = 0
""",
(self.queue_key, proposal_id),
).fetchone()
if row is None:
raise FileNotFoundError(proposal_id)
return self._row_to_proposal(row)
def list_pending_proposals(self) -> list[Proposal]:
if not self.db_path.is_file():
return []
with self._connection() as conn:
rows = conn.execute(
"""
SELECT p.* FROM supervise_proposals p
WHERE p.archived = 0
AND p.queue_key = ?
AND NOT EXISTS (
SELECT 1 FROM supervise_responses r
WHERE r.queue_key = p.queue_key
AND r.proposal_id = p.id
AND r.archived = 0
)
ORDER BY p.arrival_timestamp, p.id
""",
(self.queue_key,),
).fetchall()
return [self._row_to_proposal(row) for row in rows]
def list_all_pending_proposals(self) -> list[Proposal]:
if not self.db_path.is_file():
return []
with self._connection() as conn:
rows = conn.execute(
"""
SELECT p.* FROM supervise_proposals p
WHERE p.archived = 0
AND NOT EXISTS (
SELECT 1 FROM supervise_responses r
WHERE r.queue_key = p.queue_key
AND r.proposal_id = p.id
AND r.archived = 0
)
ORDER BY p.arrival_timestamp, p.id
"""
).fetchall()
return [self._row_to_proposal(row) for row in rows]
def write_response(self, response: Response) -> Path:
with self._connection() as conn:
conn.execute(
"""
INSERT OR REPLACE INTO supervise_responses (
queue_key, proposal_id, status, notes, final_file, archived
) VALUES (?, ?, ?, ?, ?, 0)
""",
(
self.queue_key,
response.proposal_id,
response.status,
response.notes,
response.final_file,
),
)
self._chmod()
return self.db_path
def read_response(self, proposal_id: str) -> Response:
with self._connection() as conn:
row = conn.execute(
"""
SELECT * FROM supervise_responses
WHERE queue_key = ? AND proposal_id = ? AND archived = 0
""",
(self.queue_key, proposal_id),
).fetchone()
if row is None:
raise FileNotFoundError(proposal_id)
return self._row_to_response(row)
def archive_proposal(self, proposal_id: str) -> None:
if not self.db_path.is_file():
return
with self._connection() as conn:
conn.execute(
"""
UPDATE supervise_proposals SET archived = 1
WHERE queue_key = ? AND id = ?
""",
(self.queue_key, proposal_id),
)
conn.execute(
"""
UPDATE supervise_responses SET archived = 1
WHERE queue_key = ? AND proposal_id = ?
""",
(self.queue_key, proposal_id),
)
def archive_all(self) -> None:
"""Archive every proposal + response for this queue key. Used to reap a
bottle's supervise records when it is torn down or reconciled away —
the retention path for a client that received a decision but died before
acknowledging it. Idempotent."""
if not self.db_path.is_file():
return
with self._connection() as conn:
conn.execute(
"UPDATE supervise_proposals SET archived = 1 WHERE queue_key = ?",
(self.queue_key,),
)
conn.execute(
"UPDATE supervise_responses SET archived = 1 WHERE queue_key = ?",
(self.queue_key,),
)
@staticmethod
def _row_to_proposal(row: sqlite3.Row) -> Proposal:
return Proposal(
id=row["id"],
bottle_slug=row["bottle_slug"],
tool=row["tool"],
proposed_file=row["proposed_file"],
justification=row["justification"],
arrival_timestamp=row["arrival_timestamp"],
current_file_hash=row["current_file_hash"],
)
@staticmethod
def _row_to_response(row: sqlite3.Row) -> Response:
return Response(
proposal_id=row["proposal_id"],
status=row["status"],
notes=row["notes"],
final_file=row["final_file"],
)
__all__ = ["QueueStore"]
@@ -0,0 +1,94 @@
"""Symmetric encryption for per-bottle egress secrets (PRD prd-new-secret-provider).
Each agent receives a random ENV_VAR_SECRET at startup — passed as an env var,
never logged or persisted. The host uses this key to encrypt each egress auth
token value before writing it to the bottled_agent_secrets table; the DB rows
(ciphertext, plaintext env-var name) without the key are insufficient to
recover the credentials.
On orchestrator restart the in-memory token map is lost. The host-side
reattachment path reads ENV_VAR_SECRET from the running agent container via
``docker exec … printenv ENV_VAR_SECRET`` and posts it to
``POST /bottles/<id>/reprovision_gateway``; the orchestrator decrypts the
stored rows and re-populates ``_tokens``.
Encryption scheme: HMAC-SHA256 used as a PRF in CTR mode (stdlib-only,
no external deps). Each value is encrypted independently. The output blob is
``nonce (16 bytes) || ciphertext`` encoded as URL-safe base64 (no padding).
keystream_block_i = HMAC-SHA256(key, nonce || i.to_bytes(4, "big"))
ciphertext_i = plaintext_i XOR keystream_block_i[:len(plaintext_i)]
"""
from __future__ import annotations
import base64
import hashlib
import hmac
import secrets
_KEY_BYTES = 32 # 256-bit key from ENV_VAR_SECRET
_NONCE_BYTES = 16 # 128-bit random nonce per encrypt call
_BLOCK = 32 # HMAC-SHA256 output width == one keystream block
# Env-var name the agent container receives at startup.
ENV_VAR_SECRET_NAME = "ENV_VAR_SECRET"
def new_env_var_secret() -> str:
"""Generate a fresh ENV_VAR_SECRET: 32 random bytes as URL-safe base64."""
return base64.urlsafe_b64encode(secrets.token_bytes(_KEY_BYTES)).rstrip(b"=").decode()
def _b64dec(s: str) -> bytes:
return base64.urlsafe_b64decode(s + "=" * (-len(s) % 4))
def _keystream(key: bytes, nonce: bytes, block_index: int) -> bytes:
return hmac.new(
key, nonce + block_index.to_bytes(4, "big"), hashlib.sha256
).digest()
def encrypt_value(secret_b64: str, plaintext: str) -> str:
"""Encrypt a single string value with *secret_b64* (the ENV_VAR_SECRET).
Returns a URL-safe base64 blob ``nonce || ciphertext`` suitable for
the ``bottled_agent_secrets.value`` column."""
key = _b64dec(secret_b64)
pt = plaintext.encode()
nonce = secrets.token_bytes(_NONCE_BYTES)
ct = bytearray()
for i in range(0, len(pt), _BLOCK):
chunk = pt[i : i + _BLOCK]
ks = _keystream(key, nonce, i)[: len(chunk)]
ct.extend(p ^ k for p, k in zip(chunk, ks))
return base64.urlsafe_b64encode(nonce + bytes(ct)).rstrip(b"=").decode()
def decrypt_value(secret_b64: str, blob_b64: str) -> str:
"""Decrypt a blob produced by :func:`encrypt_value`.
Returns the original plaintext string. Raises ``ValueError`` for malformed
input or a key mismatch (wrong key produces garbage, not an error, unless
the plaintext is non-UTF-8 — treat all such failures as wrong key)."""
key = _b64dec(secret_b64)
try:
blob = _b64dec(blob_b64)
except Exception as exc:
raise ValueError(f"invalid ciphertext blob: {exc}") from exc
if len(blob) < _NONCE_BYTES:
raise ValueError("ciphertext blob too short")
nonce, ciphertext = blob[:_NONCE_BYTES], blob[_NONCE_BYTES:]
pt = bytearray()
for i in range(0, len(ciphertext), _BLOCK):
chunk = ciphertext[i : i + _BLOCK]
ks = _keystream(key, nonce, i)[: len(chunk)]
pt.extend(c ^ k for c, k in zip(chunk, ks))
try:
return bytes(pt).decode()
except UnicodeDecodeError as exc:
raise ValueError(f"decryption produced non-UTF-8 output (wrong key?): {exc}") from exc
__all__ = ["ENV_VAR_SECRET_NAME", "new_env_var_secret", "encrypt_value", "decrypt_value"]
@@ -0,0 +1,61 @@
"""The orchestrator's store manager (PRD 0013 / 0070).
One store manager per service. Post-PRD-0070 the orchestrator is the sole
opener of `bot-bottle.db`, so it owns migrating every table in it: the supervise
`queue_store` (local to `orchestrator.store`) plus the `audit_store` and
`config_store` from the shared `bot_bottle.store` base. The data plane never
touches these — it reaches state over the control-plane RPC.
"""
from __future__ import annotations
from pathlib import Path
from .queue_store import QueueStore
from ...store.audit_store import AuditStore
from ...store.config_store import ConfigStore
_instance: StoreManager | None = None
class StoreManager:
"""Owns db_path and delegates migrate/is_migrated across the orchestrator's
stores.
Use instance() for normal access. Call reset(db_path) in tests to swap the
singleton to a temp path, then reset() with no args to restore the
default."""
def __init__(self, db_path: Path | None = None) -> None:
if db_path is None:
from ...paths import host_db_path
db_path = host_db_path()
self.db_path = db_path
@classmethod
def instance(cls) -> StoreManager:
global _instance
if _instance is None:
_instance = cls()
return _instance
@classmethod
def reset(cls, db_path: Path | None = None) -> None:
"""Replace the singleton. Pass db_path for test isolation; omit to restore default."""
global _instance
_instance = cls(db_path)
def is_migrated(self) -> bool:
return (
QueueStore("", self.db_path).is_migrated()
and AuditStore(self.db_path).is_migrated()
and ConfigStore(self.db_path).is_migrated()
)
def migrate(self) -> None:
QueueStore("", self.db_path).migrate()
AuditStore(self.db_path).migrate()
ConfigStore(self.db_path).migrate()
__all__ = ["StoreManager"]