diff --git a/docs/integrations/taskiq.md b/docs/integrations/taskiq.md new file mode 100644 index 0000000..7932786 --- /dev/null +++ b/docs/integrations/taskiq.md @@ -0,0 +1,85 @@ +# Usage with Taskiq + +There is no Taskiq bootstrapper: a worker uses `FreeBootstrapper`, plus Taskiq's own OpenTelemetry +instrumentation for task spans. + +## 1. Install `lite-bootstrap[free-all]` and `taskiq[opentelemetry]`: + +=== "uv" + + ```bash + uv add "lite-bootstrap[free-all]" "taskiq[opentelemetry]" + ``` + +=== "pip" + + ```bash + pip install "lite-bootstrap[free-all]" "taskiq[opentelemetry]" + ``` + +=== "poetry" + + ```bash + poetry add "lite-bootstrap[free-all]" "taskiq[opentelemetry]" + ``` + +On free-threaded CPython, install the instrument extras you need instead of `free-all`, with +`otl-http` in place of `otl`, e.g. `lite-bootstrap[sentry,logging,otl-http]`. Read more about +available extras [here](../introduction/installation.md). + +## 2. Bootstrap in the worker: + +```python +from taskiq import TaskiqEvents, TaskiqState +from taskiq.instrumentation import TaskiqInstrumentor + +from lite_bootstrap import FreeBootstrapper, FreeConfig + + +broker = ... # your broker + + +@broker.on_event(TaskiqEvents.WORKER_STARTUP) +async def bootstrap_observability(state: TaskiqState) -> None: + state.bootstrapper = FreeBootstrapper( + FreeConfig( + service_name="my-worker", + opentelemetry_endpoint="otl", + sentry_dsn="https://testdsn@localhost/1", + ) + ) + state.bootstrapper.bootstrap() + TaskiqInstrumentor().instrument_broker(broker) + + +@broker.on_event(TaskiqEvents.WORKER_SHUTDOWN) +async def teardown_observability(state: TaskiqState) -> None: + state.bootstrapper.teardown() +``` + +`instrument_broker` adds Taskiq's `OpenTelemetryMiddleware` to the broker you already built, and its +spans go to the `TracerProvider` lite-bootstrap has just installed. `TaskiqInstrumentor().instrument()` +is not a substitute: it patches `AsyncBroker.__init__`, so it only reaches brokers built after it runs. + +Read more about available configuration options [here](../introduction/configuration.md). + +## Client and worker processes + +One broker object serves both the processes that send tasks and the `taskiq worker` processes that +run them. The handlers above fire on the worker events only, so a client process is left alone: it +belongs to its own framework's bootstrapper, e.g. `FastAPIBootstrapper` in the web app that calls +`.kiq()`. Bootstrapping when the broker's module is imported would run in that process too, and set +up logging and tracing a second time. + +## Metrics + +lite-bootstrap has no Taskiq metrics. Taskiq's own `taskiq.middlewares.PrometheusMiddleware` +(`taskiq[metrics]`) provides them, with two side effects to know about: its constructor sets +`PROMETHEUS_MULTIPROC_DIR` for the whole process, and it serves metrics from its own HTTP server on +`server_port` (default `9000`) rather than on `prometheus_metrics_path`. + +## Sentry + +`sentry-python` has no Taskiq integration, and none is needed for failing tasks: Taskiq's receiver +logs a task's exception with `logger.exception`, which the `LoggingIntegration` lite-bootstrap adds +reports as an error event. Spans and breadcrumbs follow the Sentry configuration as usual. diff --git a/mkdocs.yml b/mkdocs.yml index 0044bac..414e824 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -15,6 +15,7 @@ nav: - FastStream: integrations/faststream.md - FastAPI: integrations/fastapi.md - FastMCP: integrations/fastmcp.md + - Taskiq: integrations/taskiq.md - Free bootstrapper: integrations/free.md theme: name: material diff --git a/pyproject.toml b/pyproject.toml index 33b6c2c..55f4eda 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -214,6 +214,8 @@ dev = [ # conftest imports structlog.typing at runtime, which the library does not. The suite's # floor is higher than the one the `logging` extra publishes; this is where it belongs. "structlog>=22.2", + # for the Taskiq recipe test only; Taskiq is not an extra. + "taskiq[opentelemetry]", ] lint = [ "ty", diff --git a/tests/test_taskiq_recipe.py b/tests/test_taskiq_recipe.py new file mode 100644 index 0000000..6d2a6d1 --- /dev/null +++ b/tests/test_taskiq_recipe.py @@ -0,0 +1,52 @@ +from opentelemetry import trace +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from taskiq import InMemoryBroker, TaskiqEvents, TaskiqState +from taskiq.instrumentation import TaskiqInstrumentor +from taskiq.middlewares.opentelemetry_middleware import OpenTelemetryMiddleware + +from lite_bootstrap import FreeBootstrapper, FreeConfig + + +def _register_recipe(broker: InMemoryBroker, config: FreeConfig) -> None: + @broker.on_event(TaskiqEvents.WORKER_STARTUP) + async def bootstrap_observability(state: TaskiqState) -> None: + state.bootstrapper = FreeBootstrapper(config) + state.bootstrapper.bootstrap() + TaskiqInstrumentor().instrument_broker(broker) + + @broker.on_event(TaskiqEvents.WORKER_SHUTDOWN) + async def teardown_observability(state: TaskiqState) -> None: + state.bootstrapper.teardown() + + +async def test_taskiq_recipe_traces_tasks_and_tears_down() -> None: + broker = InMemoryBroker() + + @broker.task + async def add(left: int, right: int) -> int: + return left + right + + _register_recipe(broker, FreeConfig(opentelemetry_log_traces=True, logging_buffer_capacity=0)) + + await broker.startup() + try: + bootstrapper = broker.state.bootstrapper + assert bootstrapper.is_bootstrapped + assert any(isinstance(middleware, OpenTelemetryMiddleware) for middleware in broker.middlewares) + + tracer_provider = trace.get_tracer_provider() + assert isinstance(tracer_provider, TracerProvider) + exporter = InMemorySpanExporter() + tracer_provider.add_span_processor(SimpleSpanProcessor(exporter)) + + task = await add.kiq(1, 2) + result = await task.wait_result() + assert result.return_value == 1 + 2 + span_kinds = {span.kind for span in exporter.get_finished_spans() if span.name.endswith(add.task_name)} + assert span_kinds == {trace.SpanKind.PRODUCER, trace.SpanKind.CONSUMER} + finally: + await broker.shutdown() + + assert not bootstrapper.is_bootstrapped