- Scheduler: lower stale threshold 240→30 min, batch limit 100→500, TTL 14400→3600 - Prediction snapshot: add 24h market_snapshots time-window fallback - Aggregation: add normalize_impact_scores() z-score normalization - Helm: signal-engine replicas → 0 (idle when dual pipeline disabled) - Quality gate: max_snapshot_age_hours 24→48 - Add backfill script for NULL price_at_prediction snapshots - Add PBT bug condition and preservation tests (14 tests)
467 lines
17 KiB
Python
467 lines
17 KiB
Python
"""Property-based tests for pipeline health preservation properties.
|
|
|
|
Feature: pipeline-health-fixes
|
|
|
|
These tests encode the CURRENT correct behavior for non-bug inputs.
|
|
They MUST PASS on unfixed code — confirming baseline behavior is preserved.
|
|
|
|
Preservation properties tested:
|
|
1. Normal Doc Processing — documents < threshold age untouched by recovery
|
|
2. Primary Price Path — valid primary price used without fallback
|
|
3. Balanced Sentiment — symmetric impact_scores produce neutral output
|
|
4. Quality Gate Normal — snapshots < max age meeting criteria pass
|
|
5. Failed Extraction Independence — retry_failed_extractions behavior unchanged
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import datetime, timedelta, timezone
|
|
from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
|
import pytest
|
|
from hypothesis import given, settings
|
|
from hypothesis import strategies as st
|
|
|
|
from services.aggregation.scoring import (
|
|
SignalWeight,
|
|
WeightedSignal,
|
|
weighted_sentiment_average,
|
|
)
|
|
from services.trading.model_quality_gate import (
|
|
QualityGateConfig,
|
|
evaluate_quality_gate,
|
|
)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Property 1: Normal Doc Processing Preservation
|
|
# Documents younger than the stale threshold are NOT touched by recovery.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPreservationNormalDocProcessing:
|
|
"""Preservation: Normal Document Processing
|
|
|
|
Documents that enter 'parsed' status and are younger than the current
|
|
stale threshold are never touched by recover_stale_documents().
|
|
On unfixed code, STALE_PARSED_THRESHOLD_MINUTES=240, so any document
|
|
younger than 240 min is untouched.
|
|
|
|
We test with documents < 30 min old which are well within the threshold
|
|
on both unfixed (240 min) and fixed (30 min) code.
|
|
|
|
**Validates: Requirements 3.1, 3.7**
|
|
"""
|
|
|
|
@pytest.mark.asyncio
|
|
@given(
|
|
doc_age_minutes=st.floats(min_value=0.0, max_value=29.0, allow_nan=False),
|
|
)
|
|
@settings(max_examples=100)
|
|
async def test_young_documents_not_recovered(self, doc_age_minutes: float):
|
|
"""Property: for all documents with age < 30 min, recovery does NOT enqueue them.
|
|
|
|
The SQL WHERE clause filters on `updated_at < NOW() - INTERVAL threshold`,
|
|
so documents younger than the threshold won't appear in the query results.
|
|
We mock pool.fetch to return empty (simulating the DB correctly filtering
|
|
out young documents) and verify 0 are enqueued.
|
|
"""
|
|
from services.scheduler.app import (
|
|
STALE_PARSED_THRESHOLD_MINUTES,
|
|
recover_stale_documents,
|
|
)
|
|
|
|
# The current threshold is 240 min on unfixed code.
|
|
# Documents < 30 min old are well below ANY threshold (240 or 30),
|
|
# so the DB query returns nothing for them.
|
|
assert doc_age_minutes < STALE_PARSED_THRESHOLD_MINUTES, (
|
|
f"Test assumes doc_age_minutes ({doc_age_minutes}) < threshold "
|
|
f"({STALE_PARSED_THRESHOLD_MINUTES})"
|
|
)
|
|
|
|
# Mock pool.fetch to return empty — simulating DB filtering out young docs
|
|
pool = AsyncMock()
|
|
pool.fetch = AsyncMock(return_value=[])
|
|
pool.execute = AsyncMock()
|
|
|
|
rds = AsyncMock()
|
|
rds.set = AsyncMock(return_value=True)
|
|
rds.rpush = AsyncMock()
|
|
|
|
result = await recover_stale_documents(pool, rds)
|
|
|
|
# No documents should be enqueued
|
|
assert result == 0, (
|
|
f"Young documents (age={doc_age_minutes:.1f} min) should not be "
|
|
f"recovered, but {result} were enqueued"
|
|
)
|
|
# Redis rpush should never have been called
|
|
rds.rpush.assert_not_called()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Property 2: Primary Price Path Preservation
|
|
# When fetch_latest_close_price returns a valid price, no fallback is invoked.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPreservationPrimaryPricePath:
|
|
"""Preservation: Primary Price Path
|
|
|
|
When fetch_latest_close_price() returns a valid price, that price is used
|
|
directly and no fallback (positions table or market_snapshots 24h) is
|
|
invoked.
|
|
|
|
**Validates: Requirements 3.2**
|
|
"""
|
|
|
|
@pytest.mark.asyncio
|
|
@given(
|
|
primary_price=st.floats(
|
|
min_value=1.0, max_value=5000.0, allow_nan=False, allow_infinity=False
|
|
),
|
|
)
|
|
@settings(max_examples=100)
|
|
async def test_valid_primary_price_used_directly(self, primary_price: float):
|
|
"""Property: for all tickers where primary price exists, result equals primary price.
|
|
|
|
When fetch_latest_close_price returns a valid float, the snapshot's
|
|
price_at_prediction equals that value without any fallback invoked.
|
|
"""
|
|
from services.validation.prediction_snapshot import (
|
|
create_prediction_snapshot,
|
|
)
|
|
|
|
pool = AsyncMock()
|
|
|
|
# Track fetchrow calls to verify no fallback query is made
|
|
call_count = {"fetchrow": 0}
|
|
|
|
async def mock_fetchrow(*args, **kwargs):
|
|
call_count["fetchrow"] += 1
|
|
# After the first price lookups, return None for sector ETF
|
|
return None
|
|
|
|
pool.fetchrow = AsyncMock(side_effect=mock_fetchrow)
|
|
pool.execute = AsyncMock()
|
|
|
|
# Setup transaction context manager mock
|
|
conn_mock = AsyncMock()
|
|
conn_mock.execute = AsyncMock()
|
|
tx_mock = AsyncMock()
|
|
tx_mock.__aenter__ = AsyncMock(return_value=None)
|
|
tx_mock.__aexit__ = AsyncMock(return_value=None)
|
|
conn_mock.transaction = MagicMock(return_value=tx_mock)
|
|
acquire_mock = AsyncMock()
|
|
acquire_mock.__aenter__ = AsyncMock(return_value=conn_mock)
|
|
acquire_mock.__aexit__ = AsyncMock(return_value=None)
|
|
pool.acquire = MagicMock(return_value=acquire_mock)
|
|
|
|
# Build minimal mocks
|
|
recommendation = MagicMock()
|
|
recommendation.ticker = "AAPL"
|
|
recommendation.generated_at = datetime.now(tz=timezone.utc)
|
|
recommendation.time_horizon = "7d"
|
|
recommendation.action.value = "buy"
|
|
recommendation.mode.value = "paper_eligible"
|
|
recommendation.confidence = 0.7
|
|
|
|
trend_summary = MagicMock()
|
|
trend_summary.market_context = None
|
|
trend_summary.window.value = "7d"
|
|
trend_summary.trend_direction.value = "bullish"
|
|
trend_summary.trend_strength = 0.7
|
|
trend_summary.contradiction_score = 0.1
|
|
trend_summary.p_bull = 0.7
|
|
|
|
# Patch fetch_latest_close_price to return the primary price
|
|
# This simulates the primary price lookup succeeding
|
|
async def mock_fetch_latest_close_price(pool_arg, ticker):
|
|
if ticker == "AAPL":
|
|
return primary_price
|
|
# SPY and sector ETF can also have prices
|
|
return primary_price * 0.5
|
|
|
|
with patch(
|
|
"services.validation.prediction_snapshot.fetch_latest_close_price",
|
|
new_callable=AsyncMock,
|
|
side_effect=mock_fetch_latest_close_price,
|
|
):
|
|
snapshot = await create_prediction_snapshot(
|
|
pool=pool,
|
|
recommendation=recommendation,
|
|
trend_summary=trend_summary,
|
|
evidence_signals=[],
|
|
evidence_docs=[],
|
|
)
|
|
|
|
# The snapshot price should equal the primary price exactly
|
|
assert snapshot.price_at_prediction == primary_price, (
|
|
f"Primary price path broken: expected {primary_price}, "
|
|
f"got {snapshot.price_at_prediction}"
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Property 3: Balanced Sentiment Preservation
|
|
# Tickers with symmetric impact_score distributions produce neutral signals.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPreservationBalancedSentiment:
|
|
"""Preservation: Balanced Sentiment
|
|
|
|
Tickers with symmetric impact_score distributions (mean ≈ 0) produce
|
|
neutral signals. The weighted_sentiment_average() with equal positive
|
|
and negative sentiments and impact_scores centered around 0 should
|
|
produce an output near 0.
|
|
|
|
**Validates: Requirements 3.3**
|
|
"""
|
|
|
|
@given(
|
|
# Generate symmetric impact_scores centered at 0
|
|
stddev=st.floats(min_value=0.1, max_value=1.0, allow_nan=False),
|
|
n_pairs=st.integers(min_value=5, max_value=50),
|
|
)
|
|
@settings(max_examples=100)
|
|
def test_symmetric_impacts_produce_neutral_output(
|
|
self,
|
|
stddev: float,
|
|
n_pairs: int,
|
|
):
|
|
"""Property: for all impact_score sets with mean ≈ 0, output is near 0.
|
|
|
|
Generate WeightedSignals with impact_scores drawn symmetrically around 0
|
|
(for each +x there is a -x) and equal positive/negative sentiments.
|
|
Verify weighted_sentiment_average() produces near-zero output.
|
|
"""
|
|
signals: list[WeightedSignal] = []
|
|
|
|
for i in range(n_pairs):
|
|
# Create symmetric pairs of impact_scores
|
|
# Use positive impact_score values (the scoring formula uses w = combined * impact_score)
|
|
# With symmetric sentiments and equal positive impact_scores, output should be 0
|
|
impact_val = 0.1 + (i * stddev / n_pairs) # All positive, equal for both
|
|
|
|
weight = SignalWeight(
|
|
recency=0.8,
|
|
credibility=0.7,
|
|
novelty_bonus=0.1,
|
|
confidence_gate=1.0,
|
|
market_ctx_multiplier=1.0,
|
|
combined=0.5,
|
|
)
|
|
|
|
# Positive sentiment signal
|
|
signals.append(
|
|
WeightedSignal(
|
|
document_id=f"doc_pos_{i}",
|
|
weight=weight,
|
|
sentiment_value=1.0,
|
|
impact_score=impact_val,
|
|
)
|
|
)
|
|
# Negative sentiment signal with same impact_score
|
|
signals.append(
|
|
WeightedSignal(
|
|
document_id=f"doc_neg_{i}",
|
|
weight=weight,
|
|
sentiment_value=-1.0,
|
|
impact_score=impact_val,
|
|
)
|
|
)
|
|
|
|
avg = weighted_sentiment_average(signals)
|
|
|
|
# With perfectly symmetric positive/negative sentiments and equal
|
|
# impact_scores, the output should be exactly 0
|
|
assert abs(avg) < 1e-9, (
|
|
f"Balanced sentiment broken: expected ≈ 0, got {avg:.6f} "
|
|
f"with {n_pairs} pairs and stddev={stddev:.3f}"
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Property 4: Quality Gate Normal Preservation
|
|
# Snapshots < max_snapshot_age_hours meeting thresholds pass the quality gate.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPreservationQualityGateNormal:
|
|
"""Preservation: Quality Gate Normal
|
|
|
|
Snapshots younger than the max age threshold that meet all metric criteria
|
|
continue to pass the quality gate and return passed=True.
|
|
|
|
On unfixed code, max_snapshot_age_hours=24, so we test with ages in [0, 24h).
|
|
This ensures the test passes on both unfixed and fixed code.
|
|
|
|
**Validates: Requirements 3.5, 3.6**
|
|
"""
|
|
|
|
@pytest.mark.asyncio
|
|
@given(
|
|
age_hours=st.floats(min_value=0.1, max_value=23.9, allow_nan=False),
|
|
win_rate=st.floats(min_value=0.53, max_value=0.85, allow_nan=False),
|
|
ic=st.floats(min_value=0.03, max_value=0.3, allow_nan=False),
|
|
prediction_count=st.integers(min_value=100, max_value=10000),
|
|
)
|
|
@settings(max_examples=100)
|
|
async def test_young_snapshots_meeting_criteria_pass(
|
|
self,
|
|
age_hours: float,
|
|
win_rate: float,
|
|
ic: float,
|
|
prediction_count: int,
|
|
):
|
|
"""Property: for all snapshot ages in [0, 24h) meeting metric criteria,
|
|
quality gate returns passed=True.
|
|
|
|
We use ages < 24h which is within the threshold on BOTH unfixed (24h)
|
|
and fixed (48h) code.
|
|
"""
|
|
now = datetime.now(tz=timezone.utc)
|
|
snapshot_time = now - timedelta(hours=age_hours)
|
|
|
|
# Build a snapshot row that meets all default thresholds
|
|
fake_snapshot_row = {
|
|
"id": uuid.uuid4(),
|
|
"generated_at": snapshot_time,
|
|
"prediction_count": prediction_count,
|
|
"win_rate": win_rate,
|
|
"directional_accuracy": 0.58,
|
|
"information_coefficient": ic,
|
|
"rank_information_coefficient": 0.04,
|
|
"avg_return": 0.02,
|
|
"avg_excess_return_vs_spy": 0.01, # >= 0.0 threshold
|
|
"avg_excess_return_vs_sector": 0.005,
|
|
"calibration_error": 0.10, # <= 0.15 threshold
|
|
"brier_score": 0.20,
|
|
"buy_win_rate": 0.62,
|
|
"sell_win_rate": 0.58,
|
|
"hold_win_rate": 0.55,
|
|
}
|
|
|
|
pool = AsyncMock()
|
|
pool.fetchrow = AsyncMock(return_value=fake_snapshot_row)
|
|
pool.execute = AsyncMock()
|
|
pool.fetchval = AsyncMock(return_value=None)
|
|
|
|
config = QualityGateConfig()
|
|
|
|
with patch(
|
|
"services.trading.model_quality_gate._store_gate_result",
|
|
new_callable=AsyncMock,
|
|
):
|
|
result = await evaluate_quality_gate(pool, config=config)
|
|
|
|
assert result.passed is True, (
|
|
f"Quality gate should PASS for {age_hours:.1f}h-old snapshot "
|
|
f"meeting all criteria, but got: {result.reason}"
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Property 5: Failed Extraction Independence Preservation
|
|
# retry_failed_extractions() behavior is unaffected by stale doc recovery changes.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPreservationFailedExtractionIndependence:
|
|
"""Preservation: Failed Extraction Independence
|
|
|
|
retry_failed_extractions() processes documents in 'extraction_failed' status
|
|
using its own threshold (EXTRACTION_FAILED_RETRY_MINUTES=60) and logic.
|
|
Its behavior is independent of changes to recover_stale_documents().
|
|
|
|
**Validates: Requirements 3.7**
|
|
"""
|
|
|
|
@pytest.mark.asyncio
|
|
@given(
|
|
num_failed_docs=st.integers(min_value=1, max_value=100),
|
|
)
|
|
@settings(max_examples=100)
|
|
async def test_retry_failed_extractions_processes_failed_docs(
|
|
self,
|
|
num_failed_docs: int,
|
|
):
|
|
"""Property: for all documents in extraction_failed status, retry logic unchanged.
|
|
|
|
Verify that retry_failed_extractions:
|
|
1. Queries documents with status='extraction_failed'
|
|
2. Uses EXTRACTION_FAILED_RETRY_MINUTES (60) as the age threshold
|
|
3. Enqueues them via _enqueue_if_new
|
|
4. Resets their status to 'parsed'
|
|
5. Deletes failed intelligence rows
|
|
"""
|
|
from services.scheduler.app import (
|
|
EXTRACTION_FAILED_RETRY_MINUTES,
|
|
retry_failed_extractions,
|
|
)
|
|
|
|
# Verify the retry threshold hasn't changed
|
|
assert EXTRACTION_FAILED_RETRY_MINUTES == 60, (
|
|
f"EXTRACTION_FAILED_RETRY_MINUTES changed from 60 to "
|
|
f"{EXTRACTION_FAILED_RETRY_MINUTES} — this should be preserved"
|
|
)
|
|
|
|
now = datetime.now(tz=timezone.utc)
|
|
old_time = now - timedelta(minutes=EXTRACTION_FAILED_RETRY_MINUTES + 10)
|
|
|
|
# Generate fake failed document rows
|
|
fake_rows = []
|
|
for i in range(num_failed_docs):
|
|
fake_rows.append({
|
|
"id": uuid.uuid4(),
|
|
"document_type": "news" if i % 3 != 0 else "macro_event",
|
|
"ticker": f"TICK{i}",
|
|
"updated_at": old_time,
|
|
})
|
|
|
|
pool = AsyncMock()
|
|
pool.fetch = AsyncMock(return_value=fake_rows)
|
|
pool.execute = AsyncMock()
|
|
|
|
rds = AsyncMock()
|
|
rds.set = AsyncMock(return_value=True) # All enqueues succeed
|
|
rds.rpush = AsyncMock()
|
|
|
|
result = await retry_failed_extractions(pool, rds)
|
|
|
|
# All docs should be enqueued
|
|
assert result == num_failed_docs, (
|
|
f"Expected {num_failed_docs} docs retried, got {result}"
|
|
)
|
|
|
|
# Verify the SQL query used the correct threshold
|
|
fetch_call = pool.fetch.call_args
|
|
sql_query = fetch_call[0][0]
|
|
threshold_param = fetch_call[0][1]
|
|
|
|
assert "extraction_failed" in sql_query, (
|
|
"retry_failed_extractions should query for 'extraction_failed' status"
|
|
)
|
|
assert threshold_param == EXTRACTION_FAILED_RETRY_MINUTES, (
|
|
f"Expected threshold param {EXTRACTION_FAILED_RETRY_MINUTES}, "
|
|
f"got {threshold_param}"
|
|
)
|
|
|
|
# Verify pool.execute was called for both DELETE and UPDATE
|
|
assert pool.execute.call_count == 2, (
|
|
f"Expected 2 pool.execute calls (DELETE + UPDATE), "
|
|
f"got {pool.execute.call_count}"
|
|
)
|
|
|
|
# Verify the DELETE query targets document_intelligence with failed status
|
|
delete_call = pool.execute.call_args_list[0]
|
|
assert "DELETE" in delete_call[0][0] and "document_intelligence" in delete_call[0][0], (
|
|
"First execute should DELETE from document_intelligence"
|
|
)
|
|
|
|
# Verify the UPDATE resets status to 'parsed'
|
|
update_call = pool.execute.call_args_list[1]
|
|
assert "parsed" in update_call[0][0] and "UPDATE" in update_call[0][0], (
|
|
"Second execute should UPDATE status to 'parsed'"
|
|
)
|