""" Performance benchmark tests for critical components. """ import pytest import asyncio import time import numpy as np from unittest.mock import AsyncMock, MagicMock, patch from sentiment_engine.nlp.pipeline import NLPProcessingPipeline from sentiment_engine.nlp.entity_extraction import EntityExtractor, AssetMapper from sentiment_engine.nlp.sentiment_emotion import SentimentEmotionAnalyzer from sentiment_engine.nlp.event_classification import EventClassifier from sentiment_engine.nlp.temporal import TemporalAnchorer from sentiment_engine.nlp.credibility import CredibilityScorer from sentiment_engine.signal.processor import FearGreedProcessor from sentiment_engine.signal.velocity import VelocityCalculator from sentiment_engine.signal.decay import DecayEngine from sentiment_engine.signal.fusion import MultiSourceFusion from sentiment_engine.schemas.payload import NormalizedPayload, SourceType, AssetMention from sentiment_engine.schemas.processed import ProcessedItem, SentimentScores, EmotionScores class TestPipelinePerformance: """Performance benchmarks for NLP pipeline""" @pytest.fixture def pipeline(self): return NLPProcessingPipeline() @pytest.mark.asyncio async def test_pipeline_initialization_time(self, pipeline): """Initialization should be fast""" start = time.time() await pipeline.initialize() elapsed = time.time() - start assert elapsed < 10 # 10 seconds max @pytest.mark.asyncio async def test_single_process_latency(self, pipeline): """Single process should be under 2 seconds""" await pipeline.initialize() payload = NormalizedPayload( source_id="benchmark", source_type=SourceType.NEWS, source_credibility_base=0.8, ingest_ts=1700000000.0, publish_ts=1700000000.0, content_length=200, raw_text="Bitcoin surges to $108k as institutional inflows surge. BlackRock IBIT sees record $1.2B daily inflow. BTC and ETH both hit new all-time highs.", metadata={} ) # Warm up await pipeline.process(payload) # Measure latencies = [] for _ in range(10): start = time.time() await pipeline.process(payload) latencies.append((time.time() - start) * 1000) avg_latency = sum(latencies) / len(latencies) p95_latency = sorted(latencies)[int(len(latencies) * 0.95)] print(f"Avg latency: {avg_latency:.1f}ms, P95: {p95_latency:.1f}ms") assert avg_latency < 3000 # 3 seconds average assert p95_latency < 5000 # 5 seconds P95 @pytest.mark.asyncio async def test_batch_throughput(self, pipeline): """Batch processing should achieve high throughput""" await pipeline.initialize() payloads = [ NormalizedPayload( source_id=f"bench_{i}", source_type=SourceType.NEWS, source_credibility_base=0.8, ingest_ts=1700000000.0, publish_ts=1700000000.0, content_length=100, raw_text=f"Bitcoin news item {i} with some content for processing.", metadata={} ) for i in range(50) ] start = time.time() results = await pipeline.process_batch(payloads) elapsed = time.time() - start throughput = len(results) / elapsed print(f"Throughput: {throughput:.1f} items/sec") assert len(results) == 50 assert throughput > 5 # At least 5 items/sec class TestEntityExtractorPerformance: """Performance benchmarks for EntityExtractor""" @pytest.fixture def extractor(self): return EntityExtractor(AssetMapper()) @pytest.mark.asyncio async def test_extraction_latency(self, extractor): """Entity extraction should be fast""" await extractor.initialize() texts = [ "BTC and ETH surge as Bitcoin hits new high. Vitalik says ETH to $10k.", "Major hack on exchange. SEC sues Kraken. Whale moves 10000 BTC.", "Ethereum Dencun upgrade live. PEPE and BONK listed on Coinbase." ] * 20 # 60 texts start = time.time() for text in texts: await extractor.extract_all(text) elapsed = time.time() - start throughput = len(texts) / elapsed print(f"Entity extraction: {throughput:.1f} texts/sec") assert throughput > 20 # At least 20 texts/sec class TestSentimentAnalyzerPerformance: """Performance benchmarks for SentimentEmotionAnalyzer""" @pytest.fixture def analyzer(self): return SentimentEmotionAnalyzer() @pytest.mark.asyncio async def test_sentiment_latency(self, analyzer): """Sentiment analysis should be fast""" await analyzer.initialize() texts = [ "Bitcoin surges to new all-time high!", "Market crashes as panic selling ensues.", "SEC approves Bitcoin ETF.", "Ethereum upgrade goes live.", "Whale moves 10000 BTC." ] * 50 # 250 texts asset_mentions = [{"asset_id": "BTC", "span": (0, 3)}] * len(texts) start = time.time() for text, mentions in zip(texts, asset_mentions): await analyzer.analyze(text, [mentions]) elapsed = time.time() - start throughput = len(texts) / elapsed print(f"Sentiment analysis: {throughput:.1f} texts/sec") assert throughput > 10 # At least 10 texts/sec class TestEventClassifierPerformance: """Performance benchmarks for EventClassifier""" @pytest.fixture def classifier(self): return EventClassifier() @pytest.mark.asyncio async def test_classification_latency(self, classifier): """Event classification should be fast""" await classifier.initialize() texts = [ "Bitcoin surges to $100k!", "Major hack on exchange!", "SEC sues exchange!", "Ethereum upgrade live!", "Coinbase lists new token!" ] * 50 # 250 texts assets = [["BTC"]] * len(texts) start = time.time() for text, asset_list in zip(texts, assets): await classifier.classify(text, asset_list) elapsed = time.time() - start throughput = len(texts) / elapsed print(f"Event classification: {throughput:.1f} texts/sec") assert throughput > 20 # At least 20 texts/sec class TestSignalProcessingPerformance: """Performance benchmarks for Signal Processing""" def test_fear_greed_computation(self): """Fear/Greed computation should be fast""" processor = FearGreedProcessor() items = [ ProcessedItem( payload_id=f"item_{i}", source_id="source", source_type="news", ingest_ts=1700000000.0, publish_ts=1700000000.0, entities=[], sentiment_per_asset={ "BTC": SentimentScores( polarity=0.5, confidence=0.8, positive_prob=0.7, negative_prob=0.1, neutral_prob=0.2 ) }, emotions_per_asset={}, events=[], temporal=None, credibility=None, processed_ts=1700000000.0, processing_latency_ms=100, model_versions={} ) for i in range(1000) ] start = time.time() result = processor.compute(items) elapsed = time.time() - start print(f"Fear/Greed: {1000/elapsed:.1f} items/sec") assert elapsed < 0.1 # < 100ms for 1000 items def test_velocity_computation(self): """Velocity computation should be fast""" calculator = VelocityCalculator() now = 1700000000.0 items = [ {"asset_id": "BTC", "publish_ts": now - i*60, "sentiment_polarity": 0.5 + i*0.01} for i in range(1000) ] start = time.time() for _ in range(100): calculator.compute_velocity("BTC", items) elapsed = time.time() - start throughput = 100 / elapsed print(f"Velocity: {throughput:.1f} computations/sec") assert throughput > 100 # At least 100/sec def test_decay_computation(self): """Decay computation should be fast""" engine = DecayEngine() now = 1700000000.0 timestamps = [now - i*60 for i in range(10000)] start = time.time() for ts in timestamps: engine.compute_decay(ts, now, halflife_minutes=60) elapsed = time.time() - start throughput = 10000 / elapsed print(f"Decay: {throughput:.1f} computations/sec") assert throughput > 10000 # At least 10k/sec def test_fusion_computation(self): """Fusion computation should be fast""" fusion = MultiSourceFusion() scores = {f"source_{i}": 0.5 for i in range(100)} weights = {f"source_{i}": 1.0 for i in range(100)} start = time.time() for _ in range(1000): fusion.fuse(scores, weights) elapsed = time.time() - start throughput = 1000 / elapsed print(f"Fusion: {throughput:.1f} fusions/sec") assert throughput > 1000 # At least 1000/sec class TestONNXInferencePerformance: """Performance benchmarks for ONNX inference""" @pytest.mark.asyncio async def test_onnx_finbert_inference(self): """ONNX FinBERT inference should be fast""" import onnxruntime as ort session = ort.InferenceSession( "models/onnx/finbert/model.onnx", providers=['CPUExecutionProvider'] ) input_ids = np.ones((1, 128), dtype=np.int64) attention_mask = np.ones((1, 128), dtype=np.int64) token_type_ids = np.zeros((1, 128), dtype=np.int64) # Warm up for _ in range(10): session.run(None, { "input_ids": input_ids, "attention_mask": attention_mask, "token_type_ids": token_type_ids }) start = time.time() for _ in range(100): session.run(None, { "input_ids": input_ids, "attention_mask": attention_mask, "token_type_ids": token_type_ids }) elapsed = time.time() - start throughput = 100 / elapsed print(f"ONNX FinBERT: {throughput:.1f} inferences/sec") assert throughput > 50 # At least 50/sec @pytest.mark.asyncio async def test_onnx_bert_events_inference(self): """ONNX BERT Events inference should be fast""" import onnxruntime as ort session = ort.InferenceSession( "models/onnx/bert-base-event/model.onnx", providers=['CPUExecutionProvider'] ) input_ids = np.ones((1, 128), dtype=np.int64) attention_mask = np.ones((1, 128), dtype=np.int64) token_type_ids = np.zeros((1, 128), dtype=np.int64) # Warm up for _ in range(10): session.run(None, { "input_ids": input_ids, "attention_mask": attention_mask, "token_type_ids": token_type_ids }) start = time.time() for _ in range(100): session.run(None, { "input_ids": input_ids, "attention_mask": attention_mask, "token_type_ids": token_type_ids }) elapsed = time.time() - start throughput = 100 / elapsed print(f"ONNX BERT Events: {throughput:.1f} inferences/sec") assert throughput > 50 # At least 50/sec @pytest.mark.asyncio async def test_onnx_emotion_inference(self): """ONNX Emotion inference should be fast""" import onnxruntime as ort session = ort.InferenceSession( "models/onnx/distilroberta-emotion/model.onnx", providers=['CPUExecutionProvider'] ) input_ids = np.ones((1, 128), dtype=np.int64) attention_mask = np.ones((1, 128), dtype=np.int64) # Warm up for _ in range(10): session.run(None, { "input_ids": input_ids, "attention_mask": attention_mask }) start = time.time() for _ in range(100): session.run(None, { "input_ids": input_ids, "attention_mask": attention_mask }) elapsed = time.time() - start throughput = 100 / elapsed print(f"ONNX Emotion: {throughput:.1f} inferences/sec") assert throughput > 100 # At least 100/sec class TestMemoryUsage: """Memory usage tests""" @pytest.mark.asyncio async def test_pipeline_memory_stable(self): """Pipeline memory should not grow unbounded""" import psutil import os pipeline = NLPProcessingPipeline() await pipeline.initialize() process = psutil.Process(os.getpid()) initial_memory = process.memory_info().rss / 1024 / 1024 # MB payload = NormalizedPayload( source_id="mem_test", source_type=SourceType.NEWS, source_credibility_base=0.8, ingest_ts=1700000000.0, publish_ts=1700000000.0, content_length=100, raw_text="Bitcoin surges to new high!", metadata={} ) # Process many items for i in range(100): payload.raw_text = f"Bitcoin news item {i}" await pipeline.process(payload) final_memory = process.memory_info().rss / 1024 / 1024 # MB memory_growth = final_memory - initial_memory print(f"Memory growth: {memory_growth:.1f} MB") assert memory_growth < 500 # Less than 500MB growth class TestConcurrency: """Concurrency tests""" @pytest.mark.asyncio async def test_pipeline_concurrent_requests(self): """Pipeline should handle concurrent requests""" pipeline = NLPProcessingPipeline() await pipeline.initialize() payload = NormalizedPayload( source_id="concurrent", source_type=SourceType.NEWS, source_credibility_base=0.8, ingest_ts=1700000000.0, publish_ts=1700000000.0, content_length=100, raw_text="Bitcoin surges to new high!", metadata={} ) # Run 20 concurrent requests tasks = [pipeline.process(NormalizedPayload( source_id=f"concurrent_{i}", source_type=SourceType.NEWS, source_credibility_base=0.8, ingest_ts=1700000000.0, publish_ts=1700000000.0, content_length=100, raw_text=f"Bitcoin news {i}", metadata={} )) for i in range(20)] start = time.time() results = await asyncio.gather(*tasks) elapsed = time.time() - start assert len(results) == 20 # Should be faster than sequential assert elapsed < 30 # Under 30 seconds for 20 concurrent if __name__ == "__main__": pytest.main([__file__, "-v", "-s"])