""" 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": , "unit": "%|K|pp|...", "confidence": "high|medium|low", "snippet": "" }} 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