""" Guidance Sync — crée/met à jour des market_events de type MACRO_RATE_GUIDANCE à partir des events de taux détectés dans ff_calendar. Lifecycle d'un event guidance : 1. Event de taux détecté dans ff_calendar (pas encore d'actual) → signal = 0 si pas de forecast (status quo attendu) → signal = (forecast - previous) * 100 en bps si forecast dispo 2. Guidance market_event créé, end_date = date de la réunion → causal_event_analyses créés pour chaque instrument → visible dans la Frise 3. À chaque sync, le signal est recalculé et mis à jour si nécessaire 4. Quand actual_value apparaît dans ff_calendar → guidance fermé (supprimé ou ignoré) Appelé automatiquement après chaque calendar_sync dans main.py. """ import json import logging import re from datetime import date, datetime from typing import Optional logger = logging.getLogger(__name__) # ── Patterns de noms d'events de taux par devise ────────────────────────────── # Matching insensible à la casse sur event_name de ff_calendar RATE_PATTERNS: dict[str, list[str]] = { "USD": ["fed funds rate", "federal funds rate", "fomc rate decision"], "EUR": ["ecb interest rate", "main refinancing rate", "refinancing rate", "ecb rate decision"], "GBP": ["official bank rate", "boe bank rate", "boe interest rate", "bank rate"], "JPY": ["boj policy rate", "boj rate", "overnight call rate", "monetary policy"], "AUD": ["cash rate", "rba cash rate", "rba rate decision"], "CAD": ["overnight rate", "boc rate", "bank of canada rate"], "NZD": ["official cash rate", "rbnz rate", "rbnz cash rate"], "CHF": ["snb policy rate", "snb rate", "snb interest rate"], } # Instruments affectés par devise (USD rate guidance → EURUSD, TLT, etc.) CURRENCY_INSTRUMENTS: dict[str, list[str]] = { "USD": ["EURUSD", "TLT", "USDJPY", "SP500", "XAUUSD", "GBPUSD", "EEM", "QQQ"], "EUR": ["EURUSD", "GBPUSD", "SP500", "XAUUSD", "TLT"], "GBP": ["GBPUSD", "EURUSD", "SP500"], "JPY": ["USDJPY", "EURUSD", "XAUUSD"], "AUD": ["EURUSD"], "CAD": ["EURUSD"], "NZD": ["EURUSD"], "CHF": ["EURUSD", "XAUUSD"], } GUIDANCE_ORIGIN = "guidance_sync" GUIDANCE_SUBTYPE = "rate_guidance" # sub_type dans market_events # ── Helpers ──────────────────────────────────────────────────────────────────── def _parse_rate(s: Optional[str]) -> Optional[float]: """'4.25%' ou '4.25' → 4.25, None si vide ou non parseable.""" if not s: return None cleaned = re.sub(r"[^\d.\-]", "", str(s).strip()) try: return float(cleaned) except (ValueError, TypeError): return None def _signal_bps(forecast: Optional[str], previous: Optional[str]) -> float: """ Signal en bps = (forecast - previous) * 100. Positif = hike attendu, négatif = cut attendu, 0 = status quo ou pas de forecast. """ f = _parse_rate(forecast) p = _parse_rate(previous) if f is None or p is None: return 0.0 return round((f - p) * 100, 1) def _get_guidance_template_id(conn) -> Optional[int]: row = conn.execute( "SELECT id FROM causal_graph_templates WHERE name = 'Macro Rate Guidance'" ).fetchone() return row["id"] if row else None def _find_existing(conn, currency: str, meeting_date: str) -> Optional[dict]: """Trouve le guidance market_event existant pour cette devise+date de réunion.""" row = conn.execute( """SELECT * FROM market_events WHERE origin=? AND end_date=? AND sub_type=? LIMIT 1""", (GUIDANCE_ORIGIN, meeting_date, GUIDANCE_SUBTYPE + "_" + currency) ).fetchone() return dict(row) if row else None def _upsert_analyses(conn, event_id: int, template_id: int, instruments: list[str], signal: float): """ Crée/recrée les causal_event_analyses pour que l'event apparaisse dans la Frise et dans la theoretical-curve de chaque instrument. Calcule prediction_json via evaluate_graph pour alimenter le factor engine. """ conn.execute( "DELETE FROM causal_event_analyses WHERE market_event_id=?", (event_id,) ) inputs_dict = {"guidance_signal": signal} inputs_json = json.dumps(inputs_dict) # Évalue le graphe pour obtenir les pips prédits par instrument node_values: dict = {} try: from services.causal_graphs import evaluate_graph, get_template tmpl = get_template(conn, template_id) if tmpl: node_values = evaluate_graph(tmpl["graph_json"], inputs_dict) except Exception as e: logger.warning(f"[guidance_sync] evaluate_graph failed: {e}") prediction_json = json.dumps(node_values) if node_values else None for inst in instruments: conn.execute(""" INSERT INTO causal_event_analyses (market_event_id, template_id, instrument, inputs_json, prediction_json, analyzed_at) VALUES (?, ?, ?, ?, ?, datetime('now')) """, (event_id, template_id, inst, inputs_json, prediction_json)) def _build_title(currency: str, signal: float, meeting_date: str) -> str: try: month = datetime.strptime(meeting_date, "%Y-%m-%d").strftime("%b %Y") except ValueError: month = meeting_date if signal == 0.0: action = "Status Quo" elif signal < 0: action = f"Cut {abs(signal):.0f}bps" else: action = f"Hike {signal:.0f}bps" return f"{currency} Rate Guidance: {action} ({month})" def _build_description(currency: str, signal: float, event_name: str) -> str: if signal == 0.0: sentiment = "Pas de changement attendu (status quo). Aucun forecast disponible ou consensus = pas de mouvement." elif signal < 0: sentiment = f"Marché anticipe un cut de {abs(signal):.0f}bps. Pression baissière sur {currency} (carry), haussière sur obligations." else: sentiment = f"Marché anticipe un hike de {signal:.0f}bps. Pression haussière sur {currency} (carry), baissière sur obligations." return f"[{event_name}] {sentiment}" # ── Main ─────────────────────────────────────────────────────────────────────── def sync_guidance() -> dict: """ Scanne ff_calendar pour les events de taux à venir, crée/met à jour les guidance market_events correspondants. Retourne un dict de stats. """ from services.database import get_conn from services.causal_graphs import init_tables conn = get_conn() init_tables(conn) today = date.today().isoformat() template_id = _get_guidance_template_id(conn) if template_id is None: logger.warning("[guidance_sync] Template 'Macro Rate Guidance' introuvable — lance seed_templates d'abord") conn.close() return {"error": "template_not_found"} # Tous les events de taux futurs (sans actual) dans ff_calendar rows = conn.execute(""" SELECT currency, event_date, event_name, forecast_value, previous_value, actual_value FROM ff_calendar WHERE event_date >= ? AND impact = 'high' AND (actual_value IS NULL OR actual_value = '') ORDER BY event_date ASC LIMIT 500 """, (today,)).fetchall() created = updated = closed = skipped = 0 seen: set[tuple] = set() # (currency, meeting_date) for row in rows: currency = (row["currency"] or "").strip() event_name = (row["event_name"] or "").lower() evt_date = row["event_date"] forecast = row["forecast_value"] previous = row["previous_value"] actual = row["actual_value"] # Filtre : correspond à un pattern de taux ? patterns = RATE_PATTERNS.get(currency, []) if not any(p in event_name for p in patterns): skipped += 1 continue # Évite les doublons (même devise × même date) key = (currency, evt_date) if key in seen: continue seen.add(key) # Si actual est là, l'event est résolu → on ne crée pas de guidance if actual and str(actual).strip(): # Supprimer l'ancien guidance si présent (l'event est passé) existing = _find_existing(conn, currency, evt_date) if existing: conn.execute("DELETE FROM market_events WHERE id=?", (existing["id"],)) conn.execute( "DELETE FROM causal_event_analyses WHERE market_event_id=?", (existing["id"],) ) closed += 1 continue signal = _signal_bps(forecast, previous) instruments = CURRENCY_INSTRUMENTS.get(currency, ["EURUSD"]) title = _build_title(currency, signal, evt_date) description = _build_description(currency, signal, row["event_name"] or "") impact = round(min(1.0, 0.6 + abs(signal) / 300), 2) affected = json.dumps(instruments) sub_type = GUIDANCE_SUBTYPE + "_" + currency existing = _find_existing(conn, currency, evt_date) if existing: ev_id = existing["id"] # Met à jour titre / description / signal si changé conn.execute(""" UPDATE market_events SET name=?, description=?, expected_value=?, impact_score=?, affected_assets=? WHERE id=? """, (title, description, str(signal), impact, affected, ev_id)) _upsert_analyses(conn, ev_id, template_id, instruments, signal) updated += 1 else: cur = conn.execute(""" INSERT INTO market_events (name, start_date, end_date, level, category, sub_type, description, affected_assets, impact_score, origin, expected_value, unit, source_refs) VALUES (?, ?, ?, 'high', 'central_bank', ?, ?, ?, ?, ?, ?, 'bps', '[]') """, ( title, today, evt_date, sub_type, description, affected, impact, GUIDANCE_ORIGIN, str(signal), )) ev_id = cur.lastrowid _upsert_analyses(conn, ev_id, template_id, instruments, signal) created += 1 conn.commit() conn.close() stats = {"created": created, "updated": updated, "closed": closed, "skipped": skipped} logger.info(f"[guidance_sync] {stats}") return stats