diff --git a/CHANGELOG.md b/CHANGELOG.md index f03d58e7b..e81a03f7a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,23 @@ Write the date in place of the "Unreleased" in the case a new version is release ### Added +- OpenTelemetry tracing for the server, exported over OTLP. The example monitoring + stack (`compose.monitoring.yml`) now includes an OpenTelemetry Collector, Jaeger, + and Grafana Tempo (traces are sent to both backends), and the Collector also + scrapes and re-exposes Tiled's Prometheus metrics. +- Emit OpenTelemetry spans for internal request phases (access control, read, + tokenize, pack) so they appear as child spans in a request's trace, giving a + per-request breakdown of where time is spent. +- Disable FastAPI's built-in OpenTelemetry auto-configuration + (`telemetry={"auto_configure": False}`) so that, on FastAPI >=0.142, it does + not register a second OTLP exporter alongside Tiled's own tracing pipeline and + export every span twice. +- Emit OpenTelemetry spans for PostgreSQL queries (asyncpg for the catalog and + authentication databases, ADBC for the storage database), Redis commands, and + outbound HTTP calls (httpx: OIDC, webhooks, external policy servers), so + external-service calls appear in traces. The example monitoring stack also + generates a service graph viewable in Grafana, with Tiled's separate Postgres + databases (catalog, storage, authn) and Redis shown as distinct nodes. - Expose the Deployment `strategy` in the helm chart, so that a deployment can use `Recreate` instead of the default `RollingUpdate`. - Add a `DELETE /api/v1/asset/{path}?id=N` endpoint to dissociate a single diff --git a/compose.monitoring.yml b/compose.monitoring.yml index 9f5a31a73..19cbda57a 100644 --- a/compose.monitoring.yml +++ b/compose.monitoring.yml @@ -1,5 +1,20 @@ --- services: + # Turn on OpenTelemetry tracing for the Tiled server (defined in compose.yml + # or compose.dev.yml) and point it at the collector below. These variables + # live here, alongside the collector, so that running the base compose files + # on their own leaves tracing off instead of exporting to a collector that is + # not running. + tiled: + environment: + - OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4318 + - OTEL_SERVICE_NAME=tiled + # Don't trace health checks and metrics scrapes (operational chatter). + - OTEL_PYTHON_FASTAPI_EXCLUDED_URLS=healthz,api/v1/metrics + # Sample a fraction of traces in production (default: trace everything). + # - OTEL_TRACES_SAMPLER=parentbased_traceidratio + # - OTEL_TRACES_SAMPLER_ARG=0.1 + prometheus: image: docker.io/prom/prometheus:v2.42.0 volumes: @@ -9,12 +24,14 @@ services: - '--storage.tsdb.path=/prometheus' - '--web.console.libraries=/usr/share/prometheus/console_libraries' - '--web.console.templates=/usr/share/prometheus/consoles' + # Accept remote-written metrics from Tempo's metrics generator (service graph). + - '--web.enable-remote-write-receiver' networks: - backend restart: unless-stopped grafana: - image: docker.io/grafana/grafana:8.2.6 + image: docker.io/grafana/grafana:11.3.0 depends_on: - prometheus ports: @@ -33,5 +50,50 @@ services: GF_AUTH_DISABLE_SIGNOUT_MENU: "true" GF_AUTH_DISABLE_LOGIN_FORM: "true" + # OpenTelemetry Collector: receives OTLP telemetry from apps and fans it + # out to backends (traces -> Jaeger and Tempo). It also scrapes Tiled's + # Prometheus metrics endpoint and re-exposes it for Prometheus. Apps on the + # 'backend' network export to http://otel-collector:4318; apps on the host + # export to http://localhost:4318. + otel-collector: + image: otel/opentelemetry-collector-contrib:0.130.0 + command: ["--config=/etc/otelcol/config.yaml"] + volumes: + - ./monitoring_example/otel-collector/otel-collector.yml:/etc/otelcol/config.yaml + ports: + - 4317:4317 # OTLP gRPC + - 4318:4318 # OTLP HTTP + - 8889:8889 # Prometheus exporter (re-exposed scraped metrics) + depends_on: + - jaeger + - tempo + networks: + - backend + restart: unless-stopped + + # Jaeger all-in-one: trace storage + query UI. In-memory storage (dev + # only; traces are lost on restart). Natively accepts OTLP. + jaeger: + image: jaegertracing/all-in-one:1.62.0 + environment: + COLLECTOR_OTLP_ENABLED: "true" + ports: + - 16686:16686 # Jaeger web UI + networks: + - backend + restart: unless-stopped + + # Grafana Tempo: trace storage queried from Grafana (no UI of its own). + # Stores traces on local ephemeral storage (see monitoring_example/tempo). + # Natively accepts OTLP. Explore traces in Grafana via the Tempo datasource. + tempo: + image: docker.io/grafana/tempo:2.6.1 + command: ["-config.file=/etc/tempo/tempo.yml"] + volumes: + - ./monitoring_example/tempo/tempo.yml:/etc/tempo/tempo.yml + networks: + - backend + restart: unless-stopped + networks: backend: {} diff --git a/docs/source/_toc.yml b/docs/source/_toc.yml index 98b9e5ef9..70f461128 100644 --- a/docs/source/_toc.yml +++ b/docs/source/_toc.yml @@ -42,6 +42,7 @@ subtrees: - file: user-guide/api-keys - file: user-guide/custom-clients - file: user-guide/metrics + - file: user-guide/tracing - file: user-guide/direct-client - file: user-guide/tiled-authn-database - file: user-guide/register diff --git a/docs/source/user-guide/tracing.md b/docs/source/user-guide/tracing.md new file mode 100644 index 000000000..a2df4c0af --- /dev/null +++ b/docs/source/user-guide/tracing.md @@ -0,0 +1,175 @@ +# Distributed Tracing + +In addition to [Prometheus metrics](./metrics.md), Tiled can emit +[OpenTelemetry](https://opentelemetry.io/) traces. Whereas metrics describe +aggregate behavior across many requests, a *trace* records the timeline of a +single request as a tree of *spans*. This is useful for investigating why a +particular request was slow. + +Traces are exported using the OpenTelemetry Protocol (OTLP) to an +[OpenTelemetry Collector](https://opentelemetry.io/docs/collector/), which +forwards them to one or more tracing backends for storage and visualization, +such as [Jaeger](https://www.jaegertracing.io/) or +[Grafana Tempo](https://grafana.com/oss/tempo/). + +```{mermaid} +flowchart LR + tiled["Tiled server"] + collector["OpenTelemetry
Collector"] + + subgraph backends["Storage backends"] + direction TB + jaeger["Jaeger"] + tempo["Grafana Tempo"] + prometheus["Prometheus"] + loki["Loki"] + end + + subgraph viz["Visualization"] + direction TB + jaegerui["Jaeger UI"] + grafana["Grafana"] + end + + %% Configured in the example + tiled -->|"traces (OTLP)"| collector + collector -->|OTLP| jaeger + collector -->|OTLP| tempo + tiled -->|"metrics (scrape)"| prometheus + + %% Visualization + jaeger --> jaegerui + tempo --> grafana + prometheus --> grafana + + %% Metrics and logs over OTLP: possible extension, not enabled + tiled -.->|"metrics (OTLP)"| collector + tiled -.->|"logs (OTLP)"| collector + collector -.->|metrics| prometheus + collector -.->|logs| loki + loki -.-> grafana +``` + +Solid arrows are what the example configures today: Tiled pushes **traces** over +OTLP to the Collector, which fans them out to Jaeger and Grafana Tempo, while +Prometheus scrapes Tiled's metrics endpoint. Dashed arrows show how the same +Collector could also carry OpenTelemetry's other two signals — **metrics** and +**logs** — over OTLP to backends such as Prometheus and Loki. Those paths are +not currently enabled. + +```{note} +Database-query spans are emitted only for PostgreSQL (via asyncpg and ADBC); +tracing with SQLite-backed catalogs is not supported. +``` + +## Enabling tracing + +Tracing is **disabled by default**. It is turned on by setting the standard +OpenTelemetry environment variable `OTEL_EXPORTER_OTLP_ENDPOINT` to the address +of an OTLP endpoint (an OpenTelemetry Collector, or a backend that accepts OTLP +directly). Related environment variables: + +| Variable | Purpose | +| --- | --- | +| `OTEL_EXPORTER_OTLP_ENDPOINT` | OTLP endpoint, e.g. `http://otel-collector:4318`. Tracing is off when this is unset. | +| `OTEL_SERVICE_NAME` | Name shown for the service in the tracing backend, e.g. `tiled`. | +| `OTEL_PYTHON_FASTAPI_EXCLUDED_URLS` | Comma-separated URL patterns to exclude from tracing, e.g. `healthz,api/v1/metrics` to skip health checks and metrics scrapes. | + + +## Sampling + +By default every request is traced in full. That is convenient for trying it out +but can be a lot of data in production, especially since each request emits a +span per database query. Sampling is controlled by the standard OpenTelemetry +environment variables: + +| Variable | Purpose | +| --- | --- | +| `OTEL_TRACES_SAMPLER` | Sampling strategy. Default `parentbased_always_on` (trace everything). Use `parentbased_traceidratio` to keep a fraction. | +| `OTEL_TRACES_SAMPLER_ARG` | Argument for the sampler; for the ratio samplers, the fraction of traces to keep (0.0-1.0). | + +For example, to keep 10% of traces: + +``` +OTEL_TRACES_SAMPLER=parentbased_traceidratio +OTEL_TRACES_SAMPLER_ARG=0.1 +``` + +The `parentbased_*` samplers make the decision once at the start of a trace and +apply it to all of that trace's spans, so a sampled request keeps its database +and cache spans together with the rest of the trace. + + +## How does it work? + +1. When `OTEL_EXPORTER_OTLP_ENDPOINT` is set, Tiled configures an OpenTelemetry + tracer and instruments the ASGI application, creating one span per incoming + HTTP request. + +2. Spans are exported over OTLP to the OpenTelemetry Collector. + +3. The Collector forwards traces to one or more backends (Jaeger and Grafana + Tempo in the example stack), which store them and make them available to + search and visualize. + + +## Try it with the example stack + +Tiled ships example configuration that runs an OpenTelemetry Collector and +Jaeger alongside the server, Prometheus, and Grafana. From the repository root, +start the server together with the monitoring services: + +``` +TILED_SINGLE_USER_API_KEY=secret \ + docker compose -f compose.dev.yml -f compose.monitoring.yml up --build +``` + +`compose.dev.yml` builds the Tiled image from this checkout (so it includes the +tracing support), and `compose.monitoring.yml` sets the `OTEL_*` variables above +and runs the Collector, so the server exports traces to it. (The published image +referenced by `compose.yml` may not yet include tracing.) + +Generate some activity using the Tiled Python client: + +```python +from tiled.client import from_uri + +c = from_uri("http://localhost:8000", api_key="secret") +c.create_container('test') +list(c) +``` + +The example forwards traces to two backends so you can compare their functionality: + +- **Jaeger:** open [http://localhost:16686](http://localhost:16686), select the + **tiled** service, and click **Find Traces**. Click a trace to see its span + waterfall. +- **Grafana Tempo:** open [http://localhost:3000](http://localhost:3000), go to + **Explore**, select the **Tempo** data source, and search using + [TraceQL](https://grafana.com/docs/tempo/latest/traceql/), for example + `{ resource.service.name = "tiled" }`. + +Each trace also includes spans for the **PostgreSQL** queries (against the +catalog, storage, and authentication databases), **Redis** commands (for the +streaming cache), and any **outbound HTTP** calls Tiled makes while serving the +request (OIDC authentication, webhooks, and external policy servers, via +httpx). These are client spans emitted by Tiled, so they share the `tiled` +service, but they carry a `db.system` attribute (`postgresql` or `redis`) — or, +for HTTP calls, the target host — that distinguishes them from Tiled's own +`tiled.*` spans. The Collector drops transaction-control statements +(`BEGIN`/`COMMIT`/`ROLLBACK`) to keep traces concise and readable. + +Filter spans by the `db.system` attribute (Grafana's span filters, or Jaeger's +find-within-trace box) to highlight the database and cache work. Grafana's +**Service Graph** (Explore → Tempo) also renders Tiled's dependencies as nodes: +each Postgres database (`tiled_catalog`, `tiled_storage`, authn) by `db.name`, +`redis`, and outbound HTTP — webhook deliveries grouped under one `webhooks` +node, other calls (e.g. OIDC) named by host. + +```{note} +The bundled Collector also scrapes Tiled's `/api/v1/metrics` endpoint and +re-exposes it on port 8889, in addition to Prometheus scraping it directly. +To disable this, remove the `metrics` pipeline from the Collector +configuration in `monitoring_example/otel-collector/otel-collector.yml`. +See [Prometheus Metrics](./metrics.md). +``` diff --git a/monitoring_example/grafana/provisioning/datasources/prometheus.yml b/monitoring_example/grafana/provisioning/datasources/prometheus.yml index 0cd616ec1..6900ca7ee 100644 --- a/monitoring_example/grafana/provisioning/datasources/prometheus.yml +++ b/monitoring_example/grafana/provisioning/datasources/prometheus.yml @@ -1,5 +1,7 @@ +apiVersion: 1 datasources: - name: Prometheus + uid: prometheus access: proxy type: prometheus url: http://prometheus:9090 diff --git a/monitoring_example/grafana/provisioning/datasources/tempo.yml b/monitoring_example/grafana/provisioning/datasources/tempo.yml new file mode 100644 index 000000000..a235a4cce --- /dev/null +++ b/monitoring_example/grafana/provisioning/datasources/tempo.yml @@ -0,0 +1,10 @@ +apiVersion: 1 +datasources: +- name: Tempo + access: proxy + type: tempo + url: http://tempo:3200 + uid: tempo + jsonData: + serviceMap: + datasourceUid: prometheus diff --git a/monitoring_example/otel-collector/otel-collector.yml b/monitoring_example/otel-collector/otel-collector.yml new file mode 100644 index 000000000..1503ef53a --- /dev/null +++ b/monitoring_example/otel-collector/otel-collector.yml @@ -0,0 +1,98 @@ +# OpenTelemetry Collector configuration. +# +# Traces: tiled --OTLP--> otel-collector --OTLP--> jaeger and tempo +# Metrics: otel-collector <--scrape-- tiled:8000/api/v1/metrics +# otel-collector --expose--> :8889 (Prometheus format) +# +# The collector receives OTLP traces and forwards the same spans to both Jaeger +# and Grafana Tempo. Database and cache calls stay under the "tiled" service +# (they are client spans emitted by Tiled) and are identified by their +# db.system / peer.service attributes; peer.service also lets Tempo draw the +# service graph. The collector also scrapes Tiled's Prometheus metrics endpoint +# and re-exposes those metrics on port 8889. + +receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:4317 + http: + endpoint: 0.0.0.0:4318 + + # Scrape Tiled's Prometheus metrics endpoint. This is a standard Prometheus + # scrape config. The API key must match TILED_SINGLE_USER_API_KEY (or an API + # key carrying the "metrics" scope in a multi-user deployment). + prometheus: + config: + scrape_configs: + - job_name: tiled + metrics_path: /api/v1/metrics + scrape_interval: 5s + authorization: + type: Apikey + credentials: secret + static_configs: + - targets: ['tiled:8000'] + +processors: + # Cap the batch size: Jaeger and Tempo reject OTLP/gRPC messages larger than + # 4 MiB by default, and an oversized batch is dropped, not retried. With no + # cap (`batch: {}`), a burst of traces can exceed that limit. 1024 spans fit + # if spans average under ~4 KiB. + batch: + send_batch_size: 512 + send_batch_max_size: 1024 + + # Drop noisy transaction/session-control statements so traces show the + # queries that matter, not every BEGIN/COMMIT/ROLLBACK. + filter/db_noise: + error_mode: ignore + traces: + span: + - 'attributes["db.system"] != nil and IsMatch(name, "^(BEGIN|COMMIT|ROLLBACK|SET|SHOW)")' + + # Set peer.service on database/cache client spans so Tempo's service graph + # names the dependency nodes. Postgres spans use the database name, so the + # catalog, storage, and authn databases appear as separate nodes; other + # datastores (e.g. Redis) fall back to db.system. + transform/peer_service: + error_mode: ignore + trace_statements: + - context: span + statements: + - set(attributes["peer.service"], attributes["db.name"]) + where attributes["db.system"] == "postgresql" and attributes["db.name"] != nil + - set(attributes["peer.service"], attributes["db.system"]) + where attributes["peer.service"] == nil and attributes["db.system"] != nil + +exporters: + # Forward traces to Jaeger's OTLP gRPC receiver. + otlp/jaeger: + endpoint: jaeger:4317 + tls: + insecure: true + # Forward the same traces to Grafana Tempo's OTLP gRPC receiver. + otlp/tempo: + endpoint: tempo:4317 + tls: + insecure: true + # Re-expose scraped metrics in Prometheus format for Prometheus to scrape. + prometheus: + endpoint: 0.0.0.0:8889 + # Log received spans to the collector's stdout (useful for debugging). + debug: + verbosity: detailed + +service: + pipelines: + traces: + receivers: [otlp] + processors: [filter/db_noise, transform/peer_service, batch] + exporters: [otlp/jaeger, otlp/tempo, debug] + metrics: + receivers: [prometheus] + processors: [batch] + exporters: [prometheus] + telemetry: + logs: + level: info diff --git a/monitoring_example/prometheus/prometheus.yml b/monitoring_example/prometheus/prometheus.yml index 17fc476e0..8cc092c10 100644 --- a/monitoring_example/prometheus/prometheus.yml +++ b/monitoring_example/prometheus/prometheus.yml @@ -3,6 +3,10 @@ global: scrape_interval: 5s scrape_configs: + # Scrape Tiled's metrics endpoint directly. The OpenTelemetry Collector + # independently scrapes the same endpoint and brings the metrics into its + # pipeline (see monitoring_example/otel-collector/otel-collector.yml), + # re-exposing them on otel-collector:8889 for forwarding to other backends. - job_name: 'tiled' metrics_path: /api/v1/metrics authorization: # Set Authorization header to 'Apikey secret'. diff --git a/monitoring_example/tempo/tempo.yml b/monitoring_example/tempo/tempo.yml new file mode 100644 index 000000000..7dcdc2c8f --- /dev/null +++ b/monitoring_example/tempo/tempo.yml @@ -0,0 +1,60 @@ +# Grafana Tempo configuration (single binary, for the example stack). +# +# Tempo receives traces over OTLP from the OpenTelemetry Collector and stores +# them on the local filesystem. This example writes to /tmp (ephemeral, lost on +# restart) so it needs no volume or permission setup. Traces are queried from +# Grafana via the Tempo datasource. + +server: + http_listen_port: 3200 + +distributor: + receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:4317 + http: + endpoint: 0.0.0.0:4318 + +ingester: + max_block_duration: 5m + +compactor: + compaction: + block_retention: 1h + +storage: + trace: + backend: local + local: + path: /tmp/tempo/blocks + wal: + path: /tmp/tempo/wal + +# Generate service-graph and span metrics from incoming traces and remote-write +# them to Prometheus. Grafana's Tempo "Service Graph" reads these. Client spans +# to Postgres and Redis (which have no server span) become virtual nodes named +# from the peer.service attribute the collector sets: the database name for +# Postgres (e.g. tiled_catalog, tiled_storage) and db.system for others (e.g. +# redis). +metrics_generator: + registry: + external_labels: + source: tempo + storage: + path: /tmp/tempo/generator/wal + remote_write: + - url: http://prometheus:9090/api/v1/write + send_exemplars: true + traces_storage: + path: /tmp/tempo/generator/traces + processor: + service_graphs: + # Use peer.service (set by the collector) to name dependency nodes. + peer_attributes: [peer.service] + +overrides: + defaults: + metrics_generator: + processors: [service-graphs, span-metrics] diff --git a/pyproject.toml b/pyproject.toml index 1a9603359..ae0a99ad7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -87,6 +87,13 @@ all = [ "numba >=0.59.0", # indirect, pinned to assist uv solve "obstore", "openpyxl", + "opentelemetry-exporter-otlp-proto-http", + "opentelemetry-instrumentation-asyncpg", + "opentelemetry-instrumentation-dbapi", + "opentelemetry-instrumentation-fastapi", + "opentelemetry-instrumentation-httpx", + "opentelemetry-instrumentation-redis", + "opentelemetry-sdk", "packaging", "pandas <3", "pillow", @@ -226,6 +233,13 @@ server = [ "numpy", "obstore", "openpyxl", + "opentelemetry-exporter-otlp-proto-http", + "opentelemetry-instrumentation-asyncpg", + "opentelemetry-instrumentation-dbapi", + "opentelemetry-instrumentation-fastapi", + "opentelemetry-instrumentation-httpx", + "opentelemetry-instrumentation-redis", + "opentelemetry-sdk", "packaging", "pandas", "pillow", diff --git a/tests/test_tracing.py b/tests/test_tracing.py new file mode 100644 index 000000000..d8c8b9f0e --- /dev/null +++ b/tests/test_tracing.py @@ -0,0 +1,372 @@ +"""In-process tests for Tiled's OpenTelemetry tracing. + +Two groups, both capturing spans with an in-memory exporter (no Collector, +Jaeger, or Tempo needed): + +* request-level tracing configured by + `tiled.server.app._setup_opentelemetry_tracing` -- the server span and phase + spans, excluded endpoints, off-by-default, and no duplicate export pipeline; +* external-service instrumentation -- asyncpg (catalog Postgres), ADBC (SQL + storage), Redis (streaming cache), and httpx (webhook delivery) -- exercised + against live backends, which skip when the backend is not configured via + `TILED_TEST_POSTGRESQL_URI` / `TILED_TEST_REDIS`. + +OpenTelemetry's global tracer provider can only be set once per process, so a +single provider is installed for the whole module and the exporter is cleared +between tests. +""" +import os +from urllib.parse import urlparse + +import numpy as np +import pyarrow +import pytest +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 opentelemetry.trace import SpanKind + +from tiled.catalog import in_memory +from tiled.client import Context, from_context +from tiled.config import Authentication, WebhooksConfig +from tiled.server.app import build_app, build_app_from_config +from tiled.server.schemas import WebhookRegistrationRequest + +# The FastAPI instrumentation reads OTEL_PYTHON_FASTAPI_EXCLUDED_URLS once, at +# import time, into a module-level default that `instrument_app` uses. Set it +# before that module is first imported (which the tracing hook does lazily on +# the first traced `build_app`). In deployments this variable is likewise set +# in the environment before the process starts. +os.environ["OTEL_PYTHON_FASTAPI_EXCLUDED_URLS"] = "healthz,api/v1/metrics" + +# Minimal in-memory tree used by the request-level tests. +CONFIG = { + "authentication": {"single_user_api_key": "secret"}, + "trees": [{"path": "/", "tree": "tiled.examples.generated_minimal:tree"}], +} +# A syntactically valid endpoint that is never actually contacted: the tracing +# hook reuses the in-memory provider installed below instead of creating an OTLP +# exporter, so no network traffic occurs. We only declare it to turn tracing on. +ENDPOINT = "http://otel-collector.invalid:4318" +API_KEY = "secret" +# respx mocks this, so no real delivery happens (and HTTPS keeps the default +# URL validator happy without allowing http/private targets). +WEBHOOK_URL = "https://webhook.example.com/tiled-events" + +# Global library instrumentation installed by the tracing hook, which must be +# undone so it does not leak into other test modules. +_GLOBAL_INSTRUMENTORS = [ + ("opentelemetry.instrumentation.asyncpg", "AsyncPGInstrumentor"), + ("opentelemetry.instrumentation.redis", "RedisInstrumentor"), + ("opentelemetry.instrumentation.httpx", "HTTPXClientInstrumentor"), +] + + +@pytest.fixture(scope="module") +def span_exporter(): + exporter = InMemorySpanExporter() + current = trace.get_tracer_provider() + if isinstance(current, TracerProvider): + # Another module already installed a real provider; attach to it. + provider = current + else: + provider = TracerProvider() + trace.set_tracer_provider(provider) + provider.add_span_processor(SimpleSpanProcessor(exporter)) + yield exporter + for module, cls in _GLOBAL_INSTRUMENTORS: + try: + mod = __import__(module, fromlist=[cls]) + getattr(mod, cls)().uninstrument() + except Exception: + pass + + +@pytest.fixture(autouse=True) +def _clear_spans(span_exporter): + span_exporter.clear() + yield + span_exporter.clear() + + +def _enable_tracing(monkeypatch): + # Set before build_app so the tracing hook runs and instruments the + # libraries against the in-memory provider. + monkeypatch.setenv("OTEL_EXPORTER_OTLP_ENDPOINT", ENDPOINT) + + +def _build_app(monkeypatch, *, endpoint=ENDPOINT): + if endpoint is None: + monkeypatch.delenv("OTEL_EXPORTER_OTLP_ENDPOINT", raising=False) + else: + monkeypatch.setenv("OTEL_EXPORTER_OTLP_ENDPOINT", endpoint) + return build_app_from_config(CONFIG) + + +def _is_descendant_of(span, ancestor_span_id, by_id): + seen = set() + cur = span + while cur is not None and cur.parent is not None: + parent_id = cur.parent.span_id + if parent_id == ancestor_span_id: + return True + if parent_id in seen: + break + seen.add(parent_id) + cur = by_id.get(parent_id) + return False + + +def _spans_where(spans, key, value): + return [s for s in spans if s.attributes.get(key) == value] + + +# --- request-level tracing ------------------------------------------------- + + +def test_traced_request_emits_server_and_phase_spans(monkeypatch, span_exporter): + app = _build_app(monkeypatch) + with Context.from_app(app) as context: + # Discard spans emitted while the Context was being set up, so we measure + # exactly one request. + span_exporter.clear() + response = context.http_client.get("/api/v1/metadata/") + assert response.status_code == 200 + spans = span_exporter.get_finished_spans() + + assert spans, "expected spans to be exported for a traced request" + ids = [s.context.span_id for s in spans] + assert len(ids) == len(set(ids)), "spans must not be duplicated" + assert len({s.context.trace_id for s in spans}) == 1, "one request => one trace" + + roots = [s for s in spans if s.parent is None] + assert len(roots) == 1, "expected exactly one root span" + assert roots[0].kind == SpanKind.SERVER, "root should be the FastAPI server span" + + by_id = {s.context.span_id: s for s in spans} + names = {s.name for s in spans} + assert "tiled.app" in names, "expected the per-request phase span 'tiled.app'" + app_span = next(s for s in spans if s.name == "tiled.app") + assert _is_descendant_of( + app_span, roots[0].context.span_id, by_id + ), "phase spans should be children of the request's server span" + + +def test_excluded_endpoint_emits_no_spans(monkeypatch, span_exporter): + app = _build_app(monkeypatch) + with Context.from_app(app) as context: + span_exporter.clear() + response = context.http_client.get("/healthz") + assert response.status_code == 200 + assert ( + not span_exporter.get_finished_spans() + ), "excluded endpoints must not produce spans (no orphan traces)" + + +def test_tracing_disabled_by_default(monkeypatch): + # Without OTEL_EXPORTER_OTLP_ENDPOINT the tracing hook returns early and does + # not instrument the app, so tracing is off and adds no overhead. (Asserting + # on emitted spans is not reliable here: FastAPI >=0.142 ships its own + # telemetry that emits request spans on any globally installed provider when + # the app is not instrumented by OpenTelemetry.) + app = _build_app(monkeypatch, endpoint=None) + assert not getattr(app, "_is_instrumented_by_opentelemetry", False) + + +def test_no_duplicate_export_pipeline(monkeypatch, span_exporter): + """Guard against FastAPI's built-in telemetry (>=0.142) registering a second + OTLP export pipeline, which would export every span twice.""" + from starlette.testclient import TestClient + + provider = trace.get_tracer_provider() + processors = provider._active_span_processor._span_processors + before = len(processors) + + app = _build_app(monkeypatch) + # Entering the TestClient runs the ASGI lifespan; FastAPI configures its + # built-in telemetry on `lifespan.startup`. + with TestClient(app): + pass + + after = len(provider._active_span_processor._span_processors) + assert after == before, ( + "a second span processor was registered on the global provider; " + "FastAPI's built-in OpenTelemetry auto-configuration is not disabled" + ) + + +# --- external-service spans (live backends; skip when not configured) ------ + + +def test_catalog_query_emits_asyncpg_spans( + monkeypatch, span_exporter, postgres_uri, tmp_path +): + """Catalog access over asyncpg produces postgresql client spans.""" + _enable_tracing(monkeypatch) + config = { + "authentication": {"single_user_api_key": API_KEY}, + "trees": [ + { + "tree": "catalog", + "path": "/", + "args": { + "uri": postgres_uri, + "writable_storage": [str(tmp_path / "data")], + "init_if_not_exists": True, + }, + } + ], + } + with Context.from_app(build_app_from_config(config)) as context: + client = from_context(context) + client.write_array(np.arange(5), key="arr") + span_exporter.clear() + list(client) # a search -> catalog SELECT over asyncpg + + spans = span_exporter.get_finished_spans() + pg_spans = _spans_where(spans, "db.system", "postgresql") + assert pg_spans, "expected asyncpg (postgresql) spans for the catalog query" + # The catalog database name appears on the span so it is distinguishable in the service graph + assert any(s.attributes.get("db.name") for s in pg_spans) + + +def test_sql_storage_write_emits_adbc_spans( + monkeypatch, span_exporter, sql_storage_uri, tmp_path +): + """Writing an appendable table to SQL storage (SQLite, DuckDB, or Postgres) + produces both the manual `adbc_ingest` span (the bulk write bypasses the DBAPI + `execute` path) and the DBAPI-level spans from the instrumented ADBC + connection.""" + _enable_tracing(monkeypatch) + dialect = urlparse(sql_storage_uri).scheme + config = { + "authentication": {"single_user_api_key": API_KEY}, + "trees": [ + { + "tree": "catalog", + "path": "/", + "args": { + "uri": f"sqlite:///{tmp_path / 'catalog.db'}", + "writable_storage": [sql_storage_uri], + "init_if_not_exists": True, + }, + } + ], + } + table = pyarrow.Table.from_pydict({"A": [1, 2, 3], "B": [4, 5, 6]}) + with Context.from_app(build_app_from_config(config)) as context: + client = from_context(context) + span_exporter.clear() + appendable = client.create_appendable_table(schema=table.schema, key="tab") + appendable.append_partition(0, table) + # Tracing must not break the storage connections. + assert appendable.read()["A"].tolist() == [1, 2, 3] + + spans = span_exporter.get_finished_spans() + ingest = [s for s in spans if s.name == "adbc_ingest"] + assert ingest, "expected the manual adbc_ingest span" + assert ingest[0].kind == SpanKind.CLIENT + assert ingest[0].attributes.get("db.system") == dialect + + # The ADBC connection factory is wrapped by _instrument_adbc_creator, so the + # DBAPI-level statements (e.g. the CREATE TABLE preceding the ingest) are + # also traced. + dbapi_spans = [ + s for s in _spans_where(spans, "db.system", dialect) if s.name != "adbc_ingest" + ] + assert dbapi_spans, "expected DBAPI-instrumented storage query spans" + # They carry the database name (`adbc_current_catalog`), except on DuckDB, whose + # ADBC driver does not implement it. + if dialect != "duckdb": + assert any(s.attributes.get("db.name") for s in dbapi_spans) + + +def test_streaming_emits_redis_spans(monkeypatch, span_exporter, redis_uri, tmp_path): + """Subscribing to a node's stream exercises the Redis streaming cache.""" + _enable_tracing(monkeypatch) + config = { + "authentication": {"single_user_api_key": API_KEY}, + "trees": [ + { + "tree": "catalog", + "path": "/", + "args": { + "uri": "sqlite:///:memory:", + "writable_storage": [str(tmp_path / "data")], + "init_if_not_exists": True, + }, + } + ], + "streaming_cache": { + "uri": redis_uri, + "data_ttl": 600, + "seq_ttl": 600, + "socket_timeout": 600, + "socket_connect_timeout": 10, + }, + } + with Context.from_app(build_app_from_config(config)) as context: + client = from_context(context) + test_client = context.http_client # the underlying starlette TestClient + node = client.write_array(np.arange(10), key="stream_node") + span_exporter.clear() + with test_client.websocket_connect( + "/api/v1/stream/single/stream_node?envelope_format=json", + headers={"Authorization": f"Apikey {API_KEY}"}, + ): + node.write(np.arange(10) + 1) + + spans = span_exporter.get_finished_spans() + redis_spans = _spans_where(spans, "db.system", "redis") + assert redis_spans, "expected Redis spans for the streaming subscription" + + +def test_webhook_delivery_emits_httpx_span(monkeypatch, span_exporter, tmp_path): + """Delivering a webhook goes through httpx, producing an outbound client + span. No external backend is needed; the delivery is mocked with respx.""" + respx = pytest.importorskip("respx") + from unittest.mock import patch + + from httpx import Response + + _enable_tracing(monkeypatch) + tree = in_memory(writable_storage=[f"file://localhost{tmp_path / 'data'}"]) + app = build_app( + tree, + authentication=Authentication(single_user_api_key=API_KEY), + # A non-None webhooks config enables the delivery dispatcher. + server_settings={"webhooks": WebhooksConfig(secret_keys=["test-webhook-key"])}, + ) + + with Context.from_app(app) as context: + client = from_context(context) + # respx mocks the delivery; patching the SSRF check lets us register an example.com target + with respx.mock, patch("tiled.server.webhook_router.check_url_ssrf_safety"): + respx.post(WEBHOOK_URL).mock(return_value=Response(200)) + context.http_client.post( + "/api/v1/webhooks/target/", + json=WebhookRegistrationRequest(url=WEBHOOK_URL).model_dump( + mode="json" + ), + ).raise_for_status() + span_exporter.clear() + client.create_container("triggers_webhook") + + spans = span_exporter.get_finished_spans() + httpx_spans = [ + s + for s in spans + if s.kind == SpanKind.CLIENT + and ( + s.attributes.get("http.method") == "POST" + or s.attributes.get("http.request.method") == "POST" + ) + ] + assert httpx_spans, "expected an outbound httpx client span for the webhook POST" + # The span targets the webhook URL's host. + urls = [ + str(s.attributes.get("http.url") or s.attributes.get("url.full") or "") + for s in httpx_spans + ] + assert any("webhook.example.com" in u for u in urls) diff --git a/tiled/adapters/sql.py b/tiled/adapters/sql.py index 5ff098259..b00d8d106 100644 --- a/tiled/adapters/sql.py +++ b/tiled/adapters/sql.py @@ -3,9 +3,10 @@ import copy import hashlib import logging +import os import re from collections.abc import Set -from contextlib import closing +from contextlib import AbstractContextManager, closing, nullcontext from typing import ( TYPE_CHECKING, Any, @@ -20,6 +21,7 @@ Union, cast, ) +from urllib.parse import urlparse from tiled.utils import UnsafeIdentifier @@ -49,6 +51,13 @@ from ..type_aliases import JSON from .array import ArrayAdapter +try: + from opentelemetry import trace as _otel_trace + + _STORAGE_TRACER = _otel_trace.get_tracer("tiled.storage") +except ImportError: # OpenTelemetry is an optional dependency. + _STORAGE_TRACER = None + DIALECTS = Literal["postgresql", "sqlite", "duckdb"] TABLE_NAME_PATTERN = re.compile(r"^[a-z][a-z0-9_]*$") COLUMN_NAME_PATTERN = re.compile(r"^[a-zA-Z_].*$") @@ -503,9 +512,36 @@ def append_partition( with closing(self.storage.connect()) as conn: with conn.cursor() as cursor: - cursor.adbc_ingest(self.table_name, table, mode="append") + with self._adbc_ingest_span(): + cursor.adbc_ingest(self.table_name, table, mode="append") conn.commit() + def _adbc_ingest_span(self) -> AbstractContextManager[Any]: + """OpenTelemetry tracing span for the bulk `adbc_ingest` write. + + `adbc_ingest` bypasses the DBAPI `execute()` path, so the generic + dbapi instrumentation does not trace it; emit a span explicitly with the + same db attributes as the other storage spans. A no-op if OpenTelemetry + is not installed or tracing is not configured. + """ + if _STORAGE_TRACER is None or not os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT"): + return nullcontext() + attributes = { + "db.system": self.storage.dialect, + "db.sql.table": self.table_name, + } + db_name = urlparse(self.storage.uri).path.lstrip("/") + if db_name: + attributes["db.name"] = db_name + return cast( + AbstractContextManager[Any], + _STORAGE_TRACER.start_as_current_span( + "adbc_ingest", + kind=_otel_trace.SpanKind.CLIENT, + attributes=attributes, + ), + ) + def _read_full_table_or_partition( self, fields: Optional[List[str]] = None, partition: Optional[int] = None ) -> pyarrow.Table: diff --git a/tiled/server/app.py b/tiled/server/app.py index 9e7c2d421..f12a565ed 100644 --- a/tiled/server/app.py +++ b/tiled/server/app.py @@ -1,6 +1,7 @@ import asyncio import collections import contextvars +import importlib.metadata import logging import os import secrets @@ -79,6 +80,7 @@ CSRF_QUERY_PARAMETER = "csrf" MINIMUM_SUPPORTED_PYTHON_CLIENT_VERSION = packaging.version.parse("0.1.0a104") +FASTAPI_VERSION = packaging.version.Version(importlib.metadata.version("fastapi")) logger = logging.getLogger(__name__) logger.setLevel("INFO") @@ -286,7 +288,17 @@ async def lifespan(app: FastAPI): finally: await shutdown_event() - app = FastAPI(lifespan=lifespan, strict_content_type=False) + # FastAPI >=0.142 ships built-in OpenTelemetry support that, when + # `OTEL_EXPORTER_OTLP_ENDPOINT` (or a related variable) is set, + # auto-registers its own OTLP export pipeline on the global tracer + # provider. Tiled configures and manages its own tracing pipeline (see + # `_setup_opentelemetry_tracing`), so FastAPI's auto-configuration + # would register a second exporter and emit every span twice. Opt out. + kwargs = dict(lifespan=lifespan, strict_content_type=False) + if FASTAPI_VERSION >= packaging.version.Version("0.142"): + kwargs["telemetry"] = {"auto_configure": False} + + app = FastAPI(**kwargs) # Healthcheck for deployment to containerized systems, needs to preempt other responses. # Standardized for Kubernetes, but also used by other systems. @@ -1067,9 +1079,88 @@ async def current_principal_logging_filter( generator=lambda: secrets.token_hex(8), ) + _setup_opentelemetry_tracing(app) + return app +def _setup_opentelemetry_tracing(app: FastAPI) -> None: + """Enable OpenTelemetry request tracing when an OTLP endpoint is configured. + + Tracing is activated only when the standard `OTEL_EXPORTER_OTLP_ENDPOINT` + environment variable is set, so it is off by default and adds no overhead + unless explicitly enabled. Spans are exported over OTLP/HTTP. + """ + if not os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT"): + return + try: + from opentelemetry import trace + from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( + OTLPSpanExporter, + ) + from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor + from opentelemetry.sdk.resources import Resource + from opentelemetry.sdk.trace import TracerProvider + from opentelemetry.sdk.trace.export import BatchSpanProcessor + except ImportError: + logger.warning( + "OTEL_EXPORTER_OTLP_ENDPOINT is set but the OpenTelemetry packages " + "are not installed; tracing is disabled." + ) + return + + # Configure the global tracer provider once per process. + if not isinstance(trace.get_tracer_provider(), TracerProvider): + resource = Resource.create( + {"service.name": os.getenv("OTEL_SERVICE_NAME", "tiled")} + ) + provider = TracerProvider(resource=resource) + provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter())) + trace.set_tracer_provider(provider) + FastAPIInstrumentor.instrument_app(app) + + # Emit spans for calls to the PostgreSQL driver (asyncpg), Redis, and + # outbound HTTP (httpx: OIDC, webhooks, external policy servers). These + # patch the libraries globally, so they are no-ops until a request uses them. + try: + from opentelemetry.instrumentation.asyncpg import AsyncPGInstrumentor + except ImportError: + pass + else: + AsyncPGInstrumentor().instrument() + try: + from opentelemetry.instrumentation.redis import RedisInstrumentor + except ImportError: + pass + else: + RedisInstrumentor().instrument() + try: + from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor + except ImportError: + pass + else: + + def _set_peer_service(span, request): + # Name the downstream dependency so it becomes its own node in the + # Tempo service graph (via peer.service) and is filterable in Jaeger. + # Bundle every webhook delivery under a single 'webhooks' node; name + # other outbound calls (OIDC, external policy servers) by their host. + if span is None or not span.is_recording(): + return + if "x-tiled-event-id" in request.headers: + span.set_attribute("peer.service", "webhooks") + elif request.url.host: + span.set_attribute("peer.service", request.url.host) + + async def _set_peer_service_async(span, request): + _set_peer_service(span, request) + + HTTPXClientInstrumentor().instrument( + request_hook=_set_peer_service, + async_request_hook=_set_peer_service_async, + ) + + def build_app_from_config(config: Union[Config, dict[str, Any]], scalable=False): """ Convenience function that calls build_app(...) given config as parsed Config instance diff --git a/tiled/server/utils.py b/tiled/server/utils.py index 6f1e3d634..84900eb5b 100644 --- a/tiled/server/utils.py +++ b/tiled/server/utils.py @@ -1,4 +1,5 @@ import contextlib +import importlib.util import time from collections.abc import Generator from typing import Any, Literal, Mapping, Optional, Sequence @@ -18,6 +19,26 @@ API_KEY_QUERY_PARAMETER = "api_key" CSRF_COOKIE_NAME = "tiled_csrf" +# Human-readable OpenTelemetry span names for the phases timed below. +_SPAN_NAMES = { + "app": "tiled.app", + "acl": "tiled.access_control", + "read": "tiled.read", + "tok": "tiled.tokenize", + "pack": "tiled.pack", +} + +# Enable tracing if the OpenTelemetry API is installed. `opentelemetry` is a +# namespace package shared by all `opentelemetry-*` distributions, so check for +# the `trace` module itself and its parent (first). +_tracer = None +if importlib.util.find_spec("opentelemetry") and importlib.util.find_spec( + "opentelemetry.trace" +): + from opentelemetry import trace + + _tracer = trace.get_tracer("tiled.server") + def normalize_root_path(root_path: Optional[str]) -> str: """Coerce a root_path to "" or "/prefix" (no trailing slash).""" @@ -28,10 +49,21 @@ def normalize_root_path(root_path: Optional[str]) -> str: @contextlib.contextmanager def record_timing(metrics: dict[str, Any], key: str) -> Generator[None]: """ - Set timings[key] equal to the run time (in milliseconds) of the context body. + Set timings[key] equal to the run time (in seconds) of the context body. + + When there is an active recording trace span (i.e. this request is being + traced), also open a child OpenTelemetry span around the body so these + phases appear in the request's trace. Outside a traced request (tracing + disabled, or an excluded endpoint such as health checks and metrics + scrapes) no span is created, avoiding orphaned single-span traces. """ + if _tracer is not None and trace.get_current_span().is_recording(): + span = _tracer.start_as_current_span(_SPAN_NAMES.get(key, f"tiled.{key}")) + else: + span = contextlib.nullcontext() t0 = time.perf_counter() - yield + with span: + yield metrics[key]["dur"] += time.perf_counter() - t0 # Units: seconds diff --git a/tiled/storage.py b/tiled/storage.py index c0ae483ad..106e846b8 100644 --- a/tiled/storage.py +++ b/tiled/storage.py @@ -247,7 +247,7 @@ def _adbc_connection(self) -> "adbc_driver_manager.dbapi.Connection": def _connection_pool(self) -> "sqlalchemy.pool.QueuePool": from .server.metrics import monitor_db_pool - creator = self._adbc_connection.adbc_clone + creator = self._instrument_adbc_creator(self._adbc_connection.adbc_clone) if (self.dialect == "duckdb") or (":memory:" in self.uri): pool = sqlalchemy.pool.StaticPool(creator) else: @@ -258,6 +258,45 @@ def _connection_pool(self) -> "sqlalchemy.pool.QueuePool": return pool + def _instrument_adbc_creator(self, creator): + """Wrap the ADBC connection factory to emit OpenTelemetry spans for its queries. + + The storage database is accessed via ADBC, which the asyncpg instrumentation does + not cover, so its queries would otherwise be invisible in traces. A no-op if tracing + is off or the optional dbapi instrumentation is not installed. + """ + if not os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT"): + return creator + try: + from opentelemetry.instrumentation.dbapi import instrument_connection + except ImportError: + return creator + + dialect = self.dialect + # Report the database name as `db.name`, so the storage database appears as + # its own node in the service graph. ADBC follows the SQL standard naming, + # where a "catalog" is a database (and a "schema" is a namespace inside it), + # so on Postgres `adbc_current_catalog` is e.g. "tiled_storage". This has + # nothing to do with Tiled's catalog. Not every driver implements it (DuckDB + # raises), and the instrumentation only tolerates a missing attribute, not + # an error, so check once and leave `db.name` unset if it fails. + try: + self._adbc_connection.adbc_current_catalog + except Exception: + connection_attributes = {} + else: + connection_attributes = {"database": "adbc_current_catalog"} + + def instrumented_creator(): + return instrument_connection( + "tiled.storage", + creator(), + dialect, + connection_attributes=connection_attributes, + ) + + return instrumented_creator + def connect(self) -> "adbc_driver_manager.dbapi.Connection": "Get a connection from the pool." return self._connection_pool.connect()