feat(sentiment): Telegram preview scraper + connector factory fix
This commit is contained in:
@@ -25,3 +25,4 @@ __all__ = [
|
||||
"WebCrawlConnector",
|
||||
"IngestionManager",
|
||||
]
|
||||
from sentiment_engine.ingestion.telegram_preview import TelegramPreviewConnector
|
||||
|
||||
@@ -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}")
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user