From c4c8ed7c9f4b1cdf4a508f500a11f66dc3e86b1e Mon Sep 17 00:00:00 2001 From: Codex Date: Fri, 25 Sep 2026 11:46:56 +0200 Subject: [PATCH] feat(sentiment): register_source upsert fix + connector factory fix for web:telegram: --- .../src/sentiment_engine/catalogue/manager.py | 55 +++++++++--- .../src/sentiment_engine/catalogue/store.py | 87 +++++++++++++++++++ 2 files changed, 128 insertions(+), 14 deletions(-) diff --git a/sentiment_engine/src/sentiment_engine/catalogue/manager.py b/sentiment_engine/src/sentiment_engine/catalogue/manager.py index bd1adcf..ddd906b 100644 --- a/sentiment_engine/src/sentiment_engine/catalogue/manager.py +++ b/sentiment_engine/src/sentiment_engine/catalogue/manager.py @@ -38,6 +38,12 @@ class CatalogueManager: # Load connector configs from settings await self._sync_connectors_from_settings() + + # Force checkpoint to clear WAL and avoid locking issues + try: + self.catalogue._conn.execute("PRAGMA force_checkpoint") + except Exception as e: + logger.warning(f"Failed to checkpoint database: {e}") async def _upsert_from_credibility(self, item: Dict) -> None: """Create/update source from credibility registry entry""" @@ -156,7 +162,7 @@ class CatalogueManager: cadence_seconds: int = 300, tags: List[str] = None ) -> SourceDefinition: - """Register a new source with validation""" + """Register a new source with validation (upsert if exists)""" schema = self.catalogue.get_schema(connector_type) if schema: # Validate cadence against schema @@ -166,19 +172,40 @@ class CatalogueManager: if cred not in (config.get("credentials", {}) if "credentials" in config else {}): print(f"⚠️ Missing required credential: {cred}") - source = SourceDefinition( - name=name, - connector_type=connector_type, - base_url=base_url, - config=config, - base_credibility=base_credibility, - relevance=relevance, - credentials_ref=credentials_ref, - cadence_seconds=cadence_seconds, - config_schema=schema.config_schema if schema else {}, - tags=tags or [] - ) - return self.catalogue.create_source(source) + # Check if source already exists + existing = self.catalogue.get_source(name) + if existing: + # Update existing source + updates = { + "connector_type": connector_type.value, + "base_url": base_url, + "config": config, + "base_credibility": base_credibility, + "relevance": relevance, + "credentials_ref": credentials_ref, + "cadence_seconds": cadence_seconds, + "config_schema": schema.config_schema if schema else {}, + "tags": tags or [], + "updated_ts": datetime.now().timestamp(), + } + self.catalogue.update_source(name, updates) + return self.catalogue.get_source(name) + else: + # Create new source + source = SourceDefinition( + source_id=name, + name=name, + connector_type=connector_type, + base_url=base_url, + config=config, + base_credibility=base_credibility, + relevance=relevance, + credentials_ref=credentials_ref, + cadence_seconds=cadence_seconds, + config_schema=schema.config_schema if schema else {}, + tags=tags or [] + ) + return self.catalogue.create_source(source) def record_fetch_result( self, diff --git a/sentiment_engine/src/sentiment_engine/catalogue/store.py b/sentiment_engine/src/sentiment_engine/catalogue/store.py index 909d4c0..e829fd9 100644 --- a/sentiment_engine/src/sentiment_engine/catalogue/store.py +++ b/sentiment_engine/src/sentiment_engine/catalogue/store.py @@ -20,6 +20,9 @@ class ConnectorType(str, Enum): REDDIT = "reddit" DISCORD = "discord" TELEGRAM = "telegram" + EXCHANGE_ANN = "exchange_ann" + REGULATORY = "regulatory" + CORPORATE = "corporate" WEB_CRAWL = "web_crawl" @@ -329,6 +332,90 @@ DEFAULT_SCHEMAS = { min_rate_limit_rps=0.01, max_rate_limit_rps=2.0 ), + ConnectorType.EXCHANGE_ANN: SourceSchema( + connector_type=ConnectorType.EXCHANGE_ANN, + version=1, + config_schema={ + "type": "object", + "properties": { + "feed_urls": {"type": "array", "items": {"type": "string"}}, + "max_items_per_feed": {"type": "integer"}, + "poll_interval_seconds": {"type": "integer"} + }, + "required": ["feed_urls"] + }, + payload_schema={ + "type": "object", + "properties": { + "title": {"type": "string"}, + "summary": {"type": "string"}, + "link": {"type": "string"}, + "published_parsed": {"type": "array"}, + "author": {"type": "string"} + } + }, + required_credentials=[], + min_cadence_seconds=60, + max_cadence_seconds=3600, + min_rate_limit_rps=0.01, + max_rate_limit_rps=1.0 + ), + ConnectorType.REGULATORY: SourceSchema( + connector_type=ConnectorType.REGULATORY, + version=1, + config_schema={ + "type": "object", + "properties": { + "feed_urls": {"type": "array", "items": {"type": "string"}}, + "max_items_per_feed": {"type": "integer"}, + "poll_interval_seconds": {"type": "integer"} + }, + "required": ["feed_urls"] + }, + payload_schema={ + "type": "object", + "properties": { + "title": {"type": "string"}, + "summary": {"type": "string"}, + "link": {"type": "string"}, + "published_parsed": {"type": "array"}, + "author": {"type": "string"} + } + }, + required_credentials=[], + min_cadence_seconds=60, + max_cadence_seconds=3600, + min_rate_limit_rps=0.01, + max_rate_limit_rps=1.0 + ), + ConnectorType.CORPORATE: SourceSchema( + connector_type=ConnectorType.CORPORATE, + version=1, + config_schema={ + "type": "object", + "properties": { + "feed_urls": {"type": "array", "items": {"type": "string"}}, + "max_items_per_feed": {"type": "integer"}, + "poll_interval_seconds": {"type": "integer"} + }, + "required": ["feed_urls"] + }, + payload_schema={ + "type": "object", + "properties": { + "title": {"type": "string"}, + "summary": {"type": "string"}, + "link": {"type": "string"}, + "published_parsed": {"type": "array"}, + "author": {"type": "string"} + } + }, + required_credentials=[], + min_cadence_seconds=60, + max_cadence_seconds=3600, + min_rate_limit_rps=0.01, + max_rate_limit_rps=1.0 + ), }