- Added 30 new sources (5 RSS + 25 Telegram) for previously ZERO-coverage assets - Fixed model loading priority: ONNX > LoRA v2 > PyTorch > Mock - ONNX FinBERT (pre-trained on 1.2M financial docs) now PRIMARY - best for real-world text - LoRA v2 models trained on 518 carefully labeled samples (balanced Bearish/Bullish/Neutral) - Emotion LoRA v2 trained with weighted loss (greed/fear 2x, joy 1.5x) - 30 new sources: STX, FET, XTZ, ENJ, ETC, TRX, ONG, DASH, LTC, ZIL, NEAR, APT, SUI, ICP - Early stopping (patience=3) on both LoRA trainings - Human-in-the-loop verification CLI tool created - Disk-conscious: save_total_limit=1, adapters 6-8MB each Pipeline now correctly classifies: - BTC breaks 100k → +0.54 Bullish ✅ - Major hack → -0.23 Bearish ✅ - HODL → +0.91 Bullish ✅ - Rug pull → -0.30 Bearish ✅ - SEC sues → -0.30 Bearish ✅ - ETF approval → +0.32 Bullish ✅ - Whale accumulation → +0.31 Bullish ✅ Models: ONNX FinBERT (PRIORITY 1) + LoRA v2 adapters (6-8MB each) Training data: 518 carefully labeled samples (190 real + 328 synthetic) Early stopping (patience=3) on both FinBERT and DistilRoBERTa LoRA Emotion LoRA v2: weighted loss (greed/fear 2x, joy 1.5x) + early stopping
129 lines
3.7 KiB
Python
129 lines
3.7 KiB
Python
"""Web crawl Prefect flow"""
|
|
|
|
import asyncio
|
|
import logging
|
|
import subprocess
|
|
import tempfile
|
|
from pathlib import Path
|
|
from typing import List
|
|
|
|
from prefect import flow, task
|
|
|
|
from sentiment_engine.utils.config import get_settings
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@task(retries=1, retry_delay_seconds=300)
|
|
async def run_hister_crawl(
|
|
seed_urls: List[str],
|
|
allowed_domains: List[str],
|
|
max_depth: int = 2,
|
|
job_timeout: int = 3600
|
|
) -> List[dict]:
|
|
"""Run Hister crawl job"""
|
|
with tempfile.TemporaryDirectory() as tmpdir:
|
|
seed_file = Path(tmpdir) / "seeds.txt"
|
|
seed_file.write_text("\n".join(seed_urls))
|
|
output_file = Path(tmpdir) / "output.jsonl"
|
|
|
|
cmd = [
|
|
"hister", "crawl",
|
|
"--input", str(seed_file),
|
|
"--job-id", f"prefect-crawl-{asyncio.current_task().get_name()}",
|
|
"--depth", str(max_depth),
|
|
"--delay", "1.0",
|
|
"--output", str(output_file),
|
|
"--format", "jsonl"
|
|
]
|
|
|
|
if allowed_domains:
|
|
cmd.extend(["--allowed-domain", ",".join(allowed_domains)])
|
|
|
|
try:
|
|
proc = await asyncio.create_subprocess_exec(
|
|
*cmd,
|
|
stdout=asyncio.subprocess.PIPE,
|
|
stderr=asyncio.subprocess.PIPE
|
|
)
|
|
|
|
stdout, stderr = await asyncio.wait_for(
|
|
proc.communicate(), timeout=job_timeout
|
|
)
|
|
|
|
if proc.returncode != 0:
|
|
logger.error(f"Hister failed: {stderr.decode()}")
|
|
return []
|
|
|
|
# Parse output
|
|
items = []
|
|
if output_file.exists():
|
|
import json
|
|
with open(output_file) as f:
|
|
for line in f:
|
|
line = line.strip()
|
|
if not line:
|
|
continue
|
|
try:
|
|
data = json.loads(line)
|
|
items.append(data)
|
|
except json.JSONDecodeError:
|
|
continue
|
|
|
|
return items
|
|
|
|
except asyncio.TimeoutError:
|
|
logger.error(f"Hister job timed out after {job_timeout}s")
|
|
return []
|
|
except FileNotFoundError:
|
|
logger.error("Hister not installed")
|
|
return []
|
|
|
|
|
|
@flow(
|
|
name="web_crawl",
|
|
log_prints=True
|
|
)
|
|
async def web_crawl_flow(
|
|
seed_urls: List[str] = None,
|
|
allowed_domains: List[str] = None,
|
|
max_depth: int = 2
|
|
):
|
|
"""Web crawl flow for sites without RSS/API"""
|
|
if seed_urls is None:
|
|
seed_urls = [
|
|
"https://www.coindesk.com",
|
|
"https://cointelegraph.com",
|
|
"https://www.theblock.co",
|
|
"https://decrypt.co",
|
|
"https://cryptoslate.com",
|
|
]
|
|
|
|
if allowed_domains is None:
|
|
allowed_domains = [
|
|
"coindesk.com", "cointelegraph.com", "theblock.co",
|
|
"decrypt.co", "cryptoslate.com", "bitcoinmagazine.com"
|
|
]
|
|
|
|
items = await run_hister_crawl(seed_urls, allowed_domains, max_depth)
|
|
|
|
logger.info(f"Crawled {len(items)} pages")
|
|
|
|
# Convert to normalized items
|
|
normalized = []
|
|
for item in items:
|
|
normalized.append({
|
|
"source_id": f"web:{item.get('url', '').split('/')[2] if item.get('url') else 'unknown'}",
|
|
"source_type": "news",
|
|
"raw_text": f"{item.get('title', '')}\n\n{item.get('content', item.get('text', ''))}",
|
|
"title": item.get("title"),
|
|
"url": item.get("url"),
|
|
"metadata": {"crawler": "hister", "job": "prefect"}
|
|
})
|
|
|
|
return normalized
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(web_crawl_flow())
|