""" Cognition Pipeline Launcher — standalone long-run service. Runs the cognition pipeline continuously with: - Rate-limited source fetching - Regime extraction and deduplication - Auto-add to ScenarioFactory - Persistence to ClickHouse - Metrics monitoring - Graceful shutdown """ from __future__ import annotations import json import logging import os import signal import sys import time from dataclasses import dataclass from typing import Optional from malkhut.training.cognition import CognitionPipeline, SourceCatalogue from malkhut.training.regime_expansion import RegimeExpander from malkhut.storage.ch_store import MalkhutCHStore LOGGER = logging.getLogger("malkhut.cognition.launcher") @dataclass class CognitionConfig: """Configuration for cognition pipeline launcher.""" catalogue_path: str = "source_catalogue.json" regime_db_path: str = "discovered_regimes.json" rate_limit_rpm: int = 30 fetch_interval_s: int = 60 metrics_interval_s: int = 300 max_regimes: int = 500 class CognitionLauncher: """ Standalone launcher for the cognition pipeline. Runs continuously with: - Rate-limited source fetching - Regime extraction and deduplication - Auto-add to ScenarioFactory - Persistence to CH and local DB - Metrics monitoring - Graceful shutdown """ def __init__(self, config: Optional[CognitionConfig] = None) -> None: self.config = config or CognitionConfig() self._pipeline = CognitionPipeline( catalogue_path=self.config.catalogue_path, rate_limit_rpm=self.config.rate_limit_rpm, ) self._regime_expander = RegimeExpander() self._store: Optional[MalkhutCHStore] = None self._running = False self._start_time = 0.0 self._total_fetched = 0 self._total_regimes = 0 self._last_metrics = 0.0 self._discovered_regimes: dict = {} # Load persisted regimes on init self._load_regimes() def run(self) -> None: """Run the cognition pipeline continuously.""" self._running = True self._start_time = time.time() # Setup self._pipeline.seed_default_sources() try: self._store = MalkhutCHStore() self._ensure_tables() except Exception: self._store = None # Load persisted regimes self._load_regimes() # Register signal handlers signal.signal(signal.SIGINT, self._signal_handler) signal.signal(signal.SIGTERM, self._signal_handler) print("=" * 70) print("MALKHUT COGNITION PIPELINE") print(f"Sources: {self._pipeline._catalogue.source_count}") print(f"Rate limit: {self.config.rate_limit_rpm} RPM") print(f"Fetch interval: {self.config.fetch_interval_s}s") print("=" * 70) try: while self._running: self._cycle() time.sleep(self.config.fetch_interval_s) # Periodic metrics if time.time() - self._last_metrics >= self.config.metrics_interval_s: self._log_metrics() self._last_metrics = time.time() except KeyboardInterrupt: print("\nShutdown...") finally: self._running = False self._save_regimes() self._log_final() def _cycle(self) -> None: """Run one fetch cycle.""" sources = self._pipeline._catalogue.get_enabled() for source in sources: if not self._running: break # Simulate fetching (in production, this would be HTTP) # For now, extract from source metadata new_regimes = self._pipeline.fetch_and_extract( source.source_id, f"Market conditions from {source.name}", ) if new_regimes: for regime in new_regimes: self._discovered_regimes[regime] = { "source": source.source_id, "first_seen": time.time_ns(), "fetch_count": 1, } self._total_regimes += len(new_regimes) LOGGER.info("Discovered %d new regimes: %s", len(new_regimes), new_regimes) def _save_regimes(self) -> None: """Persist discovered regimes to disk.""" try: with open(self.config.regime_db_path, "w") as f: json.dump(self._discovered_regimes, f, indent=2) except OSError as e: LOGGER.error("Failed to save regimes: %s", e) def _load_regimes(self) -> None: """Load persisted regimes from disk.""" if os.path.exists(self.config.regime_db_path): try: with open(self.config.regime_db_path) as f: self._discovered_regimes = json.load(f) self._total_regimes = len(self._discovered_regimes) except Exception: pass def _ensure_tables(self) -> None: """Ensure ClickHouse tables exist.""" if self._store: self._store.ensure_tables() def _log_metrics(self) -> None: elapsed = time.time() - self._start_time stats = self._pipeline.get_source_stats() print(f" [{elapsed:.0f}s] Sources={stats['total_sources']} " f"Fetched={stats['total_fetched']} " f"Regimes={self._total_regimes} " f"Errors={stats['total_errors']}") def _log_final(self) -> None: elapsed = time.time() - self._start_time stats = self._pipeline.get_source_stats() print() print("=" * 70) print("COGNITION PIPELINE FINAL") print("=" * 70) print(f"Duration: {elapsed:.1f}s ({elapsed/60:.1f} min)") print(f"Sources: {stats['total_sources']}") print(f"Fetched: {stats['total_fetched']}") print(f"Regimes: {self._total_regimes}") print(f"Errors: {stats['total_errors']}") print(f"Discovered: {stats['discovered_regimes']}") print("=" * 70) def _signal_handler(self, sig, frame): print("\nShutdown signal received...") self._running = False if __name__ == "__main__": logging.basicConfig(level=logging.INFO) launcher = CognitionLauncher() launcher.run()