"""Orchestrator-side broker transport (issue #468, chunk 1). The signer's half of the launch-broker transport gap. `BrokerClient` satisfies the exact `submit(token)` contract `OrchestratorCore` already depends on (see `broker.SubmitBroker`), but instead of verifying and launching in-process it POSTs the signed token to the host control server over HTTP (stdlib `urllib`, like `orchestrator/client.py`). Because it is drop-in for that interface, wiring a real out-of-process backend does not change the core: it still signs a request and calls `submit()`; only the wire is new. A provenance/schema rejection from the host controller (HTTP 401) is re-raised as the same `BrokerAuthError` the in-process broker raises, so the launch path's rollback-on-failure (`OrchestratorCore.launch_bottle`) behaves identically whether the broker is local or remote. """ from __future__ import annotations import json import urllib.error import urllib.request from .broker import BrokerAuthError, BrokerUnavailableError, LaunchRequest DEFAULT_TIMEOUT_SECONDS = 5.0 class BrokerClientError(RuntimeError): """The host control server *responded*, but with an unexpected status other than the fail-closed 401 (which surfaces as `BrokerAuthError`) — e.g. a 502 backend failure or a malformed body. A definite negative: the host processed the request and it did not launch. (A *no-response* failure — unreachable / timeout / dropped — is the ambiguous `BrokerUnavailableError` instead.)""" class BrokerClient: """Drop-in `submit(token)` that relays a signed request to the host control server. Holds no secret — provenance rides entirely in the signed token, so a caller that can reach this client still cannot forge a launch.""" def __init__(self, base_url: str, *, timeout: float = DEFAULT_TIMEOUT_SECONDS) -> None: self._base = base_url.rstrip("/") self._timeout = timeout def submit(self, token: str) -> LaunchRequest: """POST the signed token to the host controller and return the request it verified and acted on. Raises `BrokerAuthError` on a fail-closed 401 (bad provenance/schema — the same exception the in-process broker raises); `BrokerClientError` if the host *responds* with any other non-success status or a malformed body (a definite negative); or `BrokerUnavailableError` if no response is obtained (unreachable / timeout / dropped) — the **ambiguous** case, where the host may already have acted, so the caller must not roll back.""" data = json.dumps({"token": token}).encode() req = urllib.request.Request( f"{self._base}/broker", data=data, method="POST", headers={"Content-Type": "application/json"}, ) try: with urllib.request.urlopen(req, timeout=self._timeout) as resp: return _request_from(_json_object(resp.read())) except urllib.error.HTTPError as e: detail = _error_detail(e) if e.code == 401: raise BrokerAuthError( detail or "host controller rejected the request" ) from e raise BrokerClientError( f"POST /broker: HTTP {e.code} {detail}".rstrip() ) from e except (urllib.error.URLError, TimeoutError, OSError) as e: # No usable response — unreachable, timed out, or the connection # dropped mid-exchange. Ambiguous: the request may already have # launched the bottle, so this is NOT a definite failure. raise BrokerUnavailableError(f"POST /broker: {e}") from e def _json_object(raw: bytes) -> dict[str, object]: """Parse a JSON object, tolerating an empty or malformed body (→ {}), like the orchestrator client — a bad body becomes a clean 'missing field' error downstream rather than an opaque JSON crash.""" if not raw: return {} try: obj = json.loads(raw) except ValueError: return {} return obj if isinstance(obj, dict) else {} def _error_detail(e: urllib.error.HTTPError) -> str: """The `error` string from a structured error response, best-effort — an error body may be absent or unreadable, in which case there is no detail.""" try: detail = _json_object(e.read()).get("error", "") except Exception: # noqa: BLE001 — the error body is advisory only return "" return detail if isinstance(detail, str) else "" def _request_from(payload: dict[str, object]) -> LaunchRequest: """Reconstruct the verified `LaunchRequest` the controller echoed, so the returned value matches the in-process broker's (which returns the request it acted on). A missing op/bottle_id means a malformed response.""" op = payload.get("op") bottle_id = payload.get("bottle_id") if not isinstance(op, str) or not isinstance(bottle_id, str) or not bottle_id: raise BrokerClientError("host controller response missing op/bottle_id") source_ip = payload.get("source_ip") image_ref = payload.get("image_ref") slot = payload.get("slot") return LaunchRequest( op=op, bottle_id=bottle_id, source_ip=source_ip if isinstance(source_ip, str) else "", image_ref=image_ref if isinstance(image_ref, str) else "", slot=slot if isinstance(slot, int) and not isinstance(slot, bool) else None, ) __all__ = [ "BrokerClient", "BrokerClientError", "DEFAULT_TIMEOUT_SECONDS", ]