All checks were successful
CI / lint (push) Successful in 33s
CI / test (push) Successful in 2m6s
Deploy / notebooks (push) Has been skipped
CI / notebooks-smoke (push) Successful in 1m31s
Deploy / zotero (push) Has been skipped
Deploy / docs (push) Has been skipped
Deploy / llm (push) Successful in 1m21s
Deploy / mc (push) Has been skipped
Deploy / api (push) Successful in 1m46s
Infra CI / docs (push) Successful in 21s
Infra CI / llm (push) Successful in 15s
Infra CI / mc (push) Successful in 14s
Deploy / report (push) Successful in 11s
Infra CI / zotero (push) Successful in 17s
Infra CI / notebooks (push) Successful in 50s
Infra CI / api (push) Successful in 19s
The P26 telemetry path never reached Prometheus: setup_meter_provider only installs an in-process PrometheusMetricReader, nothing served the registry (api and llm answered 404 on /metrics), the images never installed the perf extra, STACK_TELEMETRY was off, and no scrape target existed — so the data-pipelines request-rate panel was empty from the day it was written. Now: perf.middleware.instrument mounts GET /metrics (prometheus_client registry), the api and llm images install --extra perf, compose sets STACK_TELEMETRY=true on both, services.yml scrapes api:8000 and llm:8000, and both dashboards' request-rate panels query http_server_duration_milliseconds_count (what the FastAPI instrumentor emits; stack_http_server_requests_total never existed). Verified live: stack_llm_dispatch_total is queryable in Prometheus with job=llm.
238 lines
7.2 KiB
Python
238 lines
7.2 KiB
Python
"""Integration tests — runner with perf collector, middleware wiring."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import pytest
|
|
from opentelemetry import trace
|
|
from opentelemetry.sdk.trace import TracerProvider
|
|
from opentelemetry.sdk.trace.export import (
|
|
SimpleSpanProcessor,
|
|
SpanExporter,
|
|
SpanExportResult,
|
|
)
|
|
|
|
|
|
class _MemoryExporter(SpanExporter):
|
|
def __init__(self):
|
|
self.spans = []
|
|
|
|
def export(self, spans):
|
|
self.spans.extend(spans)
|
|
return SpanExportResult.SUCCESS
|
|
|
|
def shutdown(self):
|
|
pass
|
|
|
|
|
|
@pytest.fixture()
|
|
def trace_setup(monkeypatch):
|
|
"""Wire up real tracer for integration tests."""
|
|
monkeypatch.setenv("STACK_TELEMETRY", "false")
|
|
original = trace.get_tracer_provider()
|
|
exporter = _MemoryExporter()
|
|
provider = TracerProvider()
|
|
provider.add_span_processor(SimpleSpanProcessor(exporter))
|
|
trace.set_tracer_provider(provider)
|
|
|
|
import perf._tracer as tmod
|
|
|
|
tmod._provider_ready = True
|
|
yield exporter
|
|
tmod._provider_ready = False
|
|
provider.shutdown()
|
|
trace._TRACER_PROVIDER = original
|
|
trace._TRACER_PROVIDER_SET_ONCE._done = False
|
|
|
|
|
|
class TestRunnerIntegration:
|
|
"""run_pipeline emits spans via PipelineCollector."""
|
|
|
|
def test_pipeline_creates_spans(self, trace_setup):
|
|
import polars as pl
|
|
|
|
from aco.pipe.runner import run_pipeline
|
|
|
|
# Simple two-step pipeline
|
|
def step_a(input_layer__eligibility):
|
|
return input_layer__eligibility.select("person_id")
|
|
|
|
def step_b(core__step_a):
|
|
return core__step_a
|
|
|
|
eligibility = pl.DataFrame({"person_id": ["P001", "P002"]})
|
|
|
|
def load(ref):
|
|
if ref == "input_layer.eligibility":
|
|
return eligibility
|
|
raise KeyError(ref)
|
|
|
|
cache = run_pipeline(
|
|
[
|
|
("core.step_a", step_a),
|
|
("core.step_b", step_b),
|
|
],
|
|
load,
|
|
)
|
|
|
|
assert "core.step_a" in cache
|
|
assert "core.step_b" in cache
|
|
assert len(cache["core.step_a"]) == 2
|
|
|
|
# Verify spans were created
|
|
names = [s.name for s in trace_setup.spans]
|
|
assert "pipeline.run" in names
|
|
assert "step.core.step_a" in names
|
|
assert "step.core.step_b" in names
|
|
|
|
def test_pipeline_records_row_counts(self, trace_setup):
|
|
import polars as pl
|
|
|
|
from aco.pipe.runner import run_pipeline
|
|
|
|
def step_one(input_layer__data):
|
|
return input_layer__data
|
|
|
|
data = pl.DataFrame({"id": list(range(50))})
|
|
|
|
run_pipeline(
|
|
[("core.step_one", step_one)],
|
|
lambda ref: data,
|
|
)
|
|
|
|
step_span = next(s for s in trace_setup.spans if s.name == "step.core.step_one")
|
|
assert step_span.attributes["rows.out"] == 50
|
|
|
|
def test_pipeline_without_perf_still_works(self, monkeypatch):
|
|
"""Verify graceful degradation if perf collector import fails."""
|
|
import polars as pl
|
|
|
|
from aco.pipe.runner import run_pipeline
|
|
|
|
def step_x(input_layer__table):
|
|
return input_layer__table
|
|
|
|
data = pl.DataFrame({"col": [1, 2, 3]})
|
|
cache = run_pipeline(
|
|
[("ns.step_x", step_x)],
|
|
lambda ref: data,
|
|
)
|
|
assert len(cache["ns.step_x"]) == 3
|
|
|
|
def test_pipeline_name_inference(self):
|
|
from aco.pipe.runner import _infer_pipeline_name
|
|
|
|
assert (
|
|
_infer_pipeline_name([("readmissions._int_enc", lambda: None)])
|
|
== "readmissions"
|
|
)
|
|
assert _infer_pipeline_name([("core.encounter", lambda: None)]) == "core"
|
|
assert _infer_pipeline_name([]) == "unknown"
|
|
|
|
|
|
class TestMiddleware:
|
|
"""perf.middleware.instrument is a safe no-op when disabled."""
|
|
|
|
def test_instrument_noop_when_disabled(self, monkeypatch):
|
|
monkeypatch.setenv("STACK_TELEMETRY", "false")
|
|
from unittest.mock import MagicMock
|
|
|
|
from perf.middleware import instrument
|
|
|
|
app = MagicMock()
|
|
instrument(app) # should not raise
|
|
|
|
def test_instrument_enabled_with_real_app(self, monkeypatch):
|
|
"""When telemetry enabled, instrument runs the OTel path."""
|
|
monkeypatch.setenv("STACK_TELEMETRY", "true")
|
|
from unittest.mock import MagicMock
|
|
|
|
from perf.middleware import instrument
|
|
|
|
app = MagicMock()
|
|
instrument(app) # exercises lines 26-31
|
|
|
|
def test_instrument_enabled_catches_import_failure(self, monkeypatch):
|
|
"""When OTel instrumentor is unavailable, instrument is a no-op."""
|
|
monkeypatch.setenv("STACK_TELEMETRY", "true")
|
|
import builtins
|
|
from unittest.mock import MagicMock
|
|
|
|
real_import = builtins.__import__
|
|
|
|
def fail_otel(name, *args, **kwargs):
|
|
if "opentelemetry.instrumentation.fastapi" in name:
|
|
raise ImportError("no otel instrumentor")
|
|
return real_import(name, *args, **kwargs)
|
|
|
|
monkeypatch.setattr(builtins, "__import__", fail_otel)
|
|
from perf.middleware import instrument
|
|
|
|
app = MagicMock()
|
|
instrument(app) # should not raise
|
|
|
|
def test_server_import(self):
|
|
"""Verify server.py imports cleanly with perf wiring."""
|
|
from api.server import app
|
|
|
|
assert app is not None
|
|
assert app.title == "stack"
|
|
|
|
|
|
class TestMetricsRoute:
|
|
"""#579: instrument() serves the in-process Prometheus registry at
|
|
/metrics when telemetry is enabled — the only way stack_* series
|
|
reach Prometheus (the meter's PrometheusMetricReader never pushes)."""
|
|
|
|
def test_metrics_route_served_when_enabled(self, monkeypatch):
|
|
monkeypatch.setenv("STACK_TELEMETRY", "true")
|
|
from fastapi import FastAPI
|
|
from fastapi.testclient import TestClient
|
|
|
|
from perf.middleware import instrument
|
|
|
|
app = FastAPI()
|
|
instrument(app)
|
|
r = TestClient(app).get("/metrics")
|
|
assert r.status_code == 200
|
|
assert r.headers["content-type"].startswith("text/plain")
|
|
assert "# HELP" in r.text or "# TYPE" in r.text
|
|
|
|
def test_metrics_route_absent_when_disabled(self, monkeypatch):
|
|
monkeypatch.setenv("STACK_TELEMETRY", "false")
|
|
from fastapi import FastAPI
|
|
from fastapi.testclient import TestClient
|
|
|
|
from perf.middleware import instrument
|
|
|
|
app = FastAPI()
|
|
instrument(app)
|
|
assert TestClient(app).get("/metrics").status_code == 404
|
|
|
|
def test_mount_metrics_tolerates_missing_prometheus_client(self, monkeypatch):
|
|
import builtins
|
|
|
|
real = builtins.__import__
|
|
|
|
def fake(name, *a, **k):
|
|
if name.startswith("prometheus_client"):
|
|
raise ImportError(name)
|
|
return real(name, *a, **k)
|
|
|
|
monkeypatch.setattr(builtins, "__import__", fake)
|
|
from unittest.mock import MagicMock
|
|
|
|
from perf.middleware import _mount_metrics
|
|
|
|
app = MagicMock()
|
|
_mount_metrics(app) # no route added, no exception
|
|
app.add_api_route.assert_not_called()
|
|
|
|
def test_mount_metrics_tolerates_route_failure(self):
|
|
from unittest.mock import MagicMock
|
|
|
|
from perf.middleware import _mount_metrics
|
|
|
|
app = MagicMock()
|
|
app.add_api_route.side_effect = RuntimeError("boom")
|
|
_mount_metrics(app) # swallowed
|