Files
Celes Renata a72f336ad1 feat: Intelligence Pipeline v3 — full implementation
Multi-stage evidence-grounded inference architecture replacing the
monolithic 9B model extraction pipeline. CPU-first specialist services
handle routine extraction while the 9B vLLM model is preserved for
semantic adjudication of ambiguous cases.

Key components:
- Capability-aware inference gateway (OpenAI-compatible + Ollama)
- Endpoint registry with DB migrations and REST API
- Sentence-aware document segmenter (property tests)
- Deterministic financial parsing with offset integrity
- Symbol resolution with ambiguity detection
- Specialist service (GLiNER2, dynamic batching, K8s deployment)
- Company-specific sentiment (FinBERT, calibration)
- Retrieval-based novelty and duplicate detection
- Confidence calibration pipeline
- Deterministic routing engine (property tests)
- 9B adjudication layer with VRAM gating
- Stock-specific impact model (features, labels, baseline, trained)
- Pipeline orchestrator (state machine, queues, leases, feature flags)
- Bounded parallelism (async workers, semaphore, load shedding)
- Observability (tracing, metrics, alerts)
- Compatibility adapter (v3→v2 golden mapping tests)
- Shadow/canary promotion framework
- Active learning and fine-tuning pipeline

Test results: 1,161 tests pass, ruff lint clean.
All 282 spec tasks completed.
2026-07-13 02:14:59 +00:00

531 lines
20 KiB
Python

"""Unit tests for latency, throughput, token, CPU, GPU, and memory metrics.
Validates: Requirements 16.3, 16.4
"""
from __future__ import annotations
import pytest
from services.intelligence_pipeline_v3.evaluation.resource_metrics import (
ResourceEvaluationReport,
StageTimingRecord,
compute_cpu_metrics,
compute_efficiency_metrics,
compute_gpu_metrics,
compute_latency_metrics,
compute_memory_metrics,
compute_percentile,
compute_throughput_metrics,
compute_token_usage_metrics,
evaluate_resources,
)
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _record(
document_id: str = "doc-1",
stage_name: str = "extraction",
start_time: float = 0.0,
end_time: float = 1.0,
input_tokens: int = 100,
output_tokens: int = 50,
gpu_memory_mb: float = 0.0,
cpu_seconds: float = 0.5,
gpu_seconds: float = 0.0,
) -> StageTimingRecord:
return StageTimingRecord(
document_id=document_id,
stage_name=stage_name,
start_time=start_time,
end_time=end_time,
input_tokens=input_tokens,
output_tokens=output_tokens,
gpu_memory_mb=gpu_memory_mb,
cpu_seconds=cpu_seconds,
gpu_seconds=gpu_seconds,
)
# ---------------------------------------------------------------------------
# Percentile Helper Tests
# ---------------------------------------------------------------------------
class TestComputePercentile:
def test_single_value(self) -> None:
assert compute_percentile([5.0], 50.0) == 5.0
assert compute_percentile([5.0], 0.0) == 5.0
assert compute_percentile([5.0], 100.0) == 5.0
def test_two_values_median(self) -> None:
result = compute_percentile([1.0, 3.0], 50.0)
assert result == 2.0
def test_known_percentiles(self) -> None:
values = [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0, 10.0]
p50 = compute_percentile(values, 50.0)
assert abs(p50 - 5.5) < 1e-9
def test_unsorted_input(self) -> None:
values = [5.0, 1.0, 3.0, 2.0, 4.0]
p50 = compute_percentile(values, 50.0)
assert p50 == 3.0
def test_p0_returns_min(self) -> None:
values = [3.0, 1.0, 2.0]
assert compute_percentile(values, 0.0) == 1.0
def test_p100_returns_max(self) -> None:
values = [3.0, 1.0, 2.0]
assert compute_percentile(values, 100.0) == 3.0
def test_empty_raises(self) -> None:
with pytest.raises(ValueError, match="empty"):
compute_percentile([], 50.0)
def test_out_of_range_raises(self) -> None:
with pytest.raises(ValueError, match="between 0 and 100"):
compute_percentile([1.0], 101.0)
with pytest.raises(ValueError, match="between 0 and 100"):
compute_percentile([1.0], -1.0)
# ---------------------------------------------------------------------------
# StageTimingRecord Tests
# ---------------------------------------------------------------------------
class TestStageTimingRecord:
def test_duration(self) -> None:
r = _record(start_time=1.0, end_time=3.5)
assert r.duration_seconds == 2.5
def test_total_tokens(self) -> None:
r = _record(input_tokens=100, output_tokens=50)
assert r.total_tokens == 150
def test_frozen(self) -> None:
r = _record()
with pytest.raises(Exception):
r.document_id = "other" # type: ignore[misc]
# ---------------------------------------------------------------------------
# Latency Metrics
# ---------------------------------------------------------------------------
class TestLatencyMetrics:
def test_empty_records(self) -> None:
overall, per_stage = compute_latency_metrics([])
assert overall.count == 0
assert overall.mean == 0.0
assert per_stage == []
def test_single_document_single_stage(self) -> None:
records = [_record(start_time=0.0, end_time=2.0)]
overall, per_stage = compute_latency_metrics(records)
assert overall.count == 1
assert overall.mean == 2.0
assert overall.max == 2.0
assert overall.p50 == 2.0
assert len(per_stage) == 1
assert per_stage[0].stage_name == "extraction"
def test_multiple_documents(self) -> None:
records = [
_record(document_id="doc-1", start_time=0.0, end_time=1.0),
_record(document_id="doc-2", start_time=0.0, end_time=3.0),
_record(document_id="doc-3", start_time=0.0, end_time=2.0),
]
overall, _ = compute_latency_metrics(records)
assert overall.count == 3
assert overall.mean == 2.0
assert overall.max == 3.0
assert overall.min == 1.0
def test_multi_stage_document(self) -> None:
"""Document duration is from earliest start to latest end."""
records = [
_record(document_id="doc-1", stage_name="segmentation", start_time=0.0, end_time=1.0),
_record(document_id="doc-1", stage_name="extraction", start_time=1.0, end_time=3.0),
_record(document_id="doc-1", stage_name="sentiment", start_time=3.0, end_time=4.0),
]
overall, per_stage = compute_latency_metrics(records)
# Total document duration: 0 -> 4 = 4 seconds
assert overall.count == 1
assert overall.mean == 4.0
assert len(per_stage) == 3
def test_per_stage_breakdown(self) -> None:
records = [
_record(document_id="doc-1", stage_name="extraction", start_time=0.0, end_time=2.0),
_record(document_id="doc-2", stage_name="extraction", start_time=0.0, end_time=4.0),
_record(document_id="doc-1", stage_name="sentiment", start_time=2.0, end_time=2.5),
]
_, per_stage = compute_latency_metrics(records)
stage_map = {s.stage_name: s for s in per_stage}
assert stage_map["extraction"].invocation_count == 2
assert stage_map["extraction"].latency.mean == 3.0
assert stage_map["sentiment"].invocation_count == 1
# ---------------------------------------------------------------------------
# Throughput Metrics
# ---------------------------------------------------------------------------
class TestThroughputMetrics:
def test_empty_records(self) -> None:
result = compute_throughput_metrics([])
assert result.total_documents == 0
assert result.documents_per_minute == 0.0
def test_single_document(self) -> None:
records = [_record(start_time=0.0, end_time=60.0)]
result = compute_throughput_metrics(records)
assert result.total_documents == 1
assert result.total_wall_seconds == 60.0
assert abs(result.documents_per_minute - 1.0) < 1e-9
assert abs(result.documents_per_hour - 60.0) < 1e-9
def test_multiple_documents(self) -> None:
records = [
_record(document_id="doc-1", start_time=0.0, end_time=10.0),
_record(document_id="doc-2", start_time=5.0, end_time=15.0),
_record(document_id="doc-3", start_time=10.0, end_time=30.0),
]
result = compute_throughput_metrics(records)
assert result.total_documents == 3
assert result.total_wall_seconds == 30.0
# 3 docs / 30 seconds = 0.1 docs/sec = 6 docs/min
assert abs(result.documents_per_minute - 6.0) < 1e-9
assert abs(result.documents_per_hour - 360.0) < 1e-9
def test_zero_duration(self) -> None:
"""All records start and end at same time."""
records = [_record(start_time=5.0, end_time=5.0)]
result = compute_throughput_metrics(records)
assert result.documents_per_minute == 0.0
# ---------------------------------------------------------------------------
# Token Usage Metrics
# ---------------------------------------------------------------------------
class TestTokenUsageMetrics:
def test_empty_records(self) -> None:
result = compute_token_usage_metrics([])
assert result.total_tokens == 0
assert result.per_stage == {}
def test_single_record(self) -> None:
records = [_record(input_tokens=200, output_tokens=80)]
result = compute_token_usage_metrics(records)
assert result.total_input_tokens == 200
assert result.total_output_tokens == 80
assert result.total_tokens == 280
assert result.mean_input_tokens_per_document == 200.0
assert result.mean_output_tokens_per_document == 80.0
assert result.mean_total_tokens_per_document == 280.0
def test_multiple_documents_and_stages(self) -> None:
records = [
_record(document_id="doc-1", stage_name="extraction", input_tokens=100, output_tokens=50),
_record(document_id="doc-1", stage_name="sentiment", input_tokens=50, output_tokens=20),
_record(document_id="doc-2", stage_name="extraction", input_tokens=150, output_tokens=60),
]
result = compute_token_usage_metrics(records)
assert result.total_input_tokens == 300
assert result.total_output_tokens == 130
assert result.total_tokens == 430
# 2 documents
assert result.mean_input_tokens_per_document == 150.0
assert result.mean_output_tokens_per_document == 65.0
def test_per_stage_breakdown(self) -> None:
records = [
_record(document_id="doc-1", stage_name="extraction", input_tokens=100, output_tokens=50),
_record(document_id="doc-2", stage_name="extraction", input_tokens=200, output_tokens=100),
_record(document_id="doc-1", stage_name="sentiment", input_tokens=30, output_tokens=10),
]
result = compute_token_usage_metrics(records)
assert "extraction" in result.per_stage
assert "sentiment" in result.per_stage
ext = result.per_stage["extraction"]
assert ext.count == 2
assert ext.total_input_tokens == 300
assert ext.mean_input_tokens == 150.0
sent = result.per_stage["sentiment"]
assert sent.count == 1
assert sent.total_tokens == 40
# ---------------------------------------------------------------------------
# CPU Metrics
# ---------------------------------------------------------------------------
class TestCpuMetrics:
def test_empty_records(self) -> None:
result = compute_cpu_metrics([])
assert result.total_cpu_seconds == 0.0
def test_single_record(self) -> None:
records = [_record(cpu_seconds=2.5)]
result = compute_cpu_metrics(records)
assert result.total_cpu_seconds == 2.5
assert result.mean_cpu_seconds_per_document == 2.5
assert result.peak_cpu_seconds == 2.5
def test_multiple_documents(self) -> None:
records = [
_record(document_id="doc-1", stage_name="extraction", cpu_seconds=1.0),
_record(document_id="doc-1", stage_name="sentiment", cpu_seconds=0.5),
_record(document_id="doc-2", stage_name="extraction", cpu_seconds=3.0),
]
result = compute_cpu_metrics(records)
assert result.total_cpu_seconds == 4.5
# doc-1: 1.5, doc-2: 3.0
assert result.mean_cpu_seconds_per_document == 2.25
assert result.peak_cpu_seconds == 3.0
# ---------------------------------------------------------------------------
# GPU Metrics
# ---------------------------------------------------------------------------
class TestGpuMetrics:
def test_empty_records(self) -> None:
result = compute_gpu_metrics([])
assert result.total_gpu_seconds == 0.0
assert result.gpu_utilization_percent == 0.0
def test_no_gpu_usage(self) -> None:
records = [_record(gpu_seconds=0.0, gpu_memory_mb=0.0)]
result = compute_gpu_metrics(records)
assert result.total_gpu_seconds == 0.0
assert result.peak_gpu_memory_mb == 0.0
assert result.mean_gpu_memory_mb == 0.0
def test_with_gpu_usage(self) -> None:
records = [
_record(
document_id="doc-1",
start_time=0.0, end_time=10.0,
gpu_seconds=5.0, gpu_memory_mb=4096.0,
),
_record(
document_id="doc-2",
start_time=10.0, end_time=20.0,
gpu_seconds=3.0, gpu_memory_mb=8192.0,
),
]
result = compute_gpu_metrics(records)
assert result.total_gpu_seconds == 8.0
assert result.mean_gpu_seconds_per_document == 4.0
assert result.peak_gpu_memory_mb == 8192.0
assert result.mean_gpu_memory_mb == 6144.0
# 8 gpu-seconds / 20 wall-seconds = 40%
assert abs(result.gpu_utilization_percent - 40.0) < 1e-9
def test_utilization_capped_at_100(self) -> None:
"""Parallel GPU stages could sum to more than wall time."""
records = [
_record(
document_id="doc-1",
start_time=0.0, end_time=1.0,
gpu_seconds=5.0, gpu_memory_mb=1000.0,
),
]
result = compute_gpu_metrics(records)
assert result.gpu_utilization_percent == 100.0
# ---------------------------------------------------------------------------
# Memory Metrics
# ---------------------------------------------------------------------------
class TestMemoryMetrics:
def test_empty_records_no_samples(self) -> None:
result = compute_memory_metrics([])
assert result.peak_rss_memory_mb == 0.0
assert result.mean_working_set_mb == 0.0
def test_with_rss_samples(self) -> None:
records = [_record(gpu_memory_mb=5000.0)]
# RSS samples take precedence
result = compute_memory_metrics(records, rss_samples_mb=[100.0, 200.0, 300.0])
assert result.peak_rss_memory_mb == 300.0
assert result.mean_working_set_mb == 200.0
def test_fallback_to_gpu_memory(self) -> None:
records = [
_record(gpu_memory_mb=4096.0),
_record(gpu_memory_mb=8192.0),
]
result = compute_memory_metrics(records)
assert result.peak_rss_memory_mb == 8192.0
assert result.mean_working_set_mb == 6144.0
def test_zero_gpu_memory_treated_as_no_data(self) -> None:
records = [_record(gpu_memory_mb=0.0)]
result = compute_memory_metrics(records)
assert result.peak_rss_memory_mb == 0.0
assert result.mean_working_set_mb == 0.0
# ---------------------------------------------------------------------------
# Efficiency Metrics
# ---------------------------------------------------------------------------
class TestEfficiencyMetrics:
def test_empty_records(self) -> None:
result = compute_efficiency_metrics([])
assert result.tokens_per_second == 0.0
assert result.documents_per_gpu_second == 0.0
assert result.fast_path_fraction == 0.0
assert result.adjudication_fraction == 0.0
def test_tokens_per_second(self) -> None:
records = [
_record(
start_time=0.0, end_time=10.0,
input_tokens=500, output_tokens=500,
),
]
result = compute_efficiency_metrics(records)
# 1000 tokens / 10 seconds = 100 tokens/sec
assert abs(result.tokens_per_second - 100.0) < 1e-9
def test_documents_per_gpu_second(self) -> None:
records = [
_record(document_id="doc-1", gpu_seconds=2.0),
_record(document_id="doc-2", gpu_seconds=3.0),
]
result = compute_efficiency_metrics(records)
# 2 docs / 5 gpu-seconds = 0.4 docs/gpu-sec
assert abs(result.documents_per_gpu_second - 0.4) < 1e-9
def test_no_gpu_usage_infinite_docs(self) -> None:
"""When no GPU time, documents_per_gpu_second should be 0 (avoid division by zero)."""
records = [_record(gpu_seconds=0.0)]
result = compute_efficiency_metrics(records)
assert result.documents_per_gpu_second == 0.0
def test_fast_path_vs_adjudication_split(self) -> None:
records = [
_record(stage_name="extraction", cpu_seconds=2.0, gpu_seconds=0.0),
_record(stage_name="sentiment", cpu_seconds=1.0, gpu_seconds=0.0),
_record(stage_name="adjudication", cpu_seconds=0.5, gpu_seconds=3.0),
]
result = compute_efficiency_metrics(records)
assert result.fast_path_cpu_seconds == 3.0
assert result.adjudication_cpu_seconds == 0.5
assert result.fast_path_gpu_seconds == 0.0
assert result.adjudication_gpu_seconds == 3.0
# Fast: 3.0, Adj: 3.5, Total: 6.5
assert abs(result.fast_path_fraction - 3.0 / 6.5) < 1e-9
assert abs(result.adjudication_fraction - 3.5 / 6.5) < 1e-9
def test_adjudication_stage_detection(self) -> None:
"""Various adjudication stage name patterns should be detected."""
records = [
_record(stage_name="9b_adjudication", cpu_seconds=1.0, gpu_seconds=1.0),
_record(stage_name="semantic_adjudication", cpu_seconds=1.0, gpu_seconds=1.0),
_record(stage_name="my_adjudicator_stage", cpu_seconds=1.0, gpu_seconds=1.0),
]
result = compute_efficiency_metrics(records)
assert result.adjudication_cpu_seconds == 3.0
assert result.adjudication_gpu_seconds == 3.0
assert result.fast_path_cpu_seconds == 0.0
# ---------------------------------------------------------------------------
# Full Evaluation Report
# ---------------------------------------------------------------------------
class TestEvaluateResources:
def test_empty_records(self) -> None:
report = evaluate_resources([])
assert report.document_count == 0
assert report.latency.count == 0
assert report.throughput.total_documents == 0
def test_complete_report(self) -> None:
records = [
_record(
document_id="doc-1", stage_name="extraction",
start_time=0.0, end_time=2.0,
input_tokens=200, output_tokens=100,
cpu_seconds=1.0, gpu_seconds=0.5, gpu_memory_mb=4096.0,
),
_record(
document_id="doc-1", stage_name="adjudication",
start_time=2.0, end_time=5.0,
input_tokens=500, output_tokens=200,
cpu_seconds=0.2, gpu_seconds=2.5, gpu_memory_mb=8000.0,
),
_record(
document_id="doc-2", stage_name="extraction",
start_time=5.0, end_time=7.0,
input_tokens=180, output_tokens=90,
cpu_seconds=0.8, gpu_seconds=0.3, gpu_memory_mb=3500.0,
),
]
report = evaluate_resources(records)
assert isinstance(report, ResourceEvaluationReport)
assert report.document_count == 2
# Latency: doc-1 = 5s, doc-2 = 2s
assert report.latency.count == 2
assert report.latency.max == 5.0
assert report.latency.min == 2.0
# Throughput: 2 docs / 7 seconds
assert report.throughput.total_documents == 2
assert report.throughput.total_wall_seconds == 7.0
# Token usage
assert report.token_usage.total_input_tokens == 880
assert report.token_usage.total_output_tokens == 390
assert report.token_usage.total_tokens == 1270
# CPU
assert report.cpu.total_cpu_seconds == 2.0
# GPU
assert report.gpu.total_gpu_seconds == 3.3
assert report.gpu.peak_gpu_memory_mb == 8000.0
# Memory (fallback to GPU memory)
assert report.memory.peak_rss_memory_mb == 8000.0
# Efficiency
assert report.efficiency.adjudication_gpu_seconds == 2.5
assert report.efficiency.fast_path_cpu_seconds == 1.8
def test_with_rss_samples(self) -> None:
records = [_record(gpu_memory_mb=5000.0)]
report = evaluate_resources(records, rss_samples_mb=[512.0, 1024.0, 768.0])
assert report.memory.peak_rss_memory_mb == 1024.0
assert abs(report.memory.mean_working_set_mb - 768.0) < 1e-9
def test_per_stage_latency_sorted(self) -> None:
records = [
_record(stage_name="z_stage", start_time=0.0, end_time=1.0),
_record(stage_name="a_stage", start_time=1.0, end_time=2.0),
]
report = evaluate_resources(records)
stage_names = [s.stage_name for s in report.per_stage_latency]
assert stage_names == ["a_stage", "z_stage"]