226 lines
8.2 KiB
Python
226 lines
8.2 KiB
Python
"""
|
|
Bank forecast scraper — fetches public research pages, extracts numeric forecasts via LLM.
|
|
Forecasts are stored in bank_forecasts and optionally pushed to macro_series_log as consensus.
|
|
"""
|
|
import json
|
|
import logging
|
|
import re
|
|
from datetime import datetime, timezone, timedelta
|
|
from typing import Optional
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Max article text sent to LLM (chars)
|
|
_MAX_TEXT = 6000
|
|
|
|
|
|
def _fetch_text(url: str, timeout: int = 15) -> str:
|
|
"""HTTP GET → plain text (strips HTML tags)."""
|
|
import urllib.request
|
|
import html
|
|
req = urllib.request.Request(url, headers={
|
|
"User-Agent": "Mozilla/5.0 (compatible; OpenFin-Research/1.0)"
|
|
})
|
|
with urllib.request.urlopen(req, timeout=timeout) as resp:
|
|
raw = resp.read().decode("utf-8", errors="replace")
|
|
# Strip HTML tags
|
|
text = re.sub(r"<[^>]+>", " ", raw)
|
|
text = html.unescape(text)
|
|
text = re.sub(r"\s+", " ", text).strip()
|
|
return text[:_MAX_TEXT]
|
|
|
|
|
|
def _upcoming_events(conn, days_ahead: int = 14) -> list[dict]:
|
|
"""Return ff_calendar events with series_id coming in the next N days."""
|
|
today = datetime.now(timezone.utc).date().isoformat()
|
|
horizon = (datetime.now(timezone.utc).date() + timedelta(days=days_ahead)).isoformat()
|
|
rows = conn.execute(
|
|
"""SELECT DISTINCT series_id, event_name, event_date, currency, impact
|
|
FROM ff_calendar
|
|
WHERE series_id IS NOT NULL AND series_id != ''
|
|
AND event_date BETWEEN ? AND ?
|
|
AND impact IN ('High', 'Medium')
|
|
ORDER BY event_date""",
|
|
(today, horizon)
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
|
|
def _llm_extract(article_text: str, upcoming: list[dict]) -> list[dict]:
|
|
"""
|
|
Call Claude to extract numeric forecasts for upcoming events from article text.
|
|
Returns list of {series_id, event_name, event_date, forecast_value, unit, confidence, snippet}.
|
|
"""
|
|
from services.database import get_config as _cfg
|
|
import anthropic
|
|
|
|
api_key = _cfg("anthropic_api_key") or ""
|
|
if not api_key:
|
|
logger.warning("[bank_forecast] No anthropic_api_key configured")
|
|
return []
|
|
|
|
if not upcoming:
|
|
return []
|
|
|
|
events_block = "\n".join(
|
|
f"- {e['event_name']} (series_id={e['series_id']}, release={e['event_date']}, currency={e['currency']})"
|
|
for e in upcoming
|
|
)
|
|
|
|
prompt = f"""You are extracting economic forecasts from a bank research note.
|
|
|
|
Upcoming economic releases:
|
|
{events_block}
|
|
|
|
From the text below, extract any numeric forecasts or expectations for the above indicators.
|
|
Return a JSON array (no commentary, no markdown). Each item:
|
|
{{ "series_id": "...", "event_name": "...", "event_date": "YYYY-MM-DD",
|
|
"forecast_value": <number>, "unit": "%|K|pp|...",
|
|
"confidence": "high|medium|low",
|
|
"snippet": "<exact quote from text, max 120 chars>" }}
|
|
|
|
Only include entries where a clear numeric forecast is stated. Return [] if nothing found.
|
|
|
|
TEXT:
|
|
{article_text}"""
|
|
|
|
client = anthropic.Anthropic(api_key=api_key)
|
|
msg = client.messages.create(
|
|
model="claude-haiku-4-5-20251001",
|
|
max_tokens=1024,
|
|
messages=[{"role": "user", "content": prompt}]
|
|
)
|
|
raw = msg.content[0].text.strip()
|
|
# Extract JSON array from response
|
|
m = re.search(r"\[.*\]", raw, re.DOTALL)
|
|
if not m:
|
|
return []
|
|
try:
|
|
return json.loads(m.group())
|
|
except Exception as e:
|
|
logger.warning(f"[bank_forecast] JSON parse error: {e} — raw: {raw[:200]}")
|
|
return []
|
|
|
|
|
|
def scrape_source(conn, source: dict) -> dict:
|
|
"""
|
|
Scrape one bank source: fetch URL, extract forecasts via LLM, save to bank_forecasts.
|
|
Returns {source_id, name, fetched, extracted, saved, error}.
|
|
"""
|
|
sid = source["id"]
|
|
url = (source.get("url") or "").strip()
|
|
if not url:
|
|
return {"source_id": sid, "name": source["name"], "fetched": False,
|
|
"extracted": 0, "saved": 0, "error": "no URL configured"}
|
|
|
|
# Fetch article text
|
|
try:
|
|
text = _fetch_text(url)
|
|
except Exception as e:
|
|
return {"source_id": sid, "name": source["name"], "fetched": False,
|
|
"extracted": 0, "saved": 0, "error": str(e)}
|
|
|
|
# Get upcoming events to focus extraction
|
|
upcoming = _upcoming_events(conn)
|
|
if not upcoming:
|
|
return {"source_id": sid, "name": source["name"], "fetched": True,
|
|
"extracted": 0, "saved": 0, "error": "no upcoming events with series_id"}
|
|
|
|
# LLM extraction
|
|
try:
|
|
forecasts = _llm_extract(text, upcoming)
|
|
except Exception as e:
|
|
return {"source_id": sid, "name": source["name"], "fetched": True,
|
|
"extracted": 0, "saved": 0, "error": f"LLM error: {e}"}
|
|
|
|
# Save results
|
|
saved = 0
|
|
now = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
|
|
for f in forecasts:
|
|
series_id = f.get("series_id", "").strip()
|
|
event_date = f.get("event_date", "").strip()
|
|
forecast_value = f.get("forecast_value")
|
|
if not series_id or not event_date or forecast_value is None:
|
|
continue
|
|
try:
|
|
conn.execute(
|
|
"""INSERT INTO bank_forecasts
|
|
(source_id, series_id, event_name, event_date, forecast_value,
|
|
unit, confidence, snippet, source_url, extracted_at)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?)
|
|
ON CONFLICT(source_id, series_id, event_date) DO UPDATE SET
|
|
forecast_value=excluded.forecast_value,
|
|
snippet=excluded.snippet,
|
|
extracted_at=excluded.extracted_at""",
|
|
(sid, series_id, f.get("event_name", ""), event_date,
|
|
float(forecast_value), f.get("unit", ""),
|
|
f.get("confidence", "medium"), f.get("snippet", "")[:200],
|
|
url, now)
|
|
)
|
|
saved += 1
|
|
except Exception as e:
|
|
logger.warning(f"[bank_forecast] save error {sid}/{series_id}: {e}")
|
|
|
|
conn.execute("UPDATE bank_forecast_sources SET last_scraped=? WHERE id=?", (now, sid))
|
|
conn.commit()
|
|
|
|
return {"source_id": sid, "name": source["name"], "fetched": True,
|
|
"extracted": len(forecasts), "saved": saved, "error": None}
|
|
|
|
|
|
def scrape_all_active(conn) -> list[dict]:
|
|
"""Scrape all active bank sources. Returns list of per-source results."""
|
|
sources = conn.execute(
|
|
"SELECT id, name, url, last_scraped FROM bank_forecast_sources WHERE active=1"
|
|
).fetchall()
|
|
results = []
|
|
for s in sources:
|
|
res = scrape_source(conn, dict(s))
|
|
results.append(res)
|
|
logger.info(f"[bank_forecast] {res['name']}: fetched={res['fetched']} saved={res['saved']} err={res.get('error')}")
|
|
return results
|
|
|
|
|
|
def push_consensus_to_log(conn, series_id: str, event_date: str) -> Optional[float]:
|
|
"""
|
|
Compute average of all bank forecasts for (series_id, event_date),
|
|
push to macro_series_log as source='bank_consensus'.
|
|
Returns consensus value or None.
|
|
"""
|
|
rows = conn.execute(
|
|
"SELECT forecast_value FROM bank_forecasts WHERE series_id=? AND event_date=? AND forecast_value IS NOT NULL",
|
|
(series_id, event_date)
|
|
).fetchall()
|
|
if not rows:
|
|
return None
|
|
|
|
values = [r[0] for r in rows]
|
|
consensus = round(sum(values) / len(values), 4)
|
|
|
|
# Get event_name from any row
|
|
meta = conn.execute(
|
|
"SELECT event_name FROM bank_forecasts WHERE series_id=? AND event_date=? LIMIT 1",
|
|
(series_id, event_date)
|
|
).fetchone()
|
|
event_name = meta[0] if meta else series_id
|
|
|
|
from services.macro_series_log import log_if_changed
|
|
log_if_changed(
|
|
conn, series_id=series_id, event_name=event_name, event_date=event_date,
|
|
actual_value=None, forecast_value=consensus, previous_value=None,
|
|
source="bank_consensus"
|
|
)
|
|
return consensus
|
|
|
|
|
|
def push_all_consensus(conn) -> list[dict]:
|
|
"""Push consensus for every (series_id, event_date) that has bank forecasts."""
|
|
pairs = conn.execute(
|
|
"SELECT DISTINCT series_id, event_date FROM bank_forecasts"
|
|
).fetchall()
|
|
results = []
|
|
for series_id, event_date in pairs:
|
|
val = push_consensus_to_log(conn, series_id, event_date)
|
|
results.append({"series_id": series_id, "event_date": event_date, "consensus": val})
|
|
return results
|