428 lines
15 KiB
Python
428 lines
15 KiB
Python
"""
|
|
Comprehensive tests for Output Sinks (Hazelcast, ClickHouse, LatticeDB).
|
|
"""
|
|
|
|
import pytest
|
|
import asyncio
|
|
from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
|
from sentiment_engine.output.hazelcast_sink import HazelcastSink
|
|
from sentiment_engine.output.clickhouse_sink import ClickHouseSink
|
|
from sentiment_engine.output.latticedb_sink import LatticeDBSink
|
|
from sentiment_engine.output.manager import OutputManager
|
|
from sentiment_engine.schemas.output import SentimentOutput, AssetSentiment
|
|
from sentiment_engine.schemas.processed import ProcessedItem, SentimentScores, EmotionScores
|
|
|
|
|
|
class TestHazelcastSink:
|
|
"""Tests for HazelcastSink"""
|
|
|
|
@pytest.fixture
|
|
def sink(self):
|
|
return HazelcastSink(
|
|
cluster_name="test",
|
|
cluster_members=["localhost:5701"],
|
|
maps={"sentiment_scores": "sentiment_scores_*", "sentiment_streams": "sentiment_streams"}
|
|
)
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_initialize_creates_client(self, sink):
|
|
"""Should initialize Hazelcast client"""
|
|
with patch('hazelcast.HazelcastClient') as mock_client:
|
|
mock_client.return_value = AsyncMock()
|
|
|
|
await sink.initialize()
|
|
|
|
assert sink._client is not None
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_sentiment_score(self, sink):
|
|
"""Should write sentiment score to map"""
|
|
with patch('hazelcast.HazelcastClient') as mock_client:
|
|
mock_map = AsyncMock()
|
|
mock_client.return_value.get_map.return_value = mock_map
|
|
mock_client.return_value = AsyncMock()
|
|
|
|
await sink.initialize()
|
|
|
|
from sentiment_engine.schemas.output import AssetSentiment
|
|
|
|
output = SentimentOutput(
|
|
timestamp=1700000000.0,
|
|
assets={
|
|
"BTC": AssetSentiment(
|
|
asset_id="BTC",
|
|
sentiment=SentimentScores(polarity=0.5, confidence=0.8, positive_prob=0.7, negative_prob=0.1, neutral_prob=0.2),
|
|
emotions=EmotionScores(joy=0.5, fear=0.1, anger=0.0, greed=0.3, sadness=0.0, intensity=0.5),
|
|
events=[],
|
|
mention_count=5
|
|
)
|
|
},
|
|
market_fear_greed=60.0,
|
|
global_sentiment=0.5
|
|
)
|
|
|
|
await sink.write(output)
|
|
|
|
mock_map.set.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_stream(self, sink):
|
|
"""Should write to stream map"""
|
|
with patch('hazelcast.HazelcastClient') as mock_client:
|
|
mock_map = AsyncMock()
|
|
mock_client.return_value.get_map.return_value = mock_map
|
|
mock_client.return_value = AsyncMock()
|
|
|
|
await sink.initialize()
|
|
|
|
await sink.write_stream("test_key", {"data": "test"})
|
|
|
|
mock_map.set.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_health_check(self, sink):
|
|
"""Health check should return status"""
|
|
with patch('hazelcast.HazelcastClient') as mock_client:
|
|
mock_client.return_value = AsyncMock()
|
|
|
|
await sink.initialize()
|
|
|
|
health = await sink.health_check()
|
|
|
|
assert "status" in health
|
|
assert health["status"] in ["healthy", "unhealthy"]
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_close_closes_client(self, sink):
|
|
"""Close should close Hazelcast client"""
|
|
with patch('hazelcast.HazelcastClient') as mock_client:
|
|
mock_client_instance = AsyncMock()
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
await sink.close()
|
|
|
|
mock_client_instance.shutdown.assert_called()
|
|
|
|
|
|
class TestClickHouseSink:
|
|
"""Tests for ClickHouseSink"""
|
|
|
|
@pytest.fixture
|
|
def sink(self):
|
|
return ClickHouseSink(
|
|
host="localhost",
|
|
port=8123,
|
|
database="test",
|
|
user="default",
|
|
password="",
|
|
tables={
|
|
"sentiment_events": "sentiment_events",
|
|
"sentiment_scores": "sentiment_scores",
|
|
"sentiment_raw_items": "sentiment_raw_items"
|
|
}
|
|
)
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_initialize_creates_pool(self, sink):
|
|
"""Should initialize connection pool"""
|
|
with patch('clickhouse_driver.Client') as mock_client:
|
|
mock_client.return_value = MagicMock()
|
|
|
|
await sink.initialize()
|
|
|
|
assert sink._client is not None
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_inserts_events(self, sink):
|
|
"""Should insert events into ClickHouse"""
|
|
with patch('clickhouse_driver.Client') as mock_client:
|
|
mock_client_instance = MagicMock()
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
|
|
from sentiment_engine.schemas.output import AssetSentiment
|
|
|
|
output = SentimentOutput(
|
|
timestamp=1700000000.0,
|
|
assets={
|
|
"BTC": AssetSentiment(
|
|
asset_id="BTC",
|
|
sentiment=SentimentScores(polarity=0.5, confidence=0.8, positive_prob=0.7, negative_prob=0.1, neutral_prob=0.2),
|
|
emotions=EmotionScores(joy=0.5, fear=0.1, anger=0.0, greed=0.3, sadness=0.0, intensity=0.5),
|
|
events=[],
|
|
mention_count=5
|
|
)
|
|
},
|
|
market_fear_greed=60.0,
|
|
global_sentiment=0.5
|
|
)
|
|
|
|
await sink.write(output)
|
|
|
|
mock_client_instance.execute.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_raw_item(self, sink):
|
|
"""Should write raw item"""
|
|
with patch('clickhouse_driver.Client') as mock_client:
|
|
mock_client_instance = MagicMock()
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
|
|
await sink.write_raw_item({
|
|
"source_id": "test",
|
|
"raw_text": "Test",
|
|
"timestamp": 1700000000.0
|
|
})
|
|
|
|
mock_client_instance.execute.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_health_check(self, sink):
|
|
"""Health check should return status"""
|
|
with patch('clickhouse_driver.Client') as mock_client:
|
|
mock_client_instance = MagicMock()
|
|
mock_client_instance.execute.return_value = [[1]]
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
|
|
health = await sink.health_check()
|
|
|
|
assert "status" in health
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_close_closes_connection(self, sink):
|
|
"""Close should close connection"""
|
|
with patch('clickhouse_driver.Client') as mock_client:
|
|
mock_client_instance = MagicMock()
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
await sink.close()
|
|
|
|
mock_client_instance.disconnect.assert_called()
|
|
|
|
|
|
class TestLatticeDBSink:
|
|
"""Tests for LatticeDBSink"""
|
|
|
|
@pytest.fixture
|
|
def sink(self):
|
|
return LatticeDBSink(
|
|
host="localhost",
|
|
port=7878
|
|
)
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_initialize_creates_connection(self, sink):
|
|
"""Should initialize connection"""
|
|
with patch('httpx.AsyncClient') as mock_client:
|
|
mock_client.return_value = AsyncMock()
|
|
|
|
await sink.initialize()
|
|
|
|
assert sink._client is not None
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_entities(self, sink):
|
|
"""Should write entity relationships"""
|
|
with patch('httpx.AsyncClient') as mock_client:
|
|
mock_client_instance = AsyncMock()
|
|
mock_client_instance.post = AsyncMock(return_value=MagicMock(status_code=200))
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
|
|
await sink.write_entities([
|
|
{"source": "BTC", "target": "ETH", "relationship": "correlated", "weight": 0.8}
|
|
])
|
|
|
|
mock_client_instance.post.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_query_neighbors(self, sink):
|
|
"""Should query neighbor entities"""
|
|
with patch('httpx.AsyncClient') as mock_client:
|
|
mock_client_instance = AsyncMock()
|
|
mock_client_instance.get = AsyncMock(return_value=MagicMock(
|
|
status_code=200,
|
|
json=lambda: {"neighbors": [{"entity": "ETH", "weight": 0.8}]}
|
|
))
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
|
|
neighbors = await sink.query_neighbors("BTC")
|
|
|
|
assert isinstance(neighbors, list)
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_health_check(self, sink):
|
|
"""Health check should return status"""
|
|
with patch('httpx.AsyncClient') as mock_client:
|
|
mock_client_instance = AsyncMock()
|
|
mock_client_instance.get = AsyncMock(return_value=MagicMock(status_code=200))
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
await sink.initialize()
|
|
|
|
health = await sink.health_check()
|
|
|
|
assert "status" in health
|
|
|
|
|
|
class TestOutputManager:
|
|
"""Tests for OutputManager"""
|
|
|
|
@pytest.fixture
|
|
def manager(self):
|
|
with patch('sentiment_engine.output.hazelcast_sink.HazelcastSink') as mock_hz, \
|
|
patch('sentiment_engine.output.clickhouse_sink.ClickHouseSink') as mock_ch, \
|
|
patch('sentiment_engine.output.latticedb_sink.LatticeDBSink') as mock_ldb:
|
|
|
|
mock_hz.return_value = AsyncMock()
|
|
mock_ch.return_value = AsyncMock()
|
|
mock_ldb.return_value = AsyncMock()
|
|
|
|
manager = OutputManager(
|
|
hazelcast_config={"cluster_name": "test"},
|
|
clickhouse_config={"host": "localhost"},
|
|
latticedb_config={"host": "localhost"}
|
|
)
|
|
yield manager
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_initialize_all_sinks(self, manager):
|
|
"""Should initialize all sinks"""
|
|
await manager.initialize()
|
|
|
|
assert manager.hazelcast_sink is not None
|
|
assert manager.clickhouse_sink is not None
|
|
assert manager.latticedb_sink is not None
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_to_all_sinks(self, manager):
|
|
"""Should write to all sinks"""
|
|
from sentiment_engine.schemas.output import AssetSentiment
|
|
|
|
output = SentimentOutput(
|
|
timestamp=1700000000.0,
|
|
assets={},
|
|
market_fear_greed=50.0,
|
|
global_sentiment=0.0
|
|
)
|
|
|
|
await manager.write(output)
|
|
|
|
manager.hazelcast_sink.write.assert_called()
|
|
manager.clickhouse_sink.write.assert_called()
|
|
manager.latticedb_sink.write.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_write_handles_sink_failure(self, manager):
|
|
"""Should handle individual sink failures gracefully"""
|
|
manager.hazelcast_sink.write = AsyncMock(side_effect=Exception("Hazelcast down"))
|
|
manager.clickhouse_sink.write = AsyncMock()
|
|
manager.latticedb_sink.write = AsyncMock()
|
|
|
|
output = SentimentOutput(timestamp=1700000000.0, assets={}, market_fear_greed=50.0, global_sentiment=0.0)
|
|
|
|
# Should not raise
|
|
await manager.write(output)
|
|
|
|
manager.clickhouse_sink.write.assert_called()
|
|
manager.latticedb_sink.write.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_close_all_sinks(self, manager):
|
|
"""Should close all sinks"""
|
|
await manager.close()
|
|
|
|
manager.hazelcast_sink.close.assert_called()
|
|
manager.clickhouse_sink.close.assert_called()
|
|
manager.latticedb_sink.close.assert_called()
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_health_check_all(self, manager):
|
|
"""Should check health of all sinks"""
|
|
manager.hazelcast_sink.health_check = AsyncMock(return_value={"status": "healthy"})
|
|
manager.clickhouse_sink.health_check = AsyncMock(return_value={"status": "healthy"})
|
|
manager.latticedb_sink.health_check = AsyncMock(return_value={"status": "healthy"})
|
|
|
|
health = await manager.health_check()
|
|
|
|
assert "hazelcast" in health
|
|
assert "clickhouse" in health
|
|
assert "latticedb" in health
|
|
|
|
|
|
class TestOutputEdgeCases:
|
|
"""Edge case tests for output sinks"""
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_hazelcast_reconnection(self):
|
|
"""Should handle reconnection"""
|
|
with patch('hazelcast.HazelcastClient') as mock_client:
|
|
mock_client.side_effect = [
|
|
Exception("Connection failed"),
|
|
AsyncMock()
|
|
]
|
|
|
|
sink = HazelcastSink(cluster_name="test", cluster_members=["localhost:5701"])
|
|
|
|
try:
|
|
await sink.initialize()
|
|
except:
|
|
pass
|
|
|
|
# Second attempt should succeed
|
|
await sink.initialize()
|
|
assert sink._client is not None
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_clickhouse_batch_insert(self):
|
|
"""Should batch inserts for efficiency"""
|
|
with patch('clickhouse_driver.Client') as mock_client:
|
|
mock_client_instance = MagicMock()
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
sink = ClickHouseSink(host="localhost")
|
|
await sink.initialize()
|
|
|
|
# Write multiple items
|
|
for i in range(10):
|
|
await sink.write_raw_item({"id": i, "data": f"test{i}"})
|
|
|
|
# Should have called execute
|
|
assert mock_client_instance.execute.call_count >= 1
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_latticedb_retry_on_failure(self):
|
|
"""Should retry on transient failures"""
|
|
with patch('httpx.AsyncClient') as mock_client:
|
|
mock_client_instance = AsyncMock()
|
|
mock_client_instance.post = AsyncMock(
|
|
side_effect=[
|
|
Exception("Transient error"),
|
|
MagicMock(status_code=200)
|
|
]
|
|
)
|
|
mock_client.return_value = mock_client_instance
|
|
|
|
sink = LatticeDBSink(host="localhost")
|
|
await sink.initialize()
|
|
|
|
# Should retry and succeed
|
|
await sink.write_entities([{"source": "BTC", "target": "ETH", "weight": 0.5}])
|
|
|
|
assert mock_client_instance.post.call_count == 2
|
|
|
|
|
|
if __name__ == "__main__":
|
|
pytest.main([__file__, "-v"])
|