Files
OpenFin/backend/services/market_event_detector.py
OpenSquared bd28b6a73a feat: origin tracing on all market_events
- DB: colonne origin (migration + UPDATE heuristique sur données legacy)
- save/update_market_event: persist origin
- Tous les points de création taguent leur origine:
    bootstrap_macro/eco/ma/legacy | detector_news/eco/technical/report | manual
- UI MarketEvents: badge d'origine avec icône + description dans le panneau détail,
  icône tooltip dans la liste gauche, message explicite si pas de source_refs

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-25 20:58:05 +02:00

530 lines
19 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
Isolated cycle action: Check New Market Events.
Scans 4 sources and creates market_events for significant findings:
- news : geopolitical/macro news (RSS feeds, rule-scored)
- eco : FRED economic releases with high surprise z-score
- technical: MA50/MA100/MA200 crossovers on key instruments
- reports : institutional reports (COT, EIA) with high importance
After each event is created, instrument impacts are evaluated immediately
via the AI (impact_service.evaluate_event_impacts).
"""
import json
import logging
from datetime import datetime, timedelta
from typing import Any, Dict, List, Optional
logger = logging.getLogger(__name__)
WATCH_INSTRUMENTS = [
"SPY", "QQQ", "IWM", "EEM", "EFA",
"GLD", "SLV", "USO", "TLT", "HYG",
"USDJPY=X", "EURUSD=X", "VXX", "NVDA", "BTC-USD",
]
SUBTYPE_FROM_SERIES = {
"UNRATE": "NFP", "PAYEMS": "NFP",
"CPIAUCSL": "CPI", "CPILFESL": "CPI",
"A191RL1Q225SBEA": "GDP", "GDP": "GDP",
"FEDFUNDS": "FOMC", "DFF": "FOMC",
"PCEPILFE": "PCE", "PCEPI": "PCE",
"BAMLH0A0HYM2": "Credit", "BAMLC0A0CM": "Credit",
}
# ── Helpers ───────────────────────────────────────────────────────────────────
def _get_api_key() -> str:
import os
key = os.environ.get("OPENAI_API_KEY", "")
if not key:
from services.database import get_config
key = get_config("openai_api_key") or ""
return key
def _existing_event_keys() -> set:
from services.database import get_all_market_events
return {ev["name"].lower()[:50] for ev in get_all_market_events()}
def _is_dup(name: str, existing: set) -> bool:
return name.lower()[:50] in existing
def _parse_date(raw: str) -> str:
if not raw:
return datetime.utcnow().strftime("%Y-%m-%d")
try:
return datetime.fromisoformat(raw[:19]).strftime("%Y-%m-%d")
except Exception:
pass
try:
from email.utils import parsedate_to_datetime
return parsedate_to_datetime(raw).strftime("%Y-%m-%d")
except Exception:
pass
return raw[:10] if len(raw) >= 10 else datetime.utcnow().strftime("%Y-%m-%d")
def _save_and_evaluate(ev: Dict, existing: set) -> Optional[Dict]:
"""Save a market_event and immediately evaluate instrument impacts. Returns created dict or None."""
from services.database import save_market_event
try:
event_id = save_market_event(ev)
existing.add(ev["name"].lower()[:50])
logger.info(f"[check_events] ✓ saved event #{event_id}: {ev['name']}")
except Exception as e:
logger.error(f"[check_events] save failed for '{ev['name']}': {e}")
return None
# Evaluate instrument impacts immediately
try:
from services.impact_service import evaluate_event_impacts
evaluate_event_impacts(event_id, force=False)
except Exception as e:
logger.warning(f"[check_events] impact eval failed for #{event_id}: {e}")
return {"name": ev["name"], "category": ev.get("category", ""), "date": ev.get("start_date", ""), "event_id": event_id}
# ── Source 1: Geopolitical / macro news ──────────────────────────────────────
def _check_news(
min_impact: float = 0.55,
lookback_hours: int = 48,
max_to_evaluate: int = 15,
) -> List[Dict[str, Any]]:
from services.data_fetcher import fetch_geo_news
api_key = _get_api_key()
if not api_key:
logger.warning("[check_events/news] no OpenAI key — skipping")
return []
try:
all_news = fetch_geo_news()
except Exception as e:
logger.warning(f"[check_events/news] fetch failed: {e}")
return []
cutoff_dt = datetime.utcnow() - timedelta(hours=lookback_hours)
candidates = []
for n in all_news:
if (n.get("impact_score") or 0) < min_impact:
continue
pub_date = _parse_date(n.get("date", ""))
try:
if datetime.fromisoformat(pub_date) < cutoff_dt:
continue
except Exception:
pass
candidates.append(n)
candidates = candidates[:max_to_evaluate]
if not candidates:
return []
existing = _existing_event_keys()
created: List[Dict] = []
try:
from openai import OpenAI
client = OpenAI(api_key=api_key)
except Exception as e:
logger.warning(f"[check_events/news] OpenAI init failed: {e}")
return []
for n in candidates:
title = n.get("title", "")
if not title or _is_dup(title, existing):
continue
prompt = f"""Tu es un analyste macro. Cette news représente-t-elle un événement marché structurant qui mérite un enregistrement permanent ?
TITRE: {title}
SOURCE: {n.get('source', '')}
DATE: {n.get('date', '')}
RÉSUMÉ: {str(n.get('summary', ''))[:400]}
SCORE IMPACT (règle): {n.get('impact_score', 0):.2f}
Réponds OUI seulement si c'est un fait avéré, pas une rumeur ou une opinion, et qu'il a un impact macro ou géopolitique mesurable.
FORMAT JSON STRICT:
{{
"qualifies": true/false,
"reason": "une phrase",
"name": "Nom court (≤ 60 chars)",
"category": "geopolitical|event_calendar|fundamental|report",
"sub_type": "ex: Conflit, Tarifs, Sanctions, OPEC+, Crise...",
"description": "1-2 phrases analytiques",
"affected_assets": ["SPY","GLD",...],
"impact_score": 0.5,
"level": "short|medium|long"
}}"""
try:
resp = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
response_format={"type": "json_object"},
temperature=0.1,
max_tokens=350,
)
parsed = json.loads(resp.choices[0].message.content)
except Exception as e:
logger.debug(f"[check_events/news] AI call failed for '{title[:40]}': {e}")
continue
if not parsed.get("qualifies"):
continue
ev_name = (parsed.get("name") or title)[:60]
if _is_dup(ev_name, existing):
continue
source_ref = {
"title": title,
"source": n.get("source", ""),
"url": n.get("url") or n.get("link", ""),
"date": _parse_date(n.get("date", "")),
"original_score": round(float(n.get("impact_score", 0)), 3),
}
ev = {
"name": ev_name,
"start_date": _parse_date(n.get("date", "")),
"level": parsed.get("level", "short"),
"category": parsed.get("category", "geopolitical"),
"sub_type": parsed.get("sub_type", ""),
"description": parsed.get("description", title),
"market_impact": "",
"affected_assets": parsed.get("affected_assets", []),
"impact_score": float(parsed.get("impact_score", 0.6)),
"source_refs": [source_ref],
"origin": "detector_news",
}
result = _save_and_evaluate(ev, existing)
if result:
result["source"] = "news"
created.append(result)
return created
# ── Source 2: Eco calendar — FRED surprises ───────────────────────────────────
def _check_eco(z_threshold: float = 1.5, days: int = 7) -> List[Dict[str, Any]]:
from services.database import get_recent_economic_surprises
try:
releases = get_recent_economic_surprises(days=days, min_zscore=z_threshold)
except Exception as e:
logger.warning(f"[check_events/eco] query failed: {e}")
return []
existing = _existing_event_keys()
created: List[Dict] = []
for rel in releases:
z = abs(rel.get("surprise_zscore") or 0)
s_pct = rel.get("surprise_pct") or 0
ev_date = (rel.get("event_date") or "")[:10]
s_id = rel.get("series_id", "")
ev_name_base = rel.get("event_name", s_id)
direction = rel.get("surprise_direction", "neutral")
sign = "+" if s_pct >= 0 else ""
ev_name = f"{ev_name_base} — Surprise {sign}{s_pct:.1f}% ({ev_date[:7]})"
if _is_dup(ev_name, existing):
continue
sub_type = SUBTYPE_FROM_SERIES.get(s_id, s_id[:10]) if s_id else ev_name_base[:10]
level = "long" if z >= 3 else ("medium" if z >= 2 else "short")
assets = rel.get("assets_impacted") or []
if isinstance(assets, str):
try:
assets = json.loads(assets)
except Exception:
assets = []
source_ref = {
"title": f"FRED release: {ev_name_base} ({ev_date})",
"source": "FRED",
"url": f"https://fred.stlouisfed.org/series/{s_id}" if s_id else "",
"date": ev_date,
"original_score": round(min(0.95, 0.35 + z * 0.15), 3),
}
ev = {
"name": ev_name,
"start_date": ev_date,
"level": level,
"category": "event_calendar",
"sub_type": sub_type,
"description": (
f"Surprise {direction} {sign}{s_pct:.1f}% vs baseline "
f"(z-score: {z:.1f}σ). "
f"Réel: {rel.get('actual_value', '?')} {rel.get('actual_unit', '')} "
f"/ Prévision: {rel.get('forecast_value', '?')}."
),
"market_impact": "",
"affected_assets": assets,
"impact_score": min(0.95, 0.35 + z * 0.15),
"actual_value": str(rel.get("actual_value", "")),
"expected_value": str(rel.get("forecast_value", "")),
"surprise_pct": float(s_pct),
"source_refs": [source_ref],
"origin": "detector_eco",
}
result = _save_and_evaluate(ev, existing)
if result:
result["source"] = "eco"
created.append(result)
return created
# ── Source 3: MA crossovers (technical) ──────────────────────────────────────
def _check_technical(instruments: List[str] = None, lookback_days: int = 7) -> List[Dict[str, Any]]:
try:
import yfinance as yf
import pandas as pd
except ImportError:
logger.warning("[check_events/technical] yfinance/pandas not available")
return []
if instruments is None:
instruments = WATCH_INSTRUMENTS
existing = _existing_event_keys()
created: List[Dict] = []
cutoff = (datetime.utcnow() - timedelta(days=lookback_days)).strftime("%Y-%m-%d")
for ticker in instruments:
try:
df = yf.download(ticker, period="1y", interval="1d", progress=False, auto_adjust=True)
if df is None or len(df) < 210:
continue
close = df["Close"].squeeze()
df["ma50"] = close.rolling(50).mean()
df["ma100"] = close.rolling(100).mean()
df["ma200"] = close.rolling(200).mean()
recent = df.tail(lookback_days + 2)
for i in range(1, len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
prev = recent.iloc[i - 1]
curr = recent.iloc[i]
def cross(fp, fc, sp, sc):
if any(pd.isna(v) for v in [fp, fc, sp, sc]):
return None
if fp < sp and fc >= sc:
return "golden"
if fp > sp and fc <= sc:
return "death"
return None
pairs = [
("MA50", "MA200", prev["ma50"], curr["ma50"], prev["ma200"], curr["ma200"]),
("MA50", "MA100", prev["ma50"], curr["ma50"], prev["ma100"], curr["ma100"]),
]
for fast_lbl, slow_lbl, fp, fc, sp, sc in pairs:
kind = cross(fp, fc, sp, sc)
if kind is None:
continue
cross_label = "Golden Cross" if kind == "golden" else "Death Cross"
ev_name = f"{ticker} {fast_lbl}/{slow_lbl} {cross_label} ({date_str[:7]})"
if _is_dup(ev_name, existing):
continue
direction = "bullish" if kind == "golden" else "bearish"
level = "medium" if slow_lbl == "MA200" else "short"
source_ref = {
"title": f"Technical signal: {ev_name}",
"source": "yfinance/computed",
"url": f"https://finance.yahoo.com/quote/{ticker}",
"date": date_str,
"original_score": 0.65 if slow_lbl == "MA200" else 0.45,
}
ev = {
"name": ev_name,
"start_date": date_str,
"level": level,
"category": "technical",
"sub_type": f"{fast_lbl}/{slow_lbl} Cross",
"description": (
f"{cross_label} : {fast_lbl} passe "
f"{'au-dessus' if kind == 'golden' else 'en-dessous'} "
f"de la {slow_lbl} sur {ticker}. "
f"Signal {direction} de tendance "
f"{'long terme' if slow_lbl == 'MA200' else 'moyen terme'}."
),
"market_impact": f"Signal {direction} sur {ticker}",
"affected_assets": [ticker],
"impact_score": 0.65 if slow_lbl == "MA200" else 0.45,
"source_refs": [source_ref],
"origin": "detector_technical",
}
result = _save_and_evaluate(ev, existing)
if result:
result["source"] = "technical"
created.append(result)
except Exception as e:
logger.debug(f"[check_events/technical] {ticker} failed: {e}")
return created
# ── Source 4: Institutional reports ──────────────────────────────────────────
def _check_reports(days: int = 7, min_importance: int = 3) -> List[Dict[str, Any]]:
from services.database import get_conn
try:
cutoff = (datetime.utcnow() - timedelta(days=days)).strftime("%Y-%m-%d")
conn = get_conn()
rows = conn.execute(
"""SELECT * FROM institutional_reports
WHERE report_date >= ? AND importance >= ?
ORDER BY importance DESC, report_date DESC
LIMIT 15""",
(cutoff, min_importance),
).fetchall()
conn.close()
reports = [dict(r) for r in rows]
except Exception as e:
logger.warning(f"[check_events/reports] query failed: {e}")
return []
existing = _existing_event_keys()
created: List[Dict] = []
for rpt in reports:
title = rpt.get("title", "")
rpt_type = rpt.get("report_type", "Report")
rpt_date = (rpt.get("report_date") or "")[:10]
if not title or _is_dup(title, existing):
continue
ev_name = title[:60]
summary = rpt.get("ai_summary") or rpt.get("trading_implications", "")
try:
kp = json.loads(rpt.get("key_points_json") or "[]")
if kp:
summary = " ".join(kp[:2]) + " " + summary
except Exception:
pass
assets: List[str] = []
for sig_col, asset_list in [
("signal_energy", ["USO", "XOM"]),
("signal_metals", ["GLD", "SLV"]),
("signal_indices", ["SPY", "QQQ"]),
("signal_forex", ["EURUSD=X", "USDJPY=X"]),
]:
if rpt.get(sig_col, "neutral") not in ("neutral", "", None):
assets.extend(asset_list)
source_ref = {
"title": title,
"source": rpt.get("source", rpt_type),
"url": "",
"date": rpt_date,
"original_score": round(min(0.9, 0.3 + rpt.get("importance", 2) * 0.12), 3),
}
ev = {
"name": ev_name,
"start_date": rpt_date,
"level": "medium" if rpt.get("importance", 2) >= 4 else "short",
"category": "report",
"sub_type": rpt_type.upper(),
"description": summary[:500] or f"Rapport {rpt_type} du {rpt_date}.",
"market_impact": rpt.get("trading_implications", ""),
"affected_assets": list(set(assets)),
"impact_score": min(0.9, 0.3 + rpt.get("importance", 2) * 0.12),
"source_refs": [source_ref],
"origin": "detector_report",
}
result = _save_and_evaluate(ev, existing)
if result:
result["source"] = "reports"
created.append(result)
return created
# ── Main entry point ──────────────────────────────────────────────────────────
def check_new_market_events(
sources: Optional[List[str]] = None,
news_impact_min: float = 0.55,
news_lookback_hours: int = 48,
eco_z_threshold: float = 1.5,
eco_days: int = 7,
technical_lookback_days: int = 7,
report_days: int = 7,
report_min_importance: int = 3,
) -> Dict[str, Any]:
"""
Isolated cycle action — scans all (or selected) sources, creates
market_events with source_refs, and immediately evaluates instrument impacts.
"""
if sources is None:
sources = ["news", "eco", "technical", "reports"]
results: Dict[str, Any] = {
"news": [], "eco": [], "technical": [], "reports": [],
"total_created": 0,
"ran_at": datetime.utcnow().isoformat(),
}
if "news" in sources:
try:
results["news"] = _check_news(min_impact=news_impact_min, lookback_hours=news_lookback_hours)
except Exception as e:
logger.error(f"[check_events] news source error: {e}")
if "eco" in sources:
try:
results["eco"] = _check_eco(z_threshold=eco_z_threshold, days=eco_days)
except Exception as e:
logger.error(f"[check_events] eco source error: {e}")
if "technical" in sources:
try:
results["technical"] = _check_technical(lookback_days=technical_lookback_days)
except Exception as e:
logger.error(f"[check_events] technical source error: {e}")
if "reports" in sources:
try:
results["reports"] = _check_reports(days=report_days, min_importance=report_min_importance)
except Exception as e:
logger.error(f"[check_events] reports source error: {e}")
results["total_created"] = sum(len(results[s]) for s in ["news", "eco", "technical", "reports"])
logger.info(
f"[check_events] Done — {results['total_created']} new events: "
f"news={len(results['news'])} eco={len(results['eco'])} "
f"technical={len(results['technical'])} reports={len(results['reports'])}"
)
return results