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

216 lines
8.1 KiB
Python

"""Tests for canary module — Tasks 47-48."""
from __future__ import annotations
from services.intelligence_pipeline_v3.canary.influence import (
DivergenceRecord,
PromotionStatus,
SignalInfluenceConfig,
SignalInfluenceTracker,
)
from services.intelligence_pipeline_v3.canary.routing import (
CanaryConfig,
CanaryRouter,
RollbackReason,
)
class TestCanaryRouter:
"""Task 47: Canary compatibility outputs."""
def test_disabled_always_v2(self):
router = CanaryRouter(config=CanaryConfig(enabled=False))
assert not router.should_use_v3("doc-001")
def test_percentage_routing_deterministic(self):
config = CanaryConfig(enabled=True, percentage=50)
router = CanaryRouter(config=config)
result1 = router.should_use_v3("doc-001")
# Reset counters to test determinism
router2 = CanaryRouter(config=CanaryConfig(enabled=True, percentage=50))
result2 = router2.should_use_v3("doc-001")
assert result1 == result2
def test_trading_excluded_by_default(self):
config = CanaryConfig(enabled=True, percentage=100, exclude_trading=True)
router = CanaryRouter(config=config)
assert not router.should_use_v3("doc-001", is_trading_consumer=True)
assert router.should_use_v3("doc-001", is_trading_consumer=False)
def test_document_type_filter(self):
config = CanaryConfig(
enabled=True, percentage=100, document_types={"news", "filing"}
)
router = CanaryRouter(config=config)
assert router.should_use_v3("doc-001", document_type="news")
assert not router.should_use_v3("doc-002", document_type="transcript")
def test_rollback_on_error_rate(self):
config = CanaryConfig(enabled=True, percentage=20, max_error_rate=0.05)
router = CanaryRouter(config=config)
event = router.check_rollback(error_rate=0.10)
assert event is not None
assert event.reason == RollbackReason.ERROR_RATE
assert router.config.percentage == 0 # Rolled back
def test_rollback_on_latency(self):
config = CanaryConfig(
enabled=True, percentage=30, max_p95_latency_ms=3000
)
router = CanaryRouter(config=config)
event = router.check_rollback(p95_latency_ms=5000)
assert event is not None
assert event.reason == RollbackReason.LATENCY_THRESHOLD
def test_rollback_on_low_availability(self):
config = CanaryConfig(
enabled=True, percentage=10, min_availability=0.95
)
router = CanaryRouter(config=config)
event = router.check_rollback(availability=0.90)
assert event is not None
assert event.reason == RollbackReason.AVAILABILITY_THRESHOLD
def test_rollback_on_low_correctness(self):
config = CanaryConfig(
enabled=True, percentage=10, min_correctness=0.90
)
router = CanaryRouter(config=config)
event = router.check_rollback(correctness=0.85)
assert event is not None
assert event.reason == RollbackReason.CORRECTNESS_THRESHOLD
def test_no_rollback_when_healthy(self):
config = CanaryConfig(enabled=True, percentage=50)
router = CanaryRouter(config=config)
event = router.check_rollback(
error_rate=0.01,
p95_latency_ms=1000,
queue_saturation=0.3,
availability=0.99,
correctness=0.95,
)
assert event is None
def test_manual_rollback(self):
config = CanaryConfig(enabled=True, percentage=25)
router = CanaryRouter(config=config)
event = router.manual_rollback("operator requested")
assert event.reason == RollbackReason.MANUAL
assert event.previous_percentage == 25
assert router.config.percentage == 0
def test_rollback_preserves_audit_records(self):
"""Rollback changes routing, not stored v3 data."""
config = CanaryConfig(enabled=True, percentage=50)
router = CanaryRouter(config=config)
# Process some docs
router.should_use_v3("doc-001")
router.should_use_v3("doc-002")
# Rollback
router.manual_rollback()
# Audit records (rollback events) are preserved
assert len(router.rollback_events) == 1
def test_traffic_ratio(self):
config = CanaryConfig(enabled=True, percentage=100)
router = CanaryRouter(config=config)
for i in range(10):
router.should_use_v3(f"doc-{i}")
assert router.v3_traffic_ratio == 1.0
class TestSignalInfluence:
"""Task 48: Canary signal influence in paper trading."""
def test_start_paper_trading(self):
tracker = SignalInfluenceTracker(
config=SignalInfluenceConfig()
)
tracker.start_paper_trading()
assert tracker.promotion_status == PromotionStatus.PAPER_TRADING
def test_record_divergence(self):
tracker = SignalInfluenceTracker(
config=SignalInfluenceConfig(enabled=True)
)
tracker.record_signal(is_v3=True)
div = DivergenceRecord.create(
document_id="doc-001",
v2_recommendation={"direction": "buy"},
v3_recommendation={"direction": "sell"},
divergence_type="direction_opposite",
)
tracker.record_divergence(div)
assert tracker.divergence_rate == 1.0
def test_extraction_and_trading_separate(self):
"""Task 48.2: Separate extraction correctness from trading outcomes."""
tracker = SignalInfluenceTracker(
config=SignalInfluenceConfig(
enabled=True,
report_extraction_separately=True,
report_trading_separately=True,
)
)
tracker.update_extraction_metrics({"entity_f1": 0.92})
tracker.update_trading_metrics({"sharpe": 1.5})
summary = tracker.summary()
assert summary["extraction_metrics"]["entity_f1"] == 0.92
assert summary["trading_metrics"]["sharpe"] == 1.5
def test_approval_requires_owner(self):
config = SignalInfluenceConfig(
enabled=True,
require_owner_approval=True,
owner_id="owner-1",
)
tracker = SignalInfluenceTracker(config=config)
# Wrong approver
assert not tracker.approve("random-person")
# Right approver
assert tracker.approve("owner-1")
assert tracker.promotion_status == PromotionStatus.APPROVED
def test_approval_requires_all_divergences_reviewed(self):
config = SignalInfluenceConfig(
enabled=True,
require_owner_approval=False,
max_divergence_rate=1.0, # Allow any rate so we test review requirement
)
tracker = SignalInfluenceTracker(config=config)
tracker.record_signal(is_v3=True)
div = DivergenceRecord.create(
"doc-001", {"d": "buy"}, {"d": "sell"}, "opposite"
)
tracker.record_divergence(div)
# Cannot approve with unreviewed divergences
assert not tracker.approve("owner")
# Mark reviewed
div.reviewed = True
assert tracker.approve("owner")
def test_reject(self):
tracker = SignalInfluenceTracker(config=SignalInfluenceConfig())
tracker.reject("too many divergences")
assert tracker.promotion_status == PromotionStatus.REJECTED
def test_trading_outcomes_dont_override_correctness(self):
"""Requirement 16.10: Trading performance cannot override failed gates."""
config = SignalInfluenceConfig(
enabled=True,
require_owner_approval=False,
max_divergence_rate=0.10,
)
tracker = SignalInfluenceTracker(config=config)
# Simulate 10 v3 signals, 5 divergences (50% rate)
for i in range(10):
tracker.record_signal(is_v3=True)
for i in range(5):
tracker.record_divergence(
DivergenceRecord.create(f"doc-{i}", {}, {}, "opposite")
)
# Even if trading metrics are good, correctness gates fail
tracker.update_trading_metrics({"sharpe": 3.0})
assert not tracker.approve("owner")