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