-
Notifications
You must be signed in to change notification settings - Fork 11
Add OpenTelemetry instrument for FastMCP #152
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -4,12 +4,22 @@ | |
| import prometheus_client | ||
| import typing_extensions | ||
| from fastmcp import FastMCP | ||
| from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware | ||
| from opentelemetry.util.http import ExcludeList, get_excluded_urls | ||
| from starlette.applications import Starlette | ||
| from starlette.responses import JSONResponse, Response | ||
| from starlette.routing import Match, Mount, Route | ||
|
|
||
| from microbootstrap.bootstrappers.base import ApplicationBootstrapper | ||
| from microbootstrap.config.fastmcp import FastMcpConfig | ||
| from microbootstrap.instruments.health_checks_instrument import HealthChecksInstrument, HealthCheckTypedDict | ||
| from microbootstrap.instruments.logging_instrument import LoggingInstrument | ||
| from microbootstrap.instruments.opentelemetry_instrument import ( | ||
| BaseOpentelemetryInstrument, | ||
| CombinedExcludeList, | ||
| OpentelemetryConfig, | ||
| build_span_name, | ||
| ) | ||
| from microbootstrap.instruments.prometheus_instrument import FastMcpPrometheusConfig, PrometheusInstrument | ||
| from microbootstrap.instruments.pyroscope_instrument import PyroscopeInstrument | ||
| from microbootstrap.instruments.sentry_instrument import SentryInstrument | ||
|
|
@@ -18,16 +28,44 @@ | |
|
|
||
|
|
||
| if typing.TYPE_CHECKING: | ||
| from fastmcp.server.http import StarletteWithLifespan | ||
| from starlette.requests import Request | ||
| from starlette.types import Scope | ||
|
|
||
|
|
||
| StarletteT = typing.TypeVar("StarletteT", bound=Starlette) | ||
|
|
||
|
|
||
| class KwargsFastMCP(FastMCP[typing.Any]): | ||
| def __init__(self, **kwargs: typing.Any) -> None: # noqa: ANN401 | ||
| super().__init__(**kwargs) | ||
| self.http_app_hooks: list[typing.Callable[[StarletteWithLifespan], StarletteWithLifespan]] = [] | ||
|
|
||
| def add_http_app_hook(self, hook: typing.Callable[[StarletteWithLifespan], StarletteWithLifespan]) -> None: | ||
| self.http_app_hooks.append(hook) | ||
|
|
||
| def http_app(self, *args: typing.Any, **kwargs: typing.Any) -> StarletteWithLifespan: # noqa: ANN401 | ||
| # ASGI application is created by the user after bootstrap, so instruments subscribe to its creation | ||
| http_application = super().http_app(*args, **kwargs) | ||
| for hook in self.http_app_hooks: | ||
| http_application = hook(http_application) | ||
| return http_application | ||
|
|
||
|
|
||
| def build_fastmcp_route_details_from_scope( | ||
| scope: Scope, | ||
| routes: typing.Iterable[typing.Any], | ||
| ) -> tuple[str, dict[str, str]]: | ||
| method: typing.Final = str(scope.get("method", "HTTP")).strip() | ||
| for route in routes: | ||
| if isinstance(route, (Route, Mount)) and route.matches(scope)[0] == Match.FULL: | ||
| return build_span_name(method, route.path), {"http.route": route.path} | ||
| # Unmatched paths get no `http.route` to keep its cardinality low | ||
| return method, {} | ||
|
|
||
|
|
||
| class FastMcpBootstrapper( | ||
| ApplicationBootstrapper[FastMcpSettings, FastMCP[typing.Any], FastMcpConfig], | ||
| ApplicationBootstrapper[FastMcpSettings, KwargsFastMCP, FastMcpConfig], | ||
| ): | ||
| application_config = FastMcpConfig() | ||
| application_type = KwargsFastMCP | ||
|
|
@@ -41,8 +79,8 @@ def bootstrap_before(self: typing_extensions.Self) -> dict[str, typing.Any]: | |
|
|
||
| def bootstrap_before_instruments_after_app_created( | ||
| self, | ||
| application: FastMCP[typing.Any], | ||
| ) -> FastMCP[typing.Any]: | ||
| application: KwargsFastMCP, | ||
| ) -> KwargsFastMCP: | ||
| self.console_writer.print_bootstrap_table() | ||
| return application | ||
|
|
||
|
|
@@ -51,6 +89,38 @@ def bootstrap_before_instruments_after_app_created( | |
| FastMcpBootstrapper.use_instrument()(PyroscopeInstrument) | ||
|
|
||
|
|
||
| @FastMcpBootstrapper.use_instrument() | ||
| class FastMcpOpentelemetryInstrument(BaseOpentelemetryInstrument[OpentelemetryConfig]): | ||
| def bootstrap_after(self, application: FastMCP[typing.Any]) -> FastMCP[typing.Any]: # type: ignore[override] | ||
| if isinstance(application, KwargsFastMCP): | ||
| application.add_http_app_hook(self.instrument_http_app) | ||
| return application | ||
|
|
||
| def instrument_http_app(self, http_application: StarletteT) -> StarletteT: | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Давай этот метод хотя бы с __ в начале сделаем, тк он не должен вызываться за пределами самого класса |
||
| # `StarletteInstrumentor` marks applications the same way, so each application is instrumented once | ||
| if getattr(http_application, "_is_instrumented_by_opentelemetry", False): | ||
| return http_application | ||
|
|
||
| def build_route_details(scope: Scope) -> tuple[str, dict[str, str]]: | ||
| return build_fastmcp_route_details_from_scope(scope, http_application.routes) | ||
|
|
||
| http_application.add_middleware( | ||
| OpenTelemetryMiddleware, | ||
| tracer_provider=self.tracer_provider, | ||
| default_span_details=build_route_details, | ||
| excluded_urls=CombinedExcludeList( | ||
| ExcludeList(self.define_exclude_urls()), | ||
| get_excluded_urls("STARLETTE"), | ||
| ), | ||
| ) | ||
| http_application._is_instrumented_by_opentelemetry = True # type: ignore[attr-defined] # noqa: SLF001 | ||
| return http_application | ||
|
|
||
| @classmethod | ||
| def get_config_type(cls) -> type[OpentelemetryConfig]: | ||
| return OpentelemetryConfig | ||
|
|
||
|
|
||
| @FastMcpBootstrapper.use_instrument() | ||
| class FastMcpLoggingInstrument(LoggingInstrument): | ||
| def bootstrap_after(self, application: FastMCP[typing.Any]) -> FastMCP[typing.Any]: # type: ignore[override] | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -11,7 +11,8 @@ | |
| from faststream.asgi import AsgiFastStream, AsgiResponse | ||
| from faststream.asgi import get as handle_get | ||
| from faststream.specification import AsyncAPI | ||
| from opentelemetry import trace | ||
| from opentelemetry.instrumentation.asgi import OpenTelemetryMiddleware | ||
| from opentelemetry.util.http import ExcludeList | ||
|
|
||
| from microbootstrap.bootstrappers.base import ApplicationBootstrapper | ||
| from microbootstrap.config.faststream import FastStreamConfig | ||
|
|
@@ -20,6 +21,7 @@ | |
| from microbootstrap.instruments.opentelemetry_instrument import ( | ||
| BaseOpentelemetryInstrument, | ||
| FastStreamOpentelemetryConfig, | ||
| build_span_name, | ||
| ) | ||
| from microbootstrap.instruments.prometheus_instrument import FastStreamPrometheusConfig, PrometheusInstrument | ||
| from microbootstrap.instruments.pyroscope_instrument import PyroscopeInstrument | ||
|
|
@@ -28,7 +30,10 @@ | |
| from microbootstrap.settings import FastStreamSettings | ||
|
|
||
|
|
||
| tracer: typing.Final = trace.get_tracer(__name__) | ||
| if typing.TYPE_CHECKING: | ||
| from faststream.asgi.types import ASGIApp, Receive, Scope, Send | ||
|
|
||
|
|
||
| MessageT = typing.TypeVar("MessageT") | ||
| ResponseT = typing.TypeVar("ResponseT") | ||
|
|
||
|
|
@@ -58,9 +63,32 @@ class KwargsAsgiFastStream(AsgiFastStream): | |
| def __init__(self, **kwargs: typing.Any) -> None: # noqa: ANN401 | ||
| # `broker` argument is positional-only | ||
| super().__init__(kwargs.pop("broker", None), **kwargs) | ||
| self.http_app: ASGIApp = super().__call__ | ||
|
|
||
| def add_http_middleware(self, build_middleware: typing.Callable[[ASGIApp], ASGIApp]) -> None: | ||
| self.http_app = build_middleware(self.http_app) | ||
|
|
||
| async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: | ||
| # Lifespan and websocket scopes bypass HTTP middlewares | ||
| if scope["type"] == "http": | ||
| await self.http_app(scope, receive, send) | ||
| return | ||
| await super().__call__(scope, receive, send) | ||
|
Comment on lines
+66
to
+76
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. А вот тут какие-то дубли из смежного ПРа пошли |
||
|
|
||
|
|
||
| class FastStreamBootstrapper(ApplicationBootstrapper[FastStreamSettings, AsgiFastStream, FastStreamConfig]): | ||
| def build_faststream_route_details_from_scope( | ||
| scope: Scope, | ||
| routes: typing.Iterable[tuple[str, ASGIApp]], | ||
| ) -> tuple[str, dict[str, str]]: | ||
| method: typing.Final = str(scope.get("method", "HTTP")).strip() | ||
| path: typing.Final = scope.get("path") | ||
| # FastStream matches ASGI routes by exact path, unmatched paths get no `http.route` to keep its cardinality low | ||
| if path is None or all(path != route_path for route_path, _ in routes): | ||
| return method, {} | ||
| return build_span_name(method, path), {"http.route": path} | ||
|
Comment on lines
+79
to
+88
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Коммент тот же, давай в фастстриме это поддерживать |
||
|
|
||
|
|
||
| class FastStreamBootstrapper(ApplicationBootstrapper[FastStreamSettings, KwargsAsgiFastStream, FastStreamConfig]): | ||
| application_config = FastStreamConfig() | ||
| application_type = KwargsAsgiFastStream | ||
|
|
||
|
|
@@ -94,9 +122,6 @@ def bootstrap_after(self, application: AsgiFastStream) -> AsgiFastStream: # typ | |
|
|
||
| @FastStreamBootstrapper.use_instrument() | ||
| class FastStreamOpentelemetryInstrument(BaseOpentelemetryInstrument[FastStreamOpentelemetryConfig]): | ||
| def is_ready(self) -> bool: | ||
| return bool(self.instrument_config.opentelemetry_middleware_cls and super().is_ready()) | ||
|
|
||
| def bootstrap_after(self, application: AsgiFastStream) -> AsgiFastStream: # type: ignore[override] | ||
| if self.instrument_config.opentelemetry_middleware_cls and application.broker: | ||
| application.broker.add_middleware( | ||
|
|
@@ -107,8 +132,23 @@ def bootstrap_after(self, application: AsgiFastStream) -> AsgiFastStream: # typ | |
| baggage_span_attributes=self.instrument_config.opentelemetry_baggage_span_attributes, | ||
| ), | ||
| ) | ||
| if isinstance(application, KwargsAsgiFastStream): | ||
| application.add_http_middleware( | ||
| functools.partial(self.create_open_telemetry_middleware, application=application), | ||
| ) | ||
| return application | ||
|
|
||
| def create_open_telemetry_middleware(self, app: ASGIApp, application: AsgiFastStream) -> ASGIApp: | ||
| def build_route_details(scope: Scope) -> tuple[str, dict[str, str]]: | ||
| return build_faststream_route_details_from_scope(scope, application.routes) | ||
|
|
||
| return OpenTelemetryMiddleware( | ||
| app=app, | ||
| default_span_details=build_route_details, | ||
| excluded_urls=ExcludeList(self.define_exclude_urls()), | ||
| tracer_provider=self.tracer_provider, | ||
| ) | ||
|
|
||
| @classmethod | ||
| def get_config_type(cls) -> type[FastStreamOpentelemetryConfig]: | ||
| return FastStreamOpentelemetryConfig | ||
|
|
@@ -175,11 +215,6 @@ async def check_health(scope: typing.Any) -> AsgiResponse: # noqa: ANN401, ARG0 | |
| else AsgiResponse(b"Service is unhealthy", 500, headers={"content-type": "application/json"}) | ||
| ) | ||
|
|
||
| if self.instrument_config.opentelemetry_generate_health_check_spans: | ||
| check_health = tracer.start_as_current_span(f"GET {self.instrument_config.health_checks_path}")( | ||
| check_health, | ||
| ) | ||
|
|
||
| return {"asgi_routes": ((self.instrument_config.health_checks_path, check_health),)} | ||
|
|
||
| async def define_health_status(self) -> bool: | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Короткое название слишком. Я так понимаю, это аналог миддлваря для обычных приложений?