From dcc8afae7274a68cda8d9863cd084391282dcd6e Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 10:33:30 +0200 Subject: [PATCH 1/9] Detect event-loop blocking in tests --- pyproject.toml | 2 ++ src/mcp/client/session.py | 11 ++++++++--- tests/conftest.py | 24 ++++++++++++++++++++++++ uv.lock | 20 ++++++++++++++++++++ 4 files changed, 54 insertions(+), 3 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 536ef5278a..044472aa30 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -77,6 +77,8 @@ dev = [ "strict-no-cover", "logfire>=3.0.0", "opentelemetry-sdk>=1.39.1", + # BlockBuster minor releases may add breaking detection rules. + "blockbuster>=1.5.26,<1.6", ] docs = [ # Zensical is the Material team's successor to MkDocs; it natively diff --git a/src/mcp/client/session.py b/src/mcp/client/session.py index a618112153..826021af47 100644 --- a/src/mcp/client/session.py +++ b/src/mcp/client/session.py @@ -12,6 +12,7 @@ import anyio import anyio.abc import anyio.lowlevel +import anyio.to_thread import mcp_types as types from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream from mcp_types import ( @@ -1133,12 +1134,16 @@ async def validate_tool_result(self, name: str, result: types.CallToolResult) -> logger.warning(f"Tool {name} not listed by server, cannot validate any structured content") if output_schema is not None: + if result.structured_content is None: + raise RuntimeError(f"Tool {name} has an output schema but did not return structured content") + validator = self._tool_output_validators.get(name) + if validator is None: + # First compilation lazily reads jsonschema's bundled schemas. + validator = await anyio.to_thread.run_sync(self._output_schema_validator, name, output_schema) + from jsonschema import exceptions as jsonschema_exceptions from referencing.exceptions import Unresolvable - if result.structured_content is None: - raise RuntimeError(f"Tool {name} has an output schema but did not return structured content") - validator = self._output_schema_validator(name, output_schema) # `best_match` picks the same error the previous `jsonschema.validate()` call raised, # so the message a caller sees is unchanged. It is untyped upstream. errors = validator.iter_errors(result.structured_content) diff --git a/tests/conftest.py b/tests/conftest.py index 9ade27e7f3..2a8331b933 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -1,7 +1,9 @@ import os from collections.abc import AsyncIterator, Iterator +import httpcore2 as _httpcore2 import pytest +from blockbuster import BlockBuster # OpenTelemetry's `set_tracer_provider` is set-once per process, so the suite # uses a single span-capture mechanism: logfire's `capfire` fixture (its @@ -17,12 +19,34 @@ import mcp.shared._otel # noqa: E402 +# Load httpx2's lazy default transport before BlockBuster starts. +del _httpcore2 + @pytest.fixture(scope="session") def anyio_backend() -> str: return "asyncio" +_BLOCKBUSTER = BlockBuster(["mcp", "mcp_types"]) +# Coverage reads source files while collecting data. +_BLOCKBUSTER.functions["os.stat"].can_block_in("coverage/python.py", "get_python_source") +_BLOCKBUSTER.functions["io.BufferedReader.read"].can_block_in("coverage/python.py", "read_python_source") +# These public synchronous conversions read the media file by design. +_BLOCKBUSTER.functions["io.BufferedReader.read"].can_block_in( + "mcp/server/mcpserver/utilities/types.py", ("to_image_content", "to_audio_content") +) + + +@pytest.fixture(autouse=True) +def blockbuster() -> Iterator[BlockBuster]: + try: + _BLOCKBUSTER.activate() + yield _BLOCKBUSTER + finally: + _BLOCKBUSTER.deactivate() + + @pytest.fixture(scope="module", autouse=True) async def _module_runner_lease(anyio_backend: str) -> AsyncIterator[None]: """Share one event loop across each module's tests instead of one per test. diff --git a/uv.lock b/uv.lock index ebc03ffa75..c99d6a0739 100644 --- a/uv.lock +++ b/uv.lock @@ -154,6 +154,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/8e/0d/52d98722666d6fc6c3dd4c76df339501d6efd40e0ff95e6186a7b7f0befd/black-26.3.1-py3-none-any.whl", hash = "sha256:2bd5aa94fc267d38bb21a70d7410a89f1a1d318841855f698746f8e7f51acd1b", size = 207542, upload-time = "2026-03-12T03:36:01.668Z" }, ] +[[package]] +name = "blockbuster" +version = "1.5.27" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "forbiddenfruit", marker = "implementation_name == 'cpython'" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/ff/c3/21678f5b979be2cbf0e68352f7330a84ea4e24674023e4de981c78a218cd/blockbuster-1.5.27.tar.gz", hash = "sha256:b8e9d988b9b91ba468c94530e219f26a00d3ff616b39ebf3da561a2a3eea9dd4", size = 96732, upload-time = "2026-08-17T23:53:13.378Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/63/c5/092e631bc1fba86f0a822be65c137c90a71b71ba0a0865e7e9a21f6ca05e/blockbuster-1.5.27-py3-none-any.whl", hash = "sha256:f0acf153d22a791bf5f142935332ef8530960ec215541b48a6037e6cea0a8645", size = 13517, upload-time = "2026-08-17T23:53:14.625Z" }, +] + [[package]] name = "certifi" version = "2025.8.3" @@ -577,6 +589,12 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/c1/ea/53f2148663b321f21b5a606bd5f191517cf40b7072c0497d3c92c4a13b1e/executing-2.2.1-py2.py3-none-any.whl", hash = "sha256:760643d3452b4d777d295bb167ccc74c64a81df23fb5e08eff250c425a4b2017", size = 28317, upload-time = "2025-09-01T09:48:08.5Z" }, ] +[[package]] +name = "forbiddenfruit" +version = "0.1.4" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/e6/79/d4f20e91327c98096d605646bdc6a5ffedae820f38d378d3515c42ec5e60/forbiddenfruit-0.1.4.tar.gz", hash = "sha256:e3f7e66561a29ae129aac139a85d610dbf3dd896128187ed5454b6421f624253", size = 43756, upload-time = "2021-01-16T21:03:35.401Z" } + [[package]] name = "genson" version = "1.3.0" @@ -1023,6 +1041,7 @@ codegen = [ { name = "datamodel-code-generator" }, ] dev = [ + { name = "blockbuster" }, { name = "coverage", extra = ["toml"] }, { name = "dirty-equals" }, { name = "inline-snapshot" }, @@ -1081,6 +1100,7 @@ provides-extras = ["cli", "rich"] [package.metadata.requires-dev] codegen = [{ name = "datamodel-code-generator", specifier = "==0.57.0" }] dev = [ + { name = "blockbuster", specifier = ">=1.5.26,<1.6" }, { name = "coverage", extras = ["toml"], specifier = ">=7.10.7,<=7.13" }, { name = "dirty-equals", specifier = ">=0.9.0" }, { name = "inline-snapshot", specifier = ">=0.23.0" }, From 3e2e5eb9eaded63c4d9fc8124c1564ffbfffd3d5 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 10:40:24 +0200 Subject: [PATCH 2/9] Keep validator cache updates on the event loop --- src/mcp/client/session.py | 9 ++--- tests/client/test_session_promotions.py | 47 +++++++++++++++++++++++++ 2 files changed, 50 insertions(+), 6 deletions(-) diff --git a/src/mcp/client/session.py b/src/mcp/client/session.py index 826021af47..e57179c41a 100644 --- a/src/mcp/client/session.py +++ b/src/mcp/client/session.py @@ -1140,6 +1140,8 @@ async def validate_tool_result(self, name: str, result: types.CallToolResult) -> if validator is None: # First compilation lazily reads jsonschema's bundled schemas. validator = await anyio.to_thread.run_sync(self._output_schema_validator, name, output_schema) + if _same_schema(self._tool_output_schemas.get(name), output_schema): + self._tool_output_validators[name] = validator from jsonschema import exceptions as jsonschema_exceptions from referencing.exceptions import Unresolvable @@ -1174,18 +1176,13 @@ def _output_schema_validator(self, name: str, output_schema: dict[str, Any]) -> from jsonschema.validators import validator_for from referencing import Registry - if (validator := self._tool_output_validators.get(name)) is not None: - return validator - validator_cls = validator_for(output_schema) try: validator_cls.check_schema(output_schema) except SchemaError as e: raise RuntimeError(f"Invalid schema for tool {name}: {e}") # An explicit empty registry: `$ref`s resolve within the schema document and the bundled metaschemas. - validator = validator_cls(output_schema, registry=Registry()) - self._tool_output_validators[name] = validator - return validator + return validator_cls(output_schema, registry=Registry()) async def list_prompts(self, *, params: types.PaginatedRequestParams | None = None) -> types.ListPromptsResult: """Send a prompts/list request. diff --git a/tests/client/test_session_promotions.py b/tests/client/test_session_promotions.py index 6d6b6bc8dc..7250a4dbc5 100644 --- a/tests/client/test_session_promotions.py +++ b/tests/client/test_session_promotions.py @@ -1,7 +1,11 @@ """`dispatch_input_request` and `validate_tool_result` are public `ClientSession` API.""" +import anyio +import anyio.from_thread import mcp_types as types import pytest +from jsonschema.protocols import Validator +from jsonschema.validators import validator_for from mcp_types import ( CallToolResult, ErrorData, @@ -125,3 +129,46 @@ async def on_list_tools(ctx: ServerRequestContext, params: PaginatedRequestParam await client.session.list_tools() with pytest.raises(RuntimeError, match="Invalid structured content returned by tool t"): await client.session.validate_tool_result("t", integer_result) + + +@pytest.mark.anyio +async def test_schema_change_during_compilation_does_not_cache_the_old_validator( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """SDK-defined: a concurrent tool relisting cannot leave a stale compiled validator cached.""" + schemas = [ + {"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}, + {"type": "object", "properties": {"x": {"type": "string"}}, "required": ["x"]}, + ] + + async def on_list_tools(ctx: ServerRequestContext, params: PaginatedRequestParams | None) -> ListToolsResult: + return ListToolsResult(tools=[Tool(name="t", input_schema={"type": "object"}, output_schema=schemas.pop(0))]) + + compilation_started = anyio.Event() + continue_compilation = anyio.Event() + + def paused_validator_for(schema: dict[str, object], default: type[Validator] | None = None) -> type[Validator]: + validator = validator_for(schema) if default is None else validator_for(schema, default=default) + if not compilation_started.is_set(): + anyio.from_thread.run_sync(compilation_started.set) + anyio.from_thread.run(continue_compilation.wait) + return validator + + monkeypatch.setattr("jsonschema.validators.validator_for", paused_validator_for) + server = Server("test-server", on_list_tools=on_list_tools) + async with Client(server) as client: + await client.session.list_tools() + with anyio.fail_after(5): + async with anyio.create_task_group() as task_group: + task_group.start_soon( + client.session.validate_tool_result, + "t", + CallToolResult(content=[], structured_content={"x": 1}), + ) + await compilation_started.wait() + try: + await client.session.list_tools() + finally: + continue_compilation.set() + + await client.session.validate_tool_result("t", CallToolResult(content=[], structured_content={"x": "yes"})) From 797d113dee8f61c3de67ca7aa6dc6bc576c2cb6d Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 10:45:17 +0200 Subject: [PATCH 3/9] Resolve Windows executables off the event loop --- src/mcp/client/stdio.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/src/mcp/client/stdio.py b/src/mcp/client/stdio.py index 3e03eef9ef..23d31e4929 100644 --- a/src/mcp/client/stdio.py +++ b/src/mcp/client/stdio.py @@ -18,6 +18,7 @@ import anyio import anyio.lowlevel +import anyio.to_thread import mcp_types as types from anyio.abc import AsyncResource, Process from anyio.streams.text import TextReceiveStream @@ -120,7 +121,7 @@ async def stdio_client( OSError: If the server process cannot be spawned. ValueError: If the spawn parameters are invalid (embedded NUL bytes). """ - command = _get_executable_command(server.command) + command = await _get_executable_command(server.command) process = await _create_platform_compatible_process( command=command, @@ -317,10 +318,10 @@ def _close_subprocess_transport(process: ServerProcess) -> None: close() -def _get_executable_command(command: str) -> str: +async def _get_executable_command(command: str) -> str: """Normalizes the command for the current platform.""" if sys.platform == "win32": # pragma: no cover - return get_windows_executable_command(command) + return await anyio.to_thread.run_sync(get_windows_executable_command, command) else: # pragma: lax no cover return command From b66cfadb684d0c36a2cd19a425eab678d2ca17d4 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 10:49:08 +0200 Subject: [PATCH 4/9] Isolate fake stdio processes from command resolution --- tests/client/test_stdio.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/tests/client/test_stdio.py b/tests/client/test_stdio.py index 91f829ff98..edbf96e322 100644 --- a/tests/client/test_stdio.py +++ b/tests/client/test_stdio.py @@ -199,6 +199,9 @@ def install_fake_process( """ terminated: list[FakeProcess] = [] + async def fake_get_executable_command(command: str) -> str: + return command + async def fake_spawn( command: str, args: list[str], @@ -212,6 +215,7 @@ async def fake_terminate_tree(proc: FakeProcess) -> None: terminated.append(proc) proc.exit(-15) + monkeypatch.setattr(stdio, "_get_executable_command", fake_get_executable_command) monkeypatch.setattr(stdio, "_create_platform_compatible_process", fake_spawn) monkeypatch.setattr(stdio, "_terminate_process_tree", fake_terminate_tree) if grace_period is not None: From 2f04ac87adba8a6d4e5e85af9126d789e3724891 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 10:59:11 +0200 Subject: [PATCH 5/9] Abandon canceled Windows command resolution --- src/mcp/client/stdio.py | 4 ++-- tests/client/test_stdio.py | 43 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+), 2 deletions(-) diff --git a/src/mcp/client/stdio.py b/src/mcp/client/stdio.py index 23d31e4929..6a3ad12111 100644 --- a/src/mcp/client/stdio.py +++ b/src/mcp/client/stdio.py @@ -320,8 +320,8 @@ def _close_subprocess_transport(process: ServerProcess) -> None: async def _get_executable_command(command: str) -> str: """Normalizes the command for the current platform.""" - if sys.platform == "win32": # pragma: no cover - return await anyio.to_thread.run_sync(get_windows_executable_command, command) + if sys.platform == "win32": + return await anyio.to_thread.run_sync(get_windows_executable_command, command, abandon_on_cancel=True) else: # pragma: lax no cover return command diff --git a/tests/client/test_stdio.py b/tests/client/test_stdio.py index edbf96e322..0101e45a76 100644 --- a/tests/client/test_stdio.py +++ b/tests/client/test_stdio.py @@ -14,14 +14,18 @@ import os import signal import sys +import threading from collections.abc import Callable from contextlib import AsyncExitStack, suppress from pathlib import Path +from types import SimpleNamespace from typing import TextIO, cast import anyio import anyio.abc +import anyio.from_thread import anyio.lowlevel +import anyio.to_thread import pytest import trio import trio.testing @@ -572,6 +576,45 @@ async def test_a_command_that_cannot_be_execed_raises_enoent() -> None: assert exc_info.value.errno == errno.ENOENT +@pytest.mark.anyio +async def test_cancellation_during_windows_command_resolution_returns_before_resolution_finishes( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Cancelling `stdio_client` does not wait for blocked Windows command resolution.""" + resolution_started = anyio.Event() + resolution_release = threading.Event() + resolution_finished = threading.Event() + + def blocking_resolver(command: str) -> str: + anyio.from_thread.run_sync(resolution_started.set) + resolution_release.wait() + resolution_finished.set() + return command + + monkeypatch.setattr(stdio, "sys", SimpleNamespace(platform="win32")) + monkeypatch.setattr(stdio, "get_windows_executable_command", blocking_resolver) + + cancel_scope = anyio.CancelScope() + client_stopped = anyio.Event() + + async def run_client() -> None: + with cancel_scope: + async with AsyncExitStack() as stack: + await stack.enter_async_context(stdio_client(FAKE_PARAMS)) + client_stopped.set() + + with anyio.fail_after(5): + async with anyio.create_task_group() as tg: + tg.start_soon(run_client) + await resolution_started.wait() + cancel_scope.cancel() + try: + await client_stopped.wait() + finally: + resolution_release.set() + await anyio.to_thread.run_sync(resolution_finished.wait) + + @pytest.mark.anyio async def test_cancellation_during_spawn_leaks_no_streams(monkeypatch: pytest.MonkeyPatch) -> None: """Cancellation while the spawn is still in flight must not leak the internal streams. From ca725437330750f1fd4f1acb92f4d268d9d35396 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 11:04:59 +0200 Subject: [PATCH 6/9] Keep validator compilation independent of shared workers --- src/mcp/client/session.py | 13 ++++++- tests/client/test_session_promotions.py | 48 +++++++++++++++++++++++++ 2 files changed, 60 insertions(+), 1 deletion(-) diff --git a/src/mcp/client/session.py b/src/mcp/client/session.py index e57179c41a..8289e32020 100644 --- a/src/mcp/client/session.py +++ b/src/mcp/client/session.py @@ -2,6 +2,7 @@ import json import logging +import sys from collections.abc import Callable, Mapping, Sequence from dataclasses import dataclass from functools import cache, reduce @@ -442,6 +443,7 @@ def __init__( # Compiled output-schema validators, derived from `_tool_output_schemas` and owned by # `_absorb_tool_listing`, which evicts a tool's entry whenever its schema changes. self._tool_output_validators: dict[str, Validator] = {} + self._tool_output_validator_limiter = anyio.CapacityLimiter(1) self._x_mcp_header_maps: dict[str, dict[tuple[str, ...], str]] = {} self._initialize_result: types.InitializeResult | None = None self._discover_result: types.DiscoverResult | None = None @@ -1139,7 +1141,16 @@ async def validate_tool_result(self, name: str, result: types.CallToolResult) -> validator = self._tool_output_validators.get(name) if validator is None: # First compilation lazily reads jsonschema's bundled schemas. - validator = await anyio.to_thread.run_sync(self._output_schema_validator, name, output_schema) + if sys.platform == "emscripten": + # Emscripten cannot start worker threads. + validator = self._output_schema_validator(name, output_schema) + else: + validator = await anyio.to_thread.run_sync( + self._output_schema_validator, + name, + output_schema, + limiter=self._tool_output_validator_limiter, + ) if _same_schema(self._tool_output_schemas.get(name), output_schema): self._tool_output_validators[name] = validator diff --git a/tests/client/test_session_promotions.py b/tests/client/test_session_promotions.py index 7250a4dbc5..6a2131725a 100644 --- a/tests/client/test_session_promotions.py +++ b/tests/client/test_session_promotions.py @@ -1,7 +1,12 @@ """`dispatch_input_request` and `validate_tool_result` are public `ClientSession` API.""" +import threading +from types import SimpleNamespace +from unittest.mock import AsyncMock + import anyio import anyio.from_thread +import anyio.to_thread import mcp_types as types import pytest from jsonschema.protocols import Validator @@ -15,6 +20,7 @@ Tool, ) +import mcp.client.session as session_module from mcp.client.client import Client from mcp.client.session import ClientRequestContext, ClientSession from mcp.server import Server, ServerRequestContext @@ -61,6 +67,48 @@ async def test_validate_tool_result_passes_a_conforming_result() -> None: await client.session.validate_tool_result("t", CallToolResult(content=[], structured_content={"x": 1})) +@pytest.mark.anyio +async def test_validate_tool_result_compiles_inline_without_worker_threads(monkeypatch: pytest.MonkeyPatch) -> None: + """Output schema validation remains available on Emscripten, which has no worker threads.""" + server = _make_server({"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}) + async with Client(server) as client: + await client.session.list_tools() + monkeypatch.setattr(session_module, "sys", SimpleNamespace(platform="emscripten")) + run_sync = AsyncMock(wraps=anyio.to_thread.run_sync) + monkeypatch.setattr(anyio.to_thread, "run_sync", run_sync) + + await client.session.validate_tool_result("t", CallToolResult(content=[], structured_content={"x": 1})) + + run_sync.assert_not_awaited() + + +@pytest.mark.anyio +async def test_validator_compilation_does_not_wait_for_default_worker_capacity() -> None: + """Schema compilation cannot deadlock behind sync handlers occupying the default worker pool.""" + server = _make_server({"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}) + async with Client(server) as client: + await client.session.list_tools() + result = CallToolResult(content=[], structured_content={"x": 1}) + default_limiter = anyio.to_thread.current_default_thread_limiter() + original_tokens = default_limiter.total_tokens + validation_finished = threading.Event() + + def validate_from_worker() -> None: + try: + anyio.from_thread.run(client.session.validate_tool_result, "t", result) + finally: + validation_finished.set() + + default_limiter.total_tokens = 1 + try: + with anyio.fail_after(5): + await anyio.to_thread.run_sync(validate_from_worker, abandon_on_cancel=True) + finally: + default_limiter.total_tokens = original_tokens + with anyio.fail_after(5): + await anyio.to_thread.run_sync(validation_finished.wait) + + @pytest.mark.anyio async def test_validate_tool_result_raises_on_schema_mismatch() -> None: server = _make_server({"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}) From 8b118f8ef32c47630a588192392e3c243d913c21 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 11:37:41 +0200 Subject: [PATCH 7/9] Match Starlette's BlockBuster fixture pattern --- pyproject.toml | 3 +-- tests/conftest.py | 26 ++++++++++++-------------- uv.lock | 2 +- 3 files changed, 14 insertions(+), 17 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 044472aa30..4a86cdb0f7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -77,8 +77,7 @@ dev = [ "strict-no-cover", "logfire>=3.0.0", "opentelemetry-sdk>=1.39.1", - # BlockBuster minor releases may add breaking detection rules. - "blockbuster>=1.5.26,<1.6", + "blockbuster>=1.5.23", ] docs = [ # Zensical is the Material team's successor to MkDocs; it natively diff --git a/tests/conftest.py b/tests/conftest.py index 2a8331b933..9bf193acae 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -28,23 +28,21 @@ def anyio_backend() -> str: return "asyncio" -_BLOCKBUSTER = BlockBuster(["mcp", "mcp_types"]) -# Coverage reads source files while collecting data. -_BLOCKBUSTER.functions["os.stat"].can_block_in("coverage/python.py", "get_python_source") -_BLOCKBUSTER.functions["io.BufferedReader.read"].can_block_in("coverage/python.py", "read_python_source") -# These public synchronous conversions read the media file by design. -_BLOCKBUSTER.functions["io.BufferedReader.read"].can_block_in( - "mcp/server/mcpserver/utilities/types.py", ("to_image_content", "to_audio_content") -) - - @pytest.fixture(autouse=True) -def blockbuster() -> Iterator[BlockBuster]: +def blockbuster() -> Iterator[None]: + bb = BlockBuster(["mcp", "mcp_types"]) + # Coverage reads source files while collecting data. + bb.functions["os.stat"].can_block_in("coverage/python.py", "get_python_source") + bb.functions["io.BufferedReader.read"].can_block_in("coverage/python.py", "read_python_source") + # These public synchronous conversions read the media file by design. + bb.functions["io.BufferedReader.read"].can_block_in( + "mcp/server/mcpserver/utilities/types.py", ("to_image_content", "to_audio_content") + ) + bb.activate() try: - _BLOCKBUSTER.activate() - yield _BLOCKBUSTER + yield finally: - _BLOCKBUSTER.deactivate() + bb.deactivate() @pytest.fixture(scope="module", autouse=True) diff --git a/uv.lock b/uv.lock index c99d6a0739..4fedcb4887 100644 --- a/uv.lock +++ b/uv.lock @@ -1100,7 +1100,7 @@ provides-extras = ["cli", "rich"] [package.metadata.requires-dev] codegen = [{ name = "datamodel-code-generator", specifier = "==0.57.0" }] dev = [ - { name = "blockbuster", specifier = ">=1.5.26,<1.6" }, + { name = "blockbuster", specifier = ">=1.5.23" }, { name = "coverage", extras = ["toml"], specifier = ">=7.10.7,<=7.13" }, { name = "dirty-equals", specifier = ">=0.9.0" }, { name = "inline-snapshot", specifier = ">=0.23.0" }, From 36fc996cd35a2b52fecc6815022285bed44c33e1 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 11:50:19 +0200 Subject: [PATCH 8/9] Allow jsonschema lazy imports in BlockBuster --- src/mcp/client/session.py | 31 +++----- tests/client/test_session_promotions.py | 95 ------------------------- tests/conftest.py | 3 + 3 files changed, 12 insertions(+), 117 deletions(-) diff --git a/src/mcp/client/session.py b/src/mcp/client/session.py index 8289e32020..a618112153 100644 --- a/src/mcp/client/session.py +++ b/src/mcp/client/session.py @@ -2,7 +2,6 @@ import json import logging -import sys from collections.abc import Callable, Mapping, Sequence from dataclasses import dataclass from functools import cache, reduce @@ -13,7 +12,6 @@ import anyio import anyio.abc import anyio.lowlevel -import anyio.to_thread import mcp_types as types from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream from mcp_types import ( @@ -443,7 +441,6 @@ def __init__( # Compiled output-schema validators, derived from `_tool_output_schemas` and owned by # `_absorb_tool_listing`, which evicts a tool's entry whenever its schema changes. self._tool_output_validators: dict[str, Validator] = {} - self._tool_output_validator_limiter = anyio.CapacityLimiter(1) self._x_mcp_header_maps: dict[str, dict[tuple[str, ...], str]] = {} self._initialize_result: types.InitializeResult | None = None self._discover_result: types.DiscoverResult | None = None @@ -1136,27 +1133,12 @@ async def validate_tool_result(self, name: str, result: types.CallToolResult) -> logger.warning(f"Tool {name} not listed by server, cannot validate any structured content") if output_schema is not None: - if result.structured_content is None: - raise RuntimeError(f"Tool {name} has an output schema but did not return structured content") - validator = self._tool_output_validators.get(name) - if validator is None: - # First compilation lazily reads jsonschema's bundled schemas. - if sys.platform == "emscripten": - # Emscripten cannot start worker threads. - validator = self._output_schema_validator(name, output_schema) - else: - validator = await anyio.to_thread.run_sync( - self._output_schema_validator, - name, - output_schema, - limiter=self._tool_output_validator_limiter, - ) - if _same_schema(self._tool_output_schemas.get(name), output_schema): - self._tool_output_validators[name] = validator - from jsonschema import exceptions as jsonschema_exceptions from referencing.exceptions import Unresolvable + if result.structured_content is None: + raise RuntimeError(f"Tool {name} has an output schema but did not return structured content") + validator = self._output_schema_validator(name, output_schema) # `best_match` picks the same error the previous `jsonschema.validate()` call raised, # so the message a caller sees is unchanged. It is untyped upstream. errors = validator.iter_errors(result.structured_content) @@ -1187,13 +1169,18 @@ def _output_schema_validator(self, name: str, output_schema: dict[str, Any]) -> from jsonschema.validators import validator_for from referencing import Registry + if (validator := self._tool_output_validators.get(name)) is not None: + return validator + validator_cls = validator_for(output_schema) try: validator_cls.check_schema(output_schema) except SchemaError as e: raise RuntimeError(f"Invalid schema for tool {name}: {e}") # An explicit empty registry: `$ref`s resolve within the schema document and the bundled metaschemas. - return validator_cls(output_schema, registry=Registry()) + validator = validator_cls(output_schema, registry=Registry()) + self._tool_output_validators[name] = validator + return validator async def list_prompts(self, *, params: types.PaginatedRequestParams | None = None) -> types.ListPromptsResult: """Send a prompts/list request. diff --git a/tests/client/test_session_promotions.py b/tests/client/test_session_promotions.py index 6a2131725a..6d6b6bc8dc 100644 --- a/tests/client/test_session_promotions.py +++ b/tests/client/test_session_promotions.py @@ -1,16 +1,7 @@ """`dispatch_input_request` and `validate_tool_result` are public `ClientSession` API.""" -import threading -from types import SimpleNamespace -from unittest.mock import AsyncMock - -import anyio -import anyio.from_thread -import anyio.to_thread import mcp_types as types import pytest -from jsonschema.protocols import Validator -from jsonschema.validators import validator_for from mcp_types import ( CallToolResult, ErrorData, @@ -20,7 +11,6 @@ Tool, ) -import mcp.client.session as session_module from mcp.client.client import Client from mcp.client.session import ClientRequestContext, ClientSession from mcp.server import Server, ServerRequestContext @@ -67,48 +57,6 @@ async def test_validate_tool_result_passes_a_conforming_result() -> None: await client.session.validate_tool_result("t", CallToolResult(content=[], structured_content={"x": 1})) -@pytest.mark.anyio -async def test_validate_tool_result_compiles_inline_without_worker_threads(monkeypatch: pytest.MonkeyPatch) -> None: - """Output schema validation remains available on Emscripten, which has no worker threads.""" - server = _make_server({"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}) - async with Client(server) as client: - await client.session.list_tools() - monkeypatch.setattr(session_module, "sys", SimpleNamespace(platform="emscripten")) - run_sync = AsyncMock(wraps=anyio.to_thread.run_sync) - monkeypatch.setattr(anyio.to_thread, "run_sync", run_sync) - - await client.session.validate_tool_result("t", CallToolResult(content=[], structured_content={"x": 1})) - - run_sync.assert_not_awaited() - - -@pytest.mark.anyio -async def test_validator_compilation_does_not_wait_for_default_worker_capacity() -> None: - """Schema compilation cannot deadlock behind sync handlers occupying the default worker pool.""" - server = _make_server({"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}) - async with Client(server) as client: - await client.session.list_tools() - result = CallToolResult(content=[], structured_content={"x": 1}) - default_limiter = anyio.to_thread.current_default_thread_limiter() - original_tokens = default_limiter.total_tokens - validation_finished = threading.Event() - - def validate_from_worker() -> None: - try: - anyio.from_thread.run(client.session.validate_tool_result, "t", result) - finally: - validation_finished.set() - - default_limiter.total_tokens = 1 - try: - with anyio.fail_after(5): - await anyio.to_thread.run_sync(validate_from_worker, abandon_on_cancel=True) - finally: - default_limiter.total_tokens = original_tokens - with anyio.fail_after(5): - await anyio.to_thread.run_sync(validation_finished.wait) - - @pytest.mark.anyio async def test_validate_tool_result_raises_on_schema_mismatch() -> None: server = _make_server({"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}) @@ -177,46 +125,3 @@ async def on_list_tools(ctx: ServerRequestContext, params: PaginatedRequestParam await client.session.list_tools() with pytest.raises(RuntimeError, match="Invalid structured content returned by tool t"): await client.session.validate_tool_result("t", integer_result) - - -@pytest.mark.anyio -async def test_schema_change_during_compilation_does_not_cache_the_old_validator( - monkeypatch: pytest.MonkeyPatch, -) -> None: - """SDK-defined: a concurrent tool relisting cannot leave a stale compiled validator cached.""" - schemas = [ - {"type": "object", "properties": {"x": {"type": "integer"}}, "required": ["x"]}, - {"type": "object", "properties": {"x": {"type": "string"}}, "required": ["x"]}, - ] - - async def on_list_tools(ctx: ServerRequestContext, params: PaginatedRequestParams | None) -> ListToolsResult: - return ListToolsResult(tools=[Tool(name="t", input_schema={"type": "object"}, output_schema=schemas.pop(0))]) - - compilation_started = anyio.Event() - continue_compilation = anyio.Event() - - def paused_validator_for(schema: dict[str, object], default: type[Validator] | None = None) -> type[Validator]: - validator = validator_for(schema) if default is None else validator_for(schema, default=default) - if not compilation_started.is_set(): - anyio.from_thread.run_sync(compilation_started.set) - anyio.from_thread.run(continue_compilation.wait) - return validator - - monkeypatch.setattr("jsonschema.validators.validator_for", paused_validator_for) - server = Server("test-server", on_list_tools=on_list_tools) - async with Client(server) as client: - await client.session.list_tools() - with anyio.fail_after(5): - async with anyio.create_task_group() as task_group: - task_group.start_soon( - client.session.validate_tool_result, - "t", - CallToolResult(content=[], structured_content={"x": 1}), - ) - await compilation_started.wait() - try: - await client.session.list_tools() - finally: - continue_compilation.set() - - await client.session.validate_tool_result("t", CallToolResult(content=[], structured_content={"x": "yes"})) diff --git a/tests/conftest.py b/tests/conftest.py index 9bf193acae..520f01c0dc 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -34,6 +34,9 @@ def blockbuster() -> Iterator[None]: # Coverage reads source files while collecting data. bb.functions["os.stat"].can_block_in("coverage/python.py", "get_python_source") bb.functions["io.BufferedReader.read"].can_block_in("coverage/python.py", "read_python_source") + # jsonschema discovers its bundled schemas during its first import. + bb.functions["os.scandir"].can_block_in("/jsonschema_specifications/_core.py", "_schemas") + bb.functions["io.TextIOWrapper.read"].can_block_in("/jsonschema_specifications/_core.py", "_schemas") # These public synchronous conversions read the media file by design. bb.functions["io.BufferedReader.read"].can_block_in( "mcp/server/mcpserver/utilities/types.py", ("to_image_content", "to_audio_content") From 4cd5f410e096a3e22476b9d07733c7035280a079 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Wed, 16 Sep 2026 11:56:08 +0200 Subject: [PATCH 9/9] Cover lazy imports across supported Python versions --- pyproject.toml | 2 +- tests/conftest.py | 1 + uv.lock | 2 +- 3 files changed, 3 insertions(+), 2 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 4a86cdb0f7..b2f26da55f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -77,7 +77,7 @@ dev = [ "strict-no-cover", "logfire>=3.0.0", "opentelemetry-sdk>=1.39.1", - "blockbuster>=1.5.23", + "blockbuster>=1.5.27", ] docs = [ # Zensical is the Material team's successor to MkDocs; it natively diff --git a/tests/conftest.py b/tests/conftest.py index 520f01c0dc..4350b76f78 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -35,6 +35,7 @@ def blockbuster() -> Iterator[None]: bb.functions["os.stat"].can_block_in("coverage/python.py", "get_python_source") bb.functions["io.BufferedReader.read"].can_block_in("coverage/python.py", "read_python_source") # jsonschema discovers its bundled schemas during its first import. + bb.functions["os.listdir"].can_block_in("/jsonschema_specifications/_core.py", "_schemas") bb.functions["os.scandir"].can_block_in("/jsonschema_specifications/_core.py", "_schemas") bb.functions["io.TextIOWrapper.read"].can_block_in("/jsonschema_specifications/_core.py", "_schemas") # These public synchronous conversions read the media file by design. diff --git a/uv.lock b/uv.lock index 4fedcb4887..40b563e974 100644 --- a/uv.lock +++ b/uv.lock @@ -1100,7 +1100,7 @@ provides-extras = ["cli", "rich"] [package.metadata.requires-dev] codegen = [{ name = "datamodel-code-generator", specifier = "==0.57.0" }] dev = [ - { name = "blockbuster", specifier = ">=1.5.23" }, + { name = "blockbuster", specifier = ">=1.5.27" }, { name = "coverage", extras = ["toml"], specifier = ">=7.10.7,<=7.13" }, { name = "dirty-equals", specifier = ">=0.9.0" }, { name = "inline-snapshot", specifier = ">=0.23.0" },