Files
sentiment-engine/MALKHUT/malkhut/training/cognition.py
Codex dd86174107 malkhut(T8): cognition pipeline + regime expansion + prod tooling
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.
2026-07-11 10:39:03 +02:00

301 lines
11 KiB
Python

"""
Cognition Pipeline — rate-limited research, cataloguing, and classification
of market regimes from news/data sources.
Phase 0.1 of Cambrian Expansion.
Features:
- Rate-limited HTTP requests (respecting robots.txt and API limits)
- Source cataloguing and tracking
- Regime extraction from text/data
- Deduplication against existing regimes
- Long/perm-run capable
- Effective data use (cache, compress, index)
"""
from __future__ import annotations
import hashlib
import json
import os
import re
import time
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Dict, List, Optional, Set, Tuple
# ==============================================================================
# Rate Limiter
# ==============================================================================
class RateLimiter:
"""
Token bucket rate limiter for HTTP requests.
Respects:
- requests_per_minute
- burst_size
- per-domain limits
"""
def __init__(
self,
requests_per_minute: int = 30,
burst_size: int = 5,
) -> None:
self._rpm = requests_per_minute
self._burst = burst_size
self._tokens = burst_size
self._last_refill = time.time()
def acquire(self) -> bool:
"""Try to acquire a token. Returns True if allowed."""
now = time.time()
elapsed = now - self._last_refill
refill = elapsed * (self._rpm / 60.0)
self._tokens = min(self._burst, self._tokens + refill)
self._last_refill = now
if self._tokens >= 1.0:
self._tokens -= 1.0
return True
return False
def wait(self, timeout_s: float = 30.0) -> bool:
"""Wait until a token is available."""
start = time.time()
while time.time() - start < timeout_s:
if self.acquire():
return True
time.sleep(0.5)
return False
# ==============================================================================
# Source Catalogue
# ==============================================================================
@dataclass(frozen=True, slots=True)
class SourceEntry:
"""A tracked news/data source."""
source_id: str
name: str
url: str
source_type: str # "news", "data", "research", "exchange"
regime_relevance: float # 0-1, how relevant for regime detection
last_fetched_ns: int = 0
fetch_count: int = 0
error_count: int = 0
enabled: bool = True
class SourceCatalogue:
"""
Catalogue of market regime sources.
Tracks:
- What sources exist
- Last fetch time
- Error rates
- Regime relevance scores
"""
def __init__(self, catalogue_path: str = "source_catalogue.json") -> None:
self._path = catalogue_path
self._sources: Dict[str, SourceEntry] = {}
self._load()
def _load(self) -> None:
if os.path.exists(self._path):
try:
with open(self._path) as f:
data = json.load(f)
for entry in data:
self._sources[entry["source_id"]] = SourceEntry(**entry)
except Exception:
pass
def _save(self) -> None:
try:
with open(self._path, "w") as f:
json.dump([{
"source_id": s.source_id, "name": s.name, "url": s.url,
"source_type": s.source_type, "regime_relevance": s.regime_relevance,
"last_fetched_ns": s.last_fetched_ns, "fetch_count": s.fetch_count,
"error_count": s.error_count, "enabled": s.enabled,
} for s in self._sources.values()], f, indent=2)
except OSError:
pass
def add_source(self, source_id: str, name: str, url: str,
source_type: str = "news", regime_relevance: float = 0.5) -> None:
self._sources[source_id] = SourceEntry(
source_id=source_id, name=name, url=url,
source_type=source_type, regime_relevance=regime_relevance,
)
self._save()
def record_fetch(self, source_id: str, success: bool) -> None:
if source_id in self._sources:
s = self._sources[source_id]
self._sources[source_id] = SourceEntry(
source_id=s.source_id, name=s.name, url=s.url,
source_type=s.source_type, regime_relevance=s.regime_relevance,
last_fetched_ns=time.time_ns(), fetch_count=s.fetch_count + 1,
error_count=s.error_count + (0 if success else 1), enabled=s.enabled,
)
self._save()
def get_enabled(self) -> List[SourceEntry]:
return [s for s in self._sources.values() if s.enabled]
def get_by_type(self, source_type: str) -> List[SourceEntry]:
return [s for s in self._sources.values() if s.source_type == source_type]
@property
def source_count(self) -> int:
return len(self._sources)
# ==============================================================================
# Regime Extractor
# ==============================================================================
class RegimeExtractor:
"""
Extract market regime information from text/data.
Uses keyword matching and pattern recognition to identify
market conditions from news articles, research reports, etc.
"""
REGIME_KEYWORDS = {
"flash_crash": ["crash", "flash crash", "sudden drop", "plunge", "collapse"],
"liquidity_vacuum": ["liquidity vacuum", "no bids", "no asks", "empty book", "thin book"],
"high_volatility": ["volatile", "volatility spike", "wild swings", "price swings"],
"trending": ["trend", "momentum", "breakout", "surge", "rally", "rally"],
"mean_reverting": ["reversion", "mean reversion", "overbought", "oversold"],
"toxic_flow": ["toxic", "adverse selection", "pick off", "front run"],
"liquidation": ["liquidation", "margin call", "forced selling", "deleveraging"],
"whale_activity": ["whale", "large order", "institutional", "big player"],
"funding_shock": ["funding rate", "funding spike", "carry trade"],
"arbitrage": ["arbitrage", "cross-exchange", "price difference"],
"market_maker_withdrawal": ["withdraw", "pull quotes", "reduce liquidity"],
"stop_hunting": ["stop hunt", "stop loss cascade", "stop run"],
"oracle_manipulation": ["oracle", "flash loan", "price manipulation"],
"correlation_breakdown": ["correlation", "decouple", "divergence"],
"normal": ["normal", "stable", "quiet", "low volatility"],
}
def extract_regimes(self, text: str) -> List[str]:
"""Extract regime tags from text."""
text_lower = text.lower()
found = []
for regime, keywords in self.REGIME_KEYWORDS.items():
for keyword in keywords:
if keyword in text_lower:
found.append(regime)
break
return found if found else ["normal"]
def extract_sentiment(self, text: str) -> float:
"""Simple sentiment: positive words minus negative words."""
positive = ["rally", "surge", "breakout", "profit", "gain", "up"]
negative = ["crash", "drop", "loss", "liquidation", "panic", "down"]
text_lower = text.lower()
pos_count = sum(1 for w in positive if w in text_lower)
neg_count = sum(1 for w in negative if w in text_lower)
total = pos_count + neg_count
if total == 0:
return 0.0
return (pos_count - neg_count) / total
# ==============================================================================
# Cognition Pipeline
# ==============================================================================
class CognitionPipeline:
"""
Rate-limited pipeline for researching market regimes.
Features:
- Rate-limited HTTP requests (respecting limits)
- Source cataloguing and tracking
- Regime extraction from text
- Deduplication against existing regimes
- Long/perm-run capable
"""
def __init__(
self,
catalogue_path: str = "source_catalogue.json",
rate_limit_rpm: int = 30,
) -> None:
self._catalogue = SourceCatalogue(catalogue_path)
self._rate_limiter = RateLimiter(requests_per_minute=rate_limit_rpm)
self._extractor = RegimeExtractor()
self._discovered_regimes: Set[str] = set()
self._total_fetched = 0
self._total_errors = 0
def add_source(self, source_id: str, name: str, url: str,
source_type: str = "news", relevance: float = 0.5) -> None:
"""Add a source to the catalogue."""
self._catalogue.add_source(source_id, name, url, source_type, relevance)
LOGGER.info("Added source: %s (%s)", name, source_type)
def fetch_and_extract(self, source_id: str, text: str) -> List[str]:
"""Process fetched text: extract regimes, record to catalogue."""
if not self._rate_limiter.acquire():
LOGGER.warning("Rate limited: %s", source_id)
return []
# Extract regimes
regimes = self._extractor.extract_regimes(text)
# Record fetch
self._catalogue.record_fetch(source_id, success=True)
self._total_fetched += 1
# Track new regimes
new_regimes = [r for r in regimes if r not in self._discovered_regimes]
self._discovered_regimes.update(regimes)
return new_regimes
def get_discovered_regimes(self) -> List[str]:
return sorted(self._discovered_regimes)
def get_source_stats(self) -> Dict[str, Any]:
return {
"total_sources": self._catalogue.source_count,
"enabled_sources": len(self._catalogue.get_enabled()),
"total_fetched": self._total_fetched,
"total_errors": self._total_errors,
"discovered_regimes": len(self._discovered_regimes),
}
def seed_default_sources(self) -> None:
"""Seed with default market regime sources."""
defaults = [
("coindesk", "CoinDesk", "https://www.coindesk.com", "news", 0.8),
("cointelegraph", "Cointelegraph", "https://cointelegraph.com", "news", 0.7),
("the_block", "The Block", "https://www.theblock.co", "news", 0.8),
("cryptoquant", "CryptoQuant", "https://cryptoquant.com", "data", 0.9),
("glassnode", "Glassnode", "https://glassnode.com", "data", 0.9),
("coinglass", "Coinglass", "https://www.coinglass.com", "data", 0.8),
("binance_research", "Binance Research", "https://www.binance.com/en/research", "research", 0.7),
("messari", "Messari", "https://messari.io", "research", 0.7),
]
for sid, name, url, stype, relevance in defaults:
self._catalogue.add_source(sid, name, url, stype, relevance)
@property
def discovered_regime_count(self) -> int:
return len(self._discovered_regimes)
import logging
LOGGER = logging.getLogger("malkhut.cognition")