Files
OpenFin/backend/services/guidance_sync.py
2026-07-02 18:07:21 +02:00

264 lines
11 KiB
Python
Raw Permalink 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.
"""
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