diff --git a/docs/handlers/cancellation.md b/docs/handlers/cancellation.md new file mode 100644 index 0000000000..2e20ccbf6d --- /dev/null +++ b/docs/handlers/cancellation.md @@ -0,0 +1,55 @@ +# Cancellation + +A client can give up on a call: the user pressed stop, or a timeout ran out. + +When it does, the SDK **cancels your handler**. The `await` it is waiting on raises, the function unwinds, and nothing it returns is sent. Most handlers need to do nothing about that. + +Two kinds do: a handler with something to clean up, and a handler that is a plain `def`. + +## Clean up in an `async def` tool + +Put the cleanup in a `finally`: + +```python title="server.py" hl_lines="23 26-28" +--8<-- "docs_src/cancellation/tutorial001.py" +``` + +* The `finally` runs however the tool ends: it returned, it raised, or it was cancelled. +* Cleanup that has to `await` needs `shield=True`. In a cancelled handler every further `await` raises too, so without the shield `release_hold` would stop at its first line. +* Nothing can cancel a shielded block, so give it a time limit. Here that is `5` seconds. + +!!! tip + Reach for `finally`, not `except`. The cancellation has to keep travelling up once your cleanup + is done, and a `finally` lets it. + +## Stop early in a plain `def` tool + +A plain `def` tool runs in a thread, and nothing can interrupt a thread from outside. The tool has to ask: + +```python title="server.py" hl_lines="22 25-26" +--8<-- "docs_src/cancellation/tutorial002.py" +``` + +* `anyio.from_thread.check_cancelled()` does nothing while the call is live, and raises once it has been cancelled. Call it between units of work. +* Cleanup goes in a `finally` here too. Nothing in a thread awaits, so it needs no shield. +* A `def` tool that never asks runs to the end, and its result is thrown away. + +## Where it applies + +Prompt and resource functions are cancelled exactly like tools. + +It works the same over stdio and Streamable HTTP. With this SDK's `Client`, giving up means cancelling the task that awaits `call_tool`, or letting its `read_timeout_seconds` run out. + +!!! warning + Two Streamable HTTP options keep the news from your handler: `json_response=True` on a + `2026-07-28` connection, and `stateless_http=True` on a legacy one. There the handler runs to + the end whatever the client did. + +## Recap + +* When the client gives up on a call, the SDK cancels the handler: tool, prompt or resource. +* `async def`: clean up in a `finally`, and put cleanup that awaits inside `anyio.move_on_after(seconds, shield=True)`. +* Plain `def`: call `anyio.from_thread.check_cancelled()` between units of work, or the tool runs to the end. A plain `finally` cleans up. +* `json_response=True` (modern connections) and `stateless_http=True` (legacy ones) switch cancellation off. + +Progress and cancellation are between a running tool and its *caller*. The lines it logs for *you*, the person operating the server, are a different channel: **[Logging](logging.md)**. diff --git a/docs/handlers/index.md b/docs/handlers/index.md index daf9fde19a..4e075a331f 100644 --- a/docs/handlers/index.md +++ b/docs/handlers/index.md @@ -22,6 +22,8 @@ What it can do while it runs: **[Sampling and roots](sampling-and-roots.md)**, deprecated but still served. * Report **[Progress](progress.md)** on something slow. +* Clean up, or stop early, when the client gives up on the call, with + **[Cancellation](cancellation.md)**. * Write logs (to standard error, for whoever operates the server) with **[Logging](logging.md)**. * Tell subscribed clients that something changed with diff --git a/docs/handlers/progress.md b/docs/handlers/progress.md index d48b4c6ab1..a5a6c123d5 100644 --- a/docs/handlers/progress.md +++ b/docs/handlers/progress.md @@ -115,4 +115,4 @@ The callback receives `total=None`. A client can still show *activity* ("3 impor * No callback on the call means `report_progress` does nothing. Report unconditionally. * Omit `total` when you don't know it; the callback gets `None`. -Progress is what a running tool shows the *user*. The lines it logs for *you*, the person operating the server, are a different channel: **[Logging](logging.md)**. +Progress is for a client that is still waiting. What your tool sees when the client stops waiting is **[Cancellation](cancellation.md)**. diff --git a/docs/servers/tools.md b/docs/servers/tools.md index f434389145..445700266e 100644 --- a/docs/servers/tools.md +++ b/docs/servers/tools.md @@ -136,7 +136,7 @@ You can mix and match: plain parameters next to model parameters, nested models, If a tool does I/O (calls an API, reads a file, queries a database), declare it `async def` and `await` inside it. The SDK awaits it. -A plain `def` tool works too: the SDK runs it in a thread so it never blocks the server. +A plain `def` tool works too: the SDK runs it in a thread so it never blocks the server. A long one can check whether the client is still waiting; see **[Cancellation](../handlers/cancellation.md)**. There is nothing else to configure. diff --git a/docs_src/cancellation/__init__.py b/docs_src/cancellation/__init__.py new file mode 100644 index 0000000000..e69de29bb2 diff --git a/docs_src/cancellation/tutorial001.py b/docs_src/cancellation/tutorial001.py new file mode 100644 index 0000000000..258d255d2f --- /dev/null +++ b/docs_src/cancellation/tutorial001.py @@ -0,0 +1,28 @@ +import anyio + +from mcp.server import MCPServer + +mcp = MCPServer("Bookshop") + +holds: set[str] = set() + + +async def take_payment(title: str) -> None: + await anyio.sleep(30) # the customer is typing a card number + + +async def release_hold(title: str) -> None: + await anyio.sleep(0.1) # a round trip to the stock system + holds.discard(title) + + +@mcp.tool() +async def order_book(title: str) -> str: + """Hold a copy of a book while the customer pays for it.""" + holds.add(title) + try: + await take_payment(title) + return f"Ordered {title!r}." + finally: + with anyio.move_on_after(5, shield=True): + await release_hold(title) diff --git a/docs_src/cancellation/tutorial002.py b/docs_src/cancellation/tutorial002.py new file mode 100644 index 0000000000..4aedec982d --- /dev/null +++ b/docs_src/cancellation/tutorial002.py @@ -0,0 +1,26 @@ +import time + +import anyio.from_thread + +from mcp.server import MCPServer + +mcp = MCPServer("Bookshop") + +offline: set[str] = set() + + +def index_book(title: str) -> None: + time.sleep(1) # slow work with nothing to await + + +@mcp.tool() +def rebuild_index(titles: list[str]) -> str: + """Take search offline and rebuild its index, one book at a time.""" + offline.add("search") + try: + for title in titles: + anyio.from_thread.check_cancelled() + index_book(title) + return f"Indexed {len(titles)} books." + finally: + offline.discard("search") diff --git a/mkdocs.yml b/mkdocs.yml index a75053326f..5d8a2b02b3 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -40,6 +40,7 @@ nav: - Multi-round-trip requests: handlers/multi-round-trip.md - Sampling and roots: handlers/sampling-and-roots.md - Progress: handlers/progress.md + - Cancellation: handlers/cancellation.md - Logging: handlers/logging.md - Subscriptions: handlers/subscriptions.md - Running your server: diff --git a/tests/docs_src/test_cancellation.py b/tests/docs_src/test_cancellation.py new file mode 100644 index 0000000000..e28cc77c11 --- /dev/null +++ b/tests/docs_src/test_cancellation.py @@ -0,0 +1,234 @@ +"""`docs/handlers/cancellation.md`: every claim the page makes, proved against the real SDK.""" + +import threading +from collections.abc import Awaitable, Callable + +import anyio +import anyio.from_thread +import pytest +from mcp_types import REQUEST_TIMEOUT + +from docs_src.cancellation import tutorial001, tutorial002 +from mcp import Client, MCPError +from mcp.client.streamable_http import streamable_http_client +from mcp.server import MCPServer +from tests.interaction._connect import BASE_URL, mounted_app + +# See test_index.py for why this is a per-module mark and not a conftest hook. +pytestmark = [pytest.mark.anyio, pytest.mark.filterwarnings("error::mcp.MCPDeprecationWarning")] + +# "auto" dispatches in process; "legacy" puts a JSON-RPC stream, and so a cancellation message, in between. +both_connections = pytest.mark.parametrize("mode", ["auto", "legacy"]) + +TITLES = ["Dune", "Emma", "Ulysses"] + + +async def abandon( + call: Callable[[], Awaitable[object]], started: anyio.Event, then: Callable[[], object] = lambda: None +) -> None: + """Start `call`, cancel the task awaiting it once `started` is set, let the server settle, then run `then`.""" + scope = anyio.CancelScope() + + async def doomed() -> None: + with scope: + await call() + raise NotImplementedError # unreachable: the call never resolves + + async with anyio.create_task_group() as tg: + tg.start_soon(doomed) + await started.wait() + scope.cancel() + await anyio.wait_all_tasks_blocked() + then() + + +@pytest.fixture +def payment_started(monkeypatch: pytest.MonkeyPatch) -> anyio.Event: + """Replace tutorial001's `take_payment` with one that says the tool reached it and then never finishes.""" + started = anyio.Event() + + async def take_payment(title: str) -> None: + assert title in tutorial001.holds + started.set() + await anyio.sleep_forever() + + monkeypatch.setattr(tutorial001, "take_payment", take_payment) + return started + + +@pytest.fixture +def hold_released(monkeypatch: pytest.MonkeyPatch) -> anyio.Event: + """Wrap tutorial001's own `release_hold` so the test can wait for it to reach its last line.""" + released = anyio.Event() + release_hold = tutorial001.release_hold + + async def announcing_release_hold(title: str) -> None: + await release_hold(title) + released.set() + + monkeypatch.setattr(tutorial001, "release_hold", announcing_release_hold) + return released + + +@both_connections +async def test_abandoning_the_call_runs_the_shielded_cleanup_to_the_end( + mode: str, payment_started: anyio.Event, hold_released: anyio.Event +) -> None: + """tutorial001: the client gives up mid-payment, and the `finally` still awaits `release_hold` to completion.""" + with anyio.fail_after(5): + async with Client(tutorial001.mcp, mode=mode) as client: + await abandon(lambda: client.call_tool("order_book", {"title": "Dune"}), payment_started) + await hold_released.wait() + assert tutorial001.holds == set() + + +@pytest.mark.parametrize("mode", ["2026-07-28", "legacy"]) +async def test_abandoning_the_call_over_streamable_http_runs_the_cleanup_too( + mode: str, payment_started: anyio.Event, hold_released: anyio.Event +) -> None: + """The last section: with the default options, either era's way of cancelling over HTTP reaches tutorial001.""" + with anyio.fail_after(5): + async with ( + mounted_app(tutorial001.mcp) as (http, _), + Client(streamable_http_client(f"{BASE_URL}/mcp", http_client=http), mode=mode) as client, + ): + await abandon(lambda: client.call_tool("order_book", {"title": "Dune"}), payment_started) + await hold_released.wait() + # Let the legacy transport's late answer to the abandoned call land while the client is still open. + await anyio.wait_all_tasks_blocked() + assert tutorial001.holds == set() + + +@both_connections +async def test_a_client_timeout_cancels_the_tool_the_same_way( + mode: str, payment_started: anyio.Event, hold_released: anyio.Event +) -> None: + """The last section: `read_timeout_seconds` running out is the other way this SDK's client gives up.""" + with anyio.fail_after(5): + async with Client(tutorial001.mcp, mode=mode) as client: + with pytest.raises(MCPError) as exc_info: + # The tool never answers, so any positive timeout expires; this one adds no wall-clock time. + await client.call_tool("order_book", {"title": "Dune"}, read_timeout_seconds=0.000001) + await hold_released.wait() + assert exc_info.value.error.code == REQUEST_TIMEOUT + assert payment_started.is_set() + assert tutorial001.holds == set() + + +@both_connections +async def test_check_cancelled_stops_a_def_tool_at_its_next_check(mode: str, monkeypatch: pytest.MonkeyPatch) -> None: + """tutorial002: cancelled during the first book, the loop raises at its next check and the `finally` cleans up.""" + started = anyio.Event() + resume = threading.Event() + indexed: list[str] = [] + + def index_book(title: str) -> None: + assert tutorial002.offline == {"search"} + indexed.append(title) + anyio.from_thread.run_sync(started.set) + assert resume.wait(5) + + monkeypatch.setattr(tutorial002, "index_book", index_book) + with anyio.fail_after(5): + # Leaving the block waits for the tool's thread, so what it left behind is final after it. + async with Client(tutorial002.mcp, mode=mode) as client: + await abandon(lambda: client.call_tool("rebuild_index", {"titles": TITLES}), started, then=resume.set) + assert indexed == ["Dune"] + assert tutorial002.offline == set() + + +async def test_a_def_tool_nobody_cancels_indexes_every_book_and_cleans_up(monkeypatch: pytest.MonkeyPatch) -> None: + """tutorial002: while the call is live the checks do nothing, and the `finally` runs on a normal finish too.""" + indexed: list[str] = [] + monkeypatch.setattr(tutorial002, "index_book", indexed.append) + async with Client(tutorial002.mcp) as client: + result = await client.call_tool("rebuild_index", {"titles": TITLES}) + assert result.structured_content == {"result": "Indexed 3 books."} + assert indexed == TITLES + assert tutorial002.offline == set() + + +async def test_a_def_tool_that_never_checks_runs_to_the_end() -> None: + """The `def` section's last bullet: nothing interrupts the thread, so the tool outlives its own cancellation.""" + started = anyio.Event() + resume = threading.Event() + finished: list[str] = [] + mcp = MCPServer("Bookshop") + + @mcp.tool() + def rebuild_index() -> str: + anyio.from_thread.run_sync(started.set) + assert resume.wait(5) + finished.append("rebuild_index") + return "Indexed 3 books." + + with anyio.fail_after(5): + async with Client(mcp, mode="legacy") as client: + await abandon(lambda: client.call_tool("rebuild_index", {}), started, then=resume.set) + assert finished == ["rebuild_index"] + + +async def test_prompt_and_resource_functions_are_cancelled_like_tools() -> None: + """The last section: a prompt or a resource function parked on an `await` is cancelled when the client gives up.""" + started = {"blurb": anyio.Event(), "stock": anyio.Event()} + cancelled = {"blurb": anyio.Event(), "stock": anyio.Event()} + mcp = MCPServer("Bookshop") + + async def park(name: str) -> str: + started[name].set() + try: + await anyio.sleep_forever() + finally: + cancelled[name].set() + raise NotImplementedError # unreachable: only cancellation ends the sleep + + @mcp.prompt() + async def blurb() -> str: + return await park("blurb") + + @mcp.resource("stock://all") + async def stock() -> str: + return await park("stock") + + with anyio.fail_after(5): + async with Client(mcp, mode="legacy") as client: + await abandon(lambda: client.get_prompt("blurb"), started["blurb"]) + await cancelled["blurb"].wait() + await abandon(lambda: client.read_resource("stock://all"), started["stock"]) + await cancelled["stock"].wait() + + +@pytest.mark.parametrize( + ("json_response", "stateless_http", "mode"), + [(True, False, "2026-07-28"), (False, True, "legacy")], + ids=["json_response-modern", "stateless_http-legacy"], +) +async def test_two_http_options_keep_the_cancellation_from_the_handler( + json_response: bool, stateless_http: bool, mode: str +) -> None: + """The `!!! warning`: on these two pairings the abandoned tool is still there once everything has settled. + + Pins known gaps. A stateless legacy server has no session in which to find the request that + `notifications/cancelled` names. In JSON-response mode the 2026-07-28 entry does not watch for + the disconnect that is that revision's cancellation signal; if that arm starts timing out, the + gap was closed and the warning should lose that half. + """ + started = anyio.Event() + resume = anyio.Event() + finished = anyio.Event() + mcp = MCPServer("Bookshop") + + @mcp.tool() + async def order_book() -> str: + started.set() + await resume.wait() + finished.set() + return "Ordered." + + with anyio.fail_after(5): + async with ( + mounted_app(mcp, json_response=json_response, stateless_http=stateless_http) as (http, _), + Client(streamable_http_client(f"{BASE_URL}/mcp", http_client=http), mode=mode) as client, + ): + await abandon(lambda: client.call_tool("order_book", {}), started, then=resume.set) + await finished.wait()