"""Integration tests for ingestion pipeline""" import pytest import asyncio import time from datetime import datetime from sentiment_engine.catalogue.manager import CatalogueManager from sentiment_engine.ingestion.router import IngestionRouter from sentiment_engine.schemas.payload import NormalizedPayload, SourceType, AssetMention, EngagementMetrics from sentiment_engine.utils.config import get_settings @pytest.fixture async def catalogue(): cat = CatalogueManager() await cat.initialize() yield cat await cat.stop() @pytest.fixture async def router(catalogue): settings = get_settings() router = IngestionRouter( nats_servers=settings.nats_servers, stream_name=settings.nats_stream_ingestion, subject_map={ "news": "sentiment.ingest.news", "social": "sentiment.ingest.social", "regulatory": "sentiment.ingest.regulatory", "exchange": "sentiment.ingest.exchange", }, catalogue=catalogue ) await router.connect() yield router await router._nc.close() class TestIngestionPipeline: """Integration tests for ingestion pipeline""" @pytest.mark.asyncio async def test_catalogue_router_integration(self, catalogue, router): """Test that catalogue and router work together""" from sentiment_engine.schemas.payload import NormalizedPayload, AssetMention from datetime import datetime payload = NormalizedPayload( source_id="rss:test.com", source_type=SourceType.NEWS, source_credibility_base=0.8, ingest_ts=datetime.now().timestamp(), publish_ts=datetime.now().timestamp(), asset_mentions=[], raw_text="BTC surges to new highs on institutional adoption", title="BTC Surges", url="https://test.com/article", author="Test Author", content_length=100, language="en" ) # Route payload result = await router.route(payload) assert result is True # Check metrics metrics = router.get_metrics() assert metrics["received"] == 1 assert metrics["routed"] == 1 assert metrics["by_source"]["rss:test.com"] == 1 # Verify catalogue recorded the fetch sources = catalogue.catalogue.get_sources() # Should have recorded the fetch attempt @pytest.mark.asyncio async def test_deduplication(self, router): """Test that duplicate payloads are rejected""" from sentiment_engine.schemas.payload import NormalizedPayload from datetime import datetime payload = NormalizedPayload( source_id="rss:test.com", source_type="news", source_credibility_base=0.8, ingest_ts=datetime.now().timestamp(), publish_ts=datetime.now().timestamp(), asset_mentions=[], raw_text="BTC surges to new highs on institutional adoption", title="BTC Surges", url="https://test.com/article", author="Test Author", content_length=100, language="en" ) # First submission result1 = await router.route(payload) assert result1 is True # Second submission - should be deduplicated result2 = await router.route(payload) assert result2 is False metrics = router.get_metrics() assert metrics["duplicates"] == 1 class TestStaleDetectionIntegration: """Test stale source detection integration""" @pytest.mark.asyncio async def test_stale_source_detection(self, catalogue): """Test that stale sources are detected correctly""" import time conn = catalogue.catalogue._get_conn() conn.execute('UPDATE sources SET last_fetch_ts = ? WHERE source_id = ?', [time.time() - 900, 'rss:binance.com']) stale = catalogue.catalogue.get_stale_sources(multiplier=2.0) assert len(stale) >= 1 assert any(s.source_id == 'rss:binance.com' for s in stale) @pytest.mark.asyncio async def test_credibility_decay_detection(self, catalogue): """Test that credibility decay is detected""" import time conn = catalogue.catalogue._get_conn() # Use an existing source with low credibility conn.execute('UPDATE sources SET current_credibility = 0.2, credibility_updated_ts = ? WHERE source_id = ?', [time.time() - 4*86400, 'reddit:CryptoCurrency']) decay = catalogue.catalogue.get_credibility_decay_candidates(threshold=0.3, window_hours=72) assert len(decay) >= 1 assert any(s.source_id == 'reddit:CryptoCurrency' for s in decay) class TestDashboardAggregation: """Test dashboard data aggregation""" @pytest.mark.asyncio async def test_dashboard_aggregation(self, catalogue): """Test dashboard data aggregation""" from sentiment_engine.catalogue.manager import CatalogueManager as CM cat_mgr = CM.__new__(CM) cat_mgr.catalogue = catalogue.catalogue dashboard = cat_mgr.get_dashboard_data() assert dashboard["total_sources"] >= 14 assert "stale_count" in dashboard assert "decay_count" in dashboard assert "avg_credibility" in dashboard assert "by_type" in dashboard assert "sources" in dashboard