Cognition pipeline (cognition.py): rate-limited, 8 sources, dedup, perm-run. Regime expansion (regime_expansion.py): 200+ regimes from 4x4x4x4 dimensions. News sources (news_sources.py): 12 industry-standard sources with ranking. Monitor (monitor.py): metrics, health scoring, alerts, JSONL logging. Cognition launcher (cognition_launcher.py): standalone long-run service. Continuous pipeline (continuous_pipeline.py): forever-loop training runner.
192 lines
6.3 KiB
Python
192 lines
6.3 KiB
Python
"""
|
|
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()
|