mirror of
https://github.com/PrefectHQ/fastmcp.git
synced 2026-08-22 05:24:18 +02:00
Compare commits
5 commits
main
...
codex/otel
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9feca45fb1 | ||
|
|
409c971493 | ||
|
|
93074de1f1 | ||
|
|
284b21ca50 | ||
|
|
e70981e37f |
14 changed files with 966 additions and 50 deletions
|
|
@ -28,6 +28,12 @@ You can change which `.env` file is loaded by setting the `FASTMCP_ENV_FILE` env
|
||||||
| `FASTMCP_ENABLE_RICH_TRACEBACKS` | `bool` | `true` | Use rich tracebacks for errors. |
|
| `FASTMCP_ENABLE_RICH_TRACEBACKS` | `bool` | `true` | Use rich tracebacks for errors. |
|
||||||
| `FASTMCP_DEPRECATION_WARNINGS` | `bool` | `true` | Show deprecation warnings. |
|
| `FASTMCP_DEPRECATION_WARNINGS` | `bool` | `true` | Show deprecation warnings. |
|
||||||
|
|
||||||
|
## Telemetry
|
||||||
|
|
||||||
|
| Environment Variable | Type | Default | Description |
|
||||||
|
|---|---|---|---|
|
||||||
|
| `FASTMCP_TELEMETRY_MODE` | `Literal["native", "propagation_only"]` | `native` | Controls FastMCP's native OpenTelemetry span creation. `native` emits FastMCP MCP spans and propagates trace context. `propagation_only` keeps `_meta` trace propagation but suppresses FastMCP's own spans so another instrumentation layer can own the MCP span hierarchy. |
|
||||||
|
|
||||||
## Transport & HTTP
|
## Transport & HTTP
|
||||||
|
|
||||||
These control how the server listens when running with an HTTP transport.
|
These control how the server listens when running with an HTTP transport.
|
||||||
|
|
|
||||||
|
|
@ -12,11 +12,13 @@ FastMCP includes native OpenTelemetry instrumentation for observability. Traces
|
||||||
|
|
||||||
FastMCP uses the OpenTelemetry API for instrumentation. This means:
|
FastMCP uses the OpenTelemetry API for instrumentation. This means:
|
||||||
|
|
||||||
- **Zero configuration required** - Instrumentation is always active
|
- **Zero configuration required** - Native FastMCP instrumentation is available by default
|
||||||
- **No overhead when unused** - Without an SDK, all operations are no-ops
|
- **No overhead when unused** - Without an SDK, all operations are no-ops
|
||||||
- **Bring your own SDK** - You control collection, export, and sampling
|
- **Bring your own SDK** - You control collection, export, and sampling
|
||||||
- **Works with any OTEL backend** - Jaeger, Zipkin, Datadog, New Relic, etc.
|
- **Works with any OTEL backend** - Jaeger, Zipkin, Datadog, New Relic, etc.
|
||||||
|
|
||||||
|
FastMCP also propagates OpenTelemetry context through MCP `params._meta`, including `traceparent`, `tracestate`, and `baggage` when present.
|
||||||
|
|
||||||
## Enabling Telemetry
|
## Enabling Telemetry
|
||||||
|
|
||||||
The easiest way to export traces is using `opentelemetry-instrument`, which configures the SDK automatically:
|
The easiest way to export traces is using `opentelemetry-instrument`, which configures the SDK automatically:
|
||||||
|
|
@ -52,11 +54,11 @@ This works with any OTLP-compatible backend (Jaeger, Zipkin, Grafana Tempo, Data
|
||||||
|
|
||||||
## Tracing
|
## Tracing
|
||||||
|
|
||||||
FastMCP creates spans for all MCP operations, providing end-to-end visibility into request handling.
|
FastMCP creates spans for the MCP operations it currently instruments, providing end-to-end visibility into request handling.
|
||||||
|
|
||||||
### Server Spans
|
### Server Spans
|
||||||
|
|
||||||
The server creates spans for each operation using [MCP semantic conventions](https://opentelemetry.io/docs/specs/semconv/gen-ai/mcp/):
|
The server creates spans for each supported operation using [MCP semantic conventions](https://opentelemetry.io/docs/specs/semconv/gen-ai/mcp/):
|
||||||
|
|
||||||
| Span Name | Description |
|
| Span Name | Description |
|
||||||
|-----------|-------------|
|
|-----------|-------------|
|
||||||
|
|
@ -119,6 +121,27 @@ def greet(name: str) -> str:
|
||||||
The SDK must be configured **before** importing FastMCP to ensure the tracer provider is set when FastMCP initializes.
|
The SDK must be configured **before** importing FastMCP to ensure the tracer provider is set when FastMCP initializes.
|
||||||
</Tip>
|
</Tip>
|
||||||
|
|
||||||
|
## Interoperability Mode
|
||||||
|
|
||||||
|
If another MCP-aware instrumentation layer should own the MCP spans, switch FastMCP to propagation-only mode:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
export FASTMCP_TELEMETRY_MODE=propagation_only
|
||||||
|
```
|
||||||
|
|
||||||
|
In this mode, FastMCP still injects and extracts trace context through MCP `_meta`, but it stops creating its own MCP request spans and `delegate` spans.
|
||||||
|
|
||||||
|
Library authors can suppress only FastMCP's native spans programmatically without disabling unrelated nested instrumentation:
|
||||||
|
|
||||||
|
```python
|
||||||
|
from fastmcp.telemetry import suppress_fastmcp_telemetry
|
||||||
|
|
||||||
|
with suppress_fastmcp_telemetry():
|
||||||
|
result = await client.call_tool("search", {"query": "otel"})
|
||||||
|
```
|
||||||
|
|
||||||
|
FastMCP server spans follow the MCP semantic conventions for mixed transport environments: when propagated MCP trace context is present in `_meta`, FastMCP uses that extracted context as the server-span parent and records any ambient transport span (for example an HTTP request span) as a span link instead.
|
||||||
|
|
||||||
### Local Development
|
### Local Development
|
||||||
|
|
||||||
For quick local trace visualization, [otel-desktop-viewer](https://github.com/CtrlSpice/otel-desktop-viewer) is a lightweight single-binary tool:
|
For quick local trace visualization, [otel-desktop-viewer](https://github.com/CtrlSpice/otel-desktop-viewer) is a lightweight single-binary tool:
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,12 @@ import httpx
|
||||||
import mcp.types
|
import mcp.types
|
||||||
from exceptiongroup import catch
|
from exceptiongroup import catch
|
||||||
from mcp import ClientSession, McpError
|
from mcp import ClientSession, McpError
|
||||||
|
from mcp.client.session import (
|
||||||
|
SUPPORTED_PROTOCOL_VERSIONS,
|
||||||
|
_default_elicitation_callback,
|
||||||
|
_default_list_roots_callback,
|
||||||
|
_default_sampling_callback,
|
||||||
|
)
|
||||||
from mcp.types import GetTaskResult, TaskStatusNotification
|
from mcp.types import GetTaskResult, TaskStatusNotification
|
||||||
from pydantic import AnyUrl
|
from pydantic import AnyUrl
|
||||||
|
|
||||||
|
|
@ -51,8 +57,10 @@ from fastmcp.client.tasks import (
|
||||||
TaskNotificationHandler,
|
TaskNotificationHandler,
|
||||||
ToolTask,
|
ToolTask,
|
||||||
)
|
)
|
||||||
|
from fastmcp.client.telemetry import client_span
|
||||||
from fastmcp.mcp_config import MCPConfig
|
from fastmcp.mcp_config import MCPConfig
|
||||||
from fastmcp.server import FastMCP
|
from fastmcp.server import FastMCP
|
||||||
|
from fastmcp.telemetry import inject_trace_context
|
||||||
from fastmcp.utilities.exceptions import get_catch_handlers
|
from fastmcp.utilities.exceptions import get_catch_handlers
|
||||||
from fastmcp.utilities.logging import get_logger
|
from fastmcp.utilities.logging import get_logger
|
||||||
from fastmcp.utilities.timeout import (
|
from fastmcp.utilities.timeout import (
|
||||||
|
|
@ -518,12 +526,93 @@ class Client(
|
||||||
timeout = normalize_timeout_to_seconds(timeout)
|
timeout = normalize_timeout_to_seconds(timeout)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
with anyio.fail_after(timeout):
|
with client_span(
|
||||||
self._session_state.initialize_result = await self.session.initialize()
|
"initialize",
|
||||||
return self._session_state.initialize_result
|
"initialize",
|
||||||
|
"",
|
||||||
|
session_id=self.transport.get_session_id(),
|
||||||
|
):
|
||||||
|
with anyio.fail_after(timeout):
|
||||||
|
propagated_meta = inject_trace_context()
|
||||||
|
if propagated_meta is None:
|
||||||
|
self._session_state.initialize_result = (
|
||||||
|
await self.session.initialize()
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
self._session_state.initialize_result = (
|
||||||
|
await self._initialize_session_with_meta(propagated_meta)
|
||||||
|
)
|
||||||
|
return self._session_state.initialize_result
|
||||||
except TimeoutError as e:
|
except TimeoutError as e:
|
||||||
raise RuntimeError("Failed to initialize server session") from e
|
raise RuntimeError("Failed to initialize server session") from e
|
||||||
|
|
||||||
|
async def _initialize_session_with_meta(
|
||||||
|
self,
|
||||||
|
meta: dict[str, Any],
|
||||||
|
) -> mcp.types.InitializeResult:
|
||||||
|
"""Send an InitializeRequest that preserves the client's capability config.
|
||||||
|
|
||||||
|
This method accesses private session attributes (``_sampling_capabilities``,
|
||||||
|
``_elicitation_callback``, etc.) to reconstruct an InitializeRequest that
|
||||||
|
preserves the client's original capability configuration. This is necessary
|
||||||
|
because the MCP SDK requires explicit capability objects at initialization
|
||||||
|
time, and there is no public API to retrieve the currently configured
|
||||||
|
capabilities from a session. This fragility is a known limitation of the
|
||||||
|
underlying MCP client SDK; any API surface to expose these would need to
|
||||||
|
be proposed to the MCP spec/SDK project.
|
||||||
|
"""
|
||||||
|
session = self.session
|
||||||
|
|
||||||
|
sampling = (
|
||||||
|
(session._sampling_capabilities or mcp.types.SamplingCapability())
|
||||||
|
if session._sampling_callback is not _default_sampling_callback
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
elicitation = (
|
||||||
|
mcp.types.ElicitationCapability(
|
||||||
|
form=mcp.types.FormElicitationCapability(),
|
||||||
|
url=mcp.types.UrlElicitationCapability(),
|
||||||
|
)
|
||||||
|
if session._elicitation_callback is not _default_elicitation_callback
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
roots = (
|
||||||
|
mcp.types.RootsCapability(listChanged=True)
|
||||||
|
if session._list_roots_callback is not _default_list_roots_callback
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
|
||||||
|
result = await session.send_request(
|
||||||
|
mcp.types.ClientRequest(
|
||||||
|
mcp.types.InitializeRequest(
|
||||||
|
params=mcp.types.InitializeRequestParams(
|
||||||
|
protocolVersion=mcp.types.LATEST_PROTOCOL_VERSION,
|
||||||
|
capabilities=mcp.types.ClientCapabilities(
|
||||||
|
sampling=sampling,
|
||||||
|
elicitation=elicitation,
|
||||||
|
experimental=None,
|
||||||
|
roots=roots,
|
||||||
|
tasks=session._task_handlers.build_capability(),
|
||||||
|
),
|
||||||
|
clientInfo=session._client_info,
|
||||||
|
_meta=meta, # type: ignore[unknown-argument] # pydantic alias # ty:ignore[unknown-argument]
|
||||||
|
)
|
||||||
|
)
|
||||||
|
),
|
||||||
|
mcp.types.InitializeResult,
|
||||||
|
)
|
||||||
|
|
||||||
|
if result.protocolVersion not in SUPPORTED_PROTOCOL_VERSIONS:
|
||||||
|
raise RuntimeError(
|
||||||
|
f"Unsupported protocol version from the server: {result.protocolVersion}"
|
||||||
|
)
|
||||||
|
|
||||||
|
session._server_capabilities = result.capabilities
|
||||||
|
await session.send_notification(
|
||||||
|
mcp.types.ClientNotification(mcp.types.InitializedNotification())
|
||||||
|
)
|
||||||
|
return result
|
||||||
|
|
||||||
async def __aenter__(self):
|
async def __aenter__(self):
|
||||||
return await self._connect()
|
return await self._connect()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -57,8 +57,16 @@ class ClientPromptsMixin:
|
||||||
):
|
):
|
||||||
logger.debug(f"[{self.name}] called list_prompts")
|
logger.debug(f"[{self.name}] called list_prompts")
|
||||||
|
|
||||||
|
propagated_meta = inject_trace_context()
|
||||||
|
params = None
|
||||||
|
if cursor is not None or propagated_meta is not None:
|
||||||
|
params = mcp.types.PaginatedRequestParams(
|
||||||
|
cursor=cursor,
|
||||||
|
_meta=propagated_meta, # type: ignore[unknown-argument] # pydantic alias # ty:ignore[unknown-argument]
|
||||||
|
)
|
||||||
|
|
||||||
result = await self._await_with_session_monitoring(
|
result = await self._await_with_session_monitoring(
|
||||||
self.session.list_prompts(cursor=cursor)
|
self.session.list_prompts(params=params)
|
||||||
)
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -56,8 +56,16 @@ class ClientResourcesMixin:
|
||||||
):
|
):
|
||||||
logger.debug(f"[{self.name}] called list_resources")
|
logger.debug(f"[{self.name}] called list_resources")
|
||||||
|
|
||||||
|
propagated_meta = inject_trace_context()
|
||||||
|
params = None
|
||||||
|
if cursor is not None or propagated_meta is not None:
|
||||||
|
params = mcp.types.PaginatedRequestParams(
|
||||||
|
cursor=cursor,
|
||||||
|
_meta=propagated_meta, # type: ignore[unknown-argument] # pydantic alias # ty:ignore[unknown-argument]
|
||||||
|
)
|
||||||
|
|
||||||
result = await self._await_with_session_monitoring(
|
result = await self._await_with_session_monitoring(
|
||||||
self.session.list_resources(cursor=cursor)
|
self.session.list_resources(params=params)
|
||||||
)
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
@ -132,8 +140,16 @@ class ClientResourcesMixin:
|
||||||
):
|
):
|
||||||
logger.debug(f"[{self.name}] called list_resource_templates")
|
logger.debug(f"[{self.name}] called list_resource_templates")
|
||||||
|
|
||||||
|
propagated_meta = inject_trace_context()
|
||||||
|
params = None
|
||||||
|
if cursor is not None or propagated_meta is not None:
|
||||||
|
params = mcp.types.PaginatedRequestParams(
|
||||||
|
cursor=cursor,
|
||||||
|
_meta=propagated_meta, # type: ignore[unknown-argument] # pydantic alias # ty:ignore[unknown-argument]
|
||||||
|
)
|
||||||
|
|
||||||
result = await self._await_with_session_monitoring(
|
result = await self._await_with_session_monitoring(
|
||||||
self.session.list_resource_templates(cursor=cursor)
|
self.session.list_resource_templates(params=params)
|
||||||
)
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -61,8 +61,16 @@ class ClientToolsMixin:
|
||||||
):
|
):
|
||||||
logger.debug(f"[{self.name}] called list_tools")
|
logger.debug(f"[{self.name}] called list_tools")
|
||||||
|
|
||||||
|
propagated_meta = inject_trace_context()
|
||||||
|
params = None
|
||||||
|
if cursor is not None or propagated_meta is not None:
|
||||||
|
params = mcp.types.PaginatedRequestParams(
|
||||||
|
cursor=cursor,
|
||||||
|
_meta=propagated_meta, # type: ignore[unknown-argument] # pydantic alias # ty:ignore[unknown-argument]
|
||||||
|
)
|
||||||
|
|
||||||
result = await self._await_with_session_monitoring(
|
result = await self._await_with_session_monitoring(
|
||||||
self.session.list_tools(cursor=cursor)
|
self.session.list_tools(params=params)
|
||||||
)
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@ from contextlib import contextmanager
|
||||||
from opentelemetry.trace import Span, SpanKind, Status, StatusCode
|
from opentelemetry.trace import Span, SpanKind, Status, StatusCode
|
||||||
|
|
||||||
from fastmcp.exceptions import ToolError as _ToolError
|
from fastmcp.exceptions import ToolError as _ToolError
|
||||||
from fastmcp.telemetry import get_tracer
|
from fastmcp.telemetry import get_noop_span, get_tracer, native_telemetry_enabled
|
||||||
|
|
||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
|
|
@ -23,6 +23,10 @@ def client_span(
|
||||||
|
|
||||||
Automatically records any exception on the span and sets error status.
|
Automatically records any exception on the span and sets error status.
|
||||||
"""
|
"""
|
||||||
|
if not native_telemetry_enabled():
|
||||||
|
yield get_noop_span()
|
||||||
|
return
|
||||||
|
|
||||||
tracer = get_tracer()
|
tracer = get_tracer()
|
||||||
with tracer.start_as_current_span(name, kind=SpanKind.CLIENT) as span:
|
with tracer.start_as_current_span(name, kind=SpanKind.CLIENT) as span:
|
||||||
if span.is_recording():
|
if span.is_recording():
|
||||||
|
|
|
||||||
|
|
@ -12,6 +12,7 @@ from mcp.server.streamable_http import (
|
||||||
EventStore,
|
EventStore,
|
||||||
)
|
)
|
||||||
from mcp.server.streamable_http_manager import StreamableHTTPSessionManager
|
from mcp.server.streamable_http_manager import StreamableHTTPSessionManager
|
||||||
|
from opentelemetry import trace
|
||||||
from starlette.applications import Starlette
|
from starlette.applications import Starlette
|
||||||
from starlette.middleware import Middleware
|
from starlette.middleware import Middleware
|
||||||
from starlette.requests import Request
|
from starlette.requests import Request
|
||||||
|
|
@ -28,6 +29,8 @@ if TYPE_CHECKING:
|
||||||
|
|
||||||
logger = get_logger(__name__)
|
logger = get_logger(__name__)
|
||||||
|
|
||||||
|
AMBIENT_SPAN_CONTEXT_SCOPE_KEY = "fastmcp.ambient_span_context"
|
||||||
|
|
||||||
|
|
||||||
class StreamableHTTPASGIApp:
|
class StreamableHTTPASGIApp:
|
||||||
"""ASGI application wrapper for Streamable HTTP server transport."""
|
"""ASGI application wrapper for Streamable HTTP server transport."""
|
||||||
|
|
@ -41,6 +44,11 @@ class StreamableHTTPASGIApp:
|
||||||
raise RuntimeError(
|
raise RuntimeError(
|
||||||
"Task group is not initialized. Make sure to use run()."
|
"Task group is not initialized. Make sure to use run()."
|
||||||
)
|
)
|
||||||
|
ambient_span_context = trace.get_current_span().get_span_context()
|
||||||
|
if ambient_span_context.is_valid:
|
||||||
|
scope[AMBIENT_SPAN_CONTEXT_SCOPE_KEY] = ambient_span_context
|
||||||
|
else:
|
||||||
|
scope.pop(AMBIENT_SPAN_CONTEXT_SCOPE_KEY, None)
|
||||||
await self.session_manager.handle_request(scope, receive, send)
|
await self.session_manager.handle_request(scope, receive, send)
|
||||||
except RuntimeError as e:
|
except RuntimeError as e:
|
||||||
if str(e) == "Task group is not initialized. Make sure to use run().":
|
if str(e) == "Task group is not initialized. Make sure to use run().":
|
||||||
|
|
|
||||||
|
|
@ -4,11 +4,27 @@ from collections.abc import Generator
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
|
|
||||||
from mcp.server.lowlevel.server import request_ctx
|
from mcp.server.lowlevel.server import request_ctx
|
||||||
|
from opentelemetry import context as otel_context
|
||||||
|
from opentelemetry import trace
|
||||||
from opentelemetry.context import Context
|
from opentelemetry.context import Context
|
||||||
from opentelemetry.trace import Span, SpanKind, Status, StatusCode
|
from opentelemetry.trace import (
|
||||||
|
Link,
|
||||||
|
Span,
|
||||||
|
SpanContext,
|
||||||
|
SpanKind,
|
||||||
|
Status,
|
||||||
|
StatusCode,
|
||||||
|
)
|
||||||
|
|
||||||
from fastmcp.exceptions import ToolError as _ToolError
|
from fastmcp.exceptions import ToolError as _ToolError
|
||||||
from fastmcp.telemetry import extract_trace_context, get_tracer
|
from fastmcp.server.http import AMBIENT_SPAN_CONTEXT_SCOPE_KEY
|
||||||
|
from fastmcp.telemetry import (
|
||||||
|
extract_propagation_keys_from_meta,
|
||||||
|
extract_trace_context,
|
||||||
|
get_noop_span,
|
||||||
|
get_tracer,
|
||||||
|
native_telemetry_enabled,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def get_auth_span_attributes() -> dict[str, str]:
|
def get_auth_span_attributes() -> dict[str, str]:
|
||||||
|
|
@ -42,15 +58,56 @@ def get_session_span_attributes() -> dict[str, str]:
|
||||||
return attrs
|
return attrs
|
||||||
|
|
||||||
|
|
||||||
def _get_parent_trace_context() -> Context | None:
|
def _get_parent_trace_context() -> tuple[Context | None, list[Link] | None]:
|
||||||
"""Get parent trace context from request meta for distributed tracing."""
|
"""Resolve MCP server parent context plus any ambient transport links."""
|
||||||
|
ambient_span_context = _get_ambient_or_current_span_context()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
req_ctx = request_ctx.get()
|
req_ctx = request_ctx.get()
|
||||||
if req_ctx and hasattr(req_ctx, "meta") and req_ctx.meta:
|
if req_ctx and hasattr(req_ctx, "meta") and req_ctx.meta:
|
||||||
return extract_trace_context(dict(req_ctx.meta))
|
meta = dict(req_ctx.meta)
|
||||||
|
if extract_propagation_keys_from_meta(meta):
|
||||||
|
parent_context = extract_trace_context(meta)
|
||||||
|
parent_span_context = trace.get_current_span(
|
||||||
|
parent_context
|
||||||
|
).get_span_context()
|
||||||
|
if (
|
||||||
|
ambient_span_context.is_valid
|
||||||
|
and parent_span_context != ambient_span_context
|
||||||
|
):
|
||||||
|
return parent_context, [Link(ambient_span_context)]
|
||||||
|
return parent_context, None
|
||||||
except LookupError:
|
except LookupError:
|
||||||
pass
|
pass
|
||||||
return None
|
|
||||||
|
if ambient_span_context.is_valid:
|
||||||
|
return otel_context.get_current(), None
|
||||||
|
|
||||||
|
return None, None
|
||||||
|
|
||||||
|
|
||||||
|
def _get_ambient_or_current_span_context() -> SpanContext:
|
||||||
|
"""Resolve the ambient transport span, falling back to the current span.
|
||||||
|
|
||||||
|
Returns the span context stored in the request scope by an outer transport
|
||||||
|
middleware, if present and valid. Otherwise returns the current active span.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
req_ctx = request_ctx.get()
|
||||||
|
except LookupError:
|
||||||
|
req_ctx = None
|
||||||
|
|
||||||
|
if req_ctx is not None:
|
||||||
|
request = getattr(req_ctx, "request", None)
|
||||||
|
if request is not None:
|
||||||
|
ambient_span_context = request.scope.get(AMBIENT_SPAN_CONTEXT_SCOPE_KEY)
|
||||||
|
if (
|
||||||
|
isinstance(ambient_span_context, SpanContext)
|
||||||
|
and ambient_span_context.is_valid
|
||||||
|
):
|
||||||
|
return ambient_span_context
|
||||||
|
|
||||||
|
return trace.get_current_span().get_span_context()
|
||||||
|
|
||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
|
|
@ -68,12 +125,24 @@ def server_span(
|
||||||
|
|
||||||
Automatically records any exception on the span and sets error status.
|
Automatically records any exception on the span and sets error status.
|
||||||
"""
|
"""
|
||||||
|
if not native_telemetry_enabled():
|
||||||
|
yield get_noop_span()
|
||||||
|
return
|
||||||
|
|
||||||
|
parent_context, links = _get_parent_trace_context()
|
||||||
tracer = get_tracer()
|
tracer = get_tracer()
|
||||||
with tracer.start_as_current_span(
|
span = tracer.start_span(
|
||||||
name,
|
name,
|
||||||
context=_get_parent_trace_context(),
|
context=parent_context,
|
||||||
kind=SpanKind.SERVER,
|
kind=SpanKind.SERVER,
|
||||||
) as span:
|
links=links,
|
||||||
|
)
|
||||||
|
current_context = trace.set_span_in_context(
|
||||||
|
span,
|
||||||
|
parent_context if parent_context is not None else otel_context.get_current(),
|
||||||
|
)
|
||||||
|
token = otel_context.attach(current_context)
|
||||||
|
try:
|
||||||
if span.is_recording():
|
if span.is_recording():
|
||||||
attrs: dict[str, str] = {
|
attrs: dict[str, str] = {
|
||||||
# MCP semantic conventions
|
# MCP semantic conventions
|
||||||
|
|
@ -103,6 +172,9 @@ def server_span(
|
||||||
span.record_exception(e)
|
span.record_exception(e)
|
||||||
span.set_status(Status(StatusCode.ERROR, str(e)))
|
span.set_status(Status(StatusCode.ERROR, str(e)))
|
||||||
raise
|
raise
|
||||||
|
finally:
|
||||||
|
otel_context.detach(token)
|
||||||
|
span.end()
|
||||||
|
|
||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
|
|
@ -117,6 +189,10 @@ def delegate_span(
|
||||||
Used by FastMCPProvider when delegating to mounted servers.
|
Used by FastMCPProvider when delegating to mounted servers.
|
||||||
Automatically records any exception on the span and sets error status.
|
Automatically records any exception on the span and sets error status.
|
||||||
"""
|
"""
|
||||||
|
if not native_telemetry_enabled():
|
||||||
|
yield get_noop_span()
|
||||||
|
return
|
||||||
|
|
||||||
tracer = get_tracer()
|
tracer = get_tracer()
|
||||||
with tracer.start_as_current_span(f"delegate {name}") as span:
|
with tracer.start_as_current_span(f"delegate {name}") as span:
|
||||||
if span.is_recording():
|
if span.is_recording():
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ logger = get_logger(__name__)
|
||||||
ENV_FILE = os.getenv("FASTMCP_ENV_FILE", ".env")
|
ENV_FILE = os.getenv("FASTMCP_ENV_FILE", ".env")
|
||||||
|
|
||||||
LOG_LEVEL = Literal["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"]
|
LOG_LEVEL = Literal["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"]
|
||||||
|
TELEMETRY_MODE = Literal["native", "propagation_only"]
|
||||||
|
|
||||||
MCP_LOG_LEVEL = Literal[
|
MCP_LOG_LEVEL = Literal[
|
||||||
"debug", "info", "notice", "warning", "error", "critical", "alert", "emergency"
|
"debug", "info", "notice", "warning", "error", "critical", "alert", "emergency"
|
||||||
|
|
@ -239,6 +240,23 @@ class Settings(BaseSettings):
|
||||||
),
|
),
|
||||||
] = True
|
] = True
|
||||||
|
|
||||||
|
telemetry_mode: Annotated[
|
||||||
|
TELEMETRY_MODE,
|
||||||
|
Field(
|
||||||
|
description=inspect.cleandoc(
|
||||||
|
"""
|
||||||
|
Controls FastMCP's native OpenTelemetry span creation.
|
||||||
|
|
||||||
|
- ``native`` (default): FastMCP creates MCP spans and propagates
|
||||||
|
trace context in request ``_meta``.
|
||||||
|
- ``propagation_only``: FastMCP still injects/extracts trace
|
||||||
|
context, but does not create its own MCP spans. Use this when
|
||||||
|
another instrumentation layer owns the MCP span hierarchy.
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
),
|
||||||
|
] = "native"
|
||||||
|
|
||||||
client_init_timeout: Annotated[
|
client_init_timeout: Annotated[
|
||||||
float | None,
|
float | None,
|
||||||
Field(
|
Field(
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,14 @@
|
||||||
"""OpenTelemetry instrumentation for FastMCP.
|
"""OpenTelemetry instrumentation for FastMCP.
|
||||||
|
|
||||||
This module provides native OpenTelemetry integration for FastMCP servers and clients.
|
This module provides native OpenTelemetry integration for FastMCP servers and
|
||||||
It uses only the opentelemetry-api package, so telemetry is a no-op unless the user
|
clients. It uses only the opentelemetry-api package, so telemetry is a no-op
|
||||||
installs an OpenTelemetry SDK and configures exporters.
|
unless the user installs an OpenTelemetry SDK and configures exporters.
|
||||||
|
|
||||||
|
FastMCP always propagates OpenTelemetry context through MCP ``params._meta``.
|
||||||
|
Native FastMCP spans can be suppressed globally via
|
||||||
|
``FASTMCP_TELEMETRY_MODE=propagation_only`` or programmatically with
|
||||||
|
``suppress_fastmcp_telemetry()`` when another instrumentation layer owns the
|
||||||
|
MCP span hierarchy.
|
||||||
|
|
||||||
Example usage with SDK:
|
Example usage with SDK:
|
||||||
```python
|
```python
|
||||||
|
|
@ -21,18 +27,69 @@ Example usage with SDK:
|
||||||
```
|
```
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
from collections.abc import Generator
|
||||||
|
from contextlib import contextmanager
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from opentelemetry import context as otel_context
|
from opentelemetry import context as otel_context
|
||||||
from opentelemetry import propagate, trace
|
from opentelemetry import propagate
|
||||||
from opentelemetry.context import Context
|
from opentelemetry.context import Context
|
||||||
from opentelemetry.trace import Span, Status, StatusCode, Tracer
|
from opentelemetry.trace import INVALID_SPAN, Span, Status, StatusCode, Tracer
|
||||||
from opentelemetry.trace import get_tracer as otel_get_tracer
|
from opentelemetry.trace import get_tracer as otel_get_tracer
|
||||||
|
|
||||||
INSTRUMENTATION_NAME = "fastmcp"
|
INSTRUMENTATION_NAME = "fastmcp"
|
||||||
|
|
||||||
TRACE_PARENT_KEY = "traceparent"
|
TRACE_PARENT_KEY = "traceparent"
|
||||||
TRACE_STATE_KEY = "tracestate"
|
TRACE_STATE_KEY = "tracestate"
|
||||||
|
BAGGAGE_KEY = "baggage"
|
||||||
|
|
||||||
|
_SUPPRESS_FASTMCP_TELEMETRY_KEY = otel_context.create_key("fastmcp_suppress_telemetry")
|
||||||
|
|
||||||
|
|
||||||
|
def _get_fastmcp_telemetry_mode() -> str:
|
||||||
|
"""Read the current FastMCP telemetry mode from settings."""
|
||||||
|
import fastmcp
|
||||||
|
|
||||||
|
return fastmcp.settings.telemetry_mode
|
||||||
|
|
||||||
|
|
||||||
|
def native_telemetry_enabled() -> bool:
|
||||||
|
"""Return whether FastMCP should create native MCP spans."""
|
||||||
|
return _get_fastmcp_telemetry_mode() == "native" and not otel_context.get_value(
|
||||||
|
_SUPPRESS_FASTMCP_TELEMETRY_KEY
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@contextmanager
|
||||||
|
def suppress_fastmcp_telemetry() -> Generator[None, None, None]:
|
||||||
|
"""Suppress native FastMCP spans while preserving context propagation.
|
||||||
|
|
||||||
|
This is narrower than OpenTelemetry's global instrumentation suppression:
|
||||||
|
it disables only FastMCP's own spans, allowing unrelated nested
|
||||||
|
instrumentations (HTTP clients, databases, etc.) to continue emitting.
|
||||||
|
"""
|
||||||
|
token = otel_context.attach(
|
||||||
|
otel_context.set_value(_SUPPRESS_FASTMCP_TELEMETRY_KEY, True)
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
yield
|
||||||
|
finally:
|
||||||
|
otel_context.detach(token)
|
||||||
|
|
||||||
|
|
||||||
|
def extract_propagation_keys_from_meta(meta: dict[str, Any] | None) -> dict[str, str]:
|
||||||
|
"""Extract trace-related propagation keys from an MCP ``_meta`` dict."""
|
||||||
|
if not meta:
|
||||||
|
return {}
|
||||||
|
|
||||||
|
carrier: dict[str, str] = {}
|
||||||
|
if TRACE_PARENT_KEY in meta:
|
||||||
|
carrier[TRACE_PARENT_KEY] = str(meta[TRACE_PARENT_KEY])
|
||||||
|
if TRACE_STATE_KEY in meta:
|
||||||
|
carrier[TRACE_STATE_KEY] = str(meta[TRACE_STATE_KEY])
|
||||||
|
if BAGGAGE_KEY in meta:
|
||||||
|
carrier[BAGGAGE_KEY] = str(meta[BAGGAGE_KEY])
|
||||||
|
return carrier
|
||||||
|
|
||||||
|
|
||||||
def get_tracer(version: str | None = None) -> Tracer:
|
def get_tracer(version: str | None = None) -> Tracer:
|
||||||
|
|
@ -63,10 +120,12 @@ def inject_trace_context(
|
||||||
propagate.inject(carrier)
|
propagate.inject(carrier)
|
||||||
|
|
||||||
trace_meta: dict[str, Any] = {}
|
trace_meta: dict[str, Any] = {}
|
||||||
if "traceparent" in carrier:
|
if TRACE_PARENT_KEY in carrier:
|
||||||
trace_meta[TRACE_PARENT_KEY] = carrier["traceparent"]
|
trace_meta[TRACE_PARENT_KEY] = carrier[TRACE_PARENT_KEY]
|
||||||
if "tracestate" in carrier:
|
if TRACE_STATE_KEY in carrier:
|
||||||
trace_meta[TRACE_STATE_KEY] = carrier["tracestate"]
|
trace_meta[TRACE_STATE_KEY] = carrier[TRACE_STATE_KEY]
|
||||||
|
if BAGGAGE_KEY in carrier:
|
||||||
|
trace_meta[BAGGAGE_KEY] = carrier[BAGGAGE_KEY]
|
||||||
|
|
||||||
if trace_meta:
|
if trace_meta:
|
||||||
return {**(meta or {}), **trace_meta}
|
return {**(meta or {}), **trace_meta}
|
||||||
|
|
@ -82,41 +141,36 @@ def record_span_error(span: Span, exception: BaseException) -> None:
|
||||||
def extract_trace_context(meta: dict[str, Any] | None) -> Context:
|
def extract_trace_context(meta: dict[str, Any] | None) -> Context:
|
||||||
"""Extract trace context from an MCP request meta dict.
|
"""Extract trace context from an MCP request meta dict.
|
||||||
|
|
||||||
If already in a valid trace (e.g., from HTTP propagation), the existing
|
|
||||||
trace context is preserved and meta is not used.
|
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
meta: The meta dict from an MCP request (ctx.request_context.meta)
|
meta: The meta dict from an MCP request (ctx.request_context.meta)
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
An OpenTelemetry Context with the extracted trace context,
|
An OpenTelemetry Context with propagated trace context and baggage
|
||||||
or the current context if no trace context found or already in a trace
|
propagated onto the current context, or the current context if no
|
||||||
|
propagation keys were present.
|
||||||
"""
|
"""
|
||||||
# Don't override existing trace context (e.g., from HTTP propagation)
|
carrier = extract_propagation_keys_from_meta(meta)
|
||||||
current_span = trace.get_current_span()
|
|
||||||
if current_span.get_span_context().is_valid:
|
|
||||||
return otel_context.get_current()
|
|
||||||
|
|
||||||
if not meta:
|
|
||||||
return otel_context.get_current()
|
|
||||||
|
|
||||||
carrier: dict[str, str] = {}
|
|
||||||
if TRACE_PARENT_KEY in meta:
|
|
||||||
carrier["traceparent"] = str(meta[TRACE_PARENT_KEY])
|
|
||||||
if TRACE_STATE_KEY in meta:
|
|
||||||
carrier["tracestate"] = str(meta[TRACE_STATE_KEY])
|
|
||||||
|
|
||||||
if carrier:
|
if carrier:
|
||||||
return propagate.extract(carrier)
|
return propagate.extract(carrier, context=otel_context.get_current())
|
||||||
return otel_context.get_current()
|
return otel_context.get_current()
|
||||||
|
|
||||||
|
|
||||||
|
def get_noop_span() -> Span:
|
||||||
|
"""Return the no-op span used when native FastMCP telemetry is suppressed."""
|
||||||
|
return INVALID_SPAN
|
||||||
|
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
|
"BAGGAGE_KEY",
|
||||||
"INSTRUMENTATION_NAME",
|
"INSTRUMENTATION_NAME",
|
||||||
"TRACE_PARENT_KEY",
|
"TRACE_PARENT_KEY",
|
||||||
"TRACE_STATE_KEY",
|
"TRACE_STATE_KEY",
|
||||||
|
"extract_propagation_keys_from_meta",
|
||||||
"extract_trace_context",
|
"extract_trace_context",
|
||||||
|
"get_noop_span",
|
||||||
"get_tracer",
|
"get_tracer",
|
||||||
"inject_trace_context",
|
"inject_trace_context",
|
||||||
|
"native_telemetry_enabled",
|
||||||
"record_span_error",
|
"record_span_error",
|
||||||
|
"suppress_fastmcp_telemetry",
|
||||||
]
|
]
|
||||||
|
|
|
||||||
180
tests/client/telemetry/test_client_propagation.py
Normal file
180
tests/client/telemetry/test_client_propagation.py
Normal file
|
|
@ -0,0 +1,180 @@
|
||||||
|
"""Tests for client trace-context propagation on initialize and list methods."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections.abc import Awaitable, Callable
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import mcp.types as mt
|
||||||
|
import pytest
|
||||||
|
from opentelemetry import baggage, trace
|
||||||
|
from opentelemetry import context as otel_context
|
||||||
|
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
|
||||||
|
|
||||||
|
from fastmcp import Client, FastMCP
|
||||||
|
|
||||||
|
|
||||||
|
def build_test_server() -> FastMCP:
|
||||||
|
"""Create a server with one component of each type."""
|
||||||
|
server = FastMCP("test-server")
|
||||||
|
|
||||||
|
@server.tool()
|
||||||
|
def greet(name: str = "world") -> str:
|
||||||
|
return f"Hello, {name}!"
|
||||||
|
|
||||||
|
@server.resource("data://config")
|
||||||
|
def get_config() -> str:
|
||||||
|
return "config"
|
||||||
|
|
||||||
|
@server.resource("users://{user_id}/profile")
|
||||||
|
def get_profile(user_id: str) -> str:
|
||||||
|
return f"profile {user_id}"
|
||||||
|
|
||||||
|
@server.prompt()
|
||||||
|
def greeting() -> str:
|
||||||
|
return "Hello!"
|
||||||
|
|
||||||
|
return server
|
||||||
|
|
||||||
|
|
||||||
|
def get_meta_dict(request: Any) -> dict[str, Any] | None:
|
||||||
|
"""Normalize ``request.params.meta`` to a plain dict for assertions."""
|
||||||
|
params = getattr(request, "params", None)
|
||||||
|
meta = getattr(params, "meta", None)
|
||||||
|
if meta is None:
|
||||||
|
return None
|
||||||
|
return meta.model_dump(exclude_none=True)
|
||||||
|
|
||||||
|
|
||||||
|
class TestClientPropagation:
|
||||||
|
async def test_initialize_creates_span_and_propagates_meta(
|
||||||
|
self,
|
||||||
|
monkeypatch,
|
||||||
|
trace_exporter: InMemorySpanExporter,
|
||||||
|
):
|
||||||
|
server = build_test_server()
|
||||||
|
client = Client(server, auto_initialize=False)
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
|
||||||
|
async with client:
|
||||||
|
captured_requests: list[Any] = []
|
||||||
|
original_send_request = client.session.send_request
|
||||||
|
|
||||||
|
async def wrapped_send_request(request: Any, result_type: Any) -> Any:
|
||||||
|
captured_requests.append(request.root)
|
||||||
|
return await original_send_request(request, result_type)
|
||||||
|
|
||||||
|
monkeypatch.setattr(client.session, "send_request", wrapped_send_request)
|
||||||
|
|
||||||
|
baggage_token = otel_context.attach(baggage.set_baggage("tenant", "acme"))
|
||||||
|
try:
|
||||||
|
with tracer.start_as_current_span("external-parent"):
|
||||||
|
await client.initialize()
|
||||||
|
finally:
|
||||||
|
otel_context.detach(baggage_token)
|
||||||
|
|
||||||
|
spans = trace_exporter.get_finished_spans()
|
||||||
|
initialize_span = next(
|
||||||
|
(
|
||||||
|
span
|
||||||
|
for span in spans
|
||||||
|
if span.name == "initialize"
|
||||||
|
and span.attributes is not None
|
||||||
|
and "fastmcp.server.name" not in span.attributes
|
||||||
|
),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
assert initialize_span is not None
|
||||||
|
assert initialize_span.attributes is not None
|
||||||
|
assert initialize_span.attributes["mcp.method.name"] == "initialize"
|
||||||
|
|
||||||
|
initialize_request = next(
|
||||||
|
request
|
||||||
|
for request in captured_requests
|
||||||
|
if isinstance(request, mt.InitializeRequest)
|
||||||
|
)
|
||||||
|
captured_meta = get_meta_dict(initialize_request)
|
||||||
|
assert captured_meta is not None
|
||||||
|
assert captured_meta["traceparent"].split("-")[2] == format(
|
||||||
|
initialize_span.get_span_context().span_id,
|
||||||
|
"016x",
|
||||||
|
)
|
||||||
|
assert "tenant=acme" in str(captured_meta["baggage"])
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("request_type", "operation", "expected_span_name"),
|
||||||
|
[
|
||||||
|
(mt.ListToolsRequest, lambda client: client.list_tools(), "tools/list"),
|
||||||
|
(
|
||||||
|
mt.ListResourcesRequest,
|
||||||
|
lambda client: client.list_resources(),
|
||||||
|
"resources/list",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
mt.ListResourceTemplatesRequest,
|
||||||
|
lambda client: client.list_resource_templates(),
|
||||||
|
"resources/templates/list",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
mt.ListPromptsRequest,
|
||||||
|
lambda client: client.list_prompts(),
|
||||||
|
"prompts/list",
|
||||||
|
),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
async def test_list_methods_propagate_meta(
|
||||||
|
self,
|
||||||
|
request_type: type[Any],
|
||||||
|
operation: Callable[[Client], Awaitable[Any]],
|
||||||
|
expected_span_name: str,
|
||||||
|
monkeypatch,
|
||||||
|
trace_exporter: InMemorySpanExporter,
|
||||||
|
):
|
||||||
|
server = build_test_server()
|
||||||
|
client = Client(server, auto_initialize=False)
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
|
||||||
|
async with client:
|
||||||
|
await client.initialize()
|
||||||
|
captured_requests: list[Any] = []
|
||||||
|
original_send_request = client.session.send_request
|
||||||
|
|
||||||
|
async def wrapped_send_request(request: Any, result_type: Any) -> Any:
|
||||||
|
captured_requests.append(request.root)
|
||||||
|
return await original_send_request(request, result_type)
|
||||||
|
|
||||||
|
monkeypatch.setattr(client.session, "send_request", wrapped_send_request)
|
||||||
|
trace_exporter.clear()
|
||||||
|
|
||||||
|
baggage_token = otel_context.attach(baggage.set_baggage("tenant", "acme"))
|
||||||
|
try:
|
||||||
|
with tracer.start_as_current_span("external-parent"):
|
||||||
|
await operation(client)
|
||||||
|
finally:
|
||||||
|
otel_context.detach(baggage_token)
|
||||||
|
|
||||||
|
spans = trace_exporter.get_finished_spans()
|
||||||
|
client_span = next(
|
||||||
|
(
|
||||||
|
span
|
||||||
|
for span in spans
|
||||||
|
if span.name == expected_span_name
|
||||||
|
and span.attributes is not None
|
||||||
|
and "fastmcp.server.name" not in span.attributes
|
||||||
|
),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
assert client_span is not None
|
||||||
|
|
||||||
|
request = next(
|
||||||
|
captured_request
|
||||||
|
for captured_request in captured_requests
|
||||||
|
if isinstance(captured_request, request_type)
|
||||||
|
)
|
||||||
|
captured_meta = get_meta_dict(request)
|
||||||
|
assert captured_meta is not None
|
||||||
|
assert captured_meta["traceparent"].split("-")[2] == format(
|
||||||
|
client_span.get_span_context().span_id,
|
||||||
|
"016x",
|
||||||
|
)
|
||||||
|
assert "tenant=acme" in str(captured_meta["baggage"])
|
||||||
379
tests/telemetry/test_interop.py
Normal file
379
tests/telemetry/test_interop.py
Normal file
|
|
@ -0,0 +1,379 @@
|
||||||
|
"""Tests for FastMCP telemetry interoperability modes."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
from typing import Any, cast
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
from mcp.server.lowlevel.server import request_ctx
|
||||||
|
from opentelemetry import baggage, trace
|
||||||
|
from opentelemetry import context as otel_context
|
||||||
|
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
|
||||||
|
from starlette.middleware import Middleware
|
||||||
|
|
||||||
|
from fastmcp import Client, FastMCP
|
||||||
|
from fastmcp.client.telemetry import client_span
|
||||||
|
from fastmcp.server.telemetry import server_span
|
||||||
|
from fastmcp.telemetry import inject_trace_context, suppress_fastmcp_telemetry
|
||||||
|
from fastmcp.utilities.tests import temporary_settings
|
||||||
|
|
||||||
|
|
||||||
|
class DummyReqCtx:
|
||||||
|
"""Minimal request context for server telemetry tests."""
|
||||||
|
|
||||||
|
def __init__(self, meta: dict[str, str]):
|
||||||
|
self.meta = meta
|
||||||
|
self.request = None
|
||||||
|
|
||||||
|
|
||||||
|
class AmbientHTTPSpanMiddleware:
|
||||||
|
"""Minimal ASGI middleware that simulates outer HTTP instrumentation."""
|
||||||
|
|
||||||
|
def __init__(self, app):
|
||||||
|
self.app = app
|
||||||
|
|
||||||
|
async def __call__(self, scope, receive, send):
|
||||||
|
if scope["type"] != "http":
|
||||||
|
await self.app(scope, receive, send)
|
||||||
|
return
|
||||||
|
|
||||||
|
tracer = trace.get_tracer("ambient-http")
|
||||||
|
with tracer.start_as_current_span("ambient-http-request"):
|
||||||
|
await self.app(scope, receive, send)
|
||||||
|
|
||||||
|
|
||||||
|
def parse_sse_response(body: str) -> dict[str, Any]:
|
||||||
|
"""Extract the first SSE data payload from a streamable HTTP response."""
|
||||||
|
for line in body.splitlines():
|
||||||
|
if line.startswith("data: "):
|
||||||
|
return json.loads(line[6:])
|
||||||
|
raise AssertionError(f"Missing SSE data payload in response: {body!r}")
|
||||||
|
|
||||||
|
|
||||||
|
class TestClientInteropMode:
|
||||||
|
async def test_propagation_only_mode_preserves_outer_client_span(
|
||||||
|
self, trace_exporter: InMemorySpanExporter
|
||||||
|
):
|
||||||
|
with temporary_settings(telemetry_mode="propagation_only"):
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
with tracer.start_as_current_span("external-client-parent") as parent_span:
|
||||||
|
with client_span(
|
||||||
|
"tools/call weather",
|
||||||
|
"tools/call",
|
||||||
|
"weather",
|
||||||
|
tool_name="weather",
|
||||||
|
) as span:
|
||||||
|
meta = inject_trace_context()
|
||||||
|
assert not span.is_recording()
|
||||||
|
|
||||||
|
spans = trace_exporter.get_finished_spans()
|
||||||
|
assert [span.name for span in spans] == ["external-client-parent"]
|
||||||
|
assert meta is not None
|
||||||
|
assert meta["traceparent"].split("-")[2] == format(
|
||||||
|
parent_span.get_span_context().span_id, "016x"
|
||||||
|
)
|
||||||
|
|
||||||
|
async def test_context_manager_suppresses_only_fastmcp_spans(
|
||||||
|
self, trace_exporter: InMemorySpanExporter
|
||||||
|
):
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
with tracer.start_as_current_span("external-client-parent") as parent_span:
|
||||||
|
with suppress_fastmcp_telemetry():
|
||||||
|
with client_span(
|
||||||
|
"tools/call weather",
|
||||||
|
"tools/call",
|
||||||
|
"weather",
|
||||||
|
tool_name="weather",
|
||||||
|
) as span:
|
||||||
|
meta = inject_trace_context()
|
||||||
|
assert not span.is_recording()
|
||||||
|
|
||||||
|
spans = trace_exporter.get_finished_spans()
|
||||||
|
assert [span.name for span in spans] == ["external-client-parent"]
|
||||||
|
assert meta is not None
|
||||||
|
assert meta["traceparent"].split("-")[2] == format(
|
||||||
|
parent_span.get_span_context().span_id, "016x"
|
||||||
|
)
|
||||||
|
|
||||||
|
async def test_end_to_end_propagation_only_suppresses_native_spans(
|
||||||
|
self, trace_exporter: InMemorySpanExporter
|
||||||
|
):
|
||||||
|
child = FastMCP("child-server")
|
||||||
|
|
||||||
|
@child.tool()
|
||||||
|
def child_tool() -> str:
|
||||||
|
return "child result"
|
||||||
|
|
||||||
|
parent = FastMCP("parent-server")
|
||||||
|
parent.mount(child, namespace="child")
|
||||||
|
|
||||||
|
with temporary_settings(telemetry_mode="propagation_only"):
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
with tracer.start_as_current_span("external-request"):
|
||||||
|
client = Client(parent)
|
||||||
|
async with client:
|
||||||
|
result = await client.call_tool("child_child_tool", {})
|
||||||
|
assert "child result" in str(result)
|
||||||
|
|
||||||
|
spans = trace_exporter.get_finished_spans()
|
||||||
|
assert [span.name for span in spans] == ["external-request"]
|
||||||
|
|
||||||
|
|
||||||
|
class TestServerInteropMode:
|
||||||
|
async def test_streamable_http_third_party_client_uses_meta_parent_and_baggage(
|
||||||
|
self,
|
||||||
|
trace_exporter: InMemorySpanExporter,
|
||||||
|
):
|
||||||
|
child = FastMCP("child-server")
|
||||||
|
|
||||||
|
@child.tool()
|
||||||
|
def tenant() -> str:
|
||||||
|
tenant_name = baggage.get_baggage("tenant")
|
||||||
|
return tenant_name if isinstance(tenant_name, str) else "missing"
|
||||||
|
|
||||||
|
parent = FastMCP("parent-server")
|
||||||
|
parent.mount(child, namespace="child")
|
||||||
|
|
||||||
|
app = parent.http_app(
|
||||||
|
transport="http",
|
||||||
|
path="/mcp",
|
||||||
|
middleware=[Middleware(AmbientHTTPSpanMiddleware)],
|
||||||
|
)
|
||||||
|
headers = {
|
||||||
|
"accept": "application/json, text/event-stream",
|
||||||
|
"content-type": "application/json",
|
||||||
|
}
|
||||||
|
|
||||||
|
async with app.router.lifespan_context(app):
|
||||||
|
transport = httpx.ASGITransport(app=app)
|
||||||
|
async with httpx.AsyncClient(
|
||||||
|
transport=transport,
|
||||||
|
base_url="http://testserver",
|
||||||
|
) as client:
|
||||||
|
init_response = await client.post(
|
||||||
|
"/mcp",
|
||||||
|
headers=headers,
|
||||||
|
json={
|
||||||
|
"jsonrpc": "2.0",
|
||||||
|
"id": 1,
|
||||||
|
"method": "initialize",
|
||||||
|
"params": {
|
||||||
|
"protocolVersion": "2025-03-26",
|
||||||
|
"capabilities": {},
|
||||||
|
"clientInfo": {
|
||||||
|
"name": "opaque-client",
|
||||||
|
"version": "0.1.0",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
session_id = init_response.headers["mcp-session-id"]
|
||||||
|
await client.post(
|
||||||
|
"/mcp",
|
||||||
|
headers={**headers, "mcp-session-id": session_id},
|
||||||
|
json={
|
||||||
|
"jsonrpc": "2.0",
|
||||||
|
"method": "notifications/initialized",
|
||||||
|
"params": {},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
trace_exporter.clear()
|
||||||
|
tracer = trace.get_tracer("opaque-client")
|
||||||
|
baggage_token = otel_context.attach(
|
||||||
|
baggage.set_baggage("tenant", "acme")
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
with tracer.start_as_current_span(
|
||||||
|
"opaque-client-root"
|
||||||
|
) as client_span_export:
|
||||||
|
meta = inject_trace_context()
|
||||||
|
tool_response = await client.post(
|
||||||
|
"/mcp",
|
||||||
|
headers={**headers, "mcp-session-id": session_id},
|
||||||
|
json={
|
||||||
|
"jsonrpc": "2.0",
|
||||||
|
"id": 2,
|
||||||
|
"method": "tools/call",
|
||||||
|
"params": {
|
||||||
|
"name": "child_tenant",
|
||||||
|
"arguments": {},
|
||||||
|
"_meta": meta,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
otel_context.detach(baggage_token)
|
||||||
|
|
||||||
|
payload = parse_sse_response(tool_response.text)
|
||||||
|
assert payload["result"]["content"][0]["text"] == "acme"
|
||||||
|
|
||||||
|
spans = {
|
||||||
|
span.name: span
|
||||||
|
for span in trace_exporter.get_finished_spans()
|
||||||
|
if span.name
|
||||||
|
in {
|
||||||
|
"ambient-http-request",
|
||||||
|
"opaque-client-root",
|
||||||
|
"tools/call child_tenant",
|
||||||
|
"delegate tenant",
|
||||||
|
"tools/call tenant",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
parent_span_export = spans["tools/call child_tenant"]
|
||||||
|
delegate_span_export = spans["delegate tenant"]
|
||||||
|
child_span_export = spans["tools/call tenant"]
|
||||||
|
ambient_span_export = spans["ambient-http-request"]
|
||||||
|
|
||||||
|
assert parent_span_export.parent is not None
|
||||||
|
assert (
|
||||||
|
parent_span_export.parent.span_id
|
||||||
|
== client_span_export.get_span_context().span_id
|
||||||
|
)
|
||||||
|
assert any(
|
||||||
|
link.context.span_id == ambient_span_export.get_span_context().span_id
|
||||||
|
for link in parent_span_export.links
|
||||||
|
)
|
||||||
|
assert child_span_export.parent is not None
|
||||||
|
assert (
|
||||||
|
child_span_export.parent.span_id
|
||||||
|
== client_span_export.get_span_context().span_id
|
||||||
|
)
|
||||||
|
assert delegate_span_export.parent is not None
|
||||||
|
assert (
|
||||||
|
delegate_span_export.parent.span_id
|
||||||
|
== parent_span_export.get_span_context().span_id
|
||||||
|
)
|
||||||
|
assert any(
|
||||||
|
link.context.span_id == ambient_span_export.get_span_context().span_id
|
||||||
|
for link in child_span_export.links
|
||||||
|
)
|
||||||
|
|
||||||
|
async def test_server_span_makes_propagated_baggage_current(
|
||||||
|
self,
|
||||||
|
monkeypatch,
|
||||||
|
trace_exporter: InMemorySpanExporter,
|
||||||
|
):
|
||||||
|
import fastmcp.server.telemetry as server_telemetry
|
||||||
|
|
||||||
|
monkeypatch.setattr(server_telemetry, "get_auth_span_attributes", lambda: {})
|
||||||
|
monkeypatch.setattr(server_telemetry, "get_session_span_attributes", lambda: {})
|
||||||
|
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
baggage_token = otel_context.attach(baggage.set_baggage("tenant", "acme"))
|
||||||
|
try:
|
||||||
|
with tracer.start_as_current_span("external-client-parent"):
|
||||||
|
meta = inject_trace_context()
|
||||||
|
finally:
|
||||||
|
otel_context.detach(baggage_token)
|
||||||
|
|
||||||
|
req_token = request_ctx.set(cast(Any, DummyReqCtx(meta or {})))
|
||||||
|
try:
|
||||||
|
with server_span(
|
||||||
|
"tools/call weather",
|
||||||
|
"tools/call",
|
||||||
|
"test-server",
|
||||||
|
"tool",
|
||||||
|
"weather",
|
||||||
|
tool_name="weather",
|
||||||
|
):
|
||||||
|
assert baggage.get_baggage("tenant") == "acme"
|
||||||
|
finally:
|
||||||
|
request_ctx.reset(req_token)
|
||||||
|
|
||||||
|
async def test_server_span_with_baggage_only_meta_keeps_ambient_parent(
|
||||||
|
self,
|
||||||
|
monkeypatch,
|
||||||
|
trace_exporter: InMemorySpanExporter,
|
||||||
|
):
|
||||||
|
import fastmcp.server.telemetry as server_telemetry
|
||||||
|
|
||||||
|
monkeypatch.setattr(server_telemetry, "get_auth_span_attributes", lambda: {})
|
||||||
|
monkeypatch.setattr(server_telemetry, "get_session_span_attributes", lambda: {})
|
||||||
|
|
||||||
|
with trace.get_tracer("external").start_as_current_span(
|
||||||
|
"ambient-http-request"
|
||||||
|
) as ambient_span:
|
||||||
|
req_token = request_ctx.set(
|
||||||
|
cast(Any, DummyReqCtx({"baggage": "userId=alice"}))
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
with server_span(
|
||||||
|
"tools/call weather",
|
||||||
|
"tools/call",
|
||||||
|
"test-server",
|
||||||
|
"tool",
|
||||||
|
"weather",
|
||||||
|
tool_name="weather",
|
||||||
|
):
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
request_ctx.reset(req_token)
|
||||||
|
|
||||||
|
spans = {
|
||||||
|
span.name: span
|
||||||
|
for span in trace_exporter.get_finished_spans()
|
||||||
|
if span.name in {"ambient-http-request", "tools/call weather"}
|
||||||
|
}
|
||||||
|
server_span_export = spans["tools/call weather"]
|
||||||
|
|
||||||
|
assert server_span_export.parent is not None
|
||||||
|
assert (
|
||||||
|
server_span_export.parent.span_id == ambient_span.get_span_context().span_id
|
||||||
|
)
|
||||||
|
assert server_span_export.links == ()
|
||||||
|
|
||||||
|
async def test_server_span_uses_meta_parent_and_links_ambient_context(
|
||||||
|
self,
|
||||||
|
monkeypatch,
|
||||||
|
trace_exporter: InMemorySpanExporter,
|
||||||
|
):
|
||||||
|
import fastmcp.server.telemetry as server_telemetry
|
||||||
|
|
||||||
|
monkeypatch.setattr(server_telemetry, "get_auth_span_attributes", lambda: {})
|
||||||
|
monkeypatch.setattr(server_telemetry, "get_session_span_attributes", lambda: {})
|
||||||
|
|
||||||
|
tracer = trace.get_tracer("external")
|
||||||
|
|
||||||
|
remote_parent = tracer.start_span("external-client-parent")
|
||||||
|
token = otel_context.attach(trace.set_span_in_context(remote_parent))
|
||||||
|
try:
|
||||||
|
meta = inject_trace_context()
|
||||||
|
finally:
|
||||||
|
otel_context.detach(token)
|
||||||
|
|
||||||
|
with tracer.start_as_current_span("ambient-http-request") as ambient_span:
|
||||||
|
req_token = request_ctx.set(cast(Any, DummyReqCtx(meta or {})))
|
||||||
|
try:
|
||||||
|
with server_span(
|
||||||
|
"tools/call weather",
|
||||||
|
"tools/call",
|
||||||
|
"test-server",
|
||||||
|
"tool",
|
||||||
|
"weather",
|
||||||
|
tool_name="weather",
|
||||||
|
):
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
request_ctx.reset(req_token)
|
||||||
|
|
||||||
|
remote_parent.end()
|
||||||
|
|
||||||
|
spans = {
|
||||||
|
span.name: span
|
||||||
|
for span in trace_exporter.get_finished_spans()
|
||||||
|
if span.name
|
||||||
|
in {"external-client-parent", "ambient-http-request", "tools/call weather"}
|
||||||
|
}
|
||||||
|
server_span_export = spans["tools/call weather"]
|
||||||
|
|
||||||
|
assert server_span_export.parent is not None
|
||||||
|
assert (
|
||||||
|
server_span_export.parent.span_id
|
||||||
|
== remote_parent.get_span_context().span_id
|
||||||
|
)
|
||||||
|
assert any(
|
||||||
|
link.context.span_id == ambient_span.get_span_context().span_id
|
||||||
|
for link in server_span_export.links
|
||||||
|
)
|
||||||
|
|
@ -2,11 +2,13 @@
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
from opentelemetry import trace
|
from opentelemetry import baggage, trace
|
||||||
|
from opentelemetry import context as otel_context
|
||||||
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
|
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
|
||||||
|
|
||||||
from fastmcp.server.telemetry import get_auth_span_attributes
|
from fastmcp.server.telemetry import get_auth_span_attributes
|
||||||
from fastmcp.telemetry import (
|
from fastmcp.telemetry import (
|
||||||
|
BAGGAGE_KEY,
|
||||||
INSTRUMENTATION_NAME,
|
INSTRUMENTATION_NAME,
|
||||||
TRACE_PARENT_KEY,
|
TRACE_PARENT_KEY,
|
||||||
extract_trace_context,
|
extract_trace_context,
|
||||||
|
|
@ -50,6 +52,19 @@ class TestInjectTraceContext:
|
||||||
assert TRACE_PARENT_KEY in meta
|
assert TRACE_PARENT_KEY in meta
|
||||||
assert meta[TRACE_PARENT_KEY].startswith("00-")
|
assert meta[TRACE_PARENT_KEY].startswith("00-")
|
||||||
|
|
||||||
|
def test_injects_baggage(self, trace_exporter: InMemorySpanExporter):
|
||||||
|
tracer = get_tracer()
|
||||||
|
baggage_token = otel_context.attach(baggage.set_baggage("userId", "alice"))
|
||||||
|
try:
|
||||||
|
with tracer.start_as_current_span("test"):
|
||||||
|
meta = inject_trace_context()
|
||||||
|
finally:
|
||||||
|
otel_context.detach(baggage_token)
|
||||||
|
|
||||||
|
assert meta is not None
|
||||||
|
assert BAGGAGE_KEY in meta
|
||||||
|
assert "userId=alice" in str(meta[BAGGAGE_KEY])
|
||||||
|
|
||||||
|
|
||||||
class TestExtractTraceContext:
|
class TestExtractTraceContext:
|
||||||
def test_bare_traceparent(self, trace_exporter: InMemorySpanExporter):
|
def test_bare_traceparent(self, trace_exporter: InMemorySpanExporter):
|
||||||
|
|
@ -69,6 +84,38 @@ class TestExtractTraceContext:
|
||||||
assert span_ctx.is_valid
|
assert span_ctx.is_valid
|
||||||
assert span_ctx.trace_state.get("congo") == "t61rcWkgMzE"
|
assert span_ctx.trace_state.get("congo") == "t61rcWkgMzE"
|
||||||
|
|
||||||
|
def test_bare_baggage(self, trace_exporter: InMemorySpanExporter):
|
||||||
|
ctx = extract_trace_context(
|
||||||
|
{
|
||||||
|
"traceparent": VALID_TRACEPARENT,
|
||||||
|
"baggage": "userId=alice",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
assert baggage.get_baggage("userId", context=ctx) == "alice"
|
||||||
|
|
||||||
|
def test_prefers_propagated_context_over_current(
|
||||||
|
self, trace_exporter: InMemorySpanExporter
|
||||||
|
):
|
||||||
|
tracer = get_tracer()
|
||||||
|
with tracer.start_as_current_span("current"):
|
||||||
|
ctx = extract_trace_context({"traceparent": VALID_TRACEPARENT})
|
||||||
|
|
||||||
|
span_ctx = trace.get_current_span(ctx).get_span_context()
|
||||||
|
assert span_ctx.is_valid
|
||||||
|
assert format(span_ctx.trace_id, "032x") == "0af7651916cd43dd8448eb211c80319c"
|
||||||
|
|
||||||
|
def test_baggage_only_meta_preserves_current_span(
|
||||||
|
self, trace_exporter: InMemorySpanExporter
|
||||||
|
):
|
||||||
|
tracer = get_tracer()
|
||||||
|
with tracer.start_as_current_span("current") as current_span:
|
||||||
|
ctx = extract_trace_context({"baggage": "userId=alice"})
|
||||||
|
|
||||||
|
span_ctx = trace.get_current_span(ctx).get_span_context()
|
||||||
|
assert span_ctx.is_valid
|
||||||
|
assert span_ctx.span_id == current_span.get_span_context().span_id
|
||||||
|
assert baggage.get_baggage("userId", context=ctx) == "alice"
|
||||||
|
|
||||||
def test_none_meta_returns_current_context(
|
def test_none_meta_returns_current_context(
|
||||||
self, trace_exporter: InMemorySpanExporter
|
self, trace_exporter: InMemorySpanExporter
|
||||||
):
|
):
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue