Add `perf` as the 12th skinny package with full OpenTelemetry instrumentation for pipeline traces, Prometheus metrics, and Loki-correlated logs. - src/perf/: TracerProvider, MeterProvider, LoggerProvider bridge, PipelineCollector, psutil system metrics, FileExporter fallback, FastAPI middleware, auto-file Gitea issues on step/test failure - infra/otel/: OTel Collector fan-out (traces→Jaeger, metrics→Prometheus, logs→Loki), all services migrated from jaeger:4317 to collector - infra/grafana/dashboards/: 8-panel pipeline performance dashboard - Pipeline runner auto-instruments all 189 steps with zero-cost MockTracer/MockSpan no-ops when telemetry is disabled - 82 new tests, 12,055 total passing Closes #203, closes #204, closes #205, closes #206, closes #207, closes #208, closes #209, closes #210, closes #211, closes #212, closes #213
188 lines
5.9 KiB
Python
188 lines
5.9 KiB
Python
"""Tests for StepCollector, PipelineCollector, and system snapshots."""
|
|
|
|
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 + mock telemetry enabled."""
|
|
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)
|
|
|
|
# Make perf.tracer() return real tracers from this 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 TestSystemSnapshot:
|
|
"""psutil-based system snapshots."""
|
|
|
|
def test_take_snapshot(self):
|
|
from perf._system import take_snapshot
|
|
|
|
snap = take_snapshot()
|
|
assert snap.rss_mb > 0
|
|
assert snap.disk_read_mb >= 0
|
|
assert snap.disk_write_mb >= 0
|
|
|
|
def test_delta(self):
|
|
from perf._system import SystemSnapshot, delta
|
|
|
|
before = SystemSnapshot(
|
|
rss_mb=100.0,
|
|
cpu_percent=10.0,
|
|
disk_read_mb=50.0,
|
|
disk_write_mb=20.0,
|
|
)
|
|
after = SystemSnapshot(
|
|
rss_mb=120.0,
|
|
cpu_percent=25.0,
|
|
disk_read_mb=55.0,
|
|
disk_write_mb=22.0,
|
|
)
|
|
d = delta(before, after)
|
|
assert d["memory.delta_mb"] == 20.0
|
|
assert d["memory.rss_mb"] == 120.0
|
|
assert d["cpu_percent"] == 25.0
|
|
assert d["disk.read_delta_mb"] == 5.0
|
|
assert d["disk.write_delta_mb"] == 2.0
|
|
|
|
|
|
class TestPipelineCollector:
|
|
"""PipelineCollector creates proper span hierarchy."""
|
|
|
|
def test_step_creates_span(self, trace_setup):
|
|
exporter = trace_setup
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("test_pipeline") as pc:
|
|
with pc.step("step_a") as s:
|
|
s.record(rows=42, columns=["a", "b"])
|
|
|
|
names = [s.name for s in exporter.spans]
|
|
assert "step.step_a" in names
|
|
assert "pipeline.run" in names
|
|
|
|
def test_step_records_attributes(self, trace_setup):
|
|
exporter = trace_setup
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("test_pipeline") as pc:
|
|
with pc.step("step_b") as s:
|
|
s.record(rows=100, columns=["x", "y", "z"], cache_hit=True)
|
|
|
|
step_span = next(s for s in exporter.spans if s.name == "step.step_b")
|
|
assert step_span.attributes["rows.out"] == 100
|
|
assert step_span.attributes["columns.count"] == 3
|
|
assert step_span.attributes["cache.hit"] is True
|
|
|
|
def test_multiple_steps(self, trace_setup):
|
|
exporter = trace_setup
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("multi") as pc:
|
|
with pc.step("first") as s:
|
|
s.record(rows=10)
|
|
with pc.step("second") as s:
|
|
s.record(rows=20)
|
|
with pc.step("third") as s:
|
|
s.record(rows=30)
|
|
|
|
root = next(s for s in exporter.spans if s.name == "pipeline.run")
|
|
assert root.attributes["steps.total"] == 3
|
|
|
|
def test_step_timing(self, trace_setup):
|
|
import time
|
|
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("timed") as pc:
|
|
with pc.step("slow") as s:
|
|
time.sleep(0.01)
|
|
s.record(rows=1)
|
|
|
|
assert s.result.duration_s >= 0.01
|
|
|
|
def test_step_system_delta(self, trace_setup):
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("sys") as pc:
|
|
with pc.step("mem_test") as s:
|
|
# Allocate some memory to create a delta
|
|
_data = bytearray(1024 * 1024) # 1MB
|
|
s.record(rows=1)
|
|
|
|
assert "memory.rss_mb" in s.result.system_delta
|
|
|
|
def test_cache_hit_flag(self, trace_setup):
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("cache_test") as pc:
|
|
with pc.step("cached") as s:
|
|
s.record(rows=50, cache_hit=True)
|
|
with pc.step("uncached") as s2:
|
|
s2.record(rows=50, cache_hit=False)
|
|
|
|
root = next(s for s in trace_setup.spans if s.name == "pipeline.run")
|
|
assert root.attributes["steps.cached"] == 1
|
|
|
|
def test_load_span(self, trace_setup):
|
|
exporter = trace_setup
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("load_test") as pc:
|
|
with pc.step("step_load") as s:
|
|
with pc.load("core.encounter") as lc:
|
|
lc.record(rows=500)
|
|
s.record(rows=500)
|
|
|
|
load_span = next(s for s in exporter.spans if s.name == "load.core.encounter")
|
|
assert load_span.attributes["rows"] == 500
|
|
assert load_span.attributes["table"] == "core.encounter"
|
|
|
|
def test_save_span(self, trace_setup):
|
|
exporter = trace_setup
|
|
from perf.collector import PipelineCollector
|
|
|
|
with PipelineCollector("save_test") as pc:
|
|
with pc.step("step_save") as s:
|
|
s.record(rows=200)
|
|
with pc.save("readmissions.output") as sc:
|
|
sc.record(rows=200)
|
|
|
|
save_span = next(
|
|
s for s in exporter.spans if s.name == "save.readmissions.output"
|
|
)
|
|
assert save_span.attributes["rows"] == 200
|