mirror of
https://github.com/tiennm99/DocsGPT.git
synced 2026-10-03 07:11:56 +00:00
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.
This commit is contained in:
1 parent
d26a973189
commit
30397a1905
9 files changed
+862
-8
No files matched your search
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 <agent-api-key>``; 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<id>[^/]+)/v(?P<version>\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)
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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.
|
||||
@@ -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
|
||||
@@ -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)
|
||||
Reference in new issue
Block a user