From 5c692cda74a25b04a82f651fb37a35d37ba980e7 Mon Sep 17 00:00:00 2001 From: Codex Date: Fri, 25 Sep 2026 11:40:58 +0200 Subject: [PATCH] feat(sentiment): Telegram preview scraper + connector factory fix --- .../sentiment_engine/ingestion/__init__.py | 1 + .../src/sentiment_engine/ingestion/manager.py | 23 ++++++++++- .../sentiment_engine/ingestion/telegram.py | 39 +++++++++++++------ 3 files changed, 50 insertions(+), 13 deletions(-) diff --git a/sentiment_engine/src/sentiment_engine/ingestion/__init__.py b/sentiment_engine/src/sentiment_engine/ingestion/__init__.py index 854b973..90ec5c2 100644 --- a/sentiment_engine/src/sentiment_engine/ingestion/__init__.py +++ b/sentiment_engine/src/sentiment_engine/ingestion/__init__.py @@ -25,3 +25,4 @@ __all__ = [ "WebCrawlConnector", "IngestionManager", ] +from sentiment_engine.ingestion.telegram_preview import TelegramPreviewConnector diff --git a/sentiment_engine/src/sentiment_engine/ingestion/manager.py b/sentiment_engine/src/sentiment_engine/ingestion/manager.py index d691ced..51acccf 100644 --- a/sentiment_engine/src/sentiment_engine/ingestion/manager.py +++ b/sentiment_engine/src/sentiment_engine/ingestion/manager.py @@ -14,6 +14,9 @@ from sentiment_engine.ingestion.base import BaseConnector, ConnectorConfig, Conn from sentiment_engine.ingestion.rss import RSSConnector from sentiment_engine.ingestion.twitter import TwitterConnector from sentiment_engine.ingestion.reddit import RedditConnector +from sentiment_engine.ingestion.telegram import TelegramConnector +from sentiment_engine.ingestion.telegram_preview import TelegramPreviewConnector +from sentiment_engine.ingestion.discord import DiscordConnector from sentiment_engine.ingestion.exchange import ExchangeConnector from sentiment_engine.ingestion.regulatory import RegulatoryConnector from sentiment_engine.ingestion.corporate import CorporateConnector @@ -109,8 +112,17 @@ class IngestionManager: connector = self._create_connector(conn_type, connector_config) if connector: self._connectors[source_id] = connector - # Register in catalogue - await self.catalogue.register_connector(source_id, connector) + # Register in catalogue using register_source + self.catalogue.register_source( + name=source_id, + connector_type=conn_type, + base_url=config.get("url", ""), + config=config.get("extra_config", {}), + base_credibility=config.get("base_credibility", 0.5), + relevance=config.get("relevance", 0.5), + cadence_seconds=config.get("cadence_seconds", 300), + tags=[conn_type.value] + ) def _create_connector(self, conn_type: ConnectorType, config: ConnectorConfig) -> Optional[BaseConnector]: """Factory method to create connector by type""" @@ -120,6 +132,10 @@ class IngestionManager: return TwitterConnector(config) elif conn_type == ConnectorType.REDDIT: return RedditConnector(config) + elif conn_type == ConnectorType.TELEGRAM: + return TelegramConnector(config) + elif conn_type == ConnectorType.DISCORD: + return DiscordConnector(config) elif conn_type == ConnectorType.EXCHANGE_ANN: return ExchangeConnector(config) elif conn_type == ConnectorType.REGULATORY: @@ -127,6 +143,9 @@ class IngestionManager: elif conn_type == ConnectorType.CORPORATE: return CorporateConnector(config) elif conn_type == ConnectorType.WEB_CRAWL: + # Check if it's a Telegram preview source + if config.source_id.startswith("telegram:") or config.source_id.startswith("web:telegram:"): + return TelegramPreviewConnector(config) return WebCrawlConnector(config) else: logger.warning(f"Unknown connector type: {conn_type}") diff --git a/sentiment_engine/src/sentiment_engine/ingestion/telegram.py b/sentiment_engine/src/sentiment_engine/ingestion/telegram.py index f1b85d5..54a1dba 100644 --- a/sentiment_engine/src/sentiment_engine/ingestion/telegram.py +++ b/sentiment_engine/src/sentiment_engine/ingestion/telegram.py @@ -11,8 +11,7 @@ from aiogram.filters import Command from aiogram.types import Message from sentiment_engine.schemas.payload import NormalizedPayload, SourceType, AssetMention, EngagementMetrics -from sentiment_engine.schemas.config import TelegramConnectorConfig -from sentiment_engine.ingestion.base import BaseConnector +from sentiment_engine.ingestion.base import BaseConnector, ConnectorConfig from sentiment_engine.utils.text import clean_html, extract_tickers, extract_cashtags, detect_language logger = logging.getLogger(__name__) @@ -21,10 +20,9 @@ logger = logging.getLogger(__name__) class TelegramConnector(BaseConnector): """Telegram bot connector for monitoring channels""" - def __init__(self, config: TelegramConnectorConfig, credibility_registry): + def __init__(self, config: ConnectorConfig): super().__init__(config) - self.credibility_registry = credibility_registry - self.channel_usernames = config.channel_usernames + self.channel_usernames = config.extra_config.get("channel_usernames", []) self._bot: Optional[Bot] = None self._dp: Optional[Dispatcher] = None self._message_queue: asyncio.Queue = asyncio.Queue() @@ -32,8 +30,8 @@ class TelegramConnector(BaseConnector): self._channel_ids: List[int] = [] # Query timing windows - self.preferred_windows = config.preferred_query_windows or [] - self.avoid_windows = config.avoid_query_windows or [] + self.preferred_windows = config.extra_config.get("preferred_query_windows", []) + self.avoid_windows = config.extra_config.get("avoid_query_windows", []) def _in_preferred_window(self) -> bool: if not self.preferred_windows: @@ -69,7 +67,7 @@ class TelegramConnector(BaseConnector): async def initialize(self) -> None: """Initialize Telegram bot""" - self._bot = Bot(token=self.config.bot_token) + self._bot = Bot(token=self.config.extra_config.get("bot_token", "")) self._dp = Dispatcher() # Resolve channel usernames to IDs @@ -147,7 +145,7 @@ class TelegramConnector(BaseConnector): publish_ts = message.date.timestamp() source_id = f"telegram:{message.chat.id}" - credibility = self.credibility_registry.get(source_id, 0.4) + credibility = self.config.base_credibility language = detect_language(raw_text) return NormalizedPayload( @@ -172,6 +170,25 @@ class TelegramConnector(BaseConnector): } ) + async def poll(self) -> List[NormalizedPayload]: + """Poll for new messages (collect from queue)""" + if not self._bot: + await self.initialize() + + payloads = [] + # Collect all available messages from queue + while not self._message_queue.empty(): + try: + message = self._message_queue.get_nowait() + payload = await self._process_message(message) + if payload: + payloads.append(payload) + except asyncio.QueueEmpty: + break + except Exception as e: + logger.error(f"Telegram message processing error: {e}") + return payloads + async def health_check(self) -> bool: try: if self._bot: @@ -181,10 +198,10 @@ class TelegramConnector(BaseConnector): pass return False - async def stop(self) -> None: + async def close(self) -> None: self._running = False if self._dp: await self._dp.stop_polling() if self._bot: await self._bot.session.close() - await super().stop() + await super().close()