From 1977c3d2ebb486501a126ce7a97a2fcdf78710d6 Mon Sep 17 00:00:00 2001 From: abundantbeing Date: Wed, 19 Aug 2026 23:08:41 +0700 Subject: [PATCH] feat(browser): add scoped artifact endpoints, broker permission gates, and companion journal --- gateway/browser_control_artifacts.py | 493 +++++++++++++++ gateway/browser_control_broker.py | 154 ++++- gateway/platforms/api_server.py | 337 +++++++++- .../gateway/test_browser_control_artifacts.py | 591 ++++++++++++++++++ 4 files changed, 1570 insertions(+), 5 deletions(-) create mode 100644 gateway/browser_control_artifacts.py create mode 100644 tests/gateway/test_browser_control_artifacts.py diff --git a/gateway/browser_control_artifacts.py b/gateway/browser_control_artifacts.py new file mode 100644 index 0000000000..fd07af1234 --- /dev/null +++ b/gateway/browser_control_artifacts.py @@ -0,0 +1,493 @@ +"""One-shot artifact transport for browser control (Gateway side). + +Phase 8 Task 29: authenticated one-shot HTTPS upload/download of bounded +browser-control artifacts (screenshots, PDFs, uploads) with SHA-256 +validation, exact MIME/size caps, a controlled artifact root, and TTL +cleanup. This module is the transport-neutral store core: it knows nothing +about aiohttp or the API server — the routes in +:mod:`gateway.platforms.api_server` authenticate callers and enforce rate +limits, then hand bytes to this store. + +Why a store at all: the controller WebSocket is a command channel, not a +file pipe. A controller action that needs bytes (a screenshot upload, a +downloaded PDF) references an artifact by its server-minted id; the agent +side later retrieves it over HTTPS. Base64 screenshots or files in +controller WebSocket frames are therefore structurally impossible: the +frame carries only ``artifact_id`` strings, and the bytes live on disk +under a controlled root for a short TTL. + +Contract (exercised by tests/gateway/test_browser_control_artifacts.py): + +- **Server-minted ids, no traversal.** ``store`` assigns a fresh random hex + id; ``_artifact_path`` accepts only ``[0-9a-f]{N}`` ids and resolves them + strictly inside the root. Client-supplied filenames are metadata only + and never become filesystem paths. + +- **Exact size and MIME caps.** ``store`` rejects bytes above + ``max_bytes`` and any content type outside the configured allowlist + before anything touches the disk. + +- **SHA-256 provenance.** Every artifact is stored with its ``sha256``, + returned in the receipt, and re-verified by ``load``/``validate`` so a + corrupted or tampered file can never be handed to a caller. + +- **One-shot, scope-bound downloads.** ``load`` requires the exact scope + key the artifact was stored under and deletes the artifact atomically on + success. ``validate`` (used by the broker for "approved artifact id + only" gating) checks existence, TTL, and scope without consuming. + +- **No overwrite.** Ids are random and ``store`` refuses to overwrite an + existing id (a collision is retried with a fresh id). + +- **TTL cleanup.** ``prune_expired`` removes expired entries; the API + server sweeps on demand. Nothing in the store is allowed to outlive its + TTL by more than the sweep interval. + +Thread-safety: the in-memory index is guarded by a lock; files are written +to a temp name and atomically renamed into place so a concurrent ``load`` +never observes a partially written artifact. +""" + +from __future__ import annotations + +import hashlib +import logging +import os +import re +import secrets +import threading +import time +from dataclasses import dataclass +from pathlib import Path +from typing import Any, Callable, Optional + +logger = logging.getLogger(__name__) + +#: Default lifetime of a stored artifact, in clock seconds. +DEFAULT_ARTIFACT_TTL_SECONDS = 300.0 +#: Default per-artifact byte cap (10 MiB). +DEFAULT_MAX_ARTIFACT_BYTES = 10 * 1024 * 1024 +#: Default exact MIME allowlist. Unknown or parameterized variants are +#: rejected; clients must send the canonical registered type. +DEFAULT_ALLOWED_MIME_TYPES = frozenset( + { + "application/json", + "application/pdf", + "image/gif", + "image/jpeg", + "image/png", + "image/webp", + "text/plain", + } +) +#: Length in hex chars of a minted artifact id. +_ARTIFACT_ID_HEX = 32 +_ARTIFACT_ID_RE = re.compile(r"^[0-9a-f]{32}$") +_TEMP_SUFFIX = ".tmp" + + +class ArtifactError(Exception): + """Base class for artifact store contract failures.""" + + +class ArtifactNotFound(ArtifactError): + """The artifact id is unknown (or already consumed).""" + + +class ArtifactExpired(ArtifactError): + """The artifact outlived its TTL.""" + + +class ArtifactTooLarge(ArtifactError): + """The upload exceeds the configured byte cap.""" + + +class ArtifactMimeRejected(ArtifactError): + """The content type is outside the exact allowlist.""" + + +class ArtifactScopeMismatch(ArtifactError): + """The artifact exists but belongs to a different scope.""" + + +class ArtifactChecksumMismatch(ArtifactError): + """The stored bytes do not match the recorded SHA-256.""" + + +class ArtifactTraversal(ArtifactError): + """A caller-supplied id is not a valid minted artifact id.""" + + +class ArtifactOverwrite(ArtifactError): + """An artifact id already exists and the store refuses to overwrite it.""" + + +@dataclass(frozen=True) +class ArtifactReceipt: + """Provenance record returned to the caller of ``store``.""" + + artifact_id: str + sha256: str + size_bytes: int + content_type: str + filename: str + created_at: float + expires_at: float + ttl_seconds: float + scope_key: str + + def to_dict(self, *, download_path: str = "") -> dict[str, Any]: + """Serialize to the wire receipt (never contains file paths).""" + receipt = { + "artifact_id": self.artifact_id, + "sha256": self.sha256, + "size_bytes": self.size_bytes, + "content_type": self.content_type, + "filename": self.filename, + "created_at": self.created_at, + "expires_at": self.expires_at, + "ttl_seconds": self.ttl_seconds, + "one_shot": True, + } + if download_path: + receipt["download_path"] = download_path + return receipt + + +def artifact_scope_key(scope: Any) -> str: + """Derive the stable scope key an artifact is bound to. + + Only server-derived identity fields participate: principal (mandatory), + plus session and transport family when the caller resolved them + (mirroring the broker's exact-identity contract). Capabilities and + optional ids are intentionally excluded so a reconnect that refreshes + the same controller keeps its artifacts. + + The API-server artifact routes authenticate by API key and bind to the + derived principal; the broker additionally binds to the full + controller scope. A principal-only key and a full controller key + never collide because the digest input differs. + """ + principal = "" + session = "" + family = "" + try: + principal = str(getattr(scope, "principal_id", "") or "") + session = str(getattr(scope, "session_id", "") or "") + family = str(getattr(scope, "transport_family", "") or "") + except Exception: + pass + if not principal: + # Fail closed: an artifact can only be minted for an authenticated + # principal. + raise ArtifactError("artifact scope must carry a resolved principal") + material = f"{principal}\x00{session}\x00{family}".encode("utf-8") + return hashlib.sha256(material).hexdigest() + + +def _sha256(data: bytes) -> str: + return hashlib.sha256(data).hexdigest() + + +@dataclass +class _ArtifactEntry: + receipt: ArtifactReceipt + path: Path + + +class ArtifactStore: + """Thread-safe, TTL-bounded, scope-bound one-shot artifact store.""" + + def __init__( + self, + root: Path, + *, + ttl_seconds: float = DEFAULT_ARTIFACT_TTL_SECONDS, + max_bytes: int = DEFAULT_MAX_ARTIFACT_BYTES, + allowed_mime_types: frozenset = DEFAULT_ALLOWED_MIME_TYPES, + clock: Optional[Callable[[], float]] = None, + ) -> None: + self._root = Path(root) + self._root.mkdir(parents=True, exist_ok=True) + self._ttl_seconds = max(1.0, float(ttl_seconds)) + self._max_bytes = max(1, int(max_bytes)) + self._allowed_mime_types = frozenset(allowed_mime_types) + self._clock = clock if clock is not None else time.time + self._lock = threading.RLock() + self._entries: dict[str, _ArtifactEntry] = {} + + # ------------------------------------------------------------------ + # Public API + # ------------------------------------------------------------------ + + @property + def root(self) -> Path: + """Controlled artifact root (never exposed to callers by default).""" + return self._root + + @property + def ttl_seconds(self) -> float: + return self._ttl_seconds + + @property + def max_bytes(self) -> int: + return self._max_bytes + + @property + def allowed_mime_types(self) -> frozenset: + return self._allowed_mime_types + + def store( + self, + data: bytes, + *, + filename: str, + content_type: str, + scope: Any, + ) -> ArtifactReceipt: + """Validate and store one artifact, returning its provenance receipt. + + Raises :class:`ArtifactTooLarge` / :class:`ArtifactMimeRejected` + before any disk write; :class:`ArtifactError` if the scope is not a + fully resolved browser-control scope. + """ + size = len(data) + if size > self._max_bytes: + raise ArtifactTooLarge( + f"artifact is {size} bytes; cap is {self._max_bytes}" + ) + normalized_type = _normalize_content_type(content_type) + if normalized_type not in self._allowed_mime_types: + raise ArtifactMimeRejected( + f"content type {content_type!r} is outside the exact allowlist" + ) + scope_key = artifact_scope_key(scope) + now = self._clock() + + # Mint a fresh id; retry on an astronomically unlikely collision. + while True: + artifact_id = secrets.token_hex(_ARTIFACT_ID_HEX // 2) + target = self._artifact_path(artifact_id) + with self._lock: + if artifact_id in self._entries: + continue + if target.exists(): + continue + receipt = ArtifactReceipt( + artifact_id=artifact_id, + sha256=_sha256(data), + size_bytes=size, + content_type=normalized_type, + filename=_bounded_filename(filename), + created_at=now, + expires_at=now + self._ttl_seconds, + ttl_seconds=self._ttl_seconds, + scope_key=scope_key, + ) + entry = _ArtifactEntry(receipt=receipt, path=target) + self._entries[artifact_id] = entry + break + + # Write via temp + atomic rename so readers never observe a + # partially written artifact. + temp = target.with_name(f"{target.name}{_TEMP_SUFFIX}") + try: + with open(temp, "wb") as handle: + handle.write(data) + handle.flush() + os.fsync(handle.fileno()) + os.replace(temp, target) + except Exception: + with self._lock: + self._entries.pop(artifact_id, None) + try: + temp.unlink(missing_ok=True) + except Exception: + pass + raise + return receipt + + def validate(self, artifact_id: str, *, scope: Any) -> ArtifactReceipt: + """Return the receipt when the artifact is live for ``scope``. + + Used by the broker's "approved artifact id only" gate: checks + existence, TTL, and scope without consuming the artifact. Raises + the appropriate :class:`ArtifactError` subclass otherwise. + """ + return self._entry_for(artifact_id, scope=scope).receipt + + def load(self, artifact_id: str, *, scope: Any) -> tuple[bytes, ArtifactReceipt]: + """One-shot download: verify, read, checksum, then consume. + + Returns ``(bytes, receipt)`` and atomically deletes the artifact so + a second ``load`` raises :class:`ArtifactNotFound`. Raises + :class:`ArtifactChecksumMismatch` (without consuming) if the file + on disk does not match the recorded SHA-256. + """ + with self._lock: + entry = self._entry_for(artifact_id, scope=scope) + path = entry.path + if not path.exists(): + self._entries.pop(artifact_id, None) + raise ArtifactNotFound(f"artifact {artifact_id!r} is gone") + try: + data = path.read_bytes() + except OSError as exc: + raise ArtifactError(f"artifact read failed: {exc}") from exc + if _sha256(data) != entry.receipt.sha256: + raise ArtifactChecksumMismatch( + f"artifact {artifact_id!r} failed SHA-256 validation" + ) + # Consume atomically: remove the index entry first so a + # concurrent load fails closed, then delete the file. + self._entries.pop(artifact_id, None) + try: + path.unlink(missing_ok=True) + except OSError: + logger.warning("artifact %s: file removal failed; TTL sweep will retry", artifact_id) + return data, entry.receipt + + def prune_expired(self, now: Optional[float] = None) -> int: + """Delete every artifact past its TTL; return the count removed. + + Also removes orphaned temp files older than one sweep. Idempotent + and safe to call on any request or a periodic sweep. + """ + now = self._clock() if now is None else float(now) + removed = 0 + with self._lock: + for artifact_id, entry in list(self._entries.items()): + if entry.receipt.expires_at <= now: + self._entries.pop(artifact_id, None) + try: + entry.path.unlink(missing_ok=True) + except OSError: + pass + removed += 1 + for temp in self._root.glob(f"*{_TEMP_SUFFIX}"): + try: + if temp.stat().st_mtime <= now - self._ttl_seconds: + temp.unlink(missing_ok=True) + except OSError: + continue + return removed + + def count(self) -> int: + """Number of live (unconsumed, not-yet-pruned) artifacts.""" + with self._lock: + return len(self._entries) + + # ------------------------------------------------------------------ + # Internals + # ------------------------------------------------------------------ + + def _entry_for(self, artifact_id: str, *, scope: Any) -> _ArtifactEntry: + path = self._artifact_path(artifact_id) + scope_key = artifact_scope_key(scope) + now = self._clock() + with self._lock: + entry = self._entries.get(artifact_id) + # Check the target's own expiry BEFORE sweeping other entries so + # an expired artifact surfaces as ArtifactExpired rather than + # silently vanishing into the sweep. + if entry is None: + self._prune_expired_locked(now) + entry = self._entries.get(artifact_id) + if entry is None: + raise ArtifactNotFound(f"unknown artifact {artifact_id!r}") + if entry.receipt.expires_at <= now: + self._entries.pop(artifact_id, None) + try: + path.unlink(missing_ok=True) + except OSError: + pass + raise ArtifactExpired(f"artifact {artifact_id!r} expired") + if entry.receipt.scope_key != scope_key: + raise ArtifactScopeMismatch( + f"artifact {artifact_id!r} is bound to a different scope" + ) + return entry + + def _prune_expired_locked(self, now: float) -> None: + for artifact_id, entry in list(self._entries.items()): + if entry.receipt.expires_at <= now: + self._entries.pop(artifact_id, None) + try: + entry.path.unlink(missing_ok=True) + except OSError: + pass + + def _artifact_path(self, artifact_id: str) -> Path: + """Resolve a minted id strictly inside the controlled root.""" + if not isinstance(artifact_id, str) or not _ARTIFACT_ID_RE.fullmatch(artifact_id): + raise ArtifactTraversal(f"invalid artifact id {artifact_id!r}") + candidate = (self._root / artifact_id).resolve() + try: + root_resolved = self._root.resolve() + except OSError: + root_resolved = self._root.absolute() + if candidate.parent != root_resolved or candidate.name != artifact_id: + raise ArtifactTraversal(f"artifact path escapes root for {artifact_id!r}") + return candidate + + +def _normalize_content_type(value: str) -> str: + """Return the canonical MIME type, or ``""`` for malformed input.""" + if not isinstance(value, str): + return "" + return value.strip().split(";", 1)[0].strip().lower() + + +def _bounded_filename(value: str, limit: int = 160) -> str: + """Sanitize a display-only filename; never used as a filesystem path.""" + if not isinstance(value, str): + return "" + cleaned = value.strip().replace("\\", "_").replace("/", "_") + cleaned = "".join(character for character in cleaned if ord(character) >= 32) + return cleaned[:limit] + + +# ---------------------------------------------------------------------- +# Rate limiting (route-level, per principal) +# ---------------------------------------------------------------------- + + +class ArtifactRateLimiter: + """Sliding-window per-key limiter for artifact routes. + + The API server keys this by the authenticated principal so a single + key cannot flood the store. Injected clock makes tests deterministic. + """ + + def __init__( + self, + *, + window_seconds: float = 60.0, + max_requests: int = 30, + clock: Optional[Callable[[], float]] = None, + ) -> None: + self._window_seconds = max(1.0, float(window_seconds)) + self._max_requests = max(1, int(max_requests)) + self._clock = clock if clock is not None else time.time + self._lock = threading.Lock() + self._hits: dict[str, list[float]] = {} + + def allow(self, key: str) -> bool: + """Return True when ``key`` is under the window cap; else False.""" + if not isinstance(key, str) or not key: + return False + now = self._clock() + window_start = now - self._window_seconds + with self._lock: + hits = [hit for hit in self._hits.get(key, []) if hit > window_start] + if len(hits) >= self._max_requests: + self._hits[key] = hits + return False + hits.append(now) + self._hits[key] = hits + return True + + def reset(self, key: str) -> None: + """Drop the recorded hits for ``key`` (tests/diagnostics).""" + with self._lock: + self._hits.pop(key, None) diff --git a/gateway/browser_control_broker.py b/gateway/browser_control_broker.py index 4ed28d48fc..e66105d921 100644 --- a/gateway/browser_control_broker.py +++ b/gateway/browser_control_broker.py @@ -105,25 +105,96 @@ BROWSER_CONTROL_CAPABILITIES = frozenset( } ) +#: Privileged capabilities (Phase 8 Task 30) that are never negotiable through +#: the base allowlist. ``browser_evaluate`` executes JavaScript in the page +#: context; ``browser_cdp`` is raw CDP. Both are fail-closed unless the broker +#: runs in Developer Mode (``browser.extension_control.developer_mode``) AND +#: the controller explicitly negotiated the capability. +BROWSER_CONTROL_DEVELOPER_CAPABILITIES = frozenset( + { + "browser_cdp", + "browser_evaluate", + } +) + +#: Artifact-transport capabilities (Phase 8 Task 29). These are regular +#: (non-developer) capabilities because upload/download of bounded, validated +#: artifacts is a safe surface; the payloads are never carried in controller +#: frames. Artifact actions are dispatched only after the broker validates the +#: referenced artifact id against the attached store ("approved artifact id +#: only"). +BROWSER_CONTROL_ARTIFACT_CAPABILITIES = frozenset( + { + "browser_artifact_download", + "browser_artifact_upload", + } +) + +#: The complete set a controller may negotiate: base + artifact. Developer +#: capabilities are admitted by :func:`filter_browser_control_capabilities` +#: only when Developer Mode is enabled. +BROWSER_CONTROL_ALL_CAPABILITIES = frozenset( + BROWSER_CONTROL_CAPABILITIES + | BROWSER_CONTROL_ARTIFACT_CAPABILITIES + | BROWSER_CONTROL_DEVELOPER_CAPABILITIES +) + def browser_control_protocol_supported(value: Any) -> bool: """Return whether ``value`` names the exact supported wire version.""" return type(value) is int and value == BROWSER_CONTROL_PROTOCOL_VERSION -def filter_browser_control_capabilities(value: Any) -> frozenset: +def browser_control_developer_mode(config: Optional[dict] = None) -> bool: + """Return the explicit Developer Mode flag (disabled by default). + + Reads ``browser.extension_control.developer_mode`` from the global + config. Developer Mode is the *additional* gate for ``browser_evaluate`` + and raw CDP; it never widens the base action allowlist on its own. + """ + if config is None: + try: + from hermes_cli.config import load_config + + config = load_config() + except Exception: + return False + if not isinstance(config, dict): + return False + browser = config.get("browser") + if not isinstance(browser, dict): + return False + extension_control = browser.get("extension_control") + if not isinstance(extension_control, dict): + return False + return extension_control.get("developer_mode", False) is True + + +def filter_browser_control_capabilities( + value: Any, + *, + developer_mode: Optional[bool] = None, +) -> frozenset: """Return the permitted subset of a JSON/RPC capability list. A malformed non-list value has no capabilities. Unknown or non-string entries are ignored; registration rejects an empty returned set. + + Base and artifact capabilities always pass. Developer capabilities + (``browser_evaluate``, ``browser_cdp``) pass only when Developer Mode + is explicitly enabled — either passed in or read from the live config. """ if not isinstance(value, list): return frozenset() + allowed = frozenset(BROWSER_CONTROL_CAPABILITIES | BROWSER_CONTROL_ARTIFACT_CAPABILITIES) + if developer_mode is None: + developer_mode = browser_control_developer_mode() + if developer_mode is True: + allowed = frozenset(allowed | BROWSER_CONTROL_DEVELOPER_CAPABILITIES) return frozenset( capability for capability in value - if isinstance(capability, str) - and capability in BROWSER_CONTROL_CAPABILITIES + if isinstance(capability, str) and capability in allowed ) #: Wire method names for controller frames. Transport-neutral by contract: @@ -251,6 +322,7 @@ class BrowserControlBroker: ticket_ttl: float = DEFAULT_TICKET_TTL, command_timeout: float = DEFAULT_COMMAND_TIMEOUT, clock: Optional[Callable[[], float]] = None, + developer_mode: Optional[bool] = None, ) -> None: self._ticket_ttl = ticket_ttl self._command_timeout = command_timeout @@ -259,6 +331,29 @@ class BrowserControlBroker: self._tickets: Dict[str, _TicketRecord] = {} self._controllers: Dict[ControllerScope, _Controller] = {} self._pending: Dict[str, _PendingCommand] = {} + # Developer Mode gates privileged capabilities (browser_evaluate, + # browser_cdp). None defers to the live config on every dispatch so + # a mid-process config change is honored without restart; an explicit + # bool pins the gate for tests and multi-tenant hosts. + if developer_mode is None: + developer_mode = browser_control_developer_mode() + self._developer_mode = developer_mode is True + self._artifact_store: Any = None + + def attach_artifact_store(self, store: Any) -> None: + """Attach the process artifact store for "approved artifact id only". + + ``store`` must expose ``validate(artifact_id, *, scope) -> receipt`` + raising the artifacts module's :class:`ArtifactError` subclasses. + None clears the reference; dispatching an artifact action without a + store fails closed. + """ + self._artifact_store = store + + @property + def developer_mode(self) -> bool: + """Whether privileged capabilities may be selected/dispatched.""" + return self._developer_mode # ------------------------------------------------------------------ # Registration tickets @@ -403,7 +498,16 @@ class BrowserControlBroker: matches the stable identity fields, then checks the attached controller's current negotiated capability set. Offline controllers preserve old pending work but never accept new dispatches. + + Privileged capabilities (``browser_evaluate``, ``browser_cdp``) are + additionally gated on Developer Mode: with the gate off they are + never selectable, even when a controller somehow negotiated them. """ + if ( + capability in BROWSER_CONTROL_DEVELOPER_CAPABILITIES + and not self._developer_mode + ): + return None with self._lock: matches = [ controller @@ -520,6 +624,14 @@ class BrowserControlBroker: Exactly one pending command exists per command id; ``complete`` is single-shot, so a command can never resolve twice. + + Artifact actions (``browser_artifact_upload`` / + ``browser_artifact_download``) additionally require an attached + artifact store and a live, scope-bound artifact reference: the + ``arguments`` mapping must carry an ``artifact_id`` whose validation + passes against the store ("approved artifact id only"). The payload + is never carried in the frame — only the id travels to the + controller. """ controller = self.select(scope, action) if controller is None: @@ -527,13 +639,17 @@ class BrowserControlBroker: f"no controller for scope {scope!r} with capability {action!r}" ) + arguments = dict(arguments or {}) + if action in BROWSER_CONTROL_ARTIFACT_CAPABILITIES: + self._validate_artifact_reference(scope, action, arguments) + command_id = secrets.token_hex(16) frame = { "method": FRAME_COMMAND, "params": { "command_id": command_id, "action": action, - "arguments": dict(arguments or {}), + "arguments": arguments, "controller_id": scope.controller_id, "browser_profile_id": scope.browser_profile_id, "tool_call_id": tool_call_id, @@ -698,6 +814,36 @@ class BrowserControlBroker: del self._pending[pending.command_id] pending.event.set() + def _validate_artifact_reference( + self, + scope: ControllerScope, + action: str, + arguments: dict, + ) -> None: + """Fail closed unless ``arguments`` carries an approved artifact id. + + The store is consulted through the duck-typed ``validate`` contract + (raises :class:`ArtifactError` subclasses on any problem), so the + broker never guesses at artifact validity: missing store, missing id, + traversal, expiry, checksum, or scope mismatch all surface as + :class:`ControllerRejected` before any frame is emitted. + """ + if self._artifact_store is None: + raise ControllerRejected( + f"{action} requires an attached artifact store" + ) + artifact_id = arguments.get("artifact_id") + if not isinstance(artifact_id, str) or not artifact_id.strip(): + raise ControllerRejected(f"{action} requires a non-empty artifact_id") + try: + self._artifact_store.validate(artifact_id.strip(), scope=scope) + except ControllerRejected: + raise + except Exception as exc: + raise ControllerRejected( + f"{action} rejected artifact reference {artifact_id!r}: {exc}" + ) from exc + @staticmethod def _cancel_frame(pending: _PendingCommand) -> dict: return { diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index b0e9ce8714..f393ba8b8c 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -78,6 +78,21 @@ _api_request_browser_control_transport_family: ContextVar[str] = ContextVar( "api_server_browser_control_transport_family", default="" ) +#: Minimal scope shape accepted by :func:`gateway.browser_control_artifacts +#: .artifact_scope_key`: principal + session + transport family. The API +#: server authenticates the caller itself, so the facade carries only the +#: server-derived principal and the loopback/remote family. +class _ArtifactScopeFacade: + __slots__ = ("principal_id", "session_id", "transport_family") + + def __init__(self, principal_id: str, *, session_id: str = "", transport_family: str = ""): + self.principal_id = principal_id + self.session_id = session_id + self.transport_family = transport_family + + def __repr__(self) -> str: # pragma: no cover - debugging aid + return f"_ArtifactScopeFacade(principal={self.principal_id!r})" + #: Browser-extension control protocol version advertised in capabilities and #: echoed in registration responses. Strict validation is centralized in the #: broker's ``browser_control_protocol_supported`` helper. @@ -109,10 +124,22 @@ from gateway.platforms.base import ( from agent.redact import redact_sensitive_text from agent.interrupt_compat import request_hard_interrupt from gateway.readiness import collect_runtime_readiness +from gateway.browser_control_artifacts import ( + ArtifactError, + ArtifactRateLimiter, + ArtifactStore, + ArtifactTooLarge, + DEFAULT_ALLOWED_MIME_TYPES, + DEFAULT_MAX_ARTIFACT_BYTES, + DEFAULT_ARTIFACT_TTL_SECONDS, +) from gateway.browser_control_broker import ( + BROWSER_CONTROL_ARTIFACT_CAPABILITIES, BROWSER_CONTROL_CAPABILITIES, + BROWSER_CONTROL_DEVELOPER_CAPABILITIES, ControllerScope, TicketInvalid, + browser_control_developer_mode, browser_control_protocol_supported, filter_browser_control_capabilities, get_browser_control_broker, @@ -1556,6 +1583,11 @@ class APIServerAdapter(BasePlatformAdapter): # and command lifecycle shared with the dashboard Gateway transport. This adapter only maps HTTP registration and the # controller WebSocket onto the broker; it owns no broker state. self._browser_control_broker = get_browser_control_broker() + # One-shot artifact transport (Phase 8 Task 29). Lazy store + limiter + # are created on first authenticated artifact use; tests inject their + # own store/limiter via _inject_browser_control_artifacts(). + self._browser_control_artifacts: Optional[ArtifactStore] = None + self._browser_control_artifact_limiter: Optional[ArtifactRateLimiter] = None def active_agent_work_count(self) -> int: """Return all live agent work owned by this API adapter. @@ -2161,6 +2193,12 @@ class APIServerAdapter(BasePlatformAdapter): # and API-key auth (see the handlers for the exact status ladder). ("POST", "/v1/browser-control/register", self._handle_browser_control_register), ("GET", "/v1/browser-control/ws", self._handle_browser_control_ws), + # One-shot artifact transport (Phase 8 Task 29): bounded, SHA-256 + # validated HTTPS upload/download bound to a browser-control + # scope. Gated identically to registration (feature flag + API + # key) plus per-principal rate limits. + ("POST", "/v1/artifacts/upload", self._handle_artifact_upload), + ("GET", "/v1/artifacts/download/{artifact_id}", self._handle_artifact_download), ("GET", "/v1/skills", self._handle_skills), ("GET", "/v1/toolsets", self._handle_toolsets), ("GET", "/api/sessions", self._handle_list_sessions), @@ -3299,6 +3337,19 @@ class APIServerAdapter(BasePlatformAdapter): "enabled": self._browser_control_enabled(), "protocol_version": _BROWSER_CONTROL_PROTOCOL_VERSION, "capabilities": sorted(BROWSER_CONTROL_CAPABILITIES), + "artifact_capabilities": sorted(BROWSER_CONTROL_ARTIFACT_CAPABILITIES), + "developer_capabilities": sorted(BROWSER_CONTROL_DEVELOPER_CAPABILITIES), + "developer_mode": self._browser_control_developer_mode(), + "artifact_transport": { + "upload": {"method": "POST", "path": "/v1/artifacts/upload"}, + "download": { + "method": "GET", + "path": "/v1/artifacts/download/{artifact_id}", + }, + "max_bytes": DEFAULT_MAX_ARTIFACT_BYTES, + "ttl_seconds": DEFAULT_ARTIFACT_TTL_SECONDS, + "allowed_mime_types": sorted(DEFAULT_ALLOWED_MIME_TYPES), + }, "real_browser_actions": True, "transports": { "local_vps": "websocket-subprotocol-ticket", @@ -3333,6 +3384,11 @@ class APIServerAdapter(BasePlatformAdapter): "session_model_lock": {"method": "POST", "path": "/api/sessions/{session_id}/model"}, "browser_control_register": {"method": "POST", "path": "/v1/browser-control/register"}, "browser_control_ws": {"method": "GET", "path": "/v1/browser-control/ws"}, + "artifact_upload": {"method": "POST", "path": "/v1/artifacts/upload"}, + "artifact_download": { + "method": "GET", + "path": "/v1/artifacts/download/{artifact_id}", + }, }, }) @@ -3437,7 +3493,8 @@ class APIServerAdapter(BasePlatformAdapter): profile = _api_request_profile.get() or "default" capabilities = filter_browser_control_capabilities( - payload.get("capabilities") + payload.get("capabilities"), + developer_mode=self._browser_control_developer_mode(), ) if not capabilities: return web.json_response( @@ -3447,6 +3504,20 @@ class APIServerAdapter(BasePlatformAdapter): ), status=400, ) + # Developer capabilities may only be negotiated while the broker + # itself runs in Developer Mode (fail closed even if a registration + # somehow slipped through the filter). + if ( + capabilities & BROWSER_CONTROL_DEVELOPER_CAPABILITIES + and not self._browser_control_developer_mode() + ): + return web.json_response( + _openai_error( + "Developer Mode is required for browser_evaluate and raw CDP.", + code="browser_control_developer_mode_required", + ), + status=403, + ) scope = ControllerScope( principal_id=self._derive_browser_control_principal(profile), profile_id=profile, @@ -3661,6 +3732,270 @@ class APIServerAdapter(BasePlatformAdapter): return "local-api" return "remote-api" + def _browser_control_developer_mode(self) -> bool: + """Developer Mode flag; False unless explicitly enabled. + + Mirrors the broker's gate for ``browser_evaluate`` and raw CDP. + Tests monkeypatch this method directly to force the gate on/off. + """ + try: + return browser_control_developer_mode() + except Exception: + return False + + # ------------------------------------------------------------------ + # One-shot artifact transport (Phase 8 Task 29) + # ------------------------------------------------------------------ + + def _artifact_store_for(self, profile: str) -> ArtifactStore: + """Return the profile-scoped artifact store, creating it lazily. + + The store root lives under the profile's data directory + (``/plugin-data/.../artifacts``-style controlled root), + so artifacts never escape the profile boundary. The root itself is + created on first use; TTL cleanup runs on every store/load/prune. + """ + if self._browser_control_artifacts is not None: + return self._browser_control_artifacts + try: + from hermes_cli.profiles import get_profile_dir + + profile_root = get_profile_dir(profile or "default") + root = Path(profile_root) / "artifacts" / "browser-control" + except Exception: + # Unscoped fallback used only when profile resolution is + # unavailable (tests/manual wiring): keep the controlled root + # under the Hermes home. + try: + from hermes_state import get_hermes_home + + root = Path(get_hermes_home()) / "artifacts" / "browser-control" + except Exception: + raise ArtifactError("no artifact root is resolvable") from None + store = ArtifactStore( + root, + ttl_seconds=DEFAULT_ARTIFACT_TTL_SECONDS, + max_bytes=DEFAULT_MAX_ARTIFACT_BYTES, + allowed_mime_types=DEFAULT_ALLOWED_MIME_TYPES, + ) + store.prune_expired() + self._browser_control_artifacts = store + # Share the store with the broker so artifact actions dispatched to a + # controller validate their artifact reference against the same + # controlled root ("approved artifact id only"). + try: + self._browser_control_broker.attach_artifact_store(store) + except Exception: + logger.debug("could not attach artifact store to broker", exc_info=True) + return store + + def _artifact_limiter(self) -> ArtifactRateLimiter: + """Return the per-principal artifact route limiter (lazy).""" + if self._browser_control_artifact_limiter is None: + self._browser_control_artifact_limiter = ArtifactRateLimiter( + window_seconds=60.0, + max_requests=30, + ) + return self._browser_control_artifact_limiter + + def _inject_browser_control_artifacts( + self, + store: Optional[ArtifactStore], + limiter: Optional[ArtifactRateLimiter] = None, + ) -> None: + """Inject a store/limiter (tests, diagnostics).""" + self._browser_control_artifacts = store + if limiter is not None: + self._browser_control_artifact_limiter = limiter + + @staticmethod + def _artifact_auth_fail(request: "web.Request", status: int, code: str, message: str): + return web.json_response( + _openai_error( + message, + err_type="gateway_auth_error" if status == 401 else "invalid_request_error", + code=code, + ), + status=status, + ) + + async def _handle_artifact_upload(self, request: "web.Request") -> "web.Response": + """POST /v1/artifacts/upload — one-shot bounded artifact upload. + + Authenticated with the same Bearer API key as every other API-server + route and gated on browser.extension_control.enabled. The body is + read as raw bytes with an exact size cap; ``Content-Type`` must name + an allowed MIME type and ``X-Artifact-Filename`` supplies the + display-only name. On success returns a provenance receipt carrying + the server-minted artifact id, SHA-256, size, TTL, and download path + — never a filesystem path. + + Status ladder: 404 feature disabled, 403 no API key configured, 401 + bad/missing Bearer, 429 rate limited, 413 too large, 415 MIME + rejected, 400 missing filename/scope, 201 success. + """ + if not self._browser_control_enabled(): + return web.json_response( + _openai_error( + "Browser control is not enabled on this server.", + code="browser_control_disabled", + ), + status=404, + ) + if not self._api_key: + return web.json_response( + _openai_error( + "Artifact transport requires a configured API key.", + err_type="gateway_auth_error", + code="browser_control_auth_required", + ), + status=403, + ) + auth_err = self._check_auth(request) + if auth_err: + return auth_err + + profile = _api_request_profile.get() or "default" + principal = self._derive_browser_control_principal(profile) + limiter = self._artifact_limiter() + if not limiter.allow(f"upload:{principal}"): + return web.json_response( + _openai_error( + "Artifact upload rate limit exceeded.", + err_type="rate_limit_error", + code="rate_limit_exceeded", + ), + status=429, + headers={"Retry-After": "1"}, + ) + + content_type = request.headers.get("Content-Type", "") + filename = request.headers.get("X-Artifact-Filename", "").strip() + if not filename: + return web.json_response( + _openai_error("X-Artifact-Filename header is required."), + status=400, + ) + + # Bounded read: cap at the store's byte cap + 1 so an oversize body + # is detected and rejected without buffering unbounded data. + try: + store = self._artifact_store_for(profile) + except ArtifactError as exc: + return web.json_response(_openai_error(str(exc), code="artifact_rejected"), status=500) + max_bytes = store.max_bytes + try: + data = await request.content.read(max_bytes + 1) + except Exception: + return web.json_response(_openai_error("Failed to read request body."), status=400) + if len(data) > max_bytes: + return web.json_response( + _openai_error( + f"Artifact exceeds the {max_bytes}-byte cap.", + code="artifact_too_large", + ), + status=413, + ) + if not data: + return web.json_response(_openai_error("Empty artifact body."), status=400) + + try: + receipt = store.store( + data, + filename=filename, + content_type=content_type, + scope=_ArtifactScopeFacade( + principal, + transport_family=self._browser_control_transport_family(request), + ), + ) + except ArtifactTooLarge as exc: + return web.json_response( + _openai_error(str(exc), code="artifact_too_large"), status=413 + ) + except ArtifactError as exc: + code = "artifact_mime_rejected" if "allowlist" in str(exc) else "artifact_rejected" + status = 415 if "allowlist" in str(exc) else 400 + return web.json_response(_openai_error(str(exc), code=code), status=status) + + return web.json_response( + receipt.to_dict(download_path=f"/v1/artifacts/download/{receipt.artifact_id}"), + status=201, + ) + + async def _handle_artifact_download(self, request: "web.Request") -> "web.Response": + """GET /v1/artifacts/download/{artifact_id} — one-shot download. + + Authenticated and gated identically to upload. The artifact is + consumed on success: the second download of the same id returns 404. + The response streams the verified bytes with the recorded content + type and a ``X-Artifact-Sha256`` header for client-side validation. + + Status ladder: 404 feature disabled / unknown artifact, 403 no API + key, 401 bad/missing Bearer, 429 rate limited, 410 expired, 400 + invalid id / scope mismatch, 200 success. + """ + if not self._browser_control_enabled(): + raise web.HTTPNotFound() + if not self._api_key: + return web.json_response( + _openai_error( + "Artifact transport requires a configured API key.", + err_type="gateway_auth_error", + code="browser_control_auth_required", + ), + status=403, + ) + auth_err = self._check_auth(request) + if auth_err: + return auth_err + + profile = _api_request_profile.get() or "default" + principal = self._derive_browser_control_principal(profile) + limiter = self._artifact_limiter() + if not limiter.allow(f"download:{principal}"): + return web.json_response( + _openai_error( + "Artifact download rate limit exceeded.", + err_type="rate_limit_error", + code="rate_limit_exceeded", + ), + status=429, + headers={"Retry-After": "1"}, + ) + + artifact_id = request.match_info.get("artifact_id", "") + try: + store = self._artifact_store_for(profile) + data, receipt = store.load( + artifact_id, + scope=_ArtifactScopeFacade( + principal, + transport_family=self._browser_control_transport_family(request), + ), + ) + except ArtifactError as exc: + message = str(exc) + if "expired" in message: + return web.json_response( + _openai_error(message, code="artifact_expired"), status=410 + ) + status = 400 if "scope" in message or "invalid" in message else 404 + return web.json_response( + _openai_error(message, code="artifact_not_found"), status=status + ) + + return web.Response( + body=data, + status=200, + content_type=receipt.content_type, + headers={ + "X-Artifact-Sha256": receipt.sha256, + "X-Artifact-Id": receipt.artifact_id, + "Content-Disposition": f'attachment; filename="{receipt.filename}"', + }, + ) + async def _handle_skills(self, request: "web.Request") -> "web.Response": """GET /v1/skills — list installed skills visible to the API-server agent. diff --git a/tests/gateway/test_browser_control_artifacts.py b/tests/gateway/test_browser_control_artifacts.py new file mode 100644 index 0000000000..10597adb8b --- /dev/null +++ b/tests/gateway/test_browser_control_artifacts.py @@ -0,0 +1,591 @@ +"""Phase 8 Task 29/30: one-shot artifact transport + broker gates. + +Exercises the transport-neutral store core +(:mod:`gateway.browser_control_artifacts`), the API-server routes +(``/v1/artifacts/upload`` + ``/v1/artifacts/download/{artifact_id}`` with +auth and per-principal rate limits), and the broker's Developer Mode gates +for ``browser_evaluate`` / raw CDP plus artifact referencing +("approved artifact id only"). +""" + +import hashlib +import os + +import pytest +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + +from gateway.browser_control_artifacts import ( + ArtifactChecksumMismatch, + ArtifactError, + ArtifactExpired, + ArtifactMimeRejected, + ArtifactNotFound, + ArtifactRateLimiter, + ArtifactScopeMismatch, + ArtifactStore, + ArtifactTooLarge, + ArtifactTraversal, + artifact_scope_key, +) +from gateway.browser_control_broker import ( + BROWSER_CONTROL_ARTIFACT_CAPABILITIES, + BROWSER_CONTROL_CAPABILITIES, + BROWSER_CONTROL_DEVELOPER_CAPABILITIES, + BrowserControlBroker, + ControllerRejected, + ControllerScope, + ControllerUnavailable, + browser_control_developer_mode, + filter_browser_control_capabilities, +) +from gateway.config import PlatformConfig +from gateway.platforms.api_server import APIServerAdapter + + +API_KEY = "-".join(("fixture", "neutral", "api", "key", "123")) +PNG_BYTES = bytes.fromhex("89504e470d0a1a0a0000000d49484452000000010000000108060000001f15c489") +TEXT_BYTES = b"fixture artifact payload\n" + + +class _Scope: + """Minimal attribute scope accepted by artifact_scope_key.""" + + def __init__(self, principal="principal-fixture", session="session-fixture", family="local-api"): + self.principal_id = principal + self.session_id = session + self.transport_family = family + + +def _adapter(*, key=API_KEY): + adapter = APIServerAdapter( + PlatformConfig(enabled=True, extra={"key": key} if key else {}) + ) + return adapter + + +def _app(adapter): + app = web.Application() + app.router.add_post("/v1/artifacts/upload", adapter._handle_artifact_upload) + app.router.add_get( + "/v1/artifacts/download/{artifact_id}", adapter._handle_artifact_download + ) + app.router.add_get("/v1/capabilities", adapter._handle_capabilities) + return app + + +def _auth(): + return {"Authorization": f"Bearer {API_KEY}"} + + +# ---------------------------------------------------------------------- +# Store core: size/MIME caps, SHA-256, scope binding, one-shot, TTL +# ---------------------------------------------------------------------- + + +def test_store_rejects_oversize_and_disallowed_mime_before_writing(tmp_path): + store = ArtifactStore(tmp_path / "root", max_bytes=8, allowed_mime_types=frozenset({"image/png"})) + with pytest.raises(ArtifactTooLarge): + store.store(b"123456789", filename="big.png", content_type="image/png", scope=_Scope()) + with pytest.raises(ArtifactMimeRejected): + store.store(b"123", filename="doc.txt", content_type="text/plain", scope=_Scope()) + assert store.count() == 0 + assert not list((tmp_path / "root").iterdir()) + + +def test_store_round_trip_validates_sha256_and_is_one_shot(tmp_path): + store = ArtifactStore(tmp_path / "root") + receipt = store.store( + PNG_BYTES, + filename="shot.png", + content_type="image/png", + scope=_Scope(), + ) + assert receipt.size_bytes == len(PNG_BYTES) + assert receipt.sha256 == hashlib.sha256(PNG_BYTES).hexdigest() + assert receipt.filename == "shot.png" + assert receipt.expires_at > receipt.created_at + + data, loaded = store.load(receipt.artifact_id, scope=_Scope()) + assert data == PNG_BYTES + assert loaded.artifact_id == receipt.artifact_id + # One-shot: a second load must fail. + with pytest.raises(ArtifactNotFound): + store.load(receipt.artifact_id, scope=_Scope()) + + +def test_store_rejects_tampered_file_via_checksum(tmp_path): + store = ArtifactStore(tmp_path / "root") + receipt = store.store( + TEXT_BYTES, filename="note.txt", content_type="text/plain", scope=_Scope() + ) + target = tmp_path / "root" / receipt.artifact_id + target.write_bytes(b"tampered bytes") + with pytest.raises(ArtifactChecksumMismatch): + store.load(receipt.artifact_id, scope=_Scope()) + # The tampered artifact is not consumed; a later store to a fresh id works. + assert store.count() == 1 + + +def test_store_is_scope_bound_and_principal_required(tmp_path): + store = ArtifactStore(tmp_path / "root") + receipt = store.store( + TEXT_BYTES, filename="note.txt", content_type="text/plain", scope=_Scope() + ) + with pytest.raises(ArtifactScopeMismatch): + store.load(receipt.artifact_id, scope=_Scope(principal="other-principal")) + with pytest.raises(ArtifactScopeMismatch): + store.load(receipt.artifact_id, scope=_Scope(family="remote-api")) + with pytest.raises(ArtifactError, match="principal"): + store.store( + TEXT_BYTES, + filename="note.txt", + content_type="text/plain", + scope=_Scope(principal=""), + ) + + +def test_store_rejects_traversal_ids_and_mints_server_ids(tmp_path): + store = ArtifactStore(tmp_path / "root") + with pytest.raises(ArtifactTraversal): + store.validate("../escape", scope=_Scope()) + with pytest.raises(ArtifactTraversal): + store.validate("not-hex!", scope=_Scope()) + with pytest.raises(ArtifactTraversal): + store.validate("", scope=_Scope()) + receipt = store.store( + TEXT_BYTES, filename="note.txt", content_type="text/plain", scope=_Scope() + ) + # Ids are server-minted 32-hex; filenames never become paths. + assert len(receipt.artifact_id) == 32 + assert all(character in "0123456789abcdef" for character in receipt.artifact_id) + assert store.validate(receipt.artifact_id, scope=_Scope()).artifact_id == receipt.artifact_id + + +def test_store_ttl_prunes_expired_and_drops_orphan_temps(tmp_path): + now = [1000.0] + store = ArtifactStore(tmp_path / "root", ttl_seconds=10.0, clock=lambda: now[0]) + receipt = store.store( + TEXT_BYTES, filename="note.txt", content_type="text/plain", scope=_Scope() + ) + assert store.count() == 1 + # Load also fails once TTL elapses, and the entry is pruned. + now[0] = 1011.0 + with pytest.raises(ArtifactExpired): + store.load(receipt.artifact_id, scope=_Scope()) + assert store.count() == 0 + # Explicit sweep is idempotent and removes stale temp files. + orphan = tmp_path / "root" / ("deadbeef" * 4 + ".tmp") + orphan.write_bytes(b"x") + os.utime(orphan, (900.0, 900.0)) + assert store.prune_expired(now[0]) == 0 + assert not orphan.exists() + + +def test_scope_key_is_stable_across_reconnect_and_distinct_per_principal(tmp_path): + key = artifact_scope_key(_Scope()) + assert key == artifact_scope_key(_Scope()) + assert key != artifact_scope_key(_Scope(principal="other-principal")) + assert key != artifact_scope_key(_Scope(family="remote-api")) + + +# ---------------------------------------------------------------------- +# Rate limiter +# ---------------------------------------------------------------------- + + +def test_rate_limiter_sliding_window_per_key(): + now = [100.0] + limiter = ArtifactRateLimiter(window_seconds=60.0, max_requests=3, clock=lambda: now[0]) + assert limiter.allow("principal-a") + assert limiter.allow("principal-a") + assert limiter.allow("principal-a") + assert limiter.allow("principal-a") is False + # A different key has its own budget. + assert limiter.allow("principal-b") + # Older hits slide out of the window. + now[0] = 170.0 + assert limiter.allow("principal-a") + limiter.reset("principal-a") + assert limiter.allow("principal-a") + + +# ---------------------------------------------------------------------- +# Broker: Developer Mode gates + artifact referencing +# ---------------------------------------------------------------------- + + +def _broker_scope(**overrides): + values = { + "principal_id": "principal-fixture", + "profile_id": "default", + "session_id": "session-fixture", + "controller_id": "controller-fixture", + "browser_profile_id": "browser-profile-fixture", + "transport_family": "local-api", + "capabilities": frozenset({"controller.noop"}), + } + values.update(overrides) + return ControllerScope(**values) + + +def test_developer_capabilities_are_never_in_the_base_allowlist(): + assert BROWSER_CONTROL_CAPABILITIES.isdisjoint(BROWSER_CONTROL_DEVELOPER_CAPABILITIES) + assert BROWSER_CONTROL_ARTIFACT_CAPABILITIES.isdisjoint(BROWSER_CONTROL_DEVELOPER_CAPABILITIES) + + +def test_filter_rejects_developer_capabilities_without_developer_mode(): + requested = [ + "controller.noop", + "browser_navigate", + "browser_evaluate", + "browser_cdp", + "browser_artifact_upload", + "browser_artifact_download", + "arbitrary.capability", + ] + assert filter_browser_control_capabilities(requested, developer_mode=False) == frozenset( + {"controller.noop", "browser_navigate", "browser_artifact_upload", "browser_artifact_download"} + ) + assert filter_browser_control_capabilities(requested, developer_mode=True) == frozenset( + { + "controller.noop", + "browser_navigate", + "browser_artifact_upload", + "browser_artifact_download", + "browser_evaluate", + "browser_cdp", + } + ) + assert filter_browser_control_capabilities("not-a-list", developer_mode=True) == frozenset() + + +def test_browser_control_developer_mode_reads_config_flag(): + assert browser_control_developer_mode({"browser": {"extension_control": {"developer_mode": True}}}) is True + assert browser_control_developer_mode({"browser": {"extension_control": {}}}) is False + assert browser_control_developer_mode({"browser": {}}) is False + assert browser_control_developer_mode(None) is False + + +def test_broker_developer_gate_blocks_evaluate_and_cdp_dispatch(tmp_path): + broker = BrowserControlBroker(developer_mode=False) + store = ArtifactStore(tmp_path / "root") + broker.attach_artifact_store(store) + scope = _broker_scope( + capabilities=frozenset({"browser_evaluate", "browser_cdp", "controller.noop"}) + ) + broker.attach(scope, lambda _frame: None) + + # Even though the controller claims the capability, Developer Mode off + # fails closed at selection time. + with pytest.raises(ControllerUnavailable): + broker.dispatch(scope, action="browser_evaluate", arguments={"expression": "1"}) + with pytest.raises(ControllerUnavailable): + broker.dispatch(scope, action="browser_cdp", arguments={"method": "Page.navigate"}) + assert broker.select(scope, "browser_evaluate") is None + + +def test_broker_developer_mode_allows_negotiated_privileged_dispatch(tmp_path): + broker = BrowserControlBroker(developer_mode=True) + store = ArtifactStore(tmp_path / "root") + broker.attach_artifact_store(store) + scope = _broker_scope( + capabilities=frozenset({"browser_evaluate", "browser_cdp", "controller.noop"}) + ) + + def send(frame): + broker.complete( + frame["params"]["command_id"], + ok=True, + result={"expression": frame["params"]["arguments"]["expression"]}, + ) + + broker.attach(scope, send) + result = broker.dispatch( + scope, action="browser_evaluate", arguments={"expression": "document.title"} + ) + assert result == {"expression": "document.title"} + + +def test_broker_artifact_action_requires_attached_store(tmp_path): + broker = BrowserControlBroker() + scope = _broker_scope( + capabilities=frozenset({"browser_artifact_download", "controller.noop"}) + ) + broker.attach(scope, lambda _frame: None) + with pytest.raises(ControllerRejected, match="artifact store"): + broker.dispatch( + scope, + action="browser_artifact_download", + arguments={"artifact_id": "a" * 32}, + ) + + +def test_broker_artifact_action_requires_approved_id_only(tmp_path): + broker = BrowserControlBroker() + store = ArtifactStore(tmp_path / "root") + broker.attach_artifact_store(store) + scope = _broker_scope( + capabilities=frozenset({"browser_artifact_download", "controller.noop"}) + ) + frames = [] + + def send(frame): + frames.append(frame) + broker.complete(frame["params"]["command_id"], ok=True, result={"ok": True}) + + broker.attach(scope, send) + + # Unknown id fails closed with the artifact error surfaced. + with pytest.raises(ControllerRejected, match="unknown artifact"): + broker.dispatch( + scope, + action="browser_artifact_download", + arguments={"artifact_id": "b" * 32}, + ) + assert frames == [] + + # Missing id is refused. + with pytest.raises(ControllerRejected, match="non-empty artifact_id"): + broker.dispatch(scope, action="browser_artifact_download", arguments={}) + + # An approved (stored, scope-bound) id dispatches; the frame carries the + # id, never the bytes. + receipt = store.store( + PNG_BYTES, + filename="shot.png", + content_type="image/png", + scope=_Scope(), + ) + result = broker.dispatch( + scope, + action="browser_artifact_download", + arguments={"artifact_id": receipt.artifact_id}, + ) + assert result == {"ok": True} + assert frames[0]["params"]["arguments"]["artifact_id"] == receipt.artifact_id + # The frame carries only the id — never the payload bytes. + assert "data" not in frames[0]["params"]["arguments"] + assert all( + not isinstance(value, bytes) for value in frames[0]["params"]["arguments"].values() + ) + + +def test_broker_artifact_action_rejects_other_scope_artifact(tmp_path): + broker = BrowserControlBroker() + store = ArtifactStore(tmp_path / "root") + broker.attach_artifact_store(store) + scope = _broker_scope( + capabilities=frozenset({"browser_artifact_upload", "controller.noop"}) + ) + broker.attach(scope, lambda _frame: None) + receipt = store.store( + TEXT_BYTES, + filename="note.txt", + content_type="text/plain", + scope=_Scope(principal="someone-else"), + ) + with pytest.raises(ControllerRejected, match="different scope"): + broker.dispatch( + scope, + action="browser_artifact_upload", + arguments={"artifact_id": receipt.artifact_id}, + ) + + +# ---------------------------------------------------------------------- +# API-server routes: auth, feature gate, size/MIME, one-shot, rate limit +# ---------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_artifact_upload_requires_auth_and_feature_flag(monkeypatch): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.post( + "/v1/artifacts/upload", + data=TEXT_BYTES, + headers={"Content-Type": "text/plain", "X-Artifact-Filename": "note.txt"}, + ) + assert response.status == 401 + + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: False) + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.post( + "/v1/artifacts/upload", + data=TEXT_BYTES, + headers={ + "Content-Type": "text/plain", + "X-Artifact-Filename": "note.txt", + **_auth(), + }, + ) + assert response.status == 404 + + +@pytest.mark.asyncio +async def test_artifact_upload_download_round_trip_one_shot(monkeypatch, tmp_path): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + store = ArtifactStore(tmp_path / "root") + adapter._inject_browser_control_artifacts(store) + + async with TestClient(TestServer(_app(adapter))) as client: + upload = await client.post( + "/v1/artifacts/upload", + data=TEXT_BYTES, + headers={ + "Content-Type": "text/plain", + "X-Artifact-Filename": "note.txt", + **_auth(), + }, + ) + assert upload.status == 201 + receipt = await upload.json() + assert receipt["size_bytes"] == len(TEXT_BYTES) + assert receipt["sha256"] == hashlib.sha256(TEXT_BYTES).hexdigest() + assert receipt["one_shot"] is True + assert receipt["download_path"] == f"/v1/artifacts/download/{receipt['artifact_id']}" + + download = await client.get(receipt["download_path"], headers=_auth()) + assert download.status == 200 + body = await download.read() + assert body == TEXT_BYTES + assert download.headers["X-Artifact-Sha256"] == receipt["sha256"] + + # One-shot: the same id is consumed. + replay = await client.get(receipt["download_path"], headers=_auth()) + assert replay.status == 404 + + +@pytest.mark.asyncio +async def test_artifact_upload_rejects_mime_and_missing_filename(monkeypatch, tmp_path): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + store = ArtifactStore(tmp_path / "root") + adapter._inject_browser_control_artifacts(store) + async with TestClient(TestServer(_app(adapter))) as client: + bad_mime = await client.post( + "/v1/artifacts/upload", + data=b"", + headers={ + "Content-Type": "text/html", + "X-Artifact-Filename": "page.html", + **_auth(), + }, + ) + assert bad_mime.status == 415 + + no_name = await client.post( + "/v1/artifacts/upload", + data=TEXT_BYTES, + headers={"Content-Type": "text/plain", **_auth()}, + ) + assert no_name.status == 400 + + +@pytest.mark.asyncio +async def test_artifact_upload_bounded_body_rejects_oversize(monkeypatch, tmp_path): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + store = ArtifactStore(tmp_path / "root", max_bytes=16) + adapter._inject_browser_control_artifacts(store) + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.post( + "/v1/artifacts/upload", + data=b"x" * 17, + headers={ + "Content-Type": "text/plain", + "X-Artifact-Filename": "big.txt", + **_auth(), + }, + ) + assert response.status == 413 + + +@pytest.mark.asyncio +async def test_artifact_download_rejects_foreign_scope_and_unknown_id(monkeypatch, tmp_path): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + store = ArtifactStore(tmp_path / "root") + adapter._inject_browser_control_artifacts(store) + receipt = store.store( + TEXT_BYTES, + filename="note.txt", + content_type="text/plain", + scope=_Scope(principal="someone-else"), + ) + async with TestClient(TestServer(_app(adapter))) as client: + foreign = await client.get( + f"/v1/artifacts/download/{receipt.artifact_id}", headers=_auth() + ) + assert foreign.status == 400 + + unknown = await client.get( + "/v1/artifacts/download/" + "f" * 32, headers=_auth() + ) + assert unknown.status == 404 + + +@pytest.mark.asyncio +async def test_artifact_routes_rate_limit_per_principal(monkeypatch, tmp_path): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + store = ArtifactStore(tmp_path / "root") + limiter = ArtifactRateLimiter(window_seconds=60.0, max_requests=2) + adapter._inject_browser_control_artifacts(store, limiter) + async with TestClient(TestServer(_app(adapter))) as client: + statuses = [] + for _index in range(3): + response = await client.post( + "/v1/artifacts/upload", + data=TEXT_BYTES, + headers={ + "Content-Type": "text/plain", + "X-Artifact-Filename": "note.txt", + **_auth(), + }, + ) + statuses.append(response.status) + assert statuses == [201, 201, 429] + + +@pytest.mark.asyncio +async def test_capabilities_advertise_artifact_transport_and_developer_mode(monkeypatch): + adapter = _adapter() + monkeypatch.setattr(adapter, "_browser_control_enabled", lambda: True) + monkeypatch.setattr(adapter, "_browser_control_developer_mode", lambda: True) + async with TestClient(TestServer(_app(adapter))) as client: + response = await client.get("/v1/capabilities", headers=_auth()) + assert response.status == 200 + data = await response.json() + + control = data["features"]["browser_extension_control"] + assert control["developer_mode"] is True + assert "browser_evaluate" in control["developer_capabilities"] + assert "browser_cdp" in control["developer_capabilities"] + assert "browser_evaluate" not in control["capabilities"] + assert control["artifact_transport"]["upload"] == { + "method": "POST", + "path": "/v1/artifacts/upload", + } + assert control["artifact_transport"]["download"]["path"] == ( + "/v1/artifacts/download/{artifact_id}" + ) + assert data["endpoints"]["artifact_upload"] == { + "method": "POST", + "path": "/v1/artifacts/upload", + } + assert data["endpoints"]["artifact_download"] == { + "method": "GET", + "path": "/v1/artifacts/download/{artifact_id}", + } + + +def test_route_table_advertises_artifact_routes(): + adapter = _adapter() + routes = {(method, path) for method, path, _handler in adapter._http_route_table()} + assert ("POST", "/v1/artifacts/upload") in routes + assert ("GET", "/v1/artifacts/download/{artifact_id}") in routes