Files
OpenFin/backend/services/market_event_detector.py
OpenSquared a12d7a1ef3 feat: 6 AI desks complets — fundamental, report, sentiment (options triggers)
- Fundamental Desk: filtre corporate news (layoffs, M&A, earnings, credit)
  avec prompt dédié + dedup sémantique
- Report Desk: wiring desk config (days, min_importance, system_prompt)
- Sentiment Desk: 5 signaux VIX/SKEW pour options lab
  (vix_level thresholds, vix_spike %, vix_term_structure, vvix_extreme, skew_extreme)
  Chaque event contient options_note actionnable (vente puts, straddles, calendar spreads)
- check_new_market_events() couvre les 6 sources, charge desk configs dynamiquement
- Signal catalog: 12 signaux (7 technical + 5 sentiment), filtrés par desk_type dans UI
- AIDesks.tsx: FundamentalConfig + ReportConfig + SignalToggle filtré par desk type
- 3 nouveaux desks seedés dans init_db (idempotent)

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-26 09:21:29 +02:00

1163 lines
46 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: configurable signal catalog driven by Technical Desk
- reports : institutional reports (COT, EIA) with high importance
Desk configs are loaded from ai_desks table at runtime.
"""
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."""
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
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}
# ── Semantic deduplication ────────────────────────────────────────────────────
def _semantic_dedup(
title: str,
source: str,
date_str: str,
summary: str,
category: str,
client: Any,
dedup_lookback_days: int = 2,
system_prompt_hint: str = "",
) -> bool:
"""
Ask the AI whether this news already exists in recent market_events.
Returns True if it's a duplicate (should be skipped).
"""
from services.database import get_market_events_near_date
dedup_categories = ["geopolitical", "fundamental", "report"]
if category and category not in dedup_categories:
dedup_categories.append(category)
recent = get_market_events_near_date(date_str, days=dedup_lookback_days, categories=dedup_categories)
if not recent:
return False
recent_block = "\n".join(
f" [{r['start_date']}] {r['name']}{(r.get('description') or '')[:80]}"
for r in recent[:15]
)
hint = f"\nNote du desk: {system_prompt_hint[:200]}" if system_prompt_hint else ""
prompt = f"""Tu es un éditeur de base de données d'événements marchés.{hint}
NOUVELLE NEWS À VÉRIFIER:
- Titre: {title}
- Source: {source}
- Date: {date_str}
- Résumé: {summary[:300]}
ÉVÉNEMENTS EXISTANTS (±{dedup_lookback_days} jours):
{recent_block}
Cette news représente-t-elle le même fait qu'un événement déjà enregistré ?
Réponds JSON: {{"is_duplicate": true/false, "reason": "courte phrase"}}"""
try:
resp = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
response_format={"type": "json_object"},
temperature=0.0,
max_tokens=100,
)
parsed = json.loads(resp.choices[0].message.content)
is_dup = bool(parsed.get("is_duplicate", False))
if is_dup:
logger.debug(f"[dedup] Skipping duplicate: '{title[:40]}'{parsed.get('reason','')}")
return is_dup
except Exception as e:
logger.debug(f"[dedup] AI check failed for '{title[:40]}': {e}")
return False
# ── Source 1: Geopolitical / macro news ──────────────────────────────────────
def _check_news(desk_cfg: Dict[str, Any]) -> List[Dict[str, Any]]:
from services.data_fetcher import fetch_geo_news
min_impact = float(desk_cfg.get("min_impact", 0.55))
lookback_hours = int(desk_cfg.get("lookback_hours", 48))
max_evaluate = int(desk_cfg.get("max_evaluate", 15))
dedup_enabled = bool(desk_cfg.get("dedup_enabled", True))
dedup_days = int(desk_cfg.get("dedup_lookback_days", 2))
system_prompt = desk_cfg.get("_system_prompt", "")
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_evaluate]
if not candidates:
return []
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 []
existing = _existing_event_keys()
created: List[Dict] = []
for n in candidates:
title = n.get("title", "")
if not title or _is_dup(title, existing):
continue
pub_date = _parse_date(n.get("date", ""))
news_summary = str(n.get("summary", ""))[:400]
source = n.get("source", "")
# Semantic dedup before expensive classification call
if dedup_enabled:
if _semantic_dedup(
title, source, pub_date, news_summary,
category="geopolitical",
client=client,
dedup_lookback_days=dedup_days,
system_prompt_hint=system_prompt,
):
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: {source}
DATE: {n.get('date', '')}
RÉSUMÉ: {news_summary}
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": source,
"url": n.get("url") or n.get("link", ""),
"date": pub_date,
"original_score": round(float(n.get("impact_score", 0)), 3),
}
ev = {
"name": ev_name,
"start_date": pub_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(desk_cfg: Dict[str, Any]) -> List[Dict[str, Any]]:
from services.database import get_recent_economic_surprises
z_threshold = float(desk_cfg.get("z_threshold", 1.5))
days = int(desk_cfg.get("days", 7))
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
# ── Technical signal detectors ────────────────────────────────────────────────
def _detect_ma_cross(ticker: str, df: Any, params: Dict, cutoff: str) -> List[Dict]:
"""Golden/Death cross detector for configured MA pairs."""
import pandas as pd
events = []
pairs_cfg = params.get("pairs", [["MA50", "MA200"], ["MA50", "MA100"]])
ma_map = {"MA20": 20, "MA50": 50, "MA100": 100, "MA200": 200}
close = df["Close"].squeeze()
# Pre-compute all required MAs
needed: set = set()
for pair in pairs_cfg:
needed.update(pair)
ma_series: Dict[str, Any] = {}
for lbl in needed:
period = ma_map.get(lbl)
if period and len(df) >= period:
ma_series[lbl] = close.rolling(period).mean()
recent = df.tail(4)
for i in range(1, len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
for fast_lbl, slow_lbl in pairs_cfg:
if fast_lbl not in ma_series or slow_lbl not in ma_series:
continue
fp = ma_series[fast_lbl].iloc[-(len(recent) - i + 1)]
fc = ma_series[fast_lbl].iloc[-(len(recent) - i)]
sp = ma_series[slow_lbl].iloc[-(len(recent) - i + 1)]
sc = ma_series[slow_lbl].iloc[-(len(recent) - i)]
if any(pd.isna(v) for v in [fp, fc, sp, sc]):
continue
if fp < sp and fc >= sc:
kind = "golden"
elif fp > sp and fc <= sc:
kind = "death"
else:
continue
cross_label = "Golden Cross" if kind == "golden" else "Death Cross"
events.append({
"name": f"{ticker} {fast_lbl}/{slow_lbl} {cross_label} ({date_str[:7]})",
"date": date_str,
"direction": "bullish" if kind == "golden" else "bearish",
"sub_type": f"{fast_lbl}/{slow_lbl} Cross",
"score": 0.65 if "MA200" in (fast_lbl, slow_lbl) else 0.45,
"level": "medium" if "MA200" in (fast_lbl, slow_lbl) else "short",
"desc": f"{cross_label}: {fast_lbl} {'au-dessus' if kind=='golden' else 'en-dessous'} de {slow_lbl} sur {ticker}.",
})
return events
def _detect_rsi_extreme(ticker: str, df: Any, params: Dict, cutoff: str) -> List[Dict]:
"""RSI oversold/overbought signal."""
import pandas as pd
period = int(params.get("period", 14))
oversold = float(params.get("oversold", 30))
overbought = float(params.get("overbought", 70))
close = df["Close"].squeeze()
if len(close) < period + 2:
return []
delta = close.diff()
gain = delta.clip(lower=0).rolling(period).mean()
loss = (-delta.clip(upper=0)).rolling(period).mean()
rs = gain / loss.replace(0, float("nan"))
rsi = 100 - (100 / (1 + rs))
events = []
recent = rsi.tail(3)
for i in range(len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
val = recent.iloc[i]
if pd.isna(val):
continue
if val <= oversold:
direction, label = "bullish", "Oversold"
elif val >= overbought:
direction, label = "bearish", "Overbought"
else:
continue
events.append({
"name": f"{ticker} RSI {label} ({date_str[:7]})",
"date": date_str,
"direction": direction,
"sub_type": f"RSI {label}",
"score": 0.50 if abs(val - 50) > 30 else 0.40,
"level": "short",
"desc": f"RSI({period}) à {val:.1f} sur {ticker} — signal {label.lower()} ({direction}).",
})
return events
def _detect_bb_squeeze(ticker: str, df: Any, params: Dict, cutoff: str) -> List[Dict]:
"""Bollinger Band squeeze detector."""
import pandas as pd
period = int(params.get("period", 20))
std_mult = float(params.get("std", 2.0))
width_threshold = float(params.get("width_threshold", 0.05))
close = df["Close"].squeeze()
if len(close) < period + 2:
return []
mid = close.rolling(period).mean()
std = close.rolling(period).std()
upper = mid + std_mult * std
lower = mid - std_mult * std
width = (upper - lower) / mid
events = []
recent = width.tail(3)
for i in range(len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
w = recent.iloc[i]
if pd.isna(w):
continue
if w <= width_threshold:
events.append({
"name": f"{ticker} BB Squeeze ({date_str[:7]})",
"date": date_str,
"direction": "neutral",
"sub_type": "BB Squeeze",
"score": 0.45,
"level": "short",
"desc": f"Bandes de Bollinger({period},{std_mult}) très resserrées sur {ticker} — width={w:.3f}. Explosion de volatilité imminente.",
})
return events
def _detect_52w_extreme(ticker: str, df: Any, params: Dict, cutoff: str) -> List[Dict]:
"""New 52-week high/low detector."""
import pandas as pd
buffer_pct = float(params.get("buffer_pct", 0.5)) / 100
close = df["Close"].squeeze()
if len(close) < 252:
return []
high_52 = close.rolling(252).max()
low_52 = close.rolling(252).min()
events = []
recent_close = close.tail(3)
for i in range(len(recent_close)):
date_str = str(recent_close.index[i])[:10]
if date_str < cutoff:
continue
c = recent_close.iloc[i]
h52 = high_52.iloc[-(3 - i)]
l52 = low_52.iloc[-(3 - i)]
if pd.isna(c) or pd.isna(h52) or pd.isna(l52):
continue
if c >= h52 * (1 - buffer_pct):
events.append({
"name": f"{ticker} Nouveau 52W High ({date_str[:7]})",
"date": date_str,
"direction": "bullish",
"sub_type": "52W High",
"score": 0.60,
"level": "medium",
"desc": f"{ticker} atteint un nouveau plus haut 52 semaines à {c:.2f} (précédent: {h52:.2f}).",
})
elif c <= l52 * (1 + buffer_pct):
events.append({
"name": f"{ticker} Nouveau 52W Low ({date_str[:7]})",
"date": date_str,
"direction": "bearish",
"sub_type": "52W Low",
"score": 0.60,
"level": "medium",
"desc": f"{ticker} atteint un nouveau plus bas 52 semaines à {c:.2f} (précédent: {l52:.2f}).",
})
return events
# ── Source 3: Technical signals ───────────────────────────────────────────────
def _check_technical(desk_cfg: Dict[str, Any]) -> 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 []
instruments = desk_cfg.get("_instruments") or WATCH_INSTRUMENTS
lookback_days = int(desk_cfg.get("lookback_days", 7))
signals_config = desk_cfg.get("signals", {})
# Determine which signals are active
def sig_cfg(sig_id: str) -> Optional[Dict]:
c = signals_config.get(sig_id, {})
return c if c.get("enabled", False) else None
ma_cross_cfg = sig_cfg("ma_cross")
rsi_cfg = sig_cfg("rsi_extreme")
bb_cfg = sig_cfg("bb_squeeze")
extreme_52w = sig_cfg("new_52w_extreme")
if not any([ma_cross_cfg, rsi_cfg, bb_cfg, extreme_52w]):
logger.info("[check_events/technical] no active signals in desk config")
return []
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="2y", interval="1d", progress=False, auto_adjust=True)
if df is None or len(df) < 20:
continue
# Flatten MultiIndex if needed (yfinance ≥ 0.2 returns MultiIndex columns)
if hasattr(df.columns, "levels"):
df.columns = df.columns.get_level_values(0)
detected: List[Dict] = []
if ma_cross_cfg and len(df) >= 210:
detected += _detect_ma_cross(ticker, df, ma_cross_cfg, cutoff)
if rsi_cfg:
detected += _detect_rsi_extreme(ticker, df, rsi_cfg, cutoff)
if bb_cfg:
detected += _detect_bb_squeeze(ticker, df, bb_cfg, cutoff)
if extreme_52w and len(df) >= 252:
detected += _detect_52w_extreme(ticker, df, extreme_52w, cutoff)
for sig in detected:
ev_name = sig["name"]
if _is_dup(ev_name, existing):
continue
source_ref = {
"title": f"Technical signal: {ev_name}",
"source": "yfinance/computed",
"url": f"https://finance.yahoo.com/quote/{ticker}",
"date": sig["date"],
"original_score": sig["score"],
}
ev = {
"name": ev_name,
"start_date": sig["date"],
"level": sig["level"],
"category": "technical",
"sub_type": sig["sub_type"],
"description": sig["desc"],
"market_impact": f"Signal {sig['direction']} sur {ticker}",
"affected_assets": [ticker],
"impact_score": sig["score"],
"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: Fundamental news ───────────────────────────────────────────────
def _check_fundamental(desk_cfg: Dict[str, Any]) -> List[Dict[str, Any]]:
"""Corporate fundamental events: layoffs, M&A, earnings, credit, regulatory."""
from services.data_fetcher import fetch_geo_news
min_impact = float(desk_cfg.get("min_impact", 0.45))
lookback_hours = int(desk_cfg.get("lookback_hours", 72))
max_evaluate = int(desk_cfg.get("max_evaluate", 20))
dedup_enabled = bool(desk_cfg.get("dedup_enabled", True))
dedup_days = int(desk_cfg.get("dedup_lookback_days", 3))
system_prompt = desk_cfg.get("_system_prompt", "")
focus_types = desk_cfg.get("focus_types", ["layoffs","earnings","ma","credit","regulatory","guidance"])
api_key = _get_api_key()
if not api_key:
return []
try:
all_news = fetch_geo_news()
except Exception as e:
logger.warning(f"[check_events/fundamental] fetch failed: {e}")
return []
cutoff_dt = datetime.utcnow() - timedelta(hours=lookback_hours)
candidates = [
n for n in all_news
if (n.get("impact_score") or 0) >= min_impact
and _parse_date(n.get("date", "")) >= cutoff_dt.strftime("%Y-%m-%d")
][:max_evaluate]
if not candidates:
return []
try:
from openai import OpenAI
client = OpenAI(api_key=api_key)
except Exception as e:
logger.warning(f"[check_events/fundamental] OpenAI init failed: {e}")
return []
existing = _existing_event_keys()
created: List[Dict] = []
focus_str = ", ".join(focus_types)
for n in candidates:
title = n.get("title", "")
if not title or _is_dup(title, existing):
continue
pub_date = _parse_date(n.get("date", ""))
news_summary = str(n.get("summary", ""))[:400]
source = n.get("source", "")
if dedup_enabled:
if _semantic_dedup(title, source, pub_date, news_summary,
category="fundamental", client=client,
dedup_lookback_days=dedup_days,
system_prompt_hint=system_prompt):
continue
prompt = f"""Tu es un analyste fondamental corporate. Cette news représente-t-elle un événement
fondamental CORPORATE structurant (types attendus: {focus_str}) ?
TITRE: {title}
SOURCE: {source}
DATE: {n.get('date', '')}
RÉSUMÉ: {news_summary}
Réponds OUI uniquement si c'est un fait avéré avec impact mesurable sur un secteur ou sur les indices.
Ignore les géopolitiques purs (guerres, sanctions) — ceux-là sont traités par le News Desk.
FORMAT JSON STRICT:
{{
"qualifies": true/false,
"fundamental_type": "layoffs|earnings|ma|credit|regulatory|guidance|other",
"reason": "une phrase",
"name": "Nom court (≤ 60 chars)",
"company_sector": "ex: Tech, Energy, Financials, ou ticker si connu",
"description": "1-2 phrases analytiques",
"affected_assets": ["QQQ","HYG",...],
"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/fundamental] AI 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": source,
"url": n.get("url") or n.get("link", ""),
"date": pub_date,
"original_score": round(float(n.get("impact_score", 0)), 3),
}
ev = {
"name": ev_name,
"start_date": pub_date,
"level": parsed.get("level", "short"),
"category": "fundamental",
"sub_type": parsed.get("fundamental_type", "other"),
"description": parsed.get("description", title),
"market_impact": f"Secteur: {parsed.get('company_sector','')}",
"affected_assets": parsed.get("affected_assets", []),
"impact_score": float(parsed.get("impact_score", 0.5)),
"source_refs": [source_ref],
"origin": "detector_fundamental",
}
result = _save_and_evaluate(ev, existing)
if result:
result["source"] = "fundamental"
created.append(result)
return created
# ── Source 5: Sentiment signals (Options Lab triggers) ────────────────────────
def _check_sentiment(desk_cfg: Dict[str, Any]) -> List[Dict[str, Any]]:
"""
VIX family + SKEW signals — generates sentiment market_events that feed
directly into the options lab as volatility regime triggers.
"""
try:
import yfinance as yf
import pandas as pd
except ImportError:
logger.warning("[check_events/sentiment] yfinance not available")
return []
lookback_days = int(desk_cfg.get("lookback_days", 5))
signals_config = desk_cfg.get("signals", {})
def sig_on(sig_id: str) -> Optional[Dict]:
c = signals_config.get(sig_id, {})
return c if c.get("enabled", True) else None
vix_level_cfg = sig_on("vix_level")
vix_spike_cfg = sig_on("vix_spike")
vix_ts_cfg = sig_on("vix_term_structure")
vvix_cfg = sig_on("vvix_extreme")
skew_cfg = sig_on("skew_extreme")
if not any([vix_level_cfg, vix_spike_cfg, vix_ts_cfg, vvix_cfg, skew_cfg]):
return []
existing = _existing_event_keys()
created: List[Dict] = []
cutoff = (datetime.utcnow() - timedelta(days=lookback_days)).strftime("%Y-%m-%d")
# ── Fetch VIX family ──────────────────────────────────────────────────────
tickers = {"VIX": "^VIX", "VIX9D": "^VIX9D", "VIX3M": "^VIX3M", "VVIX": "^VVIX", "SKEW": "^SKEW"}
series: Dict[str, Any] = {}
for lbl, sym in tickers.items():
try:
df = yf.download(sym, period="30d", interval="1d", progress=False, auto_adjust=True)
if df is not None and len(df) > 0:
if hasattr(df.columns, "levels"):
df.columns = df.columns.get_level_values(0)
series[lbl] = df["Close"].squeeze().dropna()
except Exception as e:
logger.debug(f"[sentiment] {sym} download failed: {e}")
vix = series.get("VIX")
if vix is None or len(vix) < 2:
logger.warning("[check_events/sentiment] VIX data unavailable")
return []
def _emit(name: str, date_str: str, direction: str, sub_type: str,
score: float, desc: str, assets: List[str], options_note: str = ""):
if _is_dup(name, existing):
return
ev = {
"name": name,
"start_date": date_str,
"level": "short",
"category": "sentiment",
"sub_type": sub_type,
"description": desc + (f" {options_note}" if options_note else ""),
"market_impact": options_note,
"affected_assets": assets,
"impact_score": score,
"source_refs": [{
"title": f"Sentiment signal: {name}",
"source": "CBOE/yfinance",
"url": "https://www.cboe.com/tradable_products/vix/",
"date": date_str,
"original_score": score,
}],
"origin": "detector_sentiment",
}
result = _save_and_evaluate(ev, existing)
if result:
result["source"] = "sentiment"
created.append(result)
# ── VIX level threshold crossings ─────────────────────────────────────────
if vix_level_cfg:
thresholds = vix_level_cfg.get("thresholds", [20, 25, 30, 35, 45])
recent = vix.tail(lookback_days + 2)
for i in range(1, len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
prev_v, curr_v = float(recent.iloc[i-1]), float(recent.iloc[i])
for lvl in thresholds:
name = None
if prev_v < lvl <= curr_v:
name = f"VIX franchit {lvl} à la hausse ({date_str[:7]})"
options_note = (
"Opportunité: vente de puts cash-secured sur SPY." if lvl < 25
else "Régime de peur — évaluer straddles ou risk reversals."
)
_emit(name, date_str, "bearish", f"VIX >{lvl}",
min(0.9, 0.4 + lvl * 0.01),
f"VIX dépasse {lvl} (précédent: {prev_v:.1f}{curr_v:.1f}). Entrée en régime de volatilité élevée.",
["VXX","SPY","QQQ","TLT"], options_note)
elif prev_v >= lvl > curr_v:
name = f"VIX repasse sous {lvl} ({date_str[:7]})"
_emit(name, date_str, "bullish", f"VIX <{lvl}",
0.45,
f"VIX revient sous {lvl} ({prev_v:.1f}{curr_v:.1f}). Détente de la volatilité.",
["SPY","QQQ","VXX"],
"Opportunité: rachat de protection ou fermeture de couvertures.")
# ── VIX spike intraday / daily ────────────────────────────────────────────
if vix_spike_cfg:
min_pct = float(vix_spike_cfg.get("min_pct_change", 15.0))
recent = vix.tail(lookback_days + 1)
for i in range(1, len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
prev_v, curr_v = float(recent.iloc[i-1]), float(recent.iloc[i])
if prev_v <= 0:
continue
pct_chg = (curr_v - prev_v) / prev_v * 100
if abs(pct_chg) >= min_pct:
direction = "bearish" if pct_chg > 0 else "bullish"
sign = "+" if pct_chg > 0 else ""
name = f"VIX spike {sign}{pct_chg:.0f}% ({date_str[:7]})"
_emit(name, date_str, direction, "VIX Spike",
min(0.85, 0.4 + abs(pct_chg) * 0.01),
f"VIX variation journalière de {sign}{pct_chg:.1f}% ({prev_v:.1f}{curr_v:.1f}). Choc de volatilité {'haussier' if pct_chg>0 else 'baissier'}.",
["VXX","SPY","QQQ","TLT","GLD"],
"Signal pour stratégies de vol à court terme.")
# ── VIX term structure inversion (VIX9D > VIX) ────────────────────────────
if vix_ts_cfg and "VIX9D" in series:
threshold = float(vix_ts_cfg.get("inversion_threshold", 1.05))
vix9d = series["VIX9D"]
common_idx = vix.index.intersection(vix9d.index)
if len(common_idx) >= 2:
for dt_idx in common_idx[-lookback_days:]:
date_str = str(dt_idx)[:10]
if date_str < cutoff:
continue
ratio = float(vix9d[dt_idx]) / float(vix[dt_idx]) if float(vix[dt_idx]) > 0 else 0
if ratio >= threshold:
name = f"VIX Term Structure Inversée — Peur court terme ({date_str[:7]})"
if not _is_dup(name, existing):
_emit(name, date_str, "bearish", "VIX Inversion",
0.70,
f"VIX9D ({float(vix9d[dt_idx]):.1f}) > VIX ({float(vix[dt_idx]):.1f}) — ratio {ratio:.2f}. La peur est concentrée sur le très court terme.",
["VXX","SPY","TLT"],
"Stratégie: calendar spread bear — acheter protection courte vs vendre moyenne échéance.")
# ── VVIX extreme ──────────────────────────────────────────────────────────
if vvix_cfg and "VVIX" in series:
threshold = float(vvix_cfg.get("threshold", 100.0))
vvix = series["VVIX"]
recent = vvix.tail(lookback_days)
for i in range(len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
val = float(recent.iloc[i])
if val >= threshold:
name = f"VVIX extrême {val:.0f} ({date_str[:7]})"
_emit(name, date_str, "bearish", "VVIX Extreme",
min(0.80, 0.45 + (val - threshold) * 0.005),
f"VVIX à {val:.1f} (seuil: {threshold}) — volatilité de la volatilité extrême. Marché très incertain sur la direction du VIX.",
["VXX","SPY","QQQ"],
"Éviter les positions directionnelles sur vol. Stratégies non-directionnelles.")
# ── SKEW extreme ──────────────────────────────────────────────────────────
if skew_cfg and "SKEW" in series:
low_thr = float(skew_cfg.get("low_threshold", 120.0))
high_thr = float(skew_cfg.get("high_threshold", 145.0))
skew = series["SKEW"]
recent = skew.tail(lookback_days)
for i in range(len(recent)):
date_str = str(recent.index[i])[:10]
if date_str < cutoff:
continue
val = float(recent.iloc[i])
if val >= high_thr:
name = f"SKEW extrême haussier {val:.0f} ({date_str[:7]})"
_emit(name, date_str, "bearish", "SKEW Extreme",
0.65,
f"CBOE SKEW à {val:.1f} — marché paye très cher pour les puts out-of-the-money. Couverture tail-risk forte.",
["SPY","QQQ","TLT"],
"Skew élevé → vente de put spreads attractive (prime élevée sur strikes bas).")
elif val <= low_thr:
name = f"SKEW très bas {val:.0f} — complaisance ({date_str[:7]})"
_emit(name, date_str, "bullish", "SKEW Low",
0.55,
f"CBOE SKEW à {val:.1f} — marché peu préoccupé par les risques tail. Signal de complaisance.",
["VXX","SPY"],
"Skew bas → acheter protection bon marché (puts OTM relativement peu chers).")
return created
# ── Source 6: Institutional reports ──────────────────────────────────────────
def _check_reports(desk_cfg: Dict[str, Any]) -> List[Dict[str, Any]]:
from services.database import get_conn
days = int(desk_cfg.get("days", 7))
min_importance = int(desk_cfg.get("min_importance", 3))
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,
# Legacy overrides — used when called without an active desk
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]:
"""
Scans all (or selected) sources for all 6 market_event categories.
Desk configs from ai_desks table override legacy params when available.
Sources: news, fundamental, eco, technical, reports, sentiment
"""
if sources is None:
sources = ["news", "fundamental", "eco", "technical", "reports", "sentiment"]
# Load desk configs from DB
try:
from services.database import get_ai_desk_by_type
news_desk = get_ai_desk_by_type("news")
fundamental_desk = get_ai_desk_by_type("fundamental")
tech_desk = get_ai_desk_by_type("technical")
eco_desk = get_ai_desk_by_type("eco")
report_desk = get_ai_desk_by_type("report")
sentiment_desk = get_ai_desk_by_type("sentiment")
except Exception as e:
logger.warning(f"[check_events] Could not load desk configs: {e}")
news_desk = fundamental_desk = tech_desk = eco_desk = report_desk = sentiment_desk = None
def _desk_cfg(desk: Optional[Dict], fallback: Dict) -> Dict:
if not desk:
return fallback
cfg = dict(desk.get("config") or {})
cfg["_instruments"] = desk.get("instruments") or None
cfg["_system_prompt"] = desk.get("system_prompt") or ""
return cfg
news_cfg = _desk_cfg(news_desk, {
"min_impact": news_impact_min, "lookback_hours": news_lookback_hours,
"max_evaluate": 15, "dedup_enabled": False, "dedup_lookback_days": 2,
})
fundamental_cfg = _desk_cfg(fundamental_desk, {
"min_impact": 0.45, "lookback_hours": 72, "max_evaluate": 20,
"dedup_enabled": True, "dedup_lookback_days": 3,
})
eco_cfg = _desk_cfg(eco_desk, {
"z_threshold": eco_z_threshold, "days": eco_days,
})
tech_cfg = _desk_cfg(tech_desk, {
"lookback_days": technical_lookback_days,
"signals": {
"ma_cross": {"enabled": True, "pairs": [["MA50","MA200"],["MA50","MA100"]]},
"rsi_extreme": {"enabled": True, "period": 14, "oversold": 30, "overbought": 70},
"bb_squeeze": {"enabled": True, "period": 20, "std": 2.0, "width_threshold": 0.05},
"new_52w_extreme": {"enabled": True, "buffer_pct": 0.5},
},
})
report_cfg = _desk_cfg(report_desk, {
"days": report_days, "min_importance": report_min_importance,
})
sentiment_cfg = _desk_cfg(sentiment_desk, {
"lookback_days": 5,
"signals": {
"vix_level": {"enabled": True, "thresholds": [20, 25, 30, 35, 45]},
"vix_spike": {"enabled": True, "min_pct_change": 15.0},
"vix_term_structure": {"enabled": True, "inversion_threshold": 1.05},
"vvix_extreme": {"enabled": True, "threshold": 100.0},
"skew_extreme": {"enabled": True, "low_threshold": 120.0, "high_threshold": 145.0},
},
})
ALL_SOURCES = ["news", "fundamental", "eco", "technical", "reports", "sentiment"]
results: Dict[str, Any] = {s: [] for s in ALL_SOURCES}
results["total_created"] = 0
results["ran_at"] = datetime.utcnow().isoformat()
runners = [
("news", lambda: _check_news(news_cfg)),
("fundamental", lambda: _check_fundamental(fundamental_cfg)),
("eco", lambda: _check_eco(eco_cfg)),
("technical", lambda: _check_technical(tech_cfg)),
("reports", lambda: _check_reports(report_cfg)),
("sentiment", lambda: _check_sentiment(sentiment_cfg)),
]
for src, fn in runners:
if src not in sources:
continue
try:
results[src] = fn()
except Exception as e:
logger.error(f"[check_events] {src} source error: {e}")
results["total_created"] = sum(len(results[s]) for s in ALL_SOURCES)
counts = " ".join(f"{s}={len(results[s])}" for s in ALL_SOURCES)
logger.info(f"[check_events] Done — {results['total_created']} new events: {counts}")
return results