Files
sentiment-engine/sentiment_engine/refetch_and_analyze.py

215 lines
8.9 KiB
Python
Raw Permalink Normal View History

#!/usr/bin/env python3
"""
Re-fetch news using proper entity extraction and analyze sentiment for trade assets.
"""
import asyncio
import json
import sys
from datetime import datetime
from pathlib import Path
sys.path.insert(0, 'src')
from sentiment_engine.ingestion.rss import RSSConnector
from sentiment_engine.ingestion.base import ConnectorConfig, ConnectorType
from sentiment_engine.nlp.pipeline import NLPProcessingPipeline
from sentiment_engine.nlp.entity_extraction import EntityExtractor, AssetMapper
from sentiment_engine.schemas.payload import NormalizedPayload, SourceType, AssetMention, EngagementMetrics
# Trade assets we care about
TRADE_ASSETS = ["ZIL", "ONG", "ONE", "STX", "ALGO", "DASH", "LTC", "FET", "XTZ", "LINK", "ENJ", "DOGE", "XLM", "ETC", "TRX", "BTC", "ETH", "SOL", "BNB", "XRP", "ADA", "AVAX", "DOT", "MATIC", "POL", "UNI", "ATOM", "NEAR", "ICP"]
SOURCES = [
{"source_id": "rss:coindesk", "url": "https://www.coindesk.com/arc/outboundfeeds/rss/", "cred": 0.85, "feed_urls": ["https://www.coindesk.com/arc/outboundfeeds/rss/"]},
{"source_id": "rss:cointelegraph", "url": "https://cointelegraph.com/rss", "cred": 0.75, "feed_urls": ["https://cointelegraph.com/rss"]},
{"source_id": "rss:theblock", "url": "https://www.theblock.co/rss", "cred": 0.85, "feed_urls": ["https://www.theblock.co/rss"]},
{"source_id": "rss:decrypt", "url": "https://decrypt.co/feed", "cred": 0.75, "feed_urls": ["https://decrypt.co/feed"]},
{"source_id": "rss:glassnode", "url": "https://insights.glassnode.com/rss/", "cred": 0.85, "feed_urls": ["https://insights.glassnode.com/rss/"]},
{"source_id": "rss:wsj_crypto", "url": "https://feeds.a.dj.com/rss/RSSMarketsMain.xml", "cred": 0.85, "feed_urls": ["https://feeds.a.dj.com/rss/RSSMarketsMain.xml"]},
]
async def fetch_all_articles():
all_articles = []
entity_extractor = EntityExtractor(AssetMapper())
await entity_extractor.initialize()
for src in SOURCES:
config = ConnectorConfig(
source_id=src["source_id"],
connector_type=ConnectorType.RSS,
base_url=src["url"],
cadence_seconds=120,
base_credibility=src["cred"],
relevance=0.9,
extra_config={"feed_urls": src["feed_urls"], "max_items_per_feed": 100}
)
connector = RSSConnector(config)
try:
print(f"\nPolling {src['source_id']}...")
await connector.initialize()
payloads = await connector.poll()
print(f" Got {len(payloads)} items")
for payload in payloads:
title = payload.title or ""
text = payload.raw_text or ""
full_text = f"{title}. {text}"
# Use entity extractor to find asset mentions
entities = await entity_extractor.extract_all(full_text)
asset_ids = [e.asset_id for e in entities]
# Filter for our trade assets
matched_assets = [a for a in asset_ids if a in TRADE_ASSETS]
if matched_assets:
pub_ts = payload.publish_ts or datetime.now().timestamp()
article = {
"source_id": src["source_id"],
"source_credibility": src["cred"],
"title": title,
"content": full_text[:5000],
"url": payload.url,
"publish_ts": pub_ts,
"matched_assets": matched_assets,
"all_entities": asset_ids,
}
all_articles.append(article)
print(f" MATCH: {title[:80]}... | Assets: {matched_assets}")
await connector.close()
except Exception as e:
print(f" ERROR polling {src['source_id']}: {e}")
return all_articles
async def process_through_pipeline(articles):
"""Run articles through the full NLP pipeline"""
print("\n=== INITIALIZING NLP PIPELINE ===")
pipeline = NLPProcessingPipeline()
await pipeline.initialize()
results = []
for article in articles:
# Create asset mentions for matched assets
asset_mentions = []
for asset in article["matched_assets"]:
asset_mentions.append(AssetMention(
asset_id=asset,
mention_span=(0, len(asset)),
confidence=0.9,
source_text=asset,
mention_type="ticker"
))
payload = NormalizedPayload(
source_id=article["source_id"],
source_type=SourceType.NEWS,
source_credibility_base=article["source_credibility"],
ingest_ts=datetime.now().timestamp(),
publish_ts=article["publish_ts"],
asset_mentions=asset_mentions,
raw_text=article["content"],
title=article["title"],
url=article["url"],
author=None,
engagement_metrics=EngagementMetrics(),
content_length=len(article["content"]),
language="en",
metadata={}
)
try:
processed = await pipeline.process(payload)
asset_sentiments = {}
for entity in processed.entities:
asset_key = entity.asset_id
sent = processed.sentiment_per_asset.get(asset_key)
if sent:
if sent.polarity > 0.1:
label = "POSITIVE"
elif sent.polarity < -0.1:
label = "NEGATIVE"
else:
label = "NEUTRAL"
asset_sentiments[asset_key] = {
"polarity": sent.polarity,
"confidence": sent.confidence,
"positive_prob": sent.positive_prob,
"negative_prob": sent.negative_prob,
"neutral_prob": sent.neutral_prob,
"label": label,
}
result = {
"source_id": article["source_id"],
"title": article["title"],
"url": article["url"],
"publish_ts": article["publish_ts"],
"matched_assets": article["matched_assets"],
"all_entities": article["all_entities"],
"asset_sentiments": asset_sentiments,
"events": [{"type": e.event_type.value, "assets": e.assets_involved, "confidence": e.confidence, "severity": e.severity} for e in processed.events],
"credibility": processed.credibility.composite if processed.credibility else 0,
}
results.append(result)
print(f"\n PROCESSED: {article['title'][:70]}...")
for asset, sent in asset_sentiments.items():
print(f" {asset}: polarity={sent['polarity']:.3f} conf={sent['confidence']:.3f} label={sent['label']}")
if result["events"]:
for ev in result["events"]:
print(f" EVENT: {ev['type']} on {ev['assets']} conf={ev['confidence']:.3f}")
except Exception as e:
print(f" ERROR processing: {e}")
import traceback
traceback.print_exc()
return results
async def main():
print("=== FETCHING NEWS WITH PROPER ENTITY EXTRACTION ===")
articles = await fetch_all_articles()
print(f"\n=== TOTAL ARTICLES MATCHED: {len(articles)} ===")
# Save raw articles
with open("trade_news_refetched.json", "w") as f:
json.dump(articles, f, default=str, indent=2)
# Process through pipeline
results = await process_through_pipeline(articles)
# Save results
with open("trade_news_refetched_sentiment.json", "w") as f:
json.dump(results, f, default=str, indent=2)
print(f"\n=== SENTIMENT RESULTS: {len(results)} articles processed ===")
# Summary by asset
from collections import defaultdict
asset_sentiments = defaultdict(list)
for r in results:
for asset, sent in r["asset_sentiments"].items():
asset_sentiments[asset].append(sent)
print("\n=== SENTIMENT SUMMARY BY ASSET ===")
for asset, sents in sorted(asset_sentiments.items()):
avg_pol = sum(s["polarity"] for s in sents) / len(sents)
avg_conf = sum(s["confidence"] for s in sents) / len(sents)
labels = [s["label"] for s in sents]
pos = sum(1 for s in sents if s["polarity"] > 0.1)
neg = sum(1 for s in sents if s["polarity"] < -0.1)
neu = sum(1 for s in sents if -0.1 <= s["polarity"] <= 0.1)
print(f" {asset}: {len(sents)} mentions | avg_polarity={avg_pol:.3f} avg_conf={avg_conf:.3f} | Pos:{pos} Neg:{neg} Neu:{neu}")
if __name__ == "__main__":
asyncio.run(main())