Files
sentiment-engine/MALKHUT/malkhut/training/news_sources.py

196 lines
8.6 KiB
Python
Raw Normal View History

"""
News Source Repository — industry-standard tracking for market regime sources.
Tracks:
- Source metadata (name, URL, type, relevance)
- Fetch history (timestamps, success/failure)
- Source ranking (by regime relevance, freshness, reliability)
- Source health (error rates, response times)
- Auto-discovery of new sources
Format: JSON-based, compatible with standard news aggregation pipelines.
"""
from __future__ import annotations
import hashlib
import json
import os
import time
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional, Set, Tuple
@dataclass(frozen=True, slots=True)
class SourceMetadata:
"""Full metadata for a news/data source."""
source_id: str
name: str
url: str
source_type: str # "news", "data", "research", "exchange", "social"
regime_relevance: float # 0-1
reliability: float # 0-1
freshness_hours: float # how often to fetch
last_fetched_ns: int = 0
fetch_count: int = 0
error_count: int = 0
avg_response_ms: float = 0.0
tags: Tuple[str, ...] = ()
enabled: bool = True
created_ns: int = 0
@property
def health_score(self) -> float:
if self.fetch_count == 0:
return 0.5
error_rate = self.error_count / self.fetch_count
return max(0.0, 1.0 - error_rate) * self.reliability
class NewsSourceRepository:
"""
Industry-standard repository for tracking news/data sources.
Features:
- Source registration and metadata
- Fetch history tracking
- Source ranking by relevance, freshness, reliability
- Source health monitoring
- Auto-discovery hooks
- JSON persistence
Compatible with standard news aggregation pipelines.
"""
def __init__(self, repo_path: str = "news_sources.json") -> None:
self._path = repo_path
self._sources: Dict[str, SourceMetadata] = {}
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"]] = SourceMetadata(**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,
"reliability": s.reliability, "freshness_hours": s.freshness_hours,
"last_fetched_ns": s.last_fetched_ns, "fetch_count": s.fetch_count,
"error_count": s.error_count, "avg_response_ms": s.avg_response_ms,
"tags": list(s.tags), "enabled": s.enabled, "created_ns": s.created_ns,
} for s in self._sources.values()], f, indent=2)
except OSError:
pass
def register(self, source_id: str, name: str, url: str,
source_type: str = "news", regime_relevance: float = 0.5,
reliability: float = 0.8, freshness_hours: float = 1.0,
tags: Tuple[str, ...] = ()) -> None:
"""Register a new source."""
self._sources[source_id] = SourceMetadata(
source_id=source_id, name=name, url=url,
source_type=source_type, regime_relevance=regime_relevance,
reliability=reliability, freshness_hours=freshness_hours,
tags=tags, created_ns=time.time_ns(),
)
self._save()
def record_fetch(self, source_id: str, success: bool, response_ms: float = 0.0) -> None:
"""Record a fetch attempt."""
if source_id in self._sources:
s = self._sources[source_id]
new_count = s.fetch_count + 1
new_errors = s.error_count + (0 if success else 1)
new_avg = ((s.avg_response_ms * s.fetch_count) + response_ms) / new_count
self._sources[source_id] = SourceMetadata(
source_id=s.source_id, name=s.name, url=s.url,
source_type=s.source_type, regime_relevance=s.regime_relevance,
reliability=s.reliability, freshness_hours=s.freshness_hours,
last_fetched_ns=time.time_ns(), fetch_count=new_count,
error_count=new_errors, avg_response_ms=new_avg,
tags=s.tags, enabled=s.enabled, created_ns=s.created_ns,
)
self._save()
def rank_by_relevance(self) -> List[SourceMetadata]:
"""Rank sources by regime relevance."""
return sorted(self._sources.values(), key=lambda s: s.regime_relevance, reverse=True)
def rank_by_health(self) -> List[SourceMetadata]:
"""Rank sources by health score."""
return sorted(self._sources.values(), key=lambda s: s.health_score, reverse=True)
def rank_by_freshness(self) -> List[SourceMetadata]:
"""Rank sources by last fetch time (most recent first)."""
return sorted(self._sources.values(), key=lambda s: s.last_fetched_ns, reverse=True)
def get_enabled(self) -> List[SourceMetadata]:
return [s for s in self._sources.values() if s.enabled]
def get_by_type(self, source_type: str) -> List[SourceMetadata]:
return [s for s in self._sources.values() if s.source_type == source_type]
def get_by_tag(self, tag: str) -> List[SourceMetadata]:
return [s for s in self._sources.values() if tag in s.tags]
def disable(self, source_id: str) -> None:
if source_id in self._sources:
s = self._sources[source_id]
self._sources[source_id] = SourceMetadata(
source_id=s.source_id, name=s.name, url=s.url,
source_type=s.source_type, regime_relevance=s.regime_relevance,
reliability=s.reliability, freshness_hours=s.freshness_hours,
last_fetched_ns=s.last_fetched_ns, fetch_count=s.fetch_count,
error_count=s.error_count, avg_response_ms=s.avg_response_ms,
tags=s.tags, enabled=False, created_ns=s.created_ns,
)
self._save()
def enable(self, source_id: str) -> None:
if source_id in self._sources:
s = self._sources[source_id]
self._sources[source_id] = SourceMetadata(
source_id=s.source_id, name=s.name, url=s.url,
source_type=s.source_type, regime_relevance=s.regime_relevance,
reliability=s.reliability, freshness_hours=s.freshness_hours,
last_fetched_ns=s.last_fetched_ns, fetch_count=s.fetch_count,
error_count=s.error_count, avg_response_ms=s.avg_response_ms,
tags=s.tags, enabled=True, created_ns=s.created_ns,
)
self._save()
@property
def source_count(self) -> int:
return len(self._sources)
@property
def enabled_count(self) -> int:
return len(self.get_enabled())
def seed_defaults(self) -> None:
"""Seed with industry-standard crypto news sources."""
defaults = [
("coindesk", "CoinDesk", "https://www.coindesk.com", "news", 0.8, 0.9, 1.0, ("crypto", "news")),
("cointelegraph", "Cointelegraph", "https://cointelegraph.com", "news", 0.7, 0.85, 1.0, ("crypto", "news")),
("the_block", "The Block", "https://www.theblock.co", "news", 0.8, 0.9, 0.5, ("crypto", "news", "research")),
("cryptoquant", "CryptoQuant", "https://cryptoquant.com", "data", 0.9, 0.95, 4.0, ("on_chain", "data")),
("glassnode", "Glassnode", "https://glassnode.com", "data", 0.9, 0.95, 4.0, ("on_chain", "data")),
("coinglass", "Coinglass", "https://www.coinglass.com", "data", 0.8, 0.9, 1.0, ("derivatives", "data")),
("binance_research", "Binance Research", "https://www.binance.com/en/research", "research", 0.7, 0.85, 24.0, ("exchange", "research")),
("messari", "Messari", "https://messari.io", "research", 0.7, 0.85, 24.0, ("research", "fundamentals")),
("defillama", "DefiLlama", "https://defillama.com", "data", 0.6, 0.9, 1.0, ("defi", "data")),
("dune", "Dune Analytics", "https://dune.com", "data", 0.7, 0.85, 4.0, ("on_chain", "data")),
("the_blockResearch", "The Block Research", "https://www.theblock.co/research", "research", 0.8, 0.9, 24.0, ("research", "institutional")),
("coingecko", "CoinGecko", "https://www.coingecko.com", "data", 0.6, 0.85, 1.0, ("market_data", "data")),
]
for sid, name, url, stype, relevance, reliability, freshness, tags in defaults:
self.register(sid, name, url, stype, relevance, reliability, freshness, tags)