#!/usr/bin/env python """r0726_venue_tripwire.py — un sistema di allerta precoce sul FALLIMENTO DEL VENUE. PROBLEMA (dal 26/07, `r0726_venue_risk.py`): TP01+SKH01+VRP01 stanno tutti sullo stesso conto Deribit. Sono quasi-ortogonali sui ritorni e **perfettamente correlati sul fallimento del venue**. L'operatore ha deciso (26/07) di restare concentrato al 100% fino a $20k: la riduzione dell'ESPOSIZIONE e' quindi fuori discussione fino a quella soglia. Resta una sola leva: **il TEMPO**. Il modello di rischio assumeva un salto a zero istantaneo, ma i fallimenti reali non sono istantanei — Mt.Gox gato' i prelievi fiat per MESI prima di chiudere, FTX ebbe ~3 giorni fra il primo blocco e la bancarotta, QuadrigaCX settimane. Se una parte del saldo esce dentro quella finestra, la perdita non e' totale. DOMANDA MISURABILE: esiste un segnale, calcolabile dai dati che gia' scarichiamo, che distingua "questo venue e' gated" da "il mercato e' in crash"? E con quale tasso di FALSI ALLARMI? MECCANISMO IPOTIZZATO (ed e' il punto che rende il test falsificabile): quando un exchange blocca i prelievi, **l'arbitraggio si rompe** e il suo prezzo si stacca dal consenso in modo (a) PERSISTENTE e (b) A SEGNO COSTANTE — perche' l'ostacolo e' strutturale, non di liquidita'. Un crash di mercato disloca anche lui, ma per minuti e con segno che rimbalza, perche' l'arb funziona ancora, e' solo lento. Su Mt.Gox il BTC arrivo' a un PREMIO del 10-20% (si comprava BTC per far uscire il valore); su un venue in fuga si vede lo sconto. **Il segnale e' |scarto|, non il suo segno** — entrambe le direzioni dicono la stessa cosa: l'arb non chiude piu'. COSA MISURA QUESTO SCRIPT 1. la distribuzione dello scarto Deribit-vs-consenso in bps su TUTTA la storia certificata (8 anni, BTC+ETH), consenso = venue USD indipendenti (Coinbase, Kraken); 2. gli EPISODI (run consecutivi sopra soglia a segno costante): quanti, quanto lunghi, quando — cioe' il tasso di falsi allarmi di ogni possibile (soglia, durata); 3. la FRONTIERA A ZERO FALSI ALLARMI: la coppia (bps, ore) che in 8 anni non e' mai scattata, inclusi tutti i crash veri (Mar-2020, Mag-2021, LUNA, FTX, ...); 4. il MARGINE: di quante volte la dislocazione di un fallimento storico supera quella soglia. CONTROLLO DI FALSI POSITIVI CHE VA FATTO O IL SISTEMA E' INUTILE: **la referenza puo' rompersi lei**. Se Coinbase ha un outage, Deribit sembra dislocato. Per questo il consenso richiede **almeno 2 venue indipendenti d'accordo fra loro**: se le referenze non concordano, il campione si scarta invece di generare un allarme (`REF_DISAGREE_BPS`). uv run python scripts/research/r0726_venue_tripwire.py """ from __future__ import annotations import pickle import sys import time from pathlib import Path import numpy as np import pandas as pd ROOT = Path(__file__).resolve().parents[2] sys.path.insert(0, str(ROOT)) from src.data.downloader import load_data # noqa: E402 CACHE = ROOT / "data" / "_cache" / "venue_ref_1h.pkl" ASSETS = ("BTC", "ETH") # Referenze INDIPENDENTI da Deribit e denominate in USD (non USDT: il depeg del 2022 sposta # BTC/USDT fino al 3% dal dollaro e produrrebbe falsi allarmi giganti — regola del progetto). # ⚠️ KRAKEN E' STATO RIMOSSO dopo la prima corsa: il suo endpoint OHLC pubblico ritorna solo le # ultime ~700 candele qualunque `since`, e l'inner join TRONCAVA il campione da 8 anni a 29 giorni # in silenzio (la "frontiera a zero falsi allarmi in 8 anni" era misurata su un mese). Bitstamp e # Bitfinex hanno storia profonda dal 2019. Da qui la guardia di copertura in `main`. REFS = [("coinbase", {"BTC": "BTC/USD", "ETH": "ETH/USD"}, 300), ("bitstamp", {"BTC": "BTC/USD", "ETH": "ETH/USD"}, 1000), ("bitfinex", {"BTC": "BTC/USD", "ETH": "ETH/USD"}, 1000)] # Copertura minima del campione: sotto questa quota la conclusione NON si stampa. Serve perche' # la prima corsa aveva la diagnostica giusta gia' stampata (702 ore su 69.663) e il testo di # sintesi diceva comunque "8 anni": la guardia deve essere un controllo, non una riga di log. MIN_COVERAGE = 0.50 # Se le due referenze divergono FRA LORO piu' di questo, il campione non e' utilizzabile: # il problema e' della referenza, non di Deribit. Scartare, non allertare. REF_DISAGREE_BPS = 100.0 # Griglia della frontiera: soglia in bps x durata minima in ore consecutive. BPS_GRID = (25, 50, 75, 100, 150, 200, 300, 500) HOURS_GRID = (1, 2, 3, 4, 6, 8, 12, 24) # =========================================================================== # fetch delle referenze (con cache: 8 anni x 2 venue x 2 asset e' lento) # =========================================================================== def _fetch_1h(exchange_id: str, symbol: str, start_ms: int, end_ms: int, limit: int) -> pd.Series: """OHLCV 1h paginato -> Series close indicizzata sul timestamp ms. Tollerante agli errori.""" import ccxt ex = getattr(ccxt, exchange_id)({"enableRateLimit": True}) ex.load_markets() if symbol not in ex.markets: return pd.Series(dtype=float) out: dict[int, float] = {} since, guard = start_ms, 0 while since <= end_ms and guard < 6000: guard += 1 rows = None for attempt in range(3): try: rows = ex.fetch_ohlcv(symbol, "1h", since=since, limit=limit) break except Exception: if attempt == 2: return pd.Series(out, dtype=float) time.sleep(2 ** attempt) rows = [r for r in (rows or []) if int(r[0]) >= since] if not rows: break for r in rows: if start_ms <= int(r[0]) <= end_ms and r[4]: out[int(r[0])] = float(r[4]) nxt = int(rows[-1][0]) + 3_600_000 if nxt <= since: break since = nxt return pd.Series(out, dtype=float) def load_refs(refresh: bool = False) -> dict[tuple[str, str], pd.Series]: """Cache INCREMENTALE per (asset, venue). Chiave per-venue e non globale: cambiare la lista delle referenze non deve invalidare quelle gia' scaricate ne', peggio, restituire in silenzio una cache che non contiene i venue nuovi.""" out: dict[tuple[str, str], pd.Series] = {} if CACHE.exists() and not refresh: out = pickle.loads(CACHE.read_bytes()) needed = [(a, eid, syms, lim) for a in ASSETS for eid, syms, lim in REFS if (a, eid) not in out] for asset, eid, syms, lim in needed: d = load_data(asset, "1h") s_ms, e_ms = int(d["timestamp"].iloc[0]), int(d["timestamp"].iloc[-1]) print(f" fetch {eid:<9} {asset} ...", flush=True) out[(asset, eid)] = _fetch_1h(eid, syms[asset], s_ms, e_ms, lim) print(f" -> {len(out[(asset, eid)]):,} barre", flush=True) if needed: CACHE.parent.mkdir(parents=True, exist_ok=True) CACHE.write_bytes(pickle.dumps(out)) return out # =========================================================================== # nucleo puro e testabile # =========================================================================== def dislocation(deribit: pd.Series, refs: list[pd.Series], ref_disagree_bps: float = REF_DISAGREE_BPS) -> pd.DataFrame: """Scarto FIRMATO di Deribit dal consenso delle referenze, in bps, ora per ora. Il consenso e' la MEDIANA delle referenze disponibili. Le ore in cui le referenze non concordano fra loro (spread > `ref_disagree_bps`) sono SCARTATE: li' non si puo' dire se il problema sia di Deribit o della referenza, e un allarme sarebbe indistinguibile da un guasto altrui. Serve almeno 2 referenze; con una sola tutte le ore restano ma `usable=False`. OUTER JOIN, non inner. Le referenze hanno storie di lunghezza diversa (Bitfinex: 33k ore su BTC, ZERO su ETH) e un `join="inner"` le farebbe decidere TUTTE il campione — la prima corsa perse cosi' l'87% delle ore. Il consenso si forma riga per riga sulle referenze DISPONIBILI in quell'ora, e la riga e' utilizzabile se ce ne sono almeno 2 e concordano. Stessa lezione di `combine_outer` (26/07): un outer-join con rinormalizzazione per-riga, mai un'intersezione. """ cols = {"deribit": deribit} for i, r in enumerate(refs): cols[f"ref{i}"] = r m = pd.concat(cols, axis=1, join="outer") m = m[m["deribit"].notna()] ref_cols = [c for c in m.columns if c.startswith("ref")] if not len(m) or not ref_cols: return pd.DataFrame(columns=["deribit", "consensus", "bps", "usable"]) R = m[ref_cols] n_ref = R.notna().sum(axis=1) m["consensus"] = R.median(axis=1, skipna=True) spread = (R.max(axis=1) - R.min(axis=1)) / m["consensus"] * 1e4 # utilizzabile = almeno 2 referenze presenti E d'accordo fra loro. Con una sola referenza non # si puo' distinguere "Deribit e' fuori" da "la referenza e' rotta": si scarta, non si allerta. m["usable"] = (n_ref >= 2) & (spread.fillna(np.inf) <= ref_disagree_bps) m["bps"] = (m["deribit"] - m["consensus"]) / m["consensus"] * 1e4 m = m[m["consensus"].notna()] return m[["deribit", "consensus", "bps", "usable"]] def episodes(bps: pd.Series, usable: pd.Series, thr_bps: float, min_hours: int) -> list[dict]: """Run CONSECUTIVI in cui |scarto| > soglia CON SEGNO COSTANTE, lunghi >= min_hours. Il segno costante e' cio' che separa "arb rotto" (ostacolo strutturale, lo scarto sta da un lato) da "mercato sottile" (lo scarto rimbalza fra i due lati). Le ore non utilizzabili (referenze in disaccordo) ROMPONO il run: non si accumula evidenza su dati che non parlano. """ over = (bps.abs() > thr_bps) & usable sign = np.sign(bps) out, start, cur_sign = [], None, 0 for i in range(len(bps)): if over.iloc[i] and (start is None or sign.iloc[i] == cur_sign): if start is None: start, cur_sign = i, sign.iloc[i] else: if start is not None and i - start >= min_hours: seg = bps.iloc[start:i] out.append(dict(start=bps.index[start], end=bps.index[i - 1], hours=i - start, peak_bps=float(seg.abs().max()), sign=int(cur_sign))) start = None if over.iloc[i]: start, cur_sign = i, sign.iloc[i] if start is not None and len(bps) - start >= min_hours: seg = bps.iloc[start:] out.append(dict(start=bps.index[start], end=bps.index[-1], hours=len(bps) - start, peak_bps=float(seg.abs().max()), sign=int(cur_sign))) return out def zero_fp_frontier(bps: pd.Series, usable: pd.Series, bps_grid=BPS_GRID, hours_grid=HOURS_GRID) -> list[tuple[float, int]]: """Per ogni durata, la soglia in bps PIU' BASSA che non ha mai fatto scattare l'allarme. Piu' bassa = piu' sensibile a parita' di zero falsi allarmi. La frontiera e' l'insieme delle coppie ammissibili; sceglierne una e' una decisione, non un risultato. """ front = [] for h in hours_grid: for b in bps_grid: if not episodes(bps, usable, b, h): front.append((float(b), int(h))) break return front def main() -> None: print("=" * 100) print(" r0726 — TRIPWIRE DI VENUE: lo scarto Deribit-vs-consenso come allerta precoce") print("=" * 100) print(" Ipotesi: un venue che GATA i prelievi rompe l'arbitraggio -> scarto PERSISTENTE e a") print(" SEGNO COSTANTE. Un crash disloca anche lui, ma per poco e a segno che rimbalza.") print(" Referenze USD indipendenti (Coinbase, Kraken); niente USDT (il depeg 2022 le sposta).") print("\n [1/4] referenze cross-venue (cache in data/_cache/) ...", flush=True) refs = load_refs() frames, coverage = {}, {} for a in ASSETS: d = load_data(a, "1h") der = pd.Series(d["close"].astype(float).values, index=d["timestamp"].astype(int).values) rr = [refs[(a, eid)] for eid, _, _ in REFS if len(refs.get((a, eid), []))] frames[a] = dislocation(der, rr) f = frames[a] coverage[a] = len(f) / len(d) if len(d) else 0.0 cov = f["usable"].mean() * 100 if len(f) else 0.0 flag = "" if coverage[a] >= MIN_COVERAGE else " <-- ⚠️ COPERTURA INSUFFICIENTE" print(f" {a}: {len(f):,} ore in comune su {len(d):,} certificate " f"({coverage[a]*100:.1f}% del campione), {cov:.1f}% utilizzabili, " f"{len(rr)} referenze{flag}") if min(coverage.values()) < MIN_COVERAGE: print(f"\n ✋ STOP: la copertura minima e' {min(coverage.values())*100:.1f}% < " f"{MIN_COVERAGE*100:.0f}%. Una referenza corta TRONCA il campione nell'inner join") print(" e la 'frontiera a zero falsi allarmi' descriverebbe una finestra breve, non") print(" la storia. Nessuna conclusione stampata: e' cosi' che deve fallire.") sys.exit(1) # ------------------------------------------------ 2. distribuzione print("\n [2/4] DISTRIBUZIONE dello scarto |Deribit - consenso| (solo ore utilizzabili)") print(f"\n {'asset':<6}{'ore':>9}{'mediana':>10}{'p95':>9}{'p99':>9}{'p99.9':>9}" f"{'max':>10}{'>50bps':>9}{'>100bps':>10}") for a in ASSETS: f = frames[a] v = f.loc[f["usable"], "bps"].abs() if not len(v): continue print(f" {a:<6}{len(v):>9,}{v.median():>10.1f}{v.quantile(.95):>9.1f}" f"{v.quantile(.99):>9.1f}{v.quantile(.999):>9.1f}{v.max():>10.1f}" f"{(v > 50).mean()*100:>8.2f}%{(v > 100).mean()*100:>9.3f}%") # ------------------------------------------------ 3. episodi print("\n [3/4] EPISODI = run consecutivi sopra soglia A SEGNO COSTANTE (= falsi allarmi)") for a in ASSETS: f = frames[a] print(f"\n {a} {'soglia':>8} |" + "".join(f"{h:>4}h" for h in HOURS_GRID)) for b in BPS_GRID: n = [len(episodes(f["bps"], f["usable"], b, h)) for h in HOURS_GRID] print(f" {'':>4}{b:>10} bps |" + "".join(f"{x:>5}" for x in n)) # i piu' lunghi in assoluto: sono i candidati falsi allarmi da capire uno per uno print("\n I 6 episodi PIU' LUNGHI a 50 bps (cosa fara' scattare il sistema per sbaglio):") for a in ASSETS: eps = sorted(episodes(frames[a]["bps"], frames[a]["usable"], 50, 1), key=lambda e: -e["hours"])[:6] for e in eps: t0 = pd.to_datetime(e["start"], unit="ms", utc=True) print(f" {a} {t0:%Y-%m-%d %H:%M} {e['hours']:>3}h picco {e['peak_bps']:>7.0f} bps" f" segno {e['sign']:+d}") # ------------------------------------------------ 4. frontiera print("\n [4/4] FRONTIERA A ZERO FALSI ALLARMI (soglia minima che in 8 anni non scatta mai)") print(f"\n {'durata':>8} | " + " | ".join(f"{a:>10}" for a in ASSETS)) fronts = {a: dict((h, b) for b, h in zero_fp_frontier(frames[a]["bps"], frames[a]["usable"])) for a in ASSETS} for h in HOURS_GRID: cells = [] for a in ASSETS: b = fronts[a].get(h) cells.append(f"{b:>7.0f} bps" if b else f"{'nessuna':>10}") print(f" {h:>6}h | " + " | ".join(f"{c:>10}" for c in cells)) # SCELTA OPERATIVA — il criterio va DICHIARATO, perche' due criteri ingenui sbagliano in versi # opposti (entrambi provati in sessione): # "minimi bps" -> 25bps/24h: consuma 24h delle ~72h che diede FTX; # "minime ore" -> 500bps/2h: MANCA FTX (300 bps), margine 0.6x. # Criterio: (1) zero falsi allarmi su ENTRAMBI gli asset; (2) soglia con margine >= MIN_MARGIN # sul caso storico PIU' DEBOLE; (3) a quei vincoli, minima latenza. # La latenza e' poco costosa e lo dice il controllo positivo qui sotto: gli episodi di un venue # realmente gated durano CENTINAIA di ore, non poche. MIN_MARGIN = 3.0 HIST = (("Mt.Gox 2013-14 (premio)", 1000.0), ("Mt.Gox finale", 2000.0), ("FTX nov-2022 (sconto)", 300.0), ("QuadrigaCX", 500.0)) weakest = min(b for _, b in HIST) max_thr = weakest / MIN_MARGIN cand = [(h, max(fronts[a].get(h, 1e9) for a in ASSETS)) for h in HOURS_GRID] cand = [(h, b) for h, b in cand if b < 1e9] ok = [(h, b) for h, b in cand if b <= max_thr] if ok: h_sel, b_sel = min(ok, key=lambda x: (x[0], x[1])) print(f"\n --> coppie a ZERO falsi allarmi su ENTRAMBI gli asset:") for h, b in cand: why = "" if b > max_thr: why = f" (scartata: {weakest/b:.1f}x < {MIN_MARGIN:.0f}x su FTX)" elif h == h_sel: why = " <- SCELTA" print(f" {h:>3}h -> {b:>4.0f} bps{why}") print(f"\n Criterio: margine >= {MIN_MARGIN:.0f}x sul caso storico piu' debole") print(f" ({weakest:.0f} bps = FTX) -> soglia <= {max_thr:.0f} bps; poi minima latenza.") print(f" SCELTA: {b_sel:.0f} bps persistenti {h_sel}h a segno costante.") print(f"\n MARGINE della soglia scelta vs i fallimenti storici:") for lab, bp in HIST: print(f" {lab:<26} {bp:>6.0f} bps -> {bp/b_sel:>5.1f}x la soglia") else: h_sel, b_sel = min(cand, key=lambda x: (x[0], x[1])) print(f"\n ✋ NESSUNA coppia soddisfa insieme zero-falsi-allarmi e margine {MIN_MARGIN:.0f}x:") print(f" la piu' vicina e' {b_sel:.0f}bps/{h_sel}h. Il rilevatore non e' utilizzabile") print(f" cosi' com'e' e va cambiato segnale, non tarato piu' aggressivamente.") # ------------------------------------------------ CONTROLLO POSITIVO # ⚠️ Deve stare FUORI dai due rami: nella prima stesura era finito dentro l'`else` (il ramo di # fallimento) e quindi non girava MAI nel caso normale — un controllo positivo che non gira e' # esattamente il difetto che dovrebbe prevenire. # Un rilevatore che non segnala mai nulla puo' essere semplicemente ROTTO. Qui lo si punta su # un venue che ha DAVVERO avuto un episodio di stress da prelievi: Bitfinex nel 2018-19 # (problemi bancari/Tether) trattava BTC a premio persistente sul consenso USD. # Bitfinex viene ESCLUSO dal proprio consenso, altrimenti assorbirebbe la sua dislocazione. print("\n" + "=" * 100) print(" CONTROLLO POSITIVO — il rilevatore scatta su un venue realmente in stress?") print("=" * 100) print(" Bersaglio: BITFINEX 2018-19 (problemi bancari/Tether -> premio persistente).") print(" Consenso = Coinbase + Bitstamp; Bitfinex escluso dal proprio consenso.") fired = 0 for a in ASSETS: tgt = refs.get((a, "bitfinex"), pd.Series(dtype=float)) base = [refs[(a, e)] for e in ("coinbase", "bitstamp") if len(refs.get((a, e), []))] if not len(tgt) or len(base) < 2: print(f" {a}: dati insufficienti per il controllo (Bitfinex {len(tgt)} barre)") continue fx = dislocation(tgt, base) eps = episodes(fx["bps"], fx["usable"], b_sel, h_sel) fired += len(eps) v = fx.loc[fx["usable"], "bps"] print(f"\n {a}: {len(fx):,} ore, scarto mediano {v.abs().median():.1f} bps, " f"max {v.abs().max():.0f} bps -> {len(eps)} EPISODI a {b_sel:.0f}bps/{h_sel}h") for e in sorted(eps, key=lambda x: -x["hours"])[:5]: t0 = pd.to_datetime(e["start"], unit="ms", utc=True) print(f" {t0:%Y-%m-%d %H:%M} {e['hours']:>4}h picco {e['peak_bps']:>6.0f} bps" f" segno {e['sign']:+d}") if fired: print(f"\n ✅ Il rilevatore SCATTA ({fired} episodi) su un venue realmente in stress e") print(f" ZERO volte su Deribit in 8 anni. La specificita' non e' cecita'.") print(f" E la DURATA di quegli episodi (centinaia di ore) e' la risposta alla") print(f" domanda che conta: un venue gated resta dislocato per SETTIMANE, quindi") print(f" {h_sel}h di latenza di rilevamento costano una frazione trascurabile del preavviso.") else: print("\n ✋ ZERO episodi anche sul bersaglio in stress: il rilevatore e' CIECO e la") print(" frontiera sopra non significa niente. Non usare.") if __name__ == "__main__": main()