refactor(orchestrator): replace manual HTTP dispatch with FastAPI
This commit is contained in:
@@ -44,7 +44,7 @@ if TYPE_CHECKING:
|
|||||||
from ..gateway import Gateway, GatewayError
|
from ..gateway import Gateway, GatewayError
|
||||||
from .lifecycle import Orchestrator
|
from .lifecycle import Orchestrator
|
||||||
from .service import OrchestratorCore
|
from .service import OrchestratorCore
|
||||||
from .server import OrchestratorServer, dispatch, make_server
|
from .server import OrchestratorServer, create_app, dispatch, make_server
|
||||||
|
|
||||||
|
|
||||||
# Facade name -> submodule that defines it. Lazy so importing a leaf (or the
|
# Facade name -> submodule that defines it. Lazy so importing a leaf (or the
|
||||||
@@ -67,8 +67,9 @@ _LAZY: dict[str, str] = {
|
|||||||
"GatewayError": "..gateway",
|
"GatewayError": "..gateway",
|
||||||
"Orchestrator": ".lifecycle",
|
"Orchestrator": ".lifecycle",
|
||||||
"OrchestratorCore": ".service",
|
"OrchestratorCore": ".service",
|
||||||
"OrchestratorServer": ".server",
|
"create_app": ".server",
|
||||||
"dispatch": ".server",
|
"dispatch": ".server",
|
||||||
|
"OrchestratorServer": ".server",
|
||||||
"make_server": ".server",
|
"make_server": ".server",
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -100,7 +101,8 @@ __all__ = [
|
|||||||
"GatewayError",
|
"GatewayError",
|
||||||
"Orchestrator",
|
"Orchestrator",
|
||||||
"OrchestratorCore",
|
"OrchestratorCore",
|
||||||
"OrchestratorServer",
|
"create_app",
|
||||||
"dispatch",
|
"dispatch",
|
||||||
|
"OrchestratorServer",
|
||||||
"make_server",
|
"make_server",
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -62,17 +62,14 @@ def main(argv: list[str] | None = None) -> int:
|
|||||||
orchestrator = OrchestratorCore(registry, broker, secret)
|
orchestrator = OrchestratorCore(registry, broker, secret)
|
||||||
|
|
||||||
server = make_server(orchestrator, host=args.host, port=args.port)
|
server = make_server(orchestrator, host=args.host, port=args.port)
|
||||||
bound_host, bound_port = server.server_address[0], server.server_address[1]
|
|
||||||
log.info(
|
log.info(
|
||||||
"orchestrator control plane listening",
|
"orchestrator control plane listening",
|
||||||
context={"host": bound_host, "port": bound_port, "db": str(registry.db_path)},
|
context={"host": args.host, "port": args.port, "db": str(registry.db_path)},
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
server.serve_forever()
|
server.run()
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
log.info("orchestrator shutting down")
|
log.info("orchestrator shutting down")
|
||||||
finally:
|
|
||||||
server.server_close()
|
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,311 @@
|
|||||||
|
"""FastAPI control-plane routes for the orchestrator."""
|
||||||
|
# pyright: reportUnusedFunction=false
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import math
|
||||||
|
import sys
|
||||||
|
from fastapi import FastAPI, HTTPException
|
||||||
|
from fastapi.responses import JSONResponse
|
||||||
|
from pydantic import BaseModel, ConfigDict, StrictStr
|
||||||
|
from starlette.types import ASGIApp, Message, Receive, Scope, Send
|
||||||
|
|
||||||
|
from ..orchestrator_auth import ROLE_CLI, ROLES
|
||||||
|
from ..supervisor.types import TOOLS
|
||||||
|
from ..trust_domain import CONTROL_PLANE
|
||||||
|
from .service import OrchestratorCore
|
||||||
|
|
||||||
|
ORCHESTRATOR_AUTH_HEADER = "x-bot-bottle-orchestrator-auth"
|
||||||
|
MAX_BODY_BYTES = 1 * 1024 * 1024
|
||||||
|
|
||||||
|
_GATEWAY_ROUTES = frozenset({
|
||||||
|
("POST", "/resolve"),
|
||||||
|
("POST", "/supervise/propose"),
|
||||||
|
("POST", "/supervise/poll"),
|
||||||
|
})
|
||||||
|
|
||||||
|
|
||||||
|
class _StrictModel(BaseModel):
|
||||||
|
model_config = ConfigDict(extra="ignore", strict=True)
|
||||||
|
|
||||||
|
|
||||||
|
class LaunchBody(_StrictModel):
|
||||||
|
source_ip: StrictStr
|
||||||
|
image_ref: StrictStr = ""
|
||||||
|
metadata: StrictStr = ""
|
||||||
|
policy: StrictStr = ""
|
||||||
|
tokens: dict[StrictStr, StrictStr] = {}
|
||||||
|
env_var_secret: StrictStr = ""
|
||||||
|
|
||||||
|
|
||||||
|
class PolicyBody(_StrictModel):
|
||||||
|
policy: StrictStr
|
||||||
|
|
||||||
|
|
||||||
|
class ReprovisionBody(_StrictModel):
|
||||||
|
env_var_secret: StrictStr
|
||||||
|
|
||||||
|
|
||||||
|
class ReconcileBody(_StrictModel):
|
||||||
|
live_source_ips: list[StrictStr]
|
||||||
|
grace_seconds: float | None = None
|
||||||
|
|
||||||
|
|
||||||
|
class IdentityBody(_StrictModel):
|
||||||
|
source_ip: StrictStr
|
||||||
|
identity_token: StrictStr = ""
|
||||||
|
|
||||||
|
|
||||||
|
class AttributeBody(_StrictModel):
|
||||||
|
source_ip: StrictStr
|
||||||
|
identity_token: StrictStr
|
||||||
|
|
||||||
|
|
||||||
|
class RespondBody(_StrictModel):
|
||||||
|
proposal_id: StrictStr
|
||||||
|
bottle_slug: StrictStr
|
||||||
|
decision: StrictStr
|
||||||
|
notes: StrictStr = ""
|
||||||
|
final_file: StrictStr | None = None
|
||||||
|
|
||||||
|
|
||||||
|
class ProposeBody(IdentityBody):
|
||||||
|
tool: StrictStr
|
||||||
|
proposed_file: StrictStr
|
||||||
|
justification: StrictStr
|
||||||
|
|
||||||
|
|
||||||
|
class PollBody(IdentityBody):
|
||||||
|
proposal_id: StrictStr
|
||||||
|
|
||||||
|
|
||||||
|
class ControlPlaneBoundary:
|
||||||
|
"""Reject unauthenticated and oversized requests before reading a body."""
|
||||||
|
|
||||||
|
def __init__(self, app: ASGIApp, signing_key: str) -> None:
|
||||||
|
self.app = app
|
||||||
|
self.signing_key = signing_key
|
||||||
|
|
||||||
|
async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None:
|
||||||
|
if scope["type"] != "http":
|
||||||
|
await self.app(scope, receive, send)
|
||||||
|
return
|
||||||
|
method = scope["method"]
|
||||||
|
route = scope["path"].rstrip("/") or "/"
|
||||||
|
if not (method == "GET" and route == "/health"):
|
||||||
|
headers = dict(scope["headers"])
|
||||||
|
presented = headers.get(
|
||||||
|
ORCHESTRATOR_AUTH_HEADER.encode(), b"",
|
||||||
|
).decode(errors="ignore")
|
||||||
|
role = CONTROL_PLANE.verify(presented, self.signing_key)
|
||||||
|
if role is None:
|
||||||
|
await self._reject(
|
||||||
|
scope, send, 401, "control-plane authentication required",
|
||||||
|
)
|
||||||
|
return
|
||||||
|
allowed = ROLES if (method, route) in _GATEWAY_ROUTES else {ROLE_CLI}
|
||||||
|
if role not in allowed:
|
||||||
|
await self._reject(scope, send, 403, "insufficient role for this route")
|
||||||
|
return
|
||||||
|
scope["state"]["role"] = role
|
||||||
|
raw_length = dict(scope["headers"]).get(b"content-length")
|
||||||
|
if raw_length is not None:
|
||||||
|
try:
|
||||||
|
length = int(raw_length)
|
||||||
|
except ValueError:
|
||||||
|
await self._reject(scope, send, 400, "invalid Content-Length")
|
||||||
|
return
|
||||||
|
if length < 0:
|
||||||
|
await self._reject(scope, send, 400, "invalid Content-Length")
|
||||||
|
return
|
||||||
|
if length > MAX_BODY_BYTES:
|
||||||
|
await self._reject(scope, send, 413, "request body too large")
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
await self.app(scope, self._bounded_receive(receive), send)
|
||||||
|
except Exception as exc: # noqa: BLE001 - redact control-plane failures
|
||||||
|
sys.stderr.write(
|
||||||
|
f"orchestrator: {method} {route} failed "
|
||||||
|
f"[error_type={type(exc).__name__}]\n"
|
||||||
|
)
|
||||||
|
sys.stderr.flush()
|
||||||
|
await self._reject(scope, send, 500, "internal error")
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
async def _reject(
|
||||||
|
scope: Scope, send: Send, status: int, error: str,
|
||||||
|
) -> None:
|
||||||
|
response = JSONResponse({"error": error}, status_code=status)
|
||||||
|
await response(scope, ControlPlaneBoundary._empty_receive, send)
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
async def _empty_receive() -> Message:
|
||||||
|
return {"type": "http.disconnect"}
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _bounded_receive(receive: Receive) -> Receive:
|
||||||
|
consumed = 0
|
||||||
|
|
||||||
|
async def bounded() -> Message:
|
||||||
|
nonlocal consumed
|
||||||
|
message = await receive()
|
||||||
|
if message["type"] == "http.request":
|
||||||
|
consumed += len(message.get("body", b""))
|
||||||
|
if consumed > MAX_BODY_BYTES:
|
||||||
|
raise HTTPException(413, "request body too large")
|
||||||
|
return message
|
||||||
|
|
||||||
|
return bounded
|
||||||
|
|
||||||
|
|
||||||
|
def _required(value: str, name: str) -> str:
|
||||||
|
if not value:
|
||||||
|
raise HTTPException(400, f"{name} (string) is required")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def create_app(orch: OrchestratorCore, *, signing_key: str) -> FastAPI:
|
||||||
|
"""Build the authenticated orchestrator ASGI application."""
|
||||||
|
key = signing_key.strip()
|
||||||
|
if not key:
|
||||||
|
raise ValueError(
|
||||||
|
"orchestrator control-plane signing key is required; "
|
||||||
|
"refusing to start without caller authentication"
|
||||||
|
)
|
||||||
|
app = FastAPI(
|
||||||
|
title="bot-bottle orchestrator",
|
||||||
|
docs_url=None,
|
||||||
|
redoc_url=None,
|
||||||
|
openapi_url=None,
|
||||||
|
)
|
||||||
|
app.add_middleware(ControlPlaneBoundary, signing_key=key)
|
||||||
|
|
||||||
|
@app.get("/health")
|
||||||
|
def health() -> dict[str, str]:
|
||||||
|
return {"status": "ok"}
|
||||||
|
|
||||||
|
@app.get("/gateway")
|
||||||
|
def gateway() -> dict[str, object]:
|
||||||
|
return orch.gateway_status()
|
||||||
|
|
||||||
|
@app.get("/bottles")
|
||||||
|
def bottles() -> dict[str, object]:
|
||||||
|
return {"bottles": [record.redacted() for record in orch.registry.all()]}
|
||||||
|
|
||||||
|
@app.post("/bottles", status_code=201)
|
||||||
|
def launch(body: LaunchBody) -> dict[str, str]:
|
||||||
|
rec = orch.launch_bottle(
|
||||||
|
_required(body.source_ip, "source_ip"),
|
||||||
|
image_ref=body.image_ref,
|
||||||
|
metadata=body.metadata,
|
||||||
|
policy=body.policy,
|
||||||
|
tokens=dict(body.tokens),
|
||||||
|
env_var_secret=body.env_var_secret,
|
||||||
|
)
|
||||||
|
return {"bottle_id": rec.bottle_id, "identity_token": rec.identity_token}
|
||||||
|
|
||||||
|
@app.put("/bottles/{bottle_id}/policy")
|
||||||
|
def set_policy(bottle_id: str, body: PolicyBody) -> dict[str, object]:
|
||||||
|
if orch.set_policy(bottle_id, body.policy):
|
||||||
|
return {"updated": True}
|
||||||
|
raise HTTPException(404, "no such bottle")
|
||||||
|
|
||||||
|
@app.post("/bottles/{bottle_id}/reprovision_gateway")
|
||||||
|
def reprovision(bottle_id: str, body: ReprovisionBody) -> dict[str, object]:
|
||||||
|
secret = _required(body.env_var_secret, "env_var_secret")
|
||||||
|
if orch.reprovision_from_secret(bottle_id, secret):
|
||||||
|
return {"reprovisioned": True}
|
||||||
|
raise HTTPException(404, "no stored secrets for this bottle")
|
||||||
|
|
||||||
|
@app.delete("/bottles/{bottle_id}")
|
||||||
|
def teardown(bottle_id: str) -> dict[str, object]:
|
||||||
|
if orch.teardown_bottle(bottle_id):
|
||||||
|
return {"torn_down": True}
|
||||||
|
raise HTTPException(404, "no such bottle")
|
||||||
|
|
||||||
|
@app.post("/reconcile")
|
||||||
|
def reconcile(body: ReconcileBody) -> dict[str, object]:
|
||||||
|
if any(not ip for ip in body.live_source_ips):
|
||||||
|
raise HTTPException(400, "live_source_ips must contain non-empty strings")
|
||||||
|
kwargs: dict[str, float] = {}
|
||||||
|
if body.grace_seconds is not None:
|
||||||
|
if not math.isfinite(body.grace_seconds) or body.grace_seconds < 0:
|
||||||
|
raise HTTPException(
|
||||||
|
400, "grace_seconds must be a non-negative finite number",
|
||||||
|
)
|
||||||
|
kwargs["grace_seconds"] = body.grace_seconds
|
||||||
|
return {"reaped": orch.reconcile(body.live_source_ips, **kwargs)}
|
||||||
|
|
||||||
|
@app.post("/attribute")
|
||||||
|
def attribute(body: AttributeBody) -> dict[str, str]:
|
||||||
|
rec = orch.attribute(body.source_ip, body.identity_token)
|
||||||
|
if rec is None:
|
||||||
|
raise HTTPException(403, "unattributed")
|
||||||
|
return {"bottle_id": rec.bottle_id}
|
||||||
|
|
||||||
|
@app.get("/supervise/proposals")
|
||||||
|
def proposals() -> dict[str, object]:
|
||||||
|
return {"proposals": orch.supervise_pending()}
|
||||||
|
|
||||||
|
@app.post("/supervise/respond")
|
||||||
|
def respond(body: RespondBody) -> dict[str, object]:
|
||||||
|
ok, error = orch.supervise_respond(
|
||||||
|
_required(body.proposal_id, "proposal_id"),
|
||||||
|
bottle_slug=_required(body.bottle_slug, "bottle_slug"),
|
||||||
|
decision=_required(body.decision, "decision"),
|
||||||
|
notes=body.notes,
|
||||||
|
final_file=body.final_file,
|
||||||
|
)
|
||||||
|
if not ok:
|
||||||
|
raise HTTPException(409, error)
|
||||||
|
return {"responded": True}
|
||||||
|
|
||||||
|
@app.post("/supervise/propose", status_code=201)
|
||||||
|
def propose(body: ProposeBody) -> dict[str, str]:
|
||||||
|
source_ip = _required(body.source_ip, "source_ip")
|
||||||
|
if body.tool not in TOOLS:
|
||||||
|
raise HTTPException(400, f"tool (string) must be one of {TOOLS}")
|
||||||
|
rec = orch.resolve(source_ip, body.identity_token)
|
||||||
|
if rec is None:
|
||||||
|
raise HTTPException(403, "unattributed")
|
||||||
|
proposal_id = orch.supervise_queue_proposal(
|
||||||
|
rec.bottle_id,
|
||||||
|
tool=body.tool,
|
||||||
|
proposed_file=_required(body.proposed_file, "proposed_file"),
|
||||||
|
justification=_required(body.justification, "justification"),
|
||||||
|
)
|
||||||
|
return {"proposal_id": proposal_id}
|
||||||
|
|
||||||
|
@app.post("/supervise/poll")
|
||||||
|
def poll(body: PollBody) -> dict[str, object]:
|
||||||
|
rec = orch.resolve(
|
||||||
|
_required(body.source_ip, "source_ip"), body.identity_token,
|
||||||
|
)
|
||||||
|
if rec is None:
|
||||||
|
raise HTTPException(403, "unattributed")
|
||||||
|
return orch.supervise_poll_response(
|
||||||
|
rec.bottle_id, _required(body.proposal_id, "proposal_id"),
|
||||||
|
)
|
||||||
|
|
||||||
|
@app.post("/resolve")
|
||||||
|
def resolve(body: IdentityBody) -> dict[str, object]:
|
||||||
|
rec = orch.resolve(
|
||||||
|
_required(body.source_ip, "source_ip"), body.identity_token,
|
||||||
|
)
|
||||||
|
if rec is None:
|
||||||
|
raise HTTPException(403, "unattributed")
|
||||||
|
return {
|
||||||
|
"bottle_id": rec.bottle_id,
|
||||||
|
"policy": rec.policy,
|
||||||
|
"tokens": orch.tokens_for(rec.bottle_id),
|
||||||
|
}
|
||||||
|
|
||||||
|
return app
|
||||||
|
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"ControlPlaneBoundary",
|
||||||
|
"MAX_BODY_BYTES",
|
||||||
|
"ORCHESTRATOR_AUTH_HEADER",
|
||||||
|
"create_app",
|
||||||
|
]
|
||||||
@@ -1,502 +1,85 @@
|
|||||||
"""Orchestrator HTTP control plane (PRD 0070).
|
"""Uvicorn transport for the FastAPI orchestrator control plane."""
|
||||||
|
|
||||||
The backend-agnostic control-plane RPC (CLI / console -> orchestrator) over
|
|
||||||
**HTTP** — the universal transport chosen in 0070 (works on every host; no
|
|
||||||
vsock / unix-socket portability caveats):
|
|
||||||
|
|
||||||
GET /health -> 200 {"status": "ok"}
|
|
||||||
GET /gateway -> 200 {"configured", ["name","running"]}
|
|
||||||
GET /bottles -> 200 {"bottles": [ <redacted record>, ...]}
|
|
||||||
POST /bottles -> 201 {"bottle_id","identity_token"} (launch)
|
|
||||||
body: {"source_ip", ["image_ref"],
|
|
||||||
["metadata"], ["policy"],
|
|
||||||
["tokens"], ["env_var_secret"]}
|
|
||||||
PUT /bottles/<bottle_id>/policy -> 200 {"updated": true} | 404 (live reload)
|
|
||||||
body: {"policy"}
|
|
||||||
POST /bottles/<bottle_id>/reprovision_gateway
|
|
||||||
-> 200 {"reprovisioned": true} | 404
|
|
||||||
body: {"env_var_secret"}
|
|
||||||
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"}
|
|
||||||
GET /supervise/proposals -> 200 {"proposals": [ <proposal>, ...]}
|
|
||||||
POST /supervise/respond -> 200 {"responded": true} | 409 (operator)
|
|
||||||
body: {"proposal_id","bottle_slug",
|
|
||||||
"decision", ["notes"],["final_file"]}
|
|
||||||
POST /supervise/propose -> 201 {"proposal_id"} | 403 (agent)
|
|
||||||
body: {"source_ip","identity_token",
|
|
||||||
"tool","proposed_file","justification"}
|
|
||||||
POST /supervise/poll -> 200 {"status", ["notes"],["final_file"]} | 403
|
|
||||||
body: {"source_ip","identity_token",
|
|
||||||
"proposal_id"}
|
|
||||||
|
|
||||||
The `/supervise/propose` + `/supervise/poll` pair is the **agent** half of the
|
|
||||||
supervise flow: the data plane (supervise / egress / git-gate) queues a proposal
|
|
||||||
and polls for its response over RPC instead of opening `bot-bottle.db` directly.
|
|
||||||
`poll` is idempotent — it never archives, so a dropped connection can't lose an
|
|
||||||
operator decision (the row is reaped when the bottle is torn down / reconciled).
|
|
||||||
Both attribute the caller by `(source_ip, identity_token)` exactly like
|
|
||||||
`/resolve`, so a bottle can only ever queue or read its own proposals.
|
|
||||||
|
|
||||||
`POST /bottles` / `DELETE` drive the full launch lifecycle: they mint (or
|
|
||||||
tear down) the bottle in the registry AND broker the backend-native launch
|
|
||||||
via the orchestrator. Register/deregister without a launch are internal to
|
|
||||||
`OrchestratorCore`, not exposed here.
|
|
||||||
|
|
||||||
Routing/handling is the pure function `dispatch()` so it is unit-testable
|
|
||||||
without a socket; `Handler` / `OrchestratorServer` / `make_server` are a
|
|
||||||
thin stdlib adapter around it. Listing redacts identity tokens — they are
|
|
||||||
returned only once, to the caller that launches the bottle.
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import http.server
|
|
||||||
import json
|
|
||||||
import math
|
|
||||||
import os
|
import os
|
||||||
import socket
|
import socket
|
||||||
import socketserver
|
|
||||||
import sys
|
|
||||||
import threading
|
import threading
|
||||||
import typing
|
|
||||||
from urllib.parse import urlsplit
|
|
||||||
|
|
||||||
from ..orchestrator_auth import ROLE_CLI, ROLES
|
import uvicorn
|
||||||
|
|
||||||
|
from ..orchestrator_auth import ROLE_CLI, mint
|
||||||
from ..trust_domain import CONTROL_PLANE
|
from ..trust_domain import CONTROL_PLANE
|
||||||
from ..supervisor.types import TOOLS
|
from .api import MAX_BODY_BYTES, ORCHESTRATOR_AUTH_HEADER, create_app
|
||||||
from .service import OrchestratorCore
|
from .service import OrchestratorCore
|
||||||
|
|
||||||
# JSON body payload type (parsed request / rendered response).
|
MAX_REQUESTS = 32
|
||||||
Json = dict[str, object]
|
KEEP_ALIVE_TIMEOUT_SECONDS = 10
|
||||||
|
|
||||||
# The request header carrying the caller's role-scoped control-plane token (a
|
|
||||||
# signed JWT naming the caller's role — see orchestrator_auth). The role gates which
|
|
||||||
# routes the caller may reach: the data plane holds a `gateway` token good only
|
|
||||||
# for the agent-facing lookups; the host CLI holds a `cli` token for the
|
|
||||||
# operator/mutating routes. An agent that can merely *reach* the port holds no
|
|
||||||
# token at all, and a compromised gateway holds only `gateway` — neither can
|
|
||||||
# drive the operator routes (approve proposals, rewrite policy, read tokens).
|
|
||||||
ORCHESTRATOR_AUTH_HEADER = "x-bot-bottle-orchestrator-auth"
|
|
||||||
MAX_BODY_BYTES = 1 * 1024 * 1024
|
|
||||||
REQUEST_TIMEOUT_SECONDS = 10.0
|
|
||||||
MAX_REQUEST_THREADS = 32
|
|
||||||
|
|
||||||
# The routes the data plane (role `gateway`) is allowed to reach — exactly the
|
|
||||||
# per-request lookups PolicyResolver makes. Every other authenticated route is
|
|
||||||
# operator-only. `cli` is a superset role: it may reach any route.
|
|
||||||
_GATEWAY_ROUTES: frozenset[tuple[str, str]] = frozenset({
|
|
||||||
("POST", "/resolve"),
|
|
||||||
("POST", "/supervise/propose"),
|
|
||||||
("POST", "/supervise/poll"),
|
|
||||||
})
|
|
||||||
|
|
||||||
|
|
||||||
def _allowed_roles(method: str, route: str) -> frozenset[str]:
|
def dispatch(
|
||||||
"""The roles permitted on `(method, route)`: `gateway` or `cli` on the
|
orchestrator: OrchestratorCore,
|
||||||
data-plane routes, `cli`-only everywhere else."""
|
method: str,
|
||||||
if (method, route) in _GATEWAY_ROUTES:
|
path: str,
|
||||||
return ROLES
|
body: bytes,
|
||||||
return frozenset({ROLE_CLI})
|
*,
|
||||||
|
role: str | None = ROLE_CLI,
|
||||||
|
) -> tuple[int, dict[str, object]]:
|
||||||
|
"""Socket-free compatibility adapter for route unit tests.
|
||||||
|
|
||||||
|
Production requests always enter through the FastAPI ASGI application.
|
||||||
|
"""
|
||||||
|
from fastapi.testclient import TestClient
|
||||||
|
|
||||||
|
key = "in-process-dispatch-key"
|
||||||
|
headers: dict[str, str] = {"content-type": "application/json"}
|
||||||
|
if role is not None:
|
||||||
|
headers[ORCHESTRATOR_AUTH_HEADER] = mint(role, key)
|
||||||
|
response = TestClient(create_app(orchestrator, signing_key=key)).request(
|
||||||
|
method, path, content=body, headers=headers,
|
||||||
|
)
|
||||||
|
payload = response.json()
|
||||||
|
if response.status_code == 422:
|
||||||
|
detail = payload.get("detail", []) if isinstance(payload, dict) else []
|
||||||
|
field = ""
|
||||||
|
if isinstance(detail, list) and detail and isinstance(detail[0], dict):
|
||||||
|
location = detail[0].get("loc", ())
|
||||||
|
if isinstance(location, (list, tuple)) and len(location) > 1:
|
||||||
|
field = str(location[1])
|
||||||
|
suffix = f": {field}" if field else ""
|
||||||
|
return 400, {"error": f"invalid request body{suffix}"}
|
||||||
|
if isinstance(payload, dict) and "detail" in payload and "error" not in payload:
|
||||||
|
payload = {"error": payload["detail"]}
|
||||||
|
return response.status_code, payload
|
||||||
|
|
||||||
|
|
||||||
def _parse_json_object(body: bytes) -> Json:
|
class OrchestratorServer:
|
||||||
"""Parse a JSON object body. Raises ValueError for non-objects / bad JSON."""
|
"""Small lifecycle wrapper around Uvicorn with an eagerly bound socket."""
|
||||||
if not body:
|
|
||||||
return {}
|
|
||||||
obj = json.loads(body) # raises json.JSONDecodeError (a ValueError)
|
|
||||||
if not isinstance(obj, dict):
|
|
||||||
raise ValueError("request body must be a JSON object")
|
|
||||||
return obj
|
|
||||||
|
|
||||||
|
def __init__(self, config: uvicorn.Config) -> None:
|
||||||
|
self._server = uvicorn.Server(config)
|
||||||
|
self._stopped = threading.Event()
|
||||||
|
self._socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||||
|
self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||||
|
self._socket.bind((config.host, config.port))
|
||||||
|
self._socket.listen(config.backlog)
|
||||||
|
self.server_address = self._socket.getsockname()
|
||||||
|
|
||||||
def dispatch( # pylint: disable=too-many-return-statements,too-many-branches
|
def run(self) -> None:
|
||||||
orch: OrchestratorCore, method: str, path: str, body: bytes, *, role: str | None = ROLE_CLI,
|
|
||||||
) -> tuple[int, Json]:
|
|
||||||
"""Route one control-plane request to a (status, payload) pair. Pure —
|
|
||||||
no I/O beyond the orchestrator — so it is fully testable without a socket.
|
|
||||||
|
|
||||||
`role` is the caller's verified control-plane role (`gateway` or `cli`), or
|
|
||||||
None for an unauthenticated request. Every route except `GET /health`
|
|
||||||
requires a role: a missing role is 401, and a role that
|
|
||||||
doesn't cover the route is 403 — so a `gateway` data-plane token can reach
|
|
||||||
`/resolve` + `/supervise/{propose,poll}` but not the operator routes
|
|
||||||
(rewrite policy, read injected tokens, approve its own supervise proposals).
|
|
||||||
The source-IP + identity-token checks inside `/resolve` and `/attribute`
|
|
||||||
authenticate the *bottle* a request is about, not the *caller*, so this role
|
|
||||||
gate is what protects the caller-privileged routes. Defaults `cli` so unit
|
|
||||||
tests of the routing logic don't have to thread it through."""
|
|
||||||
route = urlsplit(path).path.rstrip("/") or "/"
|
|
||||||
|
|
||||||
if method == "GET" and route == "/health":
|
|
||||||
return 200, {"status": "ok"}
|
|
||||||
|
|
||||||
# Role gate — every route below is a trusted-caller operation. Deny before
|
|
||||||
# touching the registry / broker / supervise store.
|
|
||||||
if role is None:
|
|
||||||
return 401, {"error": "control-plane authentication required"}
|
|
||||||
if role not in _allowed_roles(method, route):
|
|
||||||
return 403, {"error": "insufficient role for this route"}
|
|
||||||
|
|
||||||
if method == "GET" and route == "/gateway":
|
|
||||||
return 200, orch.gateway_status()
|
|
||||||
|
|
||||||
if method == "GET" and route == "/bottles":
|
|
||||||
return 200, {"bottles": [r.redacted() for r in orch.registry.all()]}
|
|
||||||
|
|
||||||
if method == "POST" and route == "/bottles":
|
|
||||||
try:
|
try:
|
||||||
data = _parse_json_object(body)
|
self._server.run(sockets=[self._socket])
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
source_ip = data.get("source_ip")
|
|
||||||
if not isinstance(source_ip, str) or not source_ip:
|
|
||||||
return 400, {"error": "source_ip (string) is required"}
|
|
||||||
image_ref = data.get("image_ref")
|
|
||||||
metadata = data.get("metadata")
|
|
||||||
policy = data.get("policy")
|
|
||||||
raw_tokens = data.get("tokens")
|
|
||||||
tokens = {
|
|
||||||
k: v for k, v in raw_tokens.items() if isinstance(k, str) and isinstance(v, str)
|
|
||||||
} if isinstance(raw_tokens, dict) else {}
|
|
||||||
env_var_secret = data.get("env_var_secret", "")
|
|
||||||
rec = orch.launch_bottle(
|
|
||||||
source_ip,
|
|
||||||
image_ref=image_ref if isinstance(image_ref, str) else "",
|
|
||||||
metadata=metadata if isinstance(metadata, str) else "",
|
|
||||||
policy=policy if isinstance(policy, str) else "",
|
|
||||||
tokens=tokens,
|
|
||||||
env_var_secret=env_var_secret if isinstance(env_var_secret, str) else "",
|
|
||||||
)
|
|
||||||
return 201, {"bottle_id": rec.bottle_id, "identity_token": rec.identity_token}
|
|
||||||
|
|
||||||
if method == "PUT" and route.startswith("/bottles/") and route.endswith("/policy"):
|
|
||||||
bottle_id = route[len("/bottles/"):-len("/policy")]
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
policy = data.get("policy")
|
|
||||||
if not isinstance(policy, str):
|
|
||||||
return 400, {"error": "policy (string) is required"}
|
|
||||||
if orch.set_policy(bottle_id, policy):
|
|
||||||
return 200, {"updated": True}
|
|
||||||
return 404, {"error": "no such bottle"}
|
|
||||||
|
|
||||||
if (
|
|
||||||
method == "POST"
|
|
||||||
and route.startswith("/bottles/")
|
|
||||||
and route.endswith("/reprovision_gateway")
|
|
||||||
):
|
|
||||||
bottle_id = route[len("/bottles/") : -len("/reprovision_gateway")]
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
env_var_secret = data.get("env_var_secret")
|
|
||||||
if not isinstance(env_var_secret, str) or not env_var_secret:
|
|
||||||
return 400, {"error": "env_var_secret (string) is required"}
|
|
||||||
if orch.reprovision_from_secret(bottle_id, env_var_secret):
|
|
||||||
return 200, {"reprovisioned": True}
|
|
||||||
return 404, {"error": "no stored secrets for this bottle"}
|
|
||||||
|
|
||||||
if method == "DELETE" and route.startswith("/bottles/"):
|
|
||||||
bottle_id = route[len("/bottles/"):]
|
|
||||||
if orch.teardown_bottle(bottle_id):
|
|
||||||
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"}
|
|
||||||
if any(not isinstance(ip, str) or not ip for ip in raw_ips):
|
|
||||||
return 400, {"error": "live_source_ips must contain non-empty strings"}
|
|
||||||
live = raw_ips
|
|
||||||
grace = data.get("grace_seconds")
|
|
||||||
kwargs: dict[str, float] = {}
|
|
||||||
if grace is not None:
|
|
||||||
if isinstance(grace, bool) or not isinstance(grace, (int, float)):
|
|
||||||
return 400, {"error": "grace_seconds must be a non-negative finite number"}
|
|
||||||
parsed_grace = float(grace)
|
|
||||||
if not math.isfinite(parsed_grace) or parsed_grace < 0:
|
|
||||||
return 400, {"error": "grace_seconds must be a non-negative finite number"}
|
|
||||||
kwargs["grace_seconds"] = parsed_grace
|
|
||||||
return 200, {"reaped": orch.reconcile(live, **kwargs)}
|
|
||||||
|
|
||||||
if method == "POST" and route == "/attribute":
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
source_ip = data.get("source_ip")
|
|
||||||
token = data.get("identity_token")
|
|
||||||
if not isinstance(source_ip, str) or not isinstance(token, str):
|
|
||||||
return 400, {"error": "source_ip and identity_token (strings) required"}
|
|
||||||
rec = orch.attribute(source_ip, token)
|
|
||||||
if rec is None:
|
|
||||||
return 403, {"error": "unattributed"}
|
|
||||||
return 200, {"bottle_id": rec.bottle_id}
|
|
||||||
|
|
||||||
if method == "GET" and route == "/supervise/proposals":
|
|
||||||
# Operator TUI: pending supervise proposals across all bottles.
|
|
||||||
return 200, {"proposals": orch.supervise_pending()}
|
|
||||||
|
|
||||||
if method == "POST" and route == "/supervise/respond":
|
|
||||||
# Operator decision: apply (approve/modify rewrites egress policy),
|
|
||||||
# write the queued response, audit — all server-side on the one DB.
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
proposal_id = data.get("proposal_id")
|
|
||||||
bottle_slug = data.get("bottle_slug")
|
|
||||||
decision = data.get("decision")
|
|
||||||
if not (isinstance(proposal_id, str) and proposal_id):
|
|
||||||
return 400, {"error": "proposal_id (string) is required"}
|
|
||||||
if not (isinstance(bottle_slug, str) and bottle_slug):
|
|
||||||
return 400, {"error": "bottle_slug (string) is required"}
|
|
||||||
if not (isinstance(decision, str) and decision):
|
|
||||||
return 400, {"error": "decision (string) is required"}
|
|
||||||
notes = data.get("notes")
|
|
||||||
final_file = data.get("final_file")
|
|
||||||
ok, err = orch.supervise_respond(
|
|
||||||
proposal_id,
|
|
||||||
bottle_slug=bottle_slug,
|
|
||||||
decision=decision,
|
|
||||||
notes=notes if isinstance(notes, str) else "",
|
|
||||||
final_file=final_file if isinstance(final_file, str) else None,
|
|
||||||
)
|
|
||||||
if ok:
|
|
||||||
return 200, {"responded": True}
|
|
||||||
return 409, {"error": err}
|
|
||||||
|
|
||||||
if method == "POST" and route == "/supervise/propose":
|
|
||||||
# Agent half: queue a proposal, attributed to the caller resolved from
|
|
||||||
# (source_ip, identity_token) — never a caller-supplied slug — so the
|
|
||||||
# data plane can't forge attribution. Fail-closed 403 when unattributed.
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
source_ip = data.get("source_ip")
|
|
||||||
token = data.get("identity_token")
|
|
||||||
tool = data.get("tool")
|
|
||||||
proposed_file = data.get("proposed_file")
|
|
||||||
justification = data.get("justification")
|
|
||||||
if not isinstance(source_ip, str) or not source_ip:
|
|
||||||
return 400, {"error": "source_ip (string) is required"}
|
|
||||||
if not isinstance(tool, str) or tool not in TOOLS:
|
|
||||||
return 400, {"error": f"tool (string) must be one of {TOOLS}"}
|
|
||||||
if not isinstance(proposed_file, str) or not proposed_file:
|
|
||||||
return 400, {"error": "proposed_file (string) is required"}
|
|
||||||
if not isinstance(justification, str) or not justification:
|
|
||||||
return 400, {"error": "justification (string) is required"}
|
|
||||||
rec = orch.resolve(source_ip, token if isinstance(token, str) else "")
|
|
||||||
if rec is None:
|
|
||||||
return 403, {"error": "unattributed"}
|
|
||||||
proposal_id = orch.supervise_queue_proposal(
|
|
||||||
rec.bottle_id, tool=tool, proposed_file=proposed_file,
|
|
||||||
justification=justification,
|
|
||||||
)
|
|
||||||
return 201, {"proposal_id": proposal_id}
|
|
||||||
|
|
||||||
if method == "POST" and route == "/supervise/poll":
|
|
||||||
# Agent half: non-blocking read of the caller's own proposal decision.
|
|
||||||
# Attributed like /propose, and scoped to the resolved bottle id, so a
|
|
||||||
# guessed proposal_id can never read another bottle's response.
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
source_ip = data.get("source_ip")
|
|
||||||
token = data.get("identity_token")
|
|
||||||
proposal_id = data.get("proposal_id")
|
|
||||||
if not isinstance(source_ip, str) or not source_ip:
|
|
||||||
return 400, {"error": "source_ip (string) is required"}
|
|
||||||
if not isinstance(proposal_id, str) or not proposal_id:
|
|
||||||
return 400, {"error": "proposal_id (string) is required"}
|
|
||||||
rec = orch.resolve(source_ip, token if isinstance(token, str) else "")
|
|
||||||
if rec is None:
|
|
||||||
return 403, {"error": "unattributed"}
|
|
||||||
return 200, orch.supervise_poll_response(rec.bottle_id, proposal_id)
|
|
||||||
|
|
||||||
if method == "POST" and route == "/resolve":
|
|
||||||
# The per-request lookup the multi-tenant gateway makes: returns the
|
|
||||||
# bottle's policy. Requires a matching (source_ip, identity_token)
|
|
||||||
# pair — a missing/empty/mismatched token fail-closes (403), no
|
|
||||||
# source-IP-only fallback.
|
|
||||||
try:
|
|
||||||
data = _parse_json_object(body)
|
|
||||||
except ValueError as e:
|
|
||||||
return 400, {"error": f"invalid JSON: {e}"}
|
|
||||||
source_ip = data.get("source_ip")
|
|
||||||
token = data.get("identity_token")
|
|
||||||
if not isinstance(source_ip, str) or not source_ip:
|
|
||||||
return 400, {"error": "source_ip (string) is required"}
|
|
||||||
rec = orch.resolve(source_ip, token if isinstance(token, str) else "")
|
|
||||||
if rec is None:
|
|
||||||
return 403, {"error": "unattributed"}
|
|
||||||
# tokens are the in-memory per-bottle egress auth values the gateway
|
|
||||||
# injects; served here, never persisted.
|
|
||||||
return 200, {
|
|
||||||
"bottle_id": rec.bottle_id,
|
|
||||||
"policy": rec.policy,
|
|
||||||
"tokens": orch.tokens_for(rec.bottle_id),
|
|
||||||
}
|
|
||||||
|
|
||||||
return 404, {"error": "not found"}
|
|
||||||
|
|
||||||
|
|
||||||
class Handler(http.server.BaseHTTPRequestHandler):
|
|
||||||
"""Thin stdlib adapter: read the body, call `dispatch`, write JSON."""
|
|
||||||
|
|
||||||
# Quiet by default (the orchestrator has its own logging); opt back into
|
|
||||||
# stdlib access logging with BOT_BOTTLE_ORCHESTRATOR_DEBUG.
|
|
||||||
def log_message(self, format: str, *args: typing.Any) -> None: # noqa: A002
|
|
||||||
if os.environ.get("BOT_BOTTLE_ORCHESTRATOR_DEBUG"):
|
|
||||||
super().log_message(format, *args)
|
|
||||||
|
|
||||||
def _serve(self, method: str) -> None:
|
|
||||||
"""Read the request body, dispatch it, and write the JSON reply. A
|
|
||||||
dispatch failure (e.g. a broker error) returns a 500 rather than
|
|
||||||
crashing the connection, so one bad request can't take the control
|
|
||||||
plane down for the caller."""
|
|
||||||
server = self.server
|
|
||||||
assert isinstance(server, OrchestratorServer)
|
|
||||||
role = server.role_for(self.headers.get(ORCHESTRATOR_AUTH_HEADER, ""))
|
|
||||||
route = urlsplit(self.path).path.rstrip("/") or "/"
|
|
||||||
if not (method == "GET" and route == "/health") and role is None:
|
|
||||||
self._write_json(
|
|
||||||
401, {"error": "control-plane authentication required"},
|
|
||||||
)
|
|
||||||
return
|
|
||||||
length_header = self.headers.get("Content-Length")
|
|
||||||
try:
|
|
||||||
length = int(length_header) if length_header is not None else 0
|
|
||||||
except ValueError:
|
|
||||||
self._write_json(400, {"error": "invalid Content-Length"})
|
|
||||||
return
|
|
||||||
if length < 0:
|
|
||||||
self._write_json(400, {"error": "invalid Content-Length"})
|
|
||||||
return
|
|
||||||
if length > MAX_BODY_BYTES:
|
|
||||||
self._write_json(413, {"error": "request body too large"})
|
|
||||||
return
|
|
||||||
try:
|
|
||||||
body = self.rfile.read(length) if length else b""
|
|
||||||
except (TimeoutError, socket.timeout):
|
|
||||||
self._write_json(408, {"error": "request body read timed out"})
|
|
||||||
return
|
|
||||||
try:
|
|
||||||
status: int
|
|
||||||
payload: Json
|
|
||||||
status, payload = dispatch(
|
|
||||||
server.orchestrator, method, self.path, body, role=role)
|
|
||||||
except Exception as e: # noqa: BLE001 — the control plane must stay up
|
|
||||||
# Do not echo exception messages to the caller or logs: broker and
|
|
||||||
# persistence exceptions can contain request data. The operation,
|
|
||||||
# route, and exception type are enough to correlate a traceback.
|
|
||||||
sys.stderr.write(
|
|
||||||
f"orchestrator: {method} {self.path} failed "
|
|
||||||
f"[error_type={type(e).__name__}]\n"
|
|
||||||
)
|
|
||||||
sys.stderr.flush()
|
|
||||||
status, payload = 500, {"error": "internal error"}
|
|
||||||
self._write_json(status, payload)
|
|
||||||
|
|
||||||
def _write_json(self, status: int, payload: Json) -> None:
|
|
||||||
data = json.dumps(payload).encode()
|
|
||||||
self.send_response(status)
|
|
||||||
self.send_header("Content-Type", "application/json")
|
|
||||||
self.send_header("Content-Length", str(len(data)))
|
|
||||||
self.end_headers()
|
|
||||||
self.wfile.write(data)
|
|
||||||
|
|
||||||
def do_GET(self) -> None:
|
|
||||||
self._serve("GET")
|
|
||||||
|
|
||||||
def do_POST(self) -> None:
|
|
||||||
self._serve("POST")
|
|
||||||
|
|
||||||
def do_PUT(self) -> None:
|
|
||||||
self._serve("PUT")
|
|
||||||
|
|
||||||
def do_DELETE(self) -> None:
|
|
||||||
self._serve("DELETE")
|
|
||||||
|
|
||||||
|
|
||||||
class OrchestratorServer(socketserver.ThreadingMixIn, http.server.HTTPServer):
|
|
||||||
"""Threading HTTP server that carries the orchestrator for its handlers.
|
|
||||||
|
|
||||||
Holds the per-host control-plane *signing key* (from
|
|
||||||
`$BOT_BOTTLE_ORCHESTRATOR_TOKEN`, injected by the launcher into the
|
|
||||||
orchestrator process only) and verifies each request's role-scoped token
|
|
||||||
against it. Every route but `/health` requires a valid token whose role
|
|
||||||
covers the route. Construction fails when the key is absent so a new or
|
|
||||||
misconfigured launcher cannot accidentally expose an open control plane."""
|
|
||||||
|
|
||||||
daemon_threads = True
|
|
||||||
allow_reuse_address = True
|
|
||||||
|
|
||||||
def __init__(
|
|
||||||
self,
|
|
||||||
address: tuple[str, int],
|
|
||||||
orchestrator: OrchestratorCore,
|
|
||||||
*,
|
|
||||||
signing_key: str,
|
|
||||||
) -> None:
|
|
||||||
self.orchestrator = orchestrator
|
|
||||||
self._signing_key = signing_key.strip()
|
|
||||||
if not self._signing_key:
|
|
||||||
raise ValueError(
|
|
||||||
"orchestrator control-plane signing key is required; "
|
|
||||||
"refusing to start without caller authentication"
|
|
||||||
)
|
|
||||||
self._request_slots = threading.BoundedSemaphore(MAX_REQUEST_THREADS)
|
|
||||||
super().__init__(address, Handler)
|
|
||||||
|
|
||||||
def get_request(self) -> tuple[socket.socket, typing.Any]:
|
|
||||||
request, client_address = super().get_request()
|
|
||||||
request.settimeout(REQUEST_TIMEOUT_SECONDS)
|
|
||||||
return request, client_address
|
|
||||||
|
|
||||||
def process_request(
|
|
||||||
self, request: typing.Any, client_address: typing.Any,
|
|
||||||
) -> None:
|
|
||||||
# Bound concurrency before ThreadingMixIn creates a worker. Backpressure
|
|
||||||
# stays in the accept loop instead of allocating an unbounded thread per
|
|
||||||
# slow or malicious connection.
|
|
||||||
self._request_slots.acquire()
|
|
||||||
try:
|
|
||||||
super().process_request(request, client_address)
|
|
||||||
except BaseException:
|
|
||||||
self._request_slots.release()
|
|
||||||
raise
|
|
||||||
|
|
||||||
def process_request_thread(
|
|
||||||
self, request: typing.Any, client_address: typing.Any,
|
|
||||||
) -> None:
|
|
||||||
try:
|
|
||||||
super().process_request_thread(request, client_address)
|
|
||||||
finally:
|
finally:
|
||||||
self._request_slots.release()
|
self._stopped.set()
|
||||||
|
|
||||||
def role_for(self, presented: str) -> str | None:
|
def serve_forever(self) -> None:
|
||||||
"""The verified caller role, or None for a missing/invalid token."""
|
self.run()
|
||||||
return CONTROL_PLANE.verify(presented, self._signing_key)
|
|
||||||
|
def shutdown(self) -> None:
|
||||||
|
self._server.should_exit = True
|
||||||
|
self._stopped.wait(timeout=5)
|
||||||
|
|
||||||
|
def server_close(self) -> None:
|
||||||
|
self._socket.close()
|
||||||
|
|
||||||
|
|
||||||
def make_server(
|
def make_server(
|
||||||
@@ -506,18 +89,29 @@ def make_server(
|
|||||||
*,
|
*,
|
||||||
signing_key: str | None = None,
|
signing_key: str | None = None,
|
||||||
) -> OrchestratorServer:
|
) -> OrchestratorServer:
|
||||||
"""Build an authenticated control-plane server.
|
"""Build a bounded Uvicorn server around the orchestrator application."""
|
||||||
|
|
||||||
``signing_key=None`` reads the owning process's injected environment.
|
|
||||||
Empty or missing keys are rejected by :class:`OrchestratorServer`.
|
|
||||||
"""
|
|
||||||
key = CONTROL_PLANE.key_from_env() if signing_key is None else signing_key
|
key = CONTROL_PLANE.key_from_env() if signing_key is None else signing_key
|
||||||
return OrchestratorServer(
|
app = create_app(orchestrator, signing_key=key)
|
||||||
(host, port), orchestrator, signing_key=key,
|
config = uvicorn.Config(
|
||||||
|
app,
|
||||||
|
host=host,
|
||||||
|
port=port,
|
||||||
|
access_log=bool(os.environ.get("BOT_BOTTLE_ORCHESTRATOR_DEBUG")),
|
||||||
|
log_level="info",
|
||||||
|
limit_concurrency=MAX_REQUESTS,
|
||||||
|
timeout_keep_alive=KEEP_ALIVE_TIMEOUT_SECONDS,
|
||||||
|
server_header=False,
|
||||||
)
|
)
|
||||||
|
return OrchestratorServer(config)
|
||||||
|
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"dispatch", "Handler", "OrchestratorServer", "make_server", "Json",
|
"KEEP_ALIVE_TIMEOUT_SECONDS",
|
||||||
"ORCHESTRATOR_AUTH_HEADER", "MAX_BODY_BYTES",
|
"MAX_BODY_BYTES",
|
||||||
|
"MAX_REQUESTS",
|
||||||
|
"ORCHESTRATOR_AUTH_HEADER",
|
||||||
|
"OrchestratorServer",
|
||||||
|
"create_app",
|
||||||
|
"dispatch",
|
||||||
|
"make_server",
|
||||||
]
|
]
|
||||||
|
|||||||
Reference in New Issue
Block a user