diff --git a/sentiment_engine/src/sentiment_engine/output/sinks.py b/sentiment_engine/src/sentiment_engine/output/sinks.py new file mode 100644 index 0000000..1e99ebc --- /dev/null +++ b/sentiment_engine/src/sentiment_engine/output/sinks.py @@ -0,0 +1,392 @@ +""" +Output Sinks — Hazelcast and ClickHouse persistence per spec Section 13 +""" +import asyncio +import logging +import time +from typing import Dict, List, Optional, Any +from dataclasses import dataclass, field +from datetime import datetime + +import aiohttp +import clickhouse_connect +import hazelcast + +from sentiment_engine.schemas.output import SentimentOutput, AssetSentiment, MarketSentiment, IndustrySentiment +from sentiment_engine.utils.config import get_settings + +logger = logging.getLogger(__name__) + + +@dataclass +class SinkConfig: + """Configuration for output sinks""" + # Hazelcast + hz_enabled: bool = True + hz_cluster_name: str = "dolphin" + hz_cluster_members: List[str] = field(default_factory=lambda: ["localhost:5701"]) + hz_map_name: str = "dolphin_features_sentiment" + + # ClickHouse + ch_enabled: bool = True + ch_host: str = "localhost" + ch_port: int = 8123 + ch_database: str = "dolphin" + ch_user: str = "default" + ch_password: str = "" + ch_table: str = "exf_data" + + # Output cadence + market_update_interval_seconds: int = 60 + asset_update_interval_seconds: int = 5 + cache_ttl_seconds: int = 300 + + +class HazelcastSink: + """Hazelcast ExF map sink for real-time feature serving""" + + def __init__(self, config: SinkConfig): + self.config = config + self._client: Optional[hazelcast.HazelcastClient] = None + self._map = None + + async def initialize(self) -> None: + """Initialize Hazelcast client""" + try: + self._client = await hazelcast.HazelcastClient( + cluster_name=self.config.hz_cluster_name, + cluster_members=self.config.hz_cluster_members, + ) + self._map = await self._client.get_map(self.config.hz_map_name).blocking() + logger.info(f"HazelcastSink connected to {self.config.hz_cluster_members}") + except Exception as e: + logger.warning(f"Hazelcast connection failed (will retry): {e}") + self._client = None + self._map = None + + async def write_market(self, output: SentimentOutput) -> None: + """Write market-level snapshot to HZ map""" + if not self._map: + await self.initialize() + if not self._map: + return + + key = "market_snapshot" + value = { + "timestamp": output.timestamp, + "schema_version": output.schema_version, + "engine_version": output.engine_version, + "market": { + "fear_state": output.market.fear_state, + "greed_state": output.market.greed_state, + "sentiment_index": output.market.sentiment_index, + "hype_velocity": output.market.hype_velocity, + "pub_velocity": output.market.pub_velocity, + "aggregate_pump_risk": output.market.aggregate_pump_risk, + "aggregate_dump_risk": output.market.aggregate_dump_risk, + "top_pump_assets": output.market.top_pump_assets, + "top_dump_assets": output.market.top_dump_assets, + "last_update_ts": output.market.last_update_ts, + }, + "industries": { + name: { + "industry": ind.industry, + "fear_state": ind.fear_state, + "greed_state": ind.greed_state, + "avg_polarity": ind.avg_polarity, + "pump_risk": ind.pump_risk, + "dump_risk": ind.dump_risk, + "asset_count": ind.asset_count, + "last_update_ts": ind.last_update_ts, + } + for name, ind in output.industries.items() + }, + "schema_version": output.schema_version, + "engine_version": output.engine_version, + } + + try: + self._map.put(key, value) + except Exception as e: + logger.error(f"Hazelcast write_market failed: {e}") + + async def write_asset(self, asset_id: str, asset: AssetSentiment) -> None: + """Write per-asset sentiment to HZ map""" + if not self._map: + await self.initialize() + if not self._map: + return + + key = f"asset:{asset_id}" + value = { + "asset_id": asset.asset_id, + "fear_state": asset.fear_state, + "greed_state": asset.greed_state, + "sentiment_polarity": asset.sentiment_polarity, + "emotion_profile": asset.emotion_profile, + "pump_score": asset.pump_dump.pump_score if asset.pump_dump else 0, + "dump_score": asset.pump_dump.dump_score if asset.pump_dump else 0, + "pump_confidence": asset.pump_dump.pump_confidence if asset.pump_dump else 0, + "dump_confidence": asset.pump_dump.dump_confidence if asset.pump_dump else 0, + "velocity": { + "hype_velocity": asset.velocity.hype_velocity if asset.velocity else 0, + "pub_velocity": asset.velocity.pub_velocity if asset.velocity else 0, + "direction": asset.velocity.velocity_direction if asset.velocity else "neutral" + } if asset.velocity else {}, + "event_flags": [ + { + "event_type": f.event_type, + "value": f.value, + "confidence": f.confidence, + "direction": f.direction, + "flag_type": f.flag_type.value if f.flag_type else None, + "flags": f.flags, + } + for f in asset.event_flags + ], + "last_update_ts": asset.last_update_ts, + "contributing_sources": asset.contributing_sources, + "decay_factor": asset.decay_factor, + "contributing_events": asset.contributing_events, + } + + try: + self._map.put(key, value) + except Exception as e: + logger.error(f"Hazelcast write_asset {asset_id} failed: {e}") + + async def close(self) -> None: + if self._client: + await self._client.shutdown() + + +class ClickHouseSink: + """ClickHouse sink for historical persistence""" + + def __init__(self, config: SinkConfig): + self.config = config + self._client: Optional[clickhouse_connect.Client] = None + self._buffer: List[Dict] = [] + + async def initialize(self) -> None: + """Initialize ClickHouse client""" + try: + self._client = clickhouse_connect.get_client( + host=self.config.ch_host, + port=self.config.ch_port, + database=self.config.ch_database, + user=self.config.ch_user, + password=self.config.ch_password, + ) + # Ensure table exists + await self._ensure_table() + logger.info(f"ClickHouseSink connected to {self.config.ch_host}:{self.config.ch_port}") + except Exception as e: + logger.warning(f"ClickHouse connection failed (will retry): {e}") + self._client = None + + async def _ensure_table(self) -> None: + """Create exf_data table if not exists""" + if not self._client: + return + + ddl = f""" + CREATE TABLE IF NOT EXISTS {self.config.ch_table} ( + timestamp DateTime64(3), + schema_version UInt8, + engine_version String, + + -- Market level + market_fear_state Float32, + market_greed_state Float32, + market_sentiment_index Float32, + market_hype_velocity Float32, + market_pub_velocity Float32, + market_aggregate_pump_risk Float32, + market_aggregate_dump_risk Float32, + market_top_pump_assets Array(String), + market_top_dump_assets Array(String), + + -- Asset level (flattened) + asset_id String, + asset_fear_state Float32, + asset_greed_state Float32, + asset_sentiment_polarity Float32, + asset_emotion_joy Float32, + asset_emotion_fear Float32, + asset_emotion_anger Float32, + asset_emotion_greed Float32, + asset_emotion_sadness Float32, + asset_emotion_intensity Float32, + asset_pump_score Float32, + asset_dump_score Float32, + asset_pump_confidence Float32, + asset_dump_confidence Float32, + asset_hype_velocity Float32, + asset_pub_velocity Float32, + asset_velocity_direction String, + asset_event_count UInt16, + asset_last_update_ts DateTime64(3), + + -- Event flags (flattened) + event_types Array(String), + event_values Array(Float32), + event_confidences Array(Float32), + event_directions Array(String), + + -- Metadata + schema_version UInt8, + engine_version String, + ingest_ts DateTime64(3) + ) ENGINE = MergeTree() + PARTITION BY toDate(timestamp) + ORDER BY (timestamp, asset_id) + TTL timestamp + INTERVAL 90 DAY + SETTINGS index_granularity = 8192 + """ + try: + self._client.command(ddl) + except Exception as e: + logger.warning(f"Table creation failed: {e}") + + async def write_output(self, output: SentimentOutput) -> None: + """Write full output to ClickHouse (buffered)""" + if not self._client: + return + + timestamp = datetime.fromtimestamp(output.timestamp) + ingest_ts = datetime.now() + + # Market-level row + market_row = { + "timestamp": timestamp, + "schema_version": output.schema_version, + "engine_version": output.engine_version, + "market_fear_state": output.market.fear_state, + "market_greed_state": output.market.greed_state, + "market_sentiment_index": output.market.sentiment_index, + "market_hype_velocity": output.market.hype_velocity, + "market_pub_velocity": output.market.pub_velocity, + "market_aggregate_pump_risk": output.market.aggregate_pump_risk, + "market_aggregate_dump_risk": output.market.aggregate_dump_risk, + "market_top_pump_assets": output.market.top_pump_assets, + "market_top_dump_assets": output.market.top_dump_assets, + "asset_id": "MARKET", + "asset_fear_state": output.market.fear_state, + "asset_greed_state": output.market.greed_state, + "asset_sentiment_polarity": output.market.sentiment_index, + "asset_emotion_joy": 0, + "asset_emotion_fear": 0, + "asset_emotion_anger": 0, + "asset_emotion_greed": 0, + "asset_emotion_sadness": 0, + "asset_emotion_intensity": 0, + "asset_pump_score": 0, + "asset_dump_score": 0, + "asset_pump_confidence": 0, + "asset_dump_confidence": 0, + "asset_hype_velocity": output.market.hype_velocity, + "asset_pub_velocity": output.market.pub_velocity, + "asset_velocity_direction": "neutral", + "asset_event_count": len(output.market.dominant_events), + "asset_last_update_ts": datetime.fromtimestamp(output.market.last_update_ts), + "event_types": [e.event_type for e in output.market.dominant_events], + "event_values": [e.value for e in output.market.dominant_events], + "event_confidences": [e.confidence for e in output.market.dominant_events], + "event_directions": [e.direction for e in output.market.dominant_events], + "schema_version": output.schema_version, + "engine_version": output.engine_version, + "ingest_ts": datetime.now(), + } + + # Per-asset rows + rows = [market_row] + + for asset_id, asset in output.assets.items(): + row = { + "timestamp": datetime.fromtimestamp(output.timestamp), + "schema_version": output.schema_version, + "engine_version": output.engine_version, + "market_fear_state": output.market.fear_state, + "market_greed_state": output.market.greed_state, + "market_sentiment_index": output.market.sentiment_index, + "market_hype_velocity": output.market.hype_velocity, + "market_pub_velocity": output.market.pub_velocity, + "market_aggregate_pump_risk": output.market.aggregate_pump_risk, + "market_aggregate_dump_risk": output.market.aggregate_dump_risk, + "market_top_pump_assets": output.market.top_pump_assets, + "market_top_dump_assets": output.market.top_dump_assets, + "asset_id": asset.asset_id, + "asset_fear_state": asset.fear_state, + "asset_greed_state": asset.greed_state, + "asset_sentiment_polarity": asset.sentiment_polarity, + "asset_emotion_joy": asset.emotion_profile.get("joy", 0), + "asset_emotion_fear": asset.emotion_profile.get("fear", 0), + "asset_emotion_anger": asset.emotion_profile.get("anger", 0), + "asset_emotion_greed": asset.emotion_profile.get("greed", 0), + "asset_emotion_sadness": asset.emotion_profile.get("sadness", 0), + "asset_emotion_intensity": asset.emotion_profile.get("intensity", 0), + "asset_pump_score": asset.pump_dump.pump_score if asset.pump_dump else 0, + "asset_dump_score": asset.pump_dump.dump_score if asset.pump_dump else 0, + "asset_pump_confidence": asset.pump_dump.pump_confidence if asset.pump_dump else 0, + "asset_dump_confidence": asset.pump_dump.dump_confidence if asset.pump_dump else 0, + "asset_hype_velocity": asset.velocity.hype_velocity if asset.velocity else 0, + "asset_pub_velocity": asset.velocity.pub_velocity if asset.velocity else 0, + "asset_velocity_direction": asset.velocity.velocity_direction if asset.velocity else "neutral", + "asset_event_count": len(asset.event_flags), + "asset_last_update_ts": datetime.fromtimestamp(asset.last_update_ts), + "event_types": [e.event_type for e in asset.event_flags], + "event_values": [e.value for e in asset.event_flags], + "event_confidences": [e.confidence for e in asset.event_flags], + "event_directions": [e.direction for e in asset.event_flags], + "schema_version": output.schema_version, + "engine_version": output.engine_version, + "ingest_ts": datetime.now(), + } + rows.append(row) + + # Insert all rows + if self._client and rows: + try: + self._client.insert(self.config.ch_table, rows, column_names=list(rows[0].keys())) + logger.debug(f"ClickHouse: inserted {len(rows)} rows") + except Exception as e: + logger.error(f"ClickHouse insert failed: {e}") + + async def close(self) -> None: + if self._client: + self._client.close() + + +class OutputSinkManager: + """Manages all output sinks""" + + def __init__(self, config: Optional[SinkConfig] = None): + self.config = config or SinkConfig() + self.hz_sink = HazelcastSink(self.config) + self.ch_sink = ClickHouseSink(self.config) + + async def initialize(self) -> None: + await asyncio.gather( + self.hz_sink.initialize(), + self.ch_sink.initialize(), + ) + logger.info("OutputSinkManager initialized") + + async def write_output(self, output: SentimentOutput) -> None: + """Write output to all sinks""" + await asyncio.gather( + self.hz_sink.write_market(output), + self.ch_sink.write_output(output), + return_exceptions=True + ) + # Also write per-asset to Hazelcast + for asset_id, asset in output.assets.items(): + await self.hz_sink.write_asset(asset_id, output.assets[asset_id]) + + async def close(self) -> None: + await asyncio.gather( + self.hz_sink.close(), + self.ch_sink.close(), + ) + logger.info("OutputSinkManager closed")