From 30397a19056fe69ae1d1937a753dd1c4b6098948 Mon Sep 17 00:00:00 2001 From: Alex Date: Wed, 24 Jun 2026 13:57:39 +0100 Subject: [PATCH] Expose artifacts as MCP resources and add sandbox egress policy Expose a user's artifacts through the MCP server as readable resources: each is listed under an artifact:// URI with its mime type and read on demand as inline text or a base64 blob, bounded by a size cap and scoped strictly to the owning principal resolved from the request's API key, so no artifact is served across tenants. Ship the network-level egress controls the sandbox runner needs but cannot self-apply: a Kubernetes NetworkPolicy that allows public egress while denying RFC1918, link-local, and cloud-metadata ranges, an optional docker-compose egress overlay, and runner network-hardening docs. --- application/core/settings.py | 2 + application/mcp_server.py | 79 +++++- .../services/artifact_resource_service.py | 232 ++++++++++++++++++ .../k8s/deployments/sandbox-deploy.yaml | 71 ++++++ .../sandbox-egress-policy.yaml | 85 +++++++ ...ocker-compose.optional.sandbox-egress.yaml | 71 ++++++ deployment/sandbox/README.md | 53 +++- .../test_artifact_resource_service.py | 176 +++++++++++++ tests/services/test_mcp_server.py | 101 ++++++++ 9 files changed, 862 insertions(+), 8 deletions(-) create mode 100644 application/services/artifact_resource_service.py create mode 100644 deployment/k8s/deployments/sandbox-deploy.yaml create mode 100644 deployment/k8s/network-policies/sandbox-egress-policy.yaml create mode 100644 deployment/optional/docker-compose.optional.sandbox-egress.yaml create mode 100644 tests/services/test_artifact_resource_service.py diff --git a/application/core/settings.py b/application/core/settings.py index a7375cc7..2b264d71 100644 --- a/application/core/settings.py +++ b/application/core/settings.py @@ -364,6 +364,8 @@ class Settings(BaseSettings): ARTIFACT_MAX_BYTES: int = 50 * 1024 * 1024 # cap on a single stored artifact version's bytes ARTIFACT_MAX_COUNT_PER_USER: int = 5000 # cap on artifacts a user may own ARTIFACT_MAX_TOTAL_BYTES_PER_USER: int = 5 * 1024 * 1024 * 1024 # cap on a user's total stored bytes + # Cap on bytes served per MCP ``resources/read`` so a giant artifact never streams into LLM context. + ARTIFACT_RESOURCE_READ_MAX_BYTES: int = 1 * 1024 * 1024 @field_validator("POSTGRES_URI", mode="before") @classmethod diff --git a/application/mcp_server.py b/application/mcp_server.py index 23b074f1..8b96fd20 100644 --- a/application/mcp_server.py +++ b/application/mcp_server.py @@ -1,7 +1,11 @@ -"""FastMCP server exposing DocsGPT retrieval over streamable HTTP. +"""FastMCP server exposing DocsGPT retrieval + artifacts over streamable HTTP. Mounted at ``/mcp`` by ``application/asgi.py``. Bearer tokens are the -existing DocsGPT agent API keys — no new credential surface. +existing DocsGPT agent API keys — no new credential surface. The +``search_docs`` tool searches the caller's knowledge base; the artifact +resources middleware exposes the caller's own artifacts as MCP resources +(``resources/list`` / ``resources/read`` over ``artifact://`` URIs), +scoped strictly to the Bearer key's owning principal. The tool reads the ``Authorization`` header directly via ``get_http_headers(include={"authorization"})``. The ``include`` kwarg @@ -14,11 +18,23 @@ token, we opt it back in. from __future__ import annotations import asyncio +import base64 import logging +from typing import Sequence from fastmcp import FastMCP +from fastmcp.resources import ResourceContent, ResourceResult +from fastmcp.resources.base import Resource from fastmcp.server.dependencies import get_http_headers +from fastmcp.server.middleware import CallNext, Middleware, MiddlewareContext +from application.services.artifact_resource_service import ( + ArtifactReadResult, + ResourceDenied, + ResourceNotFound, + list_artifact_resources, + read_artifact_resource, +) from application.services.search_service import ( InvalidAPIKey, SearchFailed, @@ -27,8 +43,6 @@ from application.services.search_service import ( logger = logging.getLogger(__name__) -mcp = FastMCP("docsgpt") - def _extract_bearer_token() -> str | None: auth = get_http_headers(include={"authorization"}).get("authorization", "") @@ -38,6 +52,63 @@ def _extract_bearer_token() -> str | None: return parts[1] +def _read_result_to_mcp(result: ArtifactReadResult) -> ResourceResult: + """Wrap a service read result as a FastMCP ``ResourceResult`` (text or blob).""" + if result.blob_b64 is not None: + raw = base64.b64decode(result.blob_b64) + return ResourceResult([ResourceContent(raw, mime_type=result.mime_type)]) + return ResourceResult([ResourceContent(result.text or "", mime_type=result.mime_type)]) + + +class ArtifactResourcesMiddleware(Middleware): + """Expose the calling principal's artifacts as MCP resources (read/list). + + The principal is the Bearer api_key's owner; resources are scoped to that + owner and a foreign/unauthenticated read is denied. Static resources + registered on the server (if any) are preserved by chaining ``call_next``. + """ + + async def on_list_resources( + self, + context: MiddlewareContext, + call_next: CallNext, + ) -> Sequence[Resource]: + """Append the principal's artifact resources to the static resource list.""" + existing = list(await call_next(context)) + try: + artifacts = await asyncio.to_thread( + list_artifact_resources, _extract_bearer_token() + ) + except Exception: + logger.exception("on_list_resources: artifact listing failed") + return existing + existing.extend(artifacts) + return existing + + async def on_read_resource( + self, + context: MiddlewareContext, + call_next: CallNext, + ) -> ResourceResult: + """Serve ``artifact://`` reads from the artifact store; defer others.""" + uri = str(getattr(context.message, "uri", "")) + if not uri.startswith("artifact://"): + return await call_next(context) + try: + result = await asyncio.to_thread( + read_artifact_resource, _extract_bearer_token(), uri + ) + except ResourceDenied as exc: + raise PermissionError(str(exc) or "forbidden") from exc + except ResourceNotFound as exc: + raise ValueError(str(exc) or "resource not found") from exc + return _read_result_to_mcp(result) + + +mcp = FastMCP("docsgpt") +mcp.add_middleware(ArtifactResourcesMiddleware()) + + @mcp.tool async def search_docs(query: str, chunks: int = 5) -> list[dict]: """Search the caller's DocsGPT knowledge base. diff --git a/application/services/artifact_resource_service.py b/application/services/artifact_resource_service.py new file mode 100644 index 00000000..92711aae --- /dev/null +++ b/application/services/artifact_resource_service.py @@ -0,0 +1,232 @@ +"""Flask-free service exposing a principal's artifacts as MCP Resources. + +The MCP server (``application/mcp_server.py``) authenticates a request with +``Authorization: Bearer ``; that key resolves to the owning +``user_id`` via ``AgentsRepository.find_by_key`` (the same api_key->owner path +as the HTTP artifact routes). Resources are scoped strictly to that principal: +``resources/list`` returns only the principal's owned artifacts, and +``resources/read`` re-checks ownership before serving any bytes. An +unresolvable principal yields an empty list / a denied read -- never another +principal's artifact. + +Returns plain ``mcp.types`` objects so both the FastMCP middleware adapter and +the unit tests can consume them without depending on FastMCP internals. +""" + +from __future__ import annotations + +import base64 +import logging +import re +from dataclasses import dataclass +from typing import List, Optional + +import mcp.types as mt +from sqlalchemy.exc import DataError, DBAPIError + +from application.core.settings import settings +from application.storage.db.base_repository import looks_like_uuid +from application.storage.db.repositories.agents import AgentsRepository +from application.storage.db.repositories.artifacts import ArtifactsRepository +from application.storage.db.session import db_readonly +from application.storage.storage_creator import StorageCreator + +logger = logging.getLogger(__name__) + +# artifact://{artifact_id}/v{version} +_URI_RE = re.compile(r"^artifact://(?P[^/]+)/v(?P\d+)$") + +# Max resources advertised by ``resources/list`` so the model is not flooded +# with thousands of rows even when a principal owns far more artifacts. +_RESOURCE_LIST_LIMIT = 500 + +# mime types served as inline ``text`` rather than base64 ``blob``. +_TEXT_MIME_PREFIXES = ("text/",) +_TEXT_MIME_EXACT = { + "application/json", + "application/xml", + "application/javascript", + "application/x-ndjson", + "image/svg+xml", +} +_TEXT_MIME_SUFFIXES = ("+json", "+xml") + + +class ResourceDenied(Exception): + """The principal may not access the requested resource (unauth or foreign).""" + + +class ResourceNotFound(Exception): + """The requested ``artifact://`` uri does not resolve to a stored version.""" + + +@dataclass(frozen=True) +class ArtifactReadResult: + """Materialized artifact contents for a ``resources/read`` response.""" + + uri: str + mime_type: str + text: Optional[str] = None + blob_b64: Optional[str] = None + + +def _read_cap() -> int: + """Return the per-read byte cap, or 0 (no cap) when the setting disables it.""" + return int(getattr(settings, "ARTIFACT_RESOURCE_READ_MAX_BYTES", 0) or 0) + + +def _cap_text(text: str) -> str: + """Truncate ``text`` to the read cap; unchanged when the cap is disabled.""" + cap = _read_cap() + return text[:cap] if cap > 0 else text + + +def _is_texty(mime_type: str) -> bool: + """Return True when ``mime_type`` should be served as inline UTF-8 text.""" + mime = (mime_type or "").split(";", 1)[0].strip().lower() + if mime in _TEXT_MIME_EXACT: + return True + if any(mime.startswith(p) for p in _TEXT_MIME_PREFIXES): + return True + return any(mime.endswith(s) for s in _TEXT_MIME_SUFFIXES) + + +def _resolve_principal(api_key: Optional[str]) -> Optional[str]: + """Resolve a Bearer api_key to its owning ``user_id``; None when unresolvable.""" + if not api_key: + return None + try: + with db_readonly() as conn: + agent = AgentsRepository(conn).find_by_key(api_key) + except Exception: + logger.exception("artifact resource: principal resolution failed") + return None + return agent.get("user_id") if agent else None + + +def _resource_uri(artifact_id: str, version: int) -> str: + """Build the stable ``artifact://{id}/v{version}`` resource uri.""" + return f"artifact://{artifact_id}/v{version}" + + +def list_artifact_resources(api_key: Optional[str]) -> List[mt.Resource]: + """List the calling principal's artifacts as MCP resources (empty if unresolved).""" + user_id = _resolve_principal(api_key) + if not user_id: + return [] + try: + with db_readonly() as conn: + rows = ArtifactsRepository(conn).list_artifacts(user_id=user_id) + except Exception: + logger.exception("artifact resource: list failed") + return [] + + resources: List[mt.Resource] = [] + for row in rows[:_RESOURCE_LIST_LIMIT]: + artifact_id = str(row.get("id")) + version = row.get("current_version") or 1 + title = row.get("title") or f"artifact-{artifact_id}" + resources.append( + mt.Resource( + uri=_resource_uri(artifact_id, version), + name=title, + title=title, + description=f"DocsGPT {row.get('kind') or 'file'} artifact", + mimeType=_kind_mime_hint(row.get("kind")), + ) + ) + return resources + + +def _kind_mime_hint(kind: Optional[str]) -> str: + """Concrete mime hint for a list row; ambiguous kinds fall back to octet-stream.""" + # Only kinds with a single unambiguous mime get a concrete type; everything + # else stays octet-stream so a list row never mismatches the read's bytes. + return { + "html": "text/html", + "data": "application/json", + }.get((kind or "").lower(), "application/octet-stream") + + +def read_artifact_resource(api_key: Optional[str], uri: str) -> ArtifactReadResult: + """Authorize and materialize an ``artifact://`` resource for the principal. + + Raises: + ResourceDenied: principal unresolved, or the artifact is not theirs. + ResourceNotFound: uri malformed, or the version/file does not exist. + """ + match = _URI_RE.match(uri or "") + if not match: + raise ResourceNotFound(f"unsupported resource uri: {uri!r}") + artifact_id = match.group("id") + version = int(match.group("version")) + + # A non-UUID id would reach a ``CAST(:id AS uuid)`` and raise a DB DataError; + # gate it up front like the HTTP artifact routes do. + if not looks_like_uuid(artifact_id): + raise ResourceNotFound(f"artifact {artifact_id} not found") + + user_id = _resolve_principal(api_key) + if not user_id: + raise ResourceDenied("unauthenticated") + + try: + with db_readonly() as conn: + repo = ArtifactsRepository(conn) + artifact = repo.get_artifact(artifact_id) + if artifact is None: + raise ResourceNotFound(f"artifact {artifact_id} not found") + # Ownership is the authz point: never serve another principal's artifact + # over MCP, regardless of conversation/workflow parent sharing. + if str(artifact.get("user_id")) != str(user_id): + raise ResourceDenied("forbidden") + version_row = repo.get_version(artifact_id, version) + except (DataError, DBAPIError) as exc: + raise ResourceNotFound(f"artifact {artifact_id} not found") from exc + + if version_row is None: + raise ResourceNotFound(f"version {version} of {artifact_id} not found") + + mime_type = version_row.get("mime_type") or "application/octet-stream" + uri = _resource_uri(artifact_id, version) + + # Prefer the stored preview/extracted text for texty kinds: it is already + # bounded and avoids a storage round-trip. + preview = version_row.get("preview_text") + if _is_texty(mime_type) and preview: + return ArtifactReadResult(uri=uri, mime_type=mime_type, text=_cap_text(preview)) + + storage_path = version_row.get("storage_path") + if not storage_path: + if _is_texty(mime_type) and preview is not None: + return ArtifactReadResult(uri=uri, mime_type=mime_type, text=_cap_text(preview)) + raise ResourceNotFound(f"version {version} of {artifact_id} has no stored bytes") + + try: + data = _read_capped_bytes(storage_path) + except FileNotFoundError as exc: + raise ResourceNotFound(f"version {version} of {artifact_id} has no stored bytes") from exc + + if _is_texty(mime_type): + # ``errors="ignore"`` keeps valid text texty even when the cap splits a + # multibyte char at the boundary, instead of demoting it to a blob. + return ArtifactReadResult(uri=uri, mime_type=mime_type, text=data.decode("utf-8", errors="ignore")) + return ArtifactReadResult( + uri=uri, mime_type=mime_type, blob_b64=base64.b64encode(data).decode("ascii") + ) + + +def _read_capped_bytes(storage_path: str) -> bytes: + """Read at most ``_read_cap()`` bytes of a stored artifact version (0 == all).""" + cap = _read_cap() + storage = StorageCreator.get_storage() + file_obj = storage.get_file(storage_path) + try: + return file_obj.read(cap) if cap > 0 else file_obj.read() + finally: + close = getattr(file_obj, "close", None) + if callable(close): + try: + close() + except Exception: + logger.debug("artifact resource: file close failed", exc_info=True) diff --git a/deployment/k8s/deployments/sandbox-deploy.yaml b/deployment/k8s/deployments/sandbox-deploy.yaml new file mode 100644 index 00000000..d7cc401a --- /dev/null +++ b/deployment/k8s/deployments/sandbox-deploy.yaml @@ -0,0 +1,71 @@ +# docsgpt-sandbox runner (Jupyter Kernel Gateway). Sessions are in-process +# kernels, never child containers; the Docker socket is NOT mounted. Network +# egress is constrained by the NetworkPolicy under +# deployment/k8s/network-policies/sandbox-egress-policy.yaml -- apply both. +# +# On Linux prod, schedule this onto a gVisor `runsc` RuntimeClass for kernel +# isolation (uncomment `runtimeClassName` once the node has it installed). +apiVersion: apps/v1 +kind: Deployment +metadata: + name: docsgpt-sandbox +spec: + replicas: 1 + selector: + matchLabels: + app: docsgpt-sandbox + template: + metadata: + labels: + app: docsgpt-sandbox + spec: + # runtimeClassName: gvisor + securityContext: + runAsNonRoot: true + seccompProfile: + type: RuntimeDefault + containers: + - name: docsgpt-sandbox + image: arc53/docsgpt-sandbox + ports: + - containerPort: 8888 + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: true + capabilities: + drop: + - ALL + env: + - name: JUPYTER_RUNTIME_DIR + value: /tmp/jupyter-runtime + - name: JUPYTER_DATA_DIR + value: /tmp/jupyter-data + resources: + limits: + memory: "1Gi" + cpu: "1" + requests: + memory: "256Mi" + cpu: "250m" + volumeMounts: + # Per-session workspaces + Jupyter runtime files on an emptyDir; + # the root FS stays read-only everywhere else. + - name: scratch + mountPath: /tmp + volumes: + - name: scratch + emptyDir: {} +--- +apiVersion: v1 +kind: Service +metadata: + name: docsgpt-sandbox +spec: + selector: + app: docsgpt-sandbox + ports: + - protocol: TCP + port: 8888 + targetPort: 8888 + # Cluster-internal only: the runner must never be exposed outside the cluster. + type: ClusterIP diff --git a/deployment/k8s/network-policies/sandbox-egress-policy.yaml b/deployment/k8s/network-policies/sandbox-egress-policy.yaml new file mode 100644 index 00000000..489cb6eb --- /dev/null +++ b/deployment/k8s/network-policies/sandbox-egress-policy.yaml @@ -0,0 +1,85 @@ +# Egress / SSRF NetworkPolicy for the docsgpt-sandbox runner. +# +# The sandbox executes arbitrary LLM-authored code, so app-level URL checks +# cannot contain it -- the code opens its own sockets. The hardened container +# runs without NET_ADMIN and cannot self-apply iptables, so egress is enforced +# here at the network layer: broad public-internet egress is allowed, but all +# private / link-local / metadata ranges are DENIED. Ingress is restricted to +# the backend and worker pods only. +# +# Requires a NetworkPolicy-enforcing CNI (Calico, Cilium, etc.). The default +# kube-proxy / flannel setups do NOT enforce these rules. Apply with: +# kubectl apply -f deployment/k8s/network-policies/sandbox-egress-policy.yaml +# +# Pods this selects must carry the label `app: docsgpt-sandbox`. +apiVersion: networking.k8s.io/v1 +kind: NetworkPolicy +metadata: + name: docsgpt-sandbox-egress + labels: + app: docsgpt-sandbox +spec: + podSelector: + matchLabels: + app: docsgpt-sandbox + policyTypes: + - Ingress + - Egress + ingress: + # Only the backend/worker may reach the sandbox gateway (TCP 8888). + - from: + - podSelector: + matchLabels: + app: docsgpt-api + - podSelector: + matchLabels: + app: docsgpt-worker + ports: + - protocol: TCP + port: 8888 + egress: + # DNS to cluster resolvers (kube-dns). Restricted to UDP/TCP 53 so the + # broad-egress rule below does not need to whitelist the resolver IP. + - to: + - namespaceSelector: {} + ports: + - protocol: UDP + port: 53 + - protocol: TCP + port: 53 + # Broad public-internet egress, with every private / link-local / ULA / + # carrier-grade-NAT / documentation range carved out via `except`. The + # cloud metadata IP 169.254.169.254 is inside the 169.254.0.0/16 hole, so + # it is unreachable. SSRF to internal services is therefore blocked while + # legitimate outbound calls (package installs, public APIs) still work. + - to: + - ipBlock: + cidr: 0.0.0.0/0 + except: + - 0.0.0.0/8 # "this network" (RFC1122) + - 10.0.0.0/8 # RFC1918 private + - 172.16.0.0/12 # RFC1918 private + - 192.168.0.0/16 # RFC1918 private + - 169.254.0.0/16 # link-local incl. 169.254.169.254 metadata + - 100.64.0.0/10 # RFC6598 carrier-grade NAT + - 127.0.0.0/8 # loopback + - 192.0.0.0/24 # IETF protocol assignments + - 192.0.2.0/24 # TEST-NET-1 (documentation) + - 198.18.0.0/15 # benchmarking + - 198.51.100.0/24 # TEST-NET-2 (documentation) + - 203.0.113.0/24 # TEST-NET-3 (documentation) + - 224.0.0.0/4 # multicast + - 240.0.0.0/4 # reserved + - 255.255.255.255/32 # limited broadcast (RFC1122) + # IPv6 public egress with private / ULA / link-local ranges carved out. + # The IPv6 metadata endpoint (fd00:ec2::254) sits inside the fc00::/7 ULA + # hole and is unreachable. Drop this block on IPv4-only clusters. + - to: + - ipBlock: + cidr: "::/0" + except: + - "::1/128" # loopback + - "fc00::/7" # unique local (ULA) incl. fd00:ec2::254 + - "fe80::/10" # link-local + - "ff00::/8" # multicast + - "2001:db8::/32" # documentation diff --git a/deployment/optional/docker-compose.optional.sandbox-egress.yaml b/deployment/optional/docker-compose.optional.sandbox-egress.yaml new file mode 100644 index 00000000..bc0916a9 --- /dev/null +++ b/deployment/optional/docker-compose.optional.sandbox-egress.yaml @@ -0,0 +1,71 @@ +# Optional egress-firewall overlay for the docsgpt-sandbox runner (compose). +# +# Docker Compose cannot express L3 egress filtering the way a Kubernetes +# NetworkPolicy can, so SSRF containment in compose deployments is delivered by +# routing the sandbox's outbound traffic through an egress-gateway sidecar that +# DENIES private / link-local / metadata ranges and ALLOWS the public internet. +# +# Two viable approaches; pick one: +# +# (1) Host / cloud firewall (simplest, recommended for single-host compose): +# Drop egress to RFC1918 (10/8, 172.16/12, 192.168/16), link-local +# (169.254/16, incl. the 169.254.169.254 metadata IP), and ULA on the +# docsgpt-sandbox container's interface using iptables/nftables on the +# Docker host (Docker does not do this for you). Example (host root): +# SBX=$(docker inspect -f '{{.NetworkSettings.Networks.docsgpt_sandbox-net.IPAddress}}' docsgpt-sandbox) +# iptables -I DOCKER-USER -s "$SBX" -d 169.254.0.0/16 -j DROP +# iptables -I DOCKER-USER -s "$SBX" -d 10.0.0.0/8 -j DROP +# iptables -I DOCKER-USER -s "$SBX" -d 172.16.0.0/12 -j DROP +# iptables -I DOCKER-USER -s "$SBX" -d 192.168.0.0/16 -j DROP +# (Allow established/return traffic and DNS as needed for your setup.) +# +# (2) Egress-gateway sidecar (this overlay): pin the sandbox to an +# internal-only network with NO direct internet route, and force all of +# its outbound traffic through a small forward proxy that blocks private +# destinations. The proxy is the ONLY container on both the internal and +# the egress network, so the sandbox cannot reach internal services +# directly. Configure the runner's outbound HTTP(S) via the proxy env +# vars below; lock down non-HTTP egress with the proxy's own rules. +# +# Apply alongside the base stack: +# docker compose -f deployment/docker-compose.yaml \ +# -f deployment/optional/docker-compose.optional.sandbox-egress.yaml up -d +# +# NOTE: a forward proxy only constrains traffic that honors the proxy env vars. +# Truly arbitrary sandbox code can ignore them, so on a multi-tenant or +# untrusted deployment prefer approach (1) (or the Kubernetes NetworkPolicy) +# which enforces at the network layer regardless of what the code does. +version: "3.8" + +services: + # Forward proxy that denies private/link-local/metadata destinations and + # allows the public internet. Replace the image/command with your proxy of + # choice (tinyproxy, squid, mitmproxy with a deny rule, etc.); the ACL must + # DENY 10/8, 172.16/12, 192.168/16, 169.254/16, 127/8 and ALLOW the rest. + sandbox-egress-proxy: + image: ghcr.io/example/egress-deny-private:latest # replace with your proxy image + restart: unless-stopped + networks: + - sandbox-net # reachable by the sandbox + - sandbox-egress # the only container with an internet route + + docsgpt-sandbox: + # Cut the sandbox off from the default (internet) bridge; it keeps only the + # internal sandbox-net and must egress via the proxy. + networks: + - sandbox-net + environment: + - HTTP_PROXY=http://sandbox-egress-proxy:8080 + - HTTPS_PROXY=http://sandbox-egress-proxy:8080 + - NO_PROXY=localhost,127.0.0.1 + +networks: + # Internal-only network shared by the sandbox and the proxy. ``internal: + # true`` removes the default gateway so the sandbox has NO direct internet + # route -- its only path out is through the proxy on the egress network. + sandbox-net: + driver: bridge + internal: true + # Internet-facing network for the proxy only. + sandbox-egress: + driver: bridge diff --git a/deployment/sandbox/README.md b/deployment/sandbox/README.md index 0b60b479..31286c95 100644 --- a/deployment/sandbox/README.md +++ b/deployment/sandbox/README.md @@ -62,8 +62,53 @@ Docling uses its own permissive PDF backend; **PyMuPDF (AGPL) is intentionally n installed.** If the base (non-extract) image is used, `extract_document` returns a clean "docling is not available in the sandbox runner" error rather than crashing. -## Hardening (separate slice — not in this image) +## Network egress / SSRF -Egress/SSRF blocks (drop RFC1918, link-local, `169.254.169.254`), the gVisor -`runsc` runtime, seccomp, read-only root FS, and cgroup CPU/mem/PID caps wired -from `SANDBOX_MEMORY` / `SANDBOX_CPUS` land in the hardening slice. +The runner allows **broad outbound egress** (so sandboxed code can `pip install` +and call public APIs) but private, link-local, and cloud-metadata ranges **MUST +be blocked at the network layer**. This is not optional: the sandbox executes +arbitrary LLM-authored code, which opens its own sockets — app-level URL checks +(the `mcp_tool.py` approach) cannot contain it. Without a network-layer block, +sandbox code can reach `169.254.169.254` (cloud instance metadata / credentials) +and internal services on the private network. + +The hardened container runs **without `NET_ADMIN`**, so it cannot self-apply +`iptables`. Enforcement therefore lives in deployment config: + +- **Kubernetes** — apply + [`deployment/k8s/network-policies/sandbox-egress-policy.yaml`](../k8s/network-policies/sandbox-egress-policy.yaml). + It allows `0.0.0.0/0` egress with `except` carve-outs for RFC1918 + (`10/8`, `172.16/12`, `192.168/16`), link-local (`169.254/16`, which contains + `169.254.169.254`), loopback, CGNAT, documentation/test ranges, and the IPv6 + ULA/link-local equivalents — and restricts ingress to the API/worker pods on + TCP 8888. It requires a policy-enforcing CNI (Calico, Cilium, …); plain + flannel/kube-proxy will silently not enforce it. The matching sandbox pod is + [`deployment/k8s/deployments/sandbox-deploy.yaml`](../k8s/deployments/sandbox-deploy.yaml) + (label `app: docsgpt-sandbox`). + + ```bash + kubectl apply -f deployment/k8s/deployments/sandbox-deploy.yaml + kubectl apply -f deployment/k8s/network-policies/sandbox-egress-policy.yaml + ``` + +- **docker-compose** — compose cannot express L3 egress filtering natively. The + base stack puts the runner on an `internal: true` network (no host port), but + that does not by itself block the metadata IP or RFC1918 reachable via the + default bridge. Add a **host/cloud firewall rule** (drop the four private + ranges on the sandbox container's interface) **or** route egress through an + **egress-gateway proxy** sidecar. Both are documented in + [`deployment/optional/docker-compose.optional.sandbox-egress.yaml`](../optional/docker-compose.optional.sandbox-egress.yaml). + On untrusted/multi-tenant hosts prefer the host-firewall rule — a forward + proxy only constrains code that honors `HTTP(S)_PROXY`. + +## Other hardening (deployment-level) + +The gVisor `runsc` runtime (kernel isolation for untrusted code), seccomp +profile, read-only root FS, non-root, and cgroup CPU/mem/PID caps (wired from +`SANDBOX_MEMORY` / `SANDBOX_CPUS`) are deployment-level concerns. The compose +service in `deployment/docker-compose.yaml` already sets `read_only`, +`mem_limit`, `cpus`, and `pids_limit`; the k8s `sandbox-deploy.yaml` sets the +equivalent `securityContext` + resource limits and has a commented +`runtimeClassName: gvisor` to enable on nodes with the `runsc` RuntimeClass +installed. These complement — they do not replace — the network egress policy +above. diff --git a/tests/services/test_artifact_resource_service.py b/tests/services/test_artifact_resource_service.py new file mode 100644 index 00000000..c9093f4a --- /dev/null +++ b/tests/services/test_artifact_resource_service.py @@ -0,0 +1,176 @@ +"""Tests for application/services/artifact_resource_service.py. + +The service exposes a Bearer-key principal's own artifacts as MCP resources. +These tests patch the DB/storage seams (``db_readonly``, the repositories, and +``StorageCreator``) so they run without Postgres, mirroring the light mocking +in ``tests/services/test_mcp_server.py``. They assert that: + +- ``resources/list`` returns only the principal's artifacts; +- ``resources/read`` returns ``text`` vs ``blob`` by mime with the right type; +- a foreign-owner artifact is denied (no cross-principal exposure); +- an unauthenticated / unresolved principal gets an empty list / denied read; +- a non-UUID id is rejected as not-found (no leaked DB error); +- the read is byte-capped. +""" + +from __future__ import annotations + +import base64 +import io +from contextlib import contextmanager + +import pytest + +from application.services import artifact_resource_service as svc + +OWNER = "owner-1" +STRANGER = "stranger-2" + +# Real UUIDs: the read path gates non-UUID ids before they reach the DB. +ART_TEXT = "11111111-1111-4111-8111-111111111111" +ART_BIN = "22222222-2222-4222-8222-222222222222" +ART_FOREIGN = "33333333-3333-4333-8333-333333333333" + + +@contextmanager +def _fake_conn(): + yield object() + + +class _FakeAgents: + """Stub AgentsRepository: maps api_key -> agent row (or None).""" + + _MAP = {"owner-key": {"user_id": OWNER}, "stranger-key": {"user_id": STRANGER}} + + def __init__(self, conn): + pass + + def find_by_key(self, key): + return self._MAP.get(key) + + +class _FakeArtifacts: + """Stub ArtifactsRepository backed by in-memory artifact/version dicts.""" + + artifacts: dict = {} + versions: dict = {} + + def __init__(self, conn): + pass + + def list_artifacts(self, user_id=None, conversation_id=None, workflow_run_id=None): + return [a for a in self.artifacts.values() if a["user_id"] == user_id] + + def get_artifact(self, artifact_id): + return self.artifacts.get(artifact_id) + + def get_version(self, artifact_id, version): + return self.versions.get((artifact_id, version)) + + +class _FakeStorage: + """Stub BaseStorage.get_file returning a capped BytesIO of fixed bytes.""" + + blob = b"x" * 10 + + def get_file(self, path): + return io.BytesIO(self.blob) + + +@pytest.fixture(autouse=True) +def _wire(monkeypatch): + """Point the service's DB/storage seams at the in-memory fakes.""" + _FakeArtifacts.artifacts = { + ART_TEXT: {"id": ART_TEXT, "user_id": OWNER, "kind": "data", "title": "notes", "current_version": 2}, + ART_BIN: {"id": ART_BIN, "user_id": OWNER, "kind": "image", "title": "chart", "current_version": 1}, + ART_FOREIGN: { + "id": ART_FOREIGN, + "user_id": STRANGER, + "kind": "data", + "title": "secret", + "current_version": 1, + }, + } + _FakeArtifacts.versions = { + (ART_TEXT, 2): {"mime_type": "text/csv", "storage_path": "k/text.csv", "preview_text": None}, + (ART_BIN, 1): {"mime_type": "image/png", "storage_path": "k/chart.png", "preview_text": None}, + (ART_FOREIGN, 1): {"mime_type": "text/plain", "storage_path": "k/secret.txt", "preview_text": None}, + } + monkeypatch.setattr(svc, "db_readonly", _fake_conn) + monkeypatch.setattr(svc, "AgentsRepository", _FakeAgents) + monkeypatch.setattr(svc, "ArtifactsRepository", _FakeArtifacts) + monkeypatch.setattr(svc.StorageCreator, "get_storage", staticmethod(lambda: _FakeStorage())) + + +@pytest.mark.unit +class TestListArtifactResources: + def test_lists_only_principal_artifacts(self): + out = svc.list_artifact_resources("owner-key") + uris = {str(r.uri) for r in out} + assert uris == {f"artifact://{ART_TEXT}/v2", f"artifact://{ART_BIN}/v1"} + assert f"artifact://{ART_FOREIGN}/v1" not in uris + + def test_unresolved_principal_is_empty(self): + assert svc.list_artifact_resources("bogus-key") == [] + + def test_missing_token_is_empty(self): + assert svc.list_artifact_resources(None) == [] + + def test_resource_carries_name_and_mime(self): + out = {str(r.uri): r for r in svc.list_artifact_resources("owner-key")} + assert out[f"artifact://{ART_TEXT}/v2"].name == "notes" + # The list row never advertises a wildcard/wrong type; the image kind + # falls back to the generic octet-stream hint. + assert out[f"artifact://{ART_BIN}/v1"].mimeType == "application/octet-stream" + + +@pytest.mark.unit +class TestReadArtifactResource: + def test_text_mime_returns_text(self): + res = svc.read_artifact_resource("owner-key", f"artifact://{ART_TEXT}/v2") + assert res.text == _FakeStorage.blob.decode("utf-8") + assert res.blob_b64 is None + assert res.mime_type == "text/csv" + + def test_binary_mime_returns_blob(self): + res = svc.read_artifact_resource("owner-key", f"artifact://{ART_BIN}/v1") + assert res.blob_b64 == base64.b64encode(_FakeStorage.blob).decode("ascii") + assert res.text is None + assert res.mime_type == "image/png" + + def test_prefers_preview_text_when_present(self, monkeypatch): + _FakeArtifacts.versions[(ART_TEXT, 2)]["preview_text"] = "cached preview" + res = svc.read_artifact_resource("owner-key", f"artifact://{ART_TEXT}/v2") + assert res.text == "cached preview" + + def test_foreign_owner_is_denied(self): + with pytest.raises(svc.ResourceDenied): + svc.read_artifact_resource("owner-key", f"artifact://{ART_FOREIGN}/v1") + + def test_unauthenticated_is_denied(self): + with pytest.raises(svc.ResourceDenied): + svc.read_artifact_resource(None, f"artifact://{ART_TEXT}/v2") + with pytest.raises(svc.ResourceDenied): + svc.read_artifact_resource("bogus-key", f"artifact://{ART_TEXT}/v2") + + def test_unknown_uri_scheme_not_found(self): + with pytest.raises(svc.ResourceNotFound): + svc.read_artifact_resource("owner-key", "https://example.com/x") + + def test_non_uuid_id_is_not_found(self): + # A non-UUID id must be rejected before the DB cast (no leaked DataError). + with pytest.raises(svc.ResourceNotFound): + svc.read_artifact_resource("owner-key", "artifact://not-a-uuid/v1") + + def test_missing_version_not_found(self): + with pytest.raises(svc.ResourceNotFound): + svc.read_artifact_resource("owner-key", f"artifact://{ART_TEXT}/v99") + + def test_read_is_byte_capped(self, monkeypatch): + monkeypatch.setattr(svc.settings, "ARTIFACT_RESOURCE_READ_MAX_BYTES", 3) + _FakeStorage.blob = b"abcdefghij" + try: + res = svc.read_artifact_resource("owner-key", f"artifact://{ART_BIN}/v1") + assert base64.b64decode(res.blob_b64) == b"abc" + finally: + _FakeStorage.blob = b"x" * 10 diff --git a/tests/services/test_mcp_server.py b/tests/services/test_mcp_server.py index c1da1cbd..9a12bc5f 100644 --- a/tests/services/test_mcp_server.py +++ b/tests/services/test_mcp_server.py @@ -9,6 +9,7 @@ full HTTP-layer plumbing (mount, lifespan, session handshake) is covered by ``tests/test_asgi.py``. """ +from types import SimpleNamespace from unittest.mock import patch import pytest @@ -132,3 +133,103 @@ class TestSearchDocsTool: ): await search_docs(query="q") mock_search.assert_called_once_with("lowercase-scheme", "q", 5) + + +async def _call_next_empty(context): + """Stand-in downstream handler returning no static resources.""" + return [] + + +@pytest.mark.unit +class TestArtifactResourcesMiddleware: + @pytest.mark.asyncio + async def test_list_appends_principal_artifacts(self): + import mcp.types as mt + from application.mcp_server import ArtifactResourcesMiddleware + + mw = ArtifactResourcesMiddleware() + res = mt.Resource(uri="artifact://a/v1", name="x", mimeType="text/plain") + with ( + patch("application.mcp_server._extract_bearer_token", return_value="k"), + patch( + "application.mcp_server.list_artifact_resources", return_value=[res] + ), + ): + out = await mw.on_list_resources(SimpleNamespace(), _call_next_empty) + assert [str(r.uri) for r in out] == ["artifact://a/v1"] + + @pytest.mark.asyncio + async def test_read_text_wraps_as_text_contents(self): + from application.mcp_server import ArtifactResourcesMiddleware + from application.services.artifact_resource_service import ArtifactReadResult + + mw = ArtifactResourcesMiddleware() + ctx = SimpleNamespace(message=SimpleNamespace(uri="artifact://a/v2")) + rr = ArtifactReadResult(uri="artifact://a/v2", mime_type="text/csv", text="a,b") + with ( + patch("application.mcp_server._extract_bearer_token", return_value="k"), + patch( + "application.mcp_server.read_artifact_resource", return_value=rr + ), + ): + result = await mw.on_read_resource(ctx, _call_next_empty) + out = result.to_mcp_result("artifact://a/v2").contents[0] + assert out.text == "a,b" + assert out.mimeType == "text/csv" + + @pytest.mark.asyncio + async def test_read_blob_wraps_as_blob_contents(self): + import base64 + + from application.mcp_server import ArtifactResourcesMiddleware + from application.services.artifact_resource_service import ArtifactReadResult + + mw = ArtifactResourcesMiddleware() + ctx = SimpleNamespace(message=SimpleNamespace(uri="artifact://a/v1")) + blob = base64.b64encode(b"\x89PNG").decode("ascii") + rr = ArtifactReadResult(uri="artifact://a/v1", mime_type="image/png", blob_b64=blob) + with ( + patch("application.mcp_server._extract_bearer_token", return_value="k"), + patch( + "application.mcp_server.read_artifact_resource", return_value=rr + ), + ): + result = await mw.on_read_resource(ctx, _call_next_empty) + out = result.to_mcp_result("artifact://a/v1").contents[0] + assert out.mimeType == "image/png" + assert base64.b64decode(out.blob) == b"\x89PNG" + + @pytest.mark.asyncio + async def test_non_artifact_uri_is_deferred(self): + from application.mcp_server import ArtifactResourcesMiddleware + + mw = ArtifactResourcesMiddleware() + ctx = SimpleNamespace(message=SimpleNamespace(uri="https://example.com/x")) + sentinel = object() + + async def _call_next(context): + return sentinel + + with patch( + "application.mcp_server.read_artifact_resource" + ) as mock_read: + out = await mw.on_read_resource(ctx, _call_next) + assert out is sentinel + mock_read.assert_not_called() + + @pytest.mark.asyncio + async def test_denied_read_raises_permission_error(self): + from application.mcp_server import ArtifactResourcesMiddleware + from application.services.artifact_resource_service import ResourceDenied + + mw = ArtifactResourcesMiddleware() + ctx = SimpleNamespace(message=SimpleNamespace(uri="artifact://foreign/v1")) + with ( + patch("application.mcp_server._extract_bearer_token", return_value="k"), + patch( + "application.mcp_server.read_artifact_resource", + side_effect=ResourceDenied("forbidden"), + ), + ): + with pytest.raises(PermissionError): + await mw.on_read_resource(ctx, _call_next_empty)