Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .pre-commit-config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ repos:
args: []
additional_dependencies:
- pytest
- mcp==2.0.0
- mcp==2.2.0
- httpx2
- af-credentials==0.3.1

Expand Down
25 changes: 25 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,3 +17,28 @@ servicex-mcp serve --backend <name-from-your-.servicex-file>
```

See `docs/plans/` for the design and implementation plan.

## Available tools

| Tool | Does | Read/write |
| ------------------------------- | --------------------------------------------------------------------- | ----------- |
| `servicex_list_datasets` | List datasets cached on this ServiceX instance | read-only |
| `servicex_get_dataset` | Get the full detail of one cached dataset by its dataset ID | read-only |
| `servicex_info` | Return the ServiceX server version and its advertised capabilities | read-only |
| `servicex_list_code_generators` | List the code generators deployed on this ServiceX instance | read-only |
| `servicex_list_transforms` | List transforms you have submitted, with status and completion counts | read-only |
| `servicex_get_transform_status` | Get the full status of one transform by its request ID | read-only |
| `servicex_submit_query` | Submit a new transform (query) request against a dataset | write |
| `servicex_cancel_transform` | Cancel a running transform by its request ID | destructive |
| `servicex_delete_transform` | Delete a transform record (and its cache entry) by request ID | destructive |
| `servicex_delete_dataset` | Delete a cached dataset record by its dataset ID | destructive |

The six read-only tools query the external ServiceX server
(`read_only_hint=true`, `open_world_hint=true` in their MCP tool annotations).
`servicex_submit_query` mutates ServiceX state but never deletes anything
(`read_only_hint=false`). The three destructive tools
(`servicex_cancel_transform`, `servicex_delete_transform`,
`servicex_delete_dataset`) act on real running/finished transforms and cached
dataset records and cannot be undone (`read_only_hint=false`,
`destructive_hint=true`). All four write/destructive tools are disabled when the
server is started with `--read-only`.
3,056 changes: 1,613 additions & 1,443 deletions pixi.lock

Large diffs are not rendered by default.

14 changes: 13 additions & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,18 @@ packages = ["src/servicex_mcp"]
minversion = "9.0"
addopts = ["-ra", "--showlocals", "--strict-markers", "--strict-config"]
strict = true
filterwarnings = ["error"]
filterwarnings = [
"error",
# anyio >=4.15 deprecates the anyio.abc.BlockingPortal re-export in favor of
# anyio.from_thread.BlockingPortal. starlette 1.6.0's own testclient.py
# still imports the old alias internally (`import anyio.abc` /
# `anyio.abc.BlockingPortal`), so every wire-level test that imports
# starlette.testclient trips this the moment anyio >4.14 is resolved --
# which the mcp 2.2.0 lock bump forces in this repo's dependency graph.
# Nothing in servicex-mcp touches the deprecated alias directly; this is
# starlette's fix to make, not ours.
"ignore:The anyio\\.abc\\.BlockingPortal alias is deprecated.*:DeprecationWarning",
]
log_level = "INFO"
testpaths = ["tests"]
markers = [
Expand All @@ -98,6 +109,7 @@ report.exclude_also = [

[tool.mypy]
files = ["src", "tests"]
mypy_path = ["src"]
python_version = "3.11"
warn_unused_configs = true
strict = true
Expand Down
28 changes: 20 additions & 8 deletions src/servicex_mcp/tools/_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
import itertools
from typing import TYPE_CHECKING, Any

from mcp.types import CallToolResult, TextContent

if TYPE_CHECKING:
from servicex.servicex_client import ServiceXClient

Expand Down Expand Up @@ -83,11 +85,15 @@ def build_hints(hints: list[str]) -> str:
return f"\n\n**Next steps:**\n{lines}"


def classify_error(exc: Exception) -> str:
"""Return an actionable error message with recovery guidance.
def classify_error(exc: Exception) -> CallToolResult:
"""Return an actionable ``is_error`` result with recovery guidance.

Pattern-matches on exception type name and message text to provide
specific recovery steps rather than a bare traceback string.
specific recovery steps rather than a bare traceback string. No
``structured_content`` is set: an error result carries no structured
payload (mcp SDK's ``convert_result`` only validates
``structured_content`` against the tool's output model when
``is_error`` is false).
"""
exc_type = type(exc).__name__
# A bare exception (e.g. a ProxyClient redeem call raising a timeout with
Expand Down Expand Up @@ -148,9 +154,13 @@ def classify_error(exc: Exception) -> str:
"Use `servicex_info` to check server connectivity and try again."
)
else:
return f"Error: {exc_msg}"
text = f"Error: {exc_msg}"
return CallToolResult(
content=[TextContent(type="text", text=text)], is_error=True
)

return f"Error: {exc_msg}\n\n**Recovery:** {guidance}"
text = f"Error: {exc_msg}\n\n**Recovery:** {guidance}"
return CallToolResult(content=[TextContent(type="text", text=text)], is_error=True)


_READ_ONLY_ERROR = (
Expand All @@ -159,10 +169,12 @@ def classify_error(exc: Exception) -> str:
)


def check_write_allowed(lifespan_context: dict[str, Any]) -> str | None:
"""Return an error string if write operations are disabled, else None."""
def check_write_allowed(lifespan_context: dict[str, Any]) -> CallToolResult | None:
"""Return an ``is_error`` result if write operations are disabled, else None."""
if lifespan_context.get("read_only"):
return _READ_ONLY_ERROR
return CallToolResult(
content=[TextContent(type="text", text=_READ_ONLY_ERROR)], is_error=True
)
return None


Expand Down
101 changes: 90 additions & 11 deletions src/servicex_mcp/tools/datasets.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,11 @@
from __future__ import annotations

import asyncio
from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING, Annotated, Any

from mcp.server.mcpserver import Context, MCPServer # noqa: TC002
from mcp.types import CallToolResult, TextContent, ToolAnnotations
from pydantic import BaseModel

from servicex_mcp.tools._helpers import (
build_hints,
Expand Down Expand Up @@ -54,22 +56,59 @@ def _dataset_to_dict(d: CachedDataset) -> dict[str, Any]:
}


class DatasetInfo(BaseModel):
"""Structured detail of one cached dataset, matching ``_dataset_to_dict``."""

id: int
name: str
did_finder: str | None
n_files: int | None
size: int | None
events: int | None
lookup_status: str | None
is_stale: bool | None
last_used: str | None
last_updated: str | None


class ServicexListDatasetsResult(BaseModel):
"""Structured result of ``servicex_list_datasets``."""

datasets: list[DatasetInfo]
offset: int
limit: int
truncated: bool


class ServicexDeleteDatasetResult(BaseModel):
"""Structured result of ``servicex_delete_dataset``."""

dataset_id: int
stale: bool


def register(mcp: MCPServer) -> None:
"""Register the dataset-related MCP tools.

Registers servicex_list_datasets, servicex_get_dataset, and
servicex_delete_dataset.
"""

@mcp.tool()
@mcp.tool(
annotations=ToolAnnotations(
title="List datasets",
read_only_hint=True,
open_world_hint=True,
)
)
async def servicex_list_datasets(
did_finder: str | None = None,
show_deleted: bool = False,
limit: int = 50,
offset: int = 0,
*,
ctx: Context[Any, Any],
) -> str:
) -> Annotated[CallToolResult, ServicexListDatasetsResult]:
"""List datasets cached on this ServiceX instance.

Shows file/event counts and cache status for each dataset. Use
Expand All @@ -87,21 +126,45 @@ async def servicex_list_datasets(
except Exception as exc: # noqa: BLE001
return classify_error(exc)
if not datasets:
return "No datasets found."
payload = ServicexListDatasetsResult(
datasets=[], offset=offset, limit=limit, truncated=False
)
return CallToolResult(
content=[TextContent(type="text", text="No datasets found.")],
structured_content=payload.model_dump(mode="json"),
)
rows, footer = paginate_iter(
(_dataset_to_dict(d) for d in datasets), limit, offset
)
hints = build_hints(
["Use `servicex_get_dataset` with a dataset_id for full detail"]
)
return (
text = (
format_list(rows, include_keys=_DATASET_KEYS, byte_keys=_BYTE_KEYS)
+ footer
+ hints
)
payload = ServicexListDatasetsResult(
datasets=[DatasetInfo(**row) for row in rows],
offset=offset,
limit=limit,
truncated=bool(footer),
)
return CallToolResult(
content=[TextContent(type="text", text=text)],
structured_content=payload.model_dump(mode="json"),
)

@mcp.tool()
async def servicex_get_dataset(dataset_id: int, *, ctx: Context[Any, Any]) -> str:
@mcp.tool(
annotations=ToolAnnotations(
title="Get dataset",
read_only_hint=True,
open_world_hint=True,
)
)
async def servicex_get_dataset(
dataset_id: int, *, ctx: Context[Any, Any]
) -> Annotated[CallToolResult, DatasetInfo]:
"""Get the full detail of one cached dataset by its dataset ID."""
try:
client = get_servicex_client(ctx)
Expand All @@ -113,12 +176,23 @@ async def servicex_get_dataset(dataset_id: int, *, ctx: Context[Any, Any]) -> st
hints = build_hints(
["Use `servicex_delete_dataset` to remove this dataset from the cache"]
)
return format_dict(_dataset_to_dict(d), byte_keys=_BYTE_KEYS) + hints
text = format_dict(_dataset_to_dict(d), byte_keys=_BYTE_KEYS) + hints
payload = DatasetInfo(**_dataset_to_dict(d))
return CallToolResult(
content=[TextContent(type="text", text=text)],
structured_content=payload.model_dump(mode="json"),
)

@mcp.tool()
@mcp.tool(
annotations=ToolAnnotations(
title="Delete dataset",
read_only_hint=False,
destructive_hint=True,
)
)
async def servicex_delete_dataset(
dataset_id: int, *, ctx: Context[Any, Any]
) -> str:
) -> Annotated[CallToolResult, ServicexDeleteDatasetResult]:
"""Delete a cached dataset record by its dataset ID."""
write_error = check_write_allowed(ctx.request_context.lifespan_context)
if write_error:
Expand All @@ -134,4 +208,9 @@ async def servicex_delete_dataset(
# record (servicex.servicex_adapter.ServiceXAdapter.delete_dataset),
# not an unconditional success/failure signal — surface it rather
# than assuming the delete always succeeded.
return f"Dataset {dataset_id} deleted (stale={stale})."
text = f"Dataset {dataset_id} deleted (stale={stale})."
payload = ServicexDeleteDatasetResult(dataset_id=dataset_id, stale=stale)
return CallToolResult(
content=[TextContent(type="text", text=text)],
structured_content=payload.model_dump(mode="json"),
)
64 changes: 54 additions & 10 deletions src/servicex_mcp/tools/info.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,11 @@

from __future__ import annotations

from typing import Any
from typing import Annotated, Any

from mcp.server.mcpserver import Context, MCPServer # noqa: TC002
from mcp.types import CallToolResult, TextContent, ToolAnnotations
from pydantic import BaseModel

from servicex_mcp.tools._helpers import (
build_hints,
Expand All @@ -14,11 +16,32 @@
)


class ServicexInfoResult(BaseModel):
"""Structured result of ``servicex_info``."""

app_version: str
capabilities: list[str]


class ServicexListCodeGeneratorsResult(BaseModel):
"""Structured result of ``servicex_list_code_generators``."""

generators: dict[str, str]


def register(mcp: MCPServer) -> None:
"""Register servicex_info and servicex_list_code_generators with the MCP server."""

@mcp.tool()
async def servicex_info(*, ctx: Context[Any, Any]) -> str:
@mcp.tool(
annotations=ToolAnnotations(
title="Get ServiceX info",
read_only_hint=True,
open_world_hint=True,
)
)
async def servicex_info(
*, ctx: Context[Any, Any]
) -> Annotated[CallToolResult, ServicexInfoResult]:
"""Return the ServiceX server version and its advertised capabilities.

Use this tool to verify the ServiceX backend is reachable and to see
Expand All @@ -37,10 +60,25 @@ async def servicex_info(*, ctx: Context[Any, Any]) -> str:
hints = build_hints(
["Use `servicex_list_code_generators` to see available codegens"]
)
return "\n".join(lines) + hints
text = "\n".join(lines) + hints
payload = ServicexInfoResult(
app_version=info.app_version, capabilities=list(info.capabilities)
)
return CallToolResult(
content=[TextContent(type="text", text=text)],
structured_content=payload.model_dump(mode="json"),
)

@mcp.tool()
async def servicex_list_code_generators(*, ctx: Context[Any, Any]) -> str:
@mcp.tool(
annotations=ToolAnnotations(
title="List code generators",
read_only_hint=True,
open_world_hint=True,
)
)
async def servicex_list_code_generators(
*, ctx: Context[Any, Any]
) -> Annotated[CallToolResult, ServicexListCodeGeneratorsResult]:
"""List the code generators deployed on this ServiceX instance.

Each code generator (e.g. `func_adl_uproot`, `python`, `uproot-raw`)
Expand All @@ -52,8 +90,14 @@ async def servicex_list_code_generators(*, ctx: Context[Any, Any]) -> str:
except Exception as exc: # noqa: BLE001
return classify_error(exc)
if not generators:
return "No code generators are registered on this ServiceX instance."
hints = build_hints(
["Use `servicex_submit_query` with one of these codegen names"]
text = "No code generators are registered on this ServiceX instance."
else:
hints = build_hints(
["Use `servicex_submit_query` with one of these codegen names"]
)
text = format_dict(generators) + hints
payload = ServicexListCodeGeneratorsResult(generators=dict(generators))
return CallToolResult(
content=[TextContent(type="text", text=text)],
structured_content=payload.model_dump(mode="json"),
)
return format_dict(generators) + hints
Loading
Loading