feat: instrument model

This commit is contained in:
OpenSquared
2026-07-02 22:46:55 +02:00
parent 7f7c343c9e
commit e2144d93e9
3 changed files with 411 additions and 66 deletions

View File

@@ -44,6 +44,27 @@ def list_instrument_models():
conn.close()
@router.get("/{instrument}/regime")
def get_instrument_regime(
instrument: str,
at_date: Optional[str] = Query(None),
) -> Dict[str, Any]:
"""Régime de marché courant pour cet instrument (détecté depuis events actifs)."""
from services.database import get_conn
from services.instrument_models import _compute_event_by_category, detect_regime
from datetime import datetime, date as date_type
conn = get_conn()
try:
try:
ref_date = date_type.fromisoformat(at_date) if at_date else datetime.utcnow().date()
except ValueError:
ref_date = datetime.utcnow().date()
ev_by_cat = _compute_event_by_category(conn, instrument.upper(), ref_date)
return detect_regime(ev_by_cat)
finally:
conn.close()
@router.get("/{instrument}/timeline")
def get_instrument_timeline(
instrument: str,

View File

@@ -1,22 +1,139 @@
"""
Instrument Models — Phase 1 : propagation en chaîne réelle.
Instrument Models — Phase 2 : saturation non-linéaire + régimes adaptatifs.
Architecture DAG (3 couches) :
Layer 0 — Inputs :
- input_event : valeur auto depuis causal_event_analyses (par catégorie template)
- input_manual : valeur utilisateur (unité native → pips via coefficient_to_pips)
- input_manual : valeur utilisateur → tanh(x/scale)*coeff → pips (saturation)
Layer 1 — Intermediate : formules agrégeant les inputs par domaine
Layer 2 — Output : formule sommant les nœuds intermédiaires
Layer 2 — Output : formule avec poids de régime (ex: 1.4*layer_monetary en MONETARY_DOMINANCE)
Évaluation : evaluate_graph() de causal_graphs.py (formula-based DAG).
Lifecycle : montée (rise_days) → plateau → décroissance (absorption_days).
Timeline : simulation jour par jour via simulate_timeline().
Nouveautés Phase 2 :
- _saturate_pips() : tanh(native/scale)*coeff, slope = coeff à l'origine, sature aux extrêmes
- detect_regime() : identifie le régime dominant depuis les catégories d'events actifs
- REGIME_WEIGHTS : multiplicateurs par couche selon 6 régimes de marché
- _apply_regime_weights() : réécrit la formule output avec les poids du régime courant
- simulate_timeline() inclut le régime du jour
"""
import json
import math
from datetime import datetime, timedelta, date as date_type
from typing import Optional
# ── Saturation scales (tanh) par unité native ─────────────────────────────────
# tanh(x/scale) : slope=1 à l'origine, sature asymptotiquement à ±1
# pips = coefficient_to_pips * scale * tanh(x / scale)
# → slope à x=0 : coefficient_to_pips (identique au linéaire)
# → max pips : coefficient_to_pips * scale (jamais dépassé)
_SATURATION_SCALES: dict[str, float] = {
"bps": 200.0, # différentiels de taux : sature autour ±300bps
"%": 3.0, # CPI / PIB différentiels : sature autour ±5%
"pts%": 3.0,
"pts": 15.0, # PMI écart depuis 50 : sature autour ±20pts
"score": 3.0, # scores subjectifs -5 à +5
"tonnes": 80.0, # tonnes or / CB buying
"Mds$": 40.0,
"Mds$/sem": 15.0,
"k lots": 80.0, # positions CFTC
"$/bbl": 25.0,
"x": 5.0, # multiples PE
"ratio": 1.5,
"t": 80.0,
}
def _saturate_pips(native: float, coeff: float, unit: str, scale_override: Optional[float] = None) -> float:
"""Convertit une valeur native en pips avec saturation tanh.
Comportement identique au linéaire pour de petites valeurs,
sature progressivement pour les valeurs extrêmes.
"""
scale = scale_override or _SATURATION_SCALES.get(unit)
if scale and scale > 0:
return coeff * scale * math.tanh(native / scale)
return coeff * native
# ── Régimes de marché et poids adaptatifs ─────────────────────────────────────
# Chaque régime amplifie certaines couches (>1) et atténue les autres (<1)
# Les clés couvrent tous les noms de couches intermédiaires possibles.
REGIME_WEIGHTS: dict[str, dict[str, float]] = {
"MONETARY_DOMINANCE": {
# Fed/BCE dominent tout — taux, OIS, anticipations
"layer_monetary": 1.40, "layer_rates": 1.40,
"layer_growth": 0.80, "layer_uk_macro": 0.80, "layer_macro": 0.80,
"layer_risk": 0.70, "layer_risk_credit": 0.70, "layer_refuge": 0.70,
"layer_positioning": 1.10, "layer_policy": 1.10,
"layer_flows": 1.00, "layer_demand": 1.00,
"layer_tech_fundamental": 0.85, "layer_supply": 1.00,
"layer_china": 0.90, "layer_dollar": 1.20,
"layer_global": 0.80,
},
"GEOPOLITICAL_RISK": {
# Tensions géo → flight to safety, or, JPY, USD
"layer_monetary": 0.80, "layer_rates": 0.80,
"layer_growth": 0.90, "layer_uk_macro": 0.85, "layer_macro": 0.85,
"layer_risk": 1.50, "layer_risk_credit": 1.50, "layer_refuge": 1.60,
"layer_positioning": 1.10, "layer_policy": 1.30,
"layer_flows": 1.00, "layer_demand": 1.20,
"layer_tech_fundamental": 0.80, "layer_supply": 0.90,
"layer_china": 0.80, "layer_dollar": 0.90,
"layer_global": 1.30,
},
"CREDIT_STRESS": {
# Banking stress / liquidity crunch → risk-off massif
"layer_monetary": 0.90, "layer_rates": 0.90,
"layer_growth": 1.10, "layer_uk_macro": 1.10, "layer_macro": 1.10,
"layer_risk": 1.60, "layer_risk_credit": 1.70, "layer_refuge": 1.50,
"layer_positioning": 0.80, "layer_policy": 0.80,
"layer_flows": 0.70, "layer_demand": 0.80,
"layer_tech_fundamental": 0.75, "layer_supply": 0.90,
"layer_china": 0.85, "layer_dollar": 1.10,
"layer_global": 1.20,
},
"GROWTH_SCARE": {
# Récession / choc croissance → pivots BC, safe haven, EM sell
"layer_monetary": 1.10, "layer_rates": 1.10,
"layer_growth": 1.50, "layer_uk_macro": 1.50, "layer_macro": 1.50,
"layer_risk": 1.20, "layer_risk_credit": 1.20, "layer_refuge": 1.30,
"layer_positioning": 0.90, "layer_policy": 0.90,
"layer_flows": 0.85, "layer_demand": 0.80,
"layer_tech_fundamental": 0.70, "layer_supply": 0.90,
"layer_china": 1.30, "layer_dollar": 0.90,
"layer_global": 1.10,
},
"COMMODITY_SHOCK": {
# Choc énergie/matières premières → inflation, EM, or
"layer_monetary": 1.10, "layer_rates": 1.10,
"layer_growth": 1.20, "layer_uk_macro": 1.10, "layer_macro": 1.20,
"layer_risk": 1.10, "layer_risk_credit": 1.00, "layer_refuge": 1.20,
"layer_positioning": 1.00, "layer_policy": 1.00,
"layer_flows": 1.00, "layer_demand": 1.40,
"layer_tech_fundamental": 0.90, "layer_supply": 1.20,
"layer_china": 1.20, "layer_dollar": 0.95,
"layer_global": 1.10,
},
"BALANCED": {
# Régime neutre : aucune amplification
},
}
# Mapping catégories d'events → régimes candidats
_CAT_TO_REGIME: dict[str, str] = {
"central_bank": "MONETARY_DOMINANCE",
"monetary_shock": "MONETARY_DOMINANCE",
"geopolitical": "GEOPOLITICAL_RISK",
"trade_policy": "GEOPOLITICAL_RISK",
"credit_stress": "CREDIT_STRESS",
"growth_shock": "GROWTH_SCARE",
"commodity": "COMMODITY_SHOCK",
"sentiment": "BALANCED",
"technical": "BALANCED",
"positioning": "BALANCED",
"unclassified": "BALANCED",
}
# ── Lifecycle & decay ──────────────────────────────────────────────────────────
def _decay(days: int, absorption: int, dtype: str) -> float:
@@ -402,6 +519,82 @@ INSTRUMENT_MODELS: dict[str, dict] = {
} # end INSTRUMENT_MODELS
# ── Regime detection ───────────────────────────────────────────────────────────
def detect_regime(ev_by_cat: dict[str, float]) -> dict:
"""
Détermine le régime de marché dominant depuis la pression event par catégorie.
Retourne : {regime, label, scores, dominant_cat, weights}
"""
# Score de chaque catégorie = |pips|
cat_scores = {cat: abs(v) for cat, v in ev_by_cat.items() if v != 0}
if not cat_scores:
return {
"regime": "BALANCED", "label": "Équilibré",
"scores": {}, "dominant_cat": None,
"weights": {},
}
# Cumul de score par régime candidat
regime_scores: dict[str, float] = {}
for cat, score in cat_scores.items():
r = _CAT_TO_REGIME.get(cat, "BALANCED")
regime_scores[r] = regime_scores.get(r, 0.0) + score
dominant_regime = max(regime_scores, key=lambda r: regime_scores[r])
# Si le score maximal ne dépasse pas 3 pips → BALANCED
max_score = regime_scores.get(dominant_regime, 0.0)
if max_score < 3.0 or dominant_regime == "BALANCED":
dominant_regime = "BALANCED"
dominant_cat = max(
(c for c in cat_scores if _CAT_TO_REGIME.get(c) == dominant_regime),
key=lambda c: cat_scores[c],
default=None,
) if dominant_regime != "BALANCED" else None
REGIME_LABELS = {
"MONETARY_DOMINANCE": "Dominance Monétaire",
"GEOPOLITICAL_RISK": "Risque Géopolitique",
"CREDIT_STRESS": "Stress Crédit",
"GROWTH_SCARE": "Choc Croissance",
"COMMODITY_SHOCK": "Choc Commodités",
"BALANCED": "Équilibré",
}
weights = REGIME_WEIGHTS.get(dominant_regime, {})
return {
"regime": dominant_regime,
"label": REGIME_LABELS.get(dominant_regime, dominant_regime),
"scores": {r: round(s, 1) for r, s in regime_scores.items()},
"dominant_cat": dominant_cat,
"weights": weights,
}
def _apply_regime_weights(formula: str, weights: dict[str, float]) -> str:
"""
Injecte les multiplicateurs de régime dans une formule d'output.
Ex: "layer_monetary + layer_risk" + {layer_monetary:1.4, layer_risk:1.5}
"1.40 * layer_monetary + 1.50 * layer_risk"
Les termes sans poids restent à 1.0 (non modifiés).
"""
if not weights:
return formula
terms = [t.strip() for t in formula.split('+')]
out = []
for term in terms:
# Extraire l'id de couche (premier token alphanumérique_)
layer_id = term.strip()
w = weights.get(layer_id, 1.0)
if abs(w - 1.0) < 0.01:
out.append(layer_id)
else:
out.append(f"{w:.2f} * {layer_id}")
return " + ".join(out)
# ── DB ─────────────────────────────────────────────────────────────────────────
def init_instrument_model_tables(conn):
@@ -517,13 +710,14 @@ def _compute_event_by_category(conn, instrument: str, ref_date: date_type) -> di
# ── Graph evaluation ───────────────────────────────────────────────────────────
def _build_inputs(graph_def: dict, overrides: dict, ev_by_cat: dict) -> dict:
def _build_inputs(
graph_def: dict, overrides: dict, ev_by_cat: dict,
saturation: bool = True,
) -> dict:
"""
Build the inputs dict for evaluate_graph():
- input_event nodes → value from event contributions (category mapping)
- input_manual nodes → user_value × coefficient_to_pips
Any node with an override uses that value directly (already in pips for events;
for manual nodes the override IS the native-unit value → converted to pips).
Build inputs dict for evaluate_graph().
Phase 2 : les nœuds input_manual utilisent _saturate_pips() (tanh)
au lieu d'une conversion purement linéaire.
"""
inputs: dict[str, float] = {}
for node in graph_def["nodes"]:
@@ -532,29 +726,37 @@ def _build_inputs(graph_def: dict, overrides: dict, ev_by_cat: dict) -> dict:
ov = overrides.get(nid)
if ntype == "input_event":
if ov:
inputs[nid] = float(ov["value"]) # override IS pips
else:
cat = node.get("event_category", "")
inputs[nid] = ev_by_cat.get(cat, 0.0)
# Events sont déjà en pips → pas de saturation (déjà non-linéaire via lifecycle)
inputs[nid] = float(ov["value"]) if ov else ev_by_cat.get(node.get("event_category", ""), 0.0)
elif ntype == "input_manual":
if ov:
coeff = float(node.get("coefficient_to_pips", 1.0))
inputs[nid] = float(ov["value"]) * coeff
coeff = float(node.get("coefficient_to_pips", 1.0))
native = float(ov["value"]) if ov else 0.0
if native == 0.0:
inputs[nid] = 0.0
elif saturation:
inputs[nid] = _saturate_pips(native, coeff, node.get("unit", ""), node.get("saturation_scale"))
else:
inputs[nid] = 0.0 # neutral
inputs[nid] = coeff * native
return inputs
def _graph_json_for_eval(graph_def: dict) -> dict:
"""Convert our model graph_def to the format expected by evaluate_graph()."""
def _graph_json_for_eval(graph_def: dict, regime_weights: Optional[dict] = None) -> dict:
"""
Convertit le graph_def au format attendu par evaluate_graph().
Phase 2 : si regime_weights est fourni, la formule du nœud output est
réécrite avec les multiplicateurs de régime.
"""
output_id = graph_def.get("output_node", "")
nodes = []
for n in graph_def["nodes"]:
entry: dict = {"id": n["id"]}
if n.get("formula"):
entry["formula"] = n["formula"]
formula = n.get("formula")
if formula:
if n["id"] == output_id and regime_weights:
formula = _apply_regime_weights(formula, regime_weights)
entry["formula"] = formula
nodes.append(entry)
return {"nodes": nodes, "coefficients": {}}
@@ -562,7 +764,10 @@ def _graph_json_for_eval(graph_def: dict) -> dict:
# ── Public API ─────────────────────────────────────────────────────────────────
def get_model_state(conn, instrument: str, at_date: Optional[str] = None) -> Optional[dict]:
"""Full model state: all node values computed via DAG evaluation."""
"""
Full model state : DAG evaluation avec saturation (Phase 2) + poids de régime.
Retourne aussi {regime: {regime, label, weights, scores}}.
"""
inst_upper = instrument.upper()
row = conn.execute(
"SELECT graph_json FROM instrument_models WHERE instrument=?", (inst_upper,)
@@ -583,57 +788,79 @@ def get_model_state(conn, instrument: str, at_date: Optional[str] = None) -> Opt
).fetchall()}
ev_by_cat = _compute_event_by_category(conn, inst_upper, ref_date)
inputs = _build_inputs(graph_def, overrides, ev_by_cat)
# DAG evaluation (propagates through intermediate nodes via formulas)
# Phase 2 : détection régime + poids adaptatifs
regime_info = detect_regime(ev_by_cat)
regime_weights = regime_info["weights"]
# Inputs avec saturation tanh
inputs = _build_inputs(graph_def, overrides, ev_by_cat, saturation=True)
# DAG evaluation avec formule output pondérée par régime
from services.causal_graphs import evaluate_graph
gj = _graph_json_for_eval(graph_def)
all_vals = evaluate_graph(gj, inputs)
gj = _graph_json_for_eval(graph_def, regime_weights)
all_vals = evaluate_graph(gj, inputs)
output_id = graph_def["output_node"]
net_pips = round(float(all_vals.get(output_id, 0.0)), 1)
nodes_out = []
for node in graph_def["nodes"]:
nid = node["id"]
ntype = node.get("node_type", "")
val = round(float(all_vals.get(nid, 0.0)), 1)
ov = overrides.get(nid)
state = dict(node)
state["computed_value"] = val
state["pip_contribution"] = val
nid = node["id"]
ntype = node.get("node_type", "")
val = round(float(all_vals.get(nid, 0.0)), 1)
ov = overrides.get(nid)
st = dict(node)
st["computed_value"] = val
st["pip_contribution"] = val
if ntype == "input_event":
cat = node.get("event_category", "")
state["source"] = "manual" if ov else ("events" if val != 0.0 else "neutral")
state["raw_value"] = ov["value"] if ov else round(ev_by_cat.get(cat, 0.0), 1)
st["source"] = "manual" if ov else ("events" if val != 0.0 else "neutral")
st["raw_value"] = ov["value"] if ov else round(ev_by_cat.get(cat, 0.0), 1)
if ov:
state["override_note"] = ov.get("note", "")
state["override_set_at"] = ov.get("set_at", "")
st["override_note"] = ov.get("note", "")
st["override_set_at"] = ov.get("set_at", "")
elif ntype == "input_manual":
state["source"] = "manual" if ov else "neutral"
state["raw_value"] = ov["value"] if ov else 0.0 # in native unit
coeff = float(node.get("coefficient_to_pips", 1.0))
native = float(ov["value"]) if ov else 0.0
# Retourne la valeur native sans saturation (pour affichage)
st["source"] = "manual" if ov else "neutral"
st["raw_value"] = native
# pip_contribution inclut la saturation (= all_vals[nid])
# On expose aussi la valeur linéaire pour comparaison
st["pip_linear"] = round(coeff * native, 1)
st["pip_saturated"] = val
st["saturation_pct"] = (
round((1.0 - val / (coeff * native)) * 100, 1)
if native != 0 and coeff != 0
else 0.0
)
if ov:
state["override_note"] = ov.get("note", "")
state["override_set_at"] = ov.get("set_at", "")
st["override_note"] = ov.get("note", "")
st["override_set_at"] = ov.get("set_at", "")
elif ntype in ("intermediate", "output"):
state["source"] = "computed"
st["source"] = "computed"
# Pour les intermédiaires, expose le multiplicateur de régime
if ntype == "intermediate":
st["regime_weight"] = regime_weights.get(nid, 1.0)
nodes_out.append(state)
nodes_out.append(st)
direction = "bullish" if net_pips > 5 else "bearish" if net_pips < -5 else "neutral"
return {
"instrument": inst_upper,
"name": graph_def["name"],
"instrument": inst_upper,
"name": graph_def["name"],
"description": graph_def.get("description", ""),
"at_date": str(ref_date),
"net_pips": net_pips,
"direction": direction,
"nodes": nodes_out,
"at_date": str(ref_date),
"net_pips": net_pips,
"direction": direction,
"nodes": nodes_out,
"output_node": output_id,
"regime": regime_info,
}
@@ -721,12 +948,10 @@ def simulate_timeline(
})
from services.causal_graphs import evaluate_graph
gj = _graph_json_for_eval(graph_def)
timeline = []
cur = date_from
while cur <= today:
# Event contributions for this day
ev_by_cat: dict[str, float] = {}
for ev in events:
if ev["ev_date"] > cur:
@@ -738,13 +963,17 @@ def simulate_timeline(
cat = ev["category"]
ev_by_cat[cat] = ev_by_cat.get(cat, 0.0) + round(ev["pips"] * df, 2)
inputs = _build_inputs(graph_def, overrides, ev_by_cat)
vals = evaluate_graph(gj, inputs)
net = round(float(vals.get(output_id, 0.0)), 1)
# Phase 2 : régime du jour → poids adaptatifs dans la formule output
ri = detect_regime(ev_by_cat)
gj = _graph_json_for_eval(graph_def, ri["weights"])
inputs = _build_inputs(graph_def, overrides, ev_by_cat, saturation=True)
vals = evaluate_graph(gj, inputs)
net = round(float(vals.get(output_id, 0.0)), 1)
timeline.append({
"date": str(cur),
"net_pips": net,
"regime": ri["regime"],
"nodes": {k: round(float(v), 1) for k, v in vals.items()},
})
cur += timedelta(days=1)