Files
Adriano Dal Pastro d55eb13533 feat(chain): assorbita la raccolta catena opzioni, cerbero-bite dismesso
cerbero-bite viene eliminato. L'unica sua parte irreversibile e' il DATO:
una catena opzioni non si ricostruisce a posteriori (Deribit non serve book
storici, non c'e' un secondo venue). Il codice si riscrive; le ore non
raccolte no.

ASSORBITO
- scripts/live/collect_chain.py + scripts/cron_chain.sh (cron 25 * * * *):
  raccolta propria, ~570 strumenti/giro, ~3 min.
- scripts/analysis/import_cb_archive.py: archivio 1.23M righe (2026-05-01+)
  + market_snapshots 17.402 righe (2026-03-26+: dealer gamma, gamma flip,
  rischio liquidazioni, funding cross — dati che non abbiamo altrove).
- snapshot sqlite integrale in /opt/docker/backups/manual/ (SHA256).

NON ASSORBITO, con motivo: motore credit-spread ETH (regola "niente
short-vol da modello in deploy", conto a $52 contro minimo $720), GUI, kill
switch/dead-man/audit (abbiamo venue_watch/edge_watch/monitor_health/
fee_watch), dvol_history (fetch_dvol.py ha storia PIU' LUNGA: 2020+ contro
2026-05), decisions/positions (0 posizioni).

TRE DIFETTI DI BITE NON REPLICATI, tutti misurati il 30/07:
1. una chiamata per strumento invece di due (get_order_book?depth=3 da' gia'
   quote+greche+IV+OI+book+underlying) + prefiltro OI in una chiamata sola:
   551 -> ~290 chiamate per asset;
2. pacing invece di raffica. Il carico non e' mai stato il problema: 570
   chiamate/ora = 0.16/s DISTRIBUITE; bite le sparava in 26s (~44/s) e si
   auto-saturava il rate limit per-IP (12.186 risposte 429 in 26h, 96% al
   minuto :00). Primo giro reale: 574 chiamate, 0 risposte 429. Il minuto :25
   e' scelto: :00 era la raffica, :07 e' cron_book (feed 5m di SKH01).
3. quote_status esplicito {ok, no_quote, error} e book_depth NULL su errore
   mai 0. "Book vuoto" e "chiamata fallita" sono cose diverse: e' per questo
   che il guasto del 29/07 (50% di quote perse, 38 ore) non produsse alcun
   segnale. Le righe ereditate restano 'unknown': bite non lo registrava e a
   posteriori non e' ricostruibile.

Battuta di cuore in data/chain_collect/runs.jsonl anche a giro fallito,
sorvegliata da monitor_health (1h, max_age 3h): un collettore fermo non
produce niente, e il niente si legge come "nessun dato quel giorno".

Difetto trovato per strada: due formati ISO nella stessa colonna (92 righe
di backfill senza microsecondi). pd.to_datetime senza `format` ne inferisce
uno solo e manda gli altri a NaT -> il dropna a valle li toglieva in
silenzio, e la serie di contesto perdeva 5 settimane slittando dal 26/03 al
01/05. Corretto con format="ISO8601" e scarto RUMOROSO.

Book, pesi, config, strategia INVARIATI. 537 test verdi.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-30 20:14:34 +00:00

262 lines
12 KiB
Python

"""COLLETTORE CATENA OPZIONI — successore di cerbero-bite, dentro PythagorasGoal (2026-07-30).
Perche' esiste: cerbero-bite viene eliminato, e con esso si fermerebbe l'unica raccolta di prezzi
opzioni REALI del progetto. Una catena non e' ricostruibile a posteriori (Deribit non serve book
storici, non c'e' un secondo venue) -> l'ora non raccolta e' persa per sempre. Il valore
dell'archivio sta negli eventi RARI: e' la rete stesa in attesa del regime di vol alta che
promuovera' o uccidera' VRP01 (criterio dichiarato il 19/06, gate IV-rank>0.30 mai attivo nella
finestra raccolta finora).
TRE DIFFERENZE DA BITE, tutte misurate il 30/07 e tutte deliberate:
1. UNA CHIAMATA PER STRUMENTO, non due. `public/get_order_book?depth=3` restituisce gia' quote,
greche, IV, open interest, volume, book E underlying_price. Bite chiamava ticker + orderbook
separatamente: doppio costo e possibilita' di disallineamento fra i due (quote di un istante,
book di un altro).
2. PACING, non raffica. Il carico non e' mai stato il problema: ~570 chiamate/ora = 0.16/s se
distribuite. Bite le sparava in ~26 secondi (~44/s) e si auto-saturava il rate limit per-IP
(12.186 risposte 429 in 26 ore, 96% nel minuto :00), con l'effetto collaterale di disturbare
il feed 5m del book live sulla stessa VPS. Qui: token bucket a `--rps` (default 4/s, ~2.5
minuti per giro) + backoff sul 429. Un giro lento non costa nulla; una raffica costa il dato.
3. STATO ESPLICITO DELLA QUOTA. Bite persisteva la riga anche quando la chiamata falliva, con
bid/ask NULL: il conteggio righe restava identico e nessun controllo di copertura se ne
accorgeva (il 29/07 il 50% delle quote e' diventato vuoto per 38 ore senza un segnale). Qui
ogni riga porta `quote_status` in {ok, no_quote, error}:
ok = il venue ha risposto e c'e' almeno un lato del book
no_quote = il venue ha risposto e il book e' vuoto da entrambi i lati (fatto di mercato)
error = la chiamata e' fallita (fatto di infrastruttura)
Sono cose diverse e non vanno mai confuse. Per lo stesso motivo `book_depth_top3` e' NULL su
errore, MAI 0: bite scriveva 0 e "chiamata fallita" diventava indistinguibile da "book vuoto".
uv run python scripts/live/collect_chain.py # un giro, entrambi gli asset
uv run python scripts/live/collect_chain.py --assets ETH --rps 8
uv run python scripts/live/collect_chain.py --dry-run # non scrive
"""
from __future__ import annotations
import argparse
import json
import sys
import time
from dataclasses import dataclass, field
from datetime import UTC, datetime
from pathlib import Path
import pandas as pd
import requests
PROJECT_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(PROJECT_ROOT))
API = "https://www.deribit.com/api/v2/public"
STORE = PROJECT_ROOT / "data" / "raw" / "cb_chain"
ASSETS = ("BTC", "ETH")
EXPIRY_MAX_DAYS = 95 # 1g..3mesi, come la finestra di bite (continuita' della serie)
OI_MIN = 100.0 # come bite: sotto questa soglia lo strumento e' rumore
DEFAULT_RPS = 4.0
TIMEOUT = 15
@dataclass
class Budget:
"""Token bucket + contabilita' del giro. Il 429 non e' un dettaglio: e' la cosa da non fare."""
rps: float
_next: float = 0.0
calls: int = 0
errors: int = 0
rate_limited: int = 0
waited_s: float = 0.0
err_samples: list[str] = field(default_factory=list)
def wait(self) -> None:
now = time.monotonic()
if now < self._next:
time.sleep(self._next - now)
self.waited_s += self._next - now
self._next = max(now, self._next) + 1.0 / self.rps
def note_error(self, msg: str) -> None:
self.errors += 1
if len(self.err_samples) < 5:
self.err_samples.append(msg[:160])
def _get(path: str, params: dict, budget: Budget, tries: int = 3) -> dict | None:
"""GET con pacing e backoff. Ritorna None se la chiamata non e' andata a buon fine."""
for k in range(tries):
budget.wait()
budget.calls += 1
try:
r = requests.get(f"{API}/{path}", params=params, timeout=TIMEOUT)
except Exception as exc: # rete: si registra QUI, non dopo
budget.note_error(f"{path}: {type(exc).__name__}: {exc}")
time.sleep(1.5 * (k + 1))
continue
if r.status_code == 429:
budget.rate_limited += 1
time.sleep(2.0 * (k + 1)) # backoff: il venue ha detto di rallentare
continue
if r.status_code != 200:
budget.note_error(f"{path}: HTTP {r.status_code}")
time.sleep(1.0 * (k + 1))
continue
try:
return r.json()["result"]
except Exception as exc:
budget.note_error(f"{path}: payload illeggibile: {exc}")
return None
return None
def instruments(asset: str, budget: Budget, now: datetime) -> list[dict]:
res = _get("get_instruments", {"currency": asset, "kind": "option", "expired": "false"}, budget)
if not res:
return []
horizon = now.timestamp() * 1000 + EXPIRY_MAX_DAYS * 86400_000
return [i for i in res if i.get("expiration_timestamp", 0) <= horizon]
def _depth_top3(side: list) -> float | None:
if side is None:
return None
return float(sum(row[1] for row in side[:3] if isinstance(row, (list, tuple)) and len(row) >= 2))
def snapshot_row(inst: dict, ob: dict | None, ts: datetime) -> dict:
"""Una riga per strumento — SEMPRE, ma con lo stato della quota dichiarato."""
name = inst["instrument_name"]
base = {
"ts": ts, "asset": inst["base_currency"], "instrument_name": name,
"strike": float(inst["strike"]), "option_type": "P" if inst["option_type"] == "put" else "C",
"exp": pd.Timestamp(inst["expiration_timestamp"], unit="ms", tz="UTC"),
"bid": None, "ask": None, "mid": None, "iv": None,
"delta": None, "gamma": None, "theta": None, "vega": None,
"open_interest": None, "volume_24h": None, "book_depth_top3": None,
"underlying_price": None, "index_price": None,
"quote_status": "error", "source": "pyg",
}
if ob is None:
return base # errore: depth resta NULL, mai 0
g = ob.get("greeks") or {}
stats = ob.get("stats") or {}
bid, ask = ob.get("best_bid_price"), ob.get("best_ask_price")
bid = float(bid) if bid else None # Deribit manda 0.0 per "nessun lato"
ask = float(ask) if ask else None
db, da = _depth_top3(ob.get("bids")), _depth_top3(ob.get("asks"))
base.update({
"bid": bid, "ask": ask,
"mid": (bid + ask) / 2 if (bid is not None and ask is not None) else None,
"iv": float(ob["mark_iv"]) if ob.get("mark_iv") is not None else None,
"delta": g.get("delta"), "gamma": g.get("gamma"),
"theta": g.get("theta"), "vega": g.get("vega"),
"open_interest": ob.get("open_interest"), "volume_24h": stats.get("volume"),
"book_depth_top3": (db or 0.0) + (da or 0.0),
"underlying_price": ob.get("underlying_price"), "index_price": ob.get("index_price"),
# il venue ha risposto: se non c'e' nessun lato e' un fatto di MERCATO, non un guasto
"quote_status": "ok" if (bid is not None or ask is not None) else "no_quote",
})
return base
def open_interest_map(asset: str, budget: Budget) -> dict[str, float] | None:
"""OI di TUTTA la catena in UNA chiamata (`get_book_summary_by_currency`).
Serve a non spendere una chiamata per scoprire che uno strumento e' sotto soglia: il prefiltro
dimezza il giro (551 -> ~300 chiamate su ETH). None = la chiamata e' fallita, e allora si
raccoglie TUTTO invece di filtrare su un dato che non si ha: un filtro su dati mancanti
scarterebbe strumenti buoni fingendo che fossero illiquidi.
"""
res = _get("get_book_summary_by_currency", {"currency": asset, "kind": "option"}, budget)
if not res:
return None
return {r["instrument_name"]: float(r.get("open_interest") or 0.0) for r in res}
def sweep(asset: str, budget: Budget, now: datetime) -> pd.DataFrame:
insts = [i for i in instruments(asset, budget, now) if float(i.get("strike") or 0) > 0]
oi = open_interest_map(asset, budget)
if oi is not None:
insts = [i for i in insts if oi.get(i["instrument_name"], 0.0) >= OI_MIN]
rows = []
for inst in insts:
ob = _get("get_order_book", {"instrument_name": inst["instrument_name"], "depth": 3}, budget)
rows.append(snapshot_row(inst, ob, now))
return pd.DataFrame(rows)
def heartbeat(now: datetime, df: pd.DataFrame, budget: Budget) -> None:
"""Una riga per giro in `data/chain_collect/runs.jsonl`, letta da `monitor_health`.
Serve perche' un collettore fermo non produce NIENTE, e il niente si legge come "nessun dato
quel giorno" invece che come "raccolta rotta" — con una serie irrecuperabile e' il modo piu'
caro di sbagliare. La battuta di cuore esiste anche quando il giro fallisce.
"""
d = PROJECT_ROOT / "data" / "chain_collect"
d.mkdir(parents=True, exist_ok=True)
st = df["quote_status"].value_counts().to_dict() if not df.empty else {}
riga = {"ts": int(now.timestamp() * 1000), "righe": int(len(df)),
"ok": int(st.get("ok", 0)), "no_quote": int(st.get("no_quote", 0)),
"error": int(st.get("error", 0)), "chiamate": budget.calls,
"rate_limited": budget.rate_limited}
with (d / "runs.jsonl").open("a") as fh:
fh.write(json.dumps(riga) + "\n")
def write(df: pd.DataFrame, day: datetime) -> Path:
"""Un parquet per giorno: append-friendly e nessun file che cresce senza fine."""
STORE.mkdir(parents=True, exist_ok=True)
p = STORE / f"{day:%Y-%m-%d}.parquet"
if p.exists():
df = pd.concat([pd.read_parquet(p), df], ignore_index=True)
df = df.drop_duplicates(subset=["ts", "instrument_name"], keep="last")
df.to_parquet(p, index=False)
return p
def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--assets", nargs="+", default=list(ASSETS))
ap.add_argument("--rps", type=float, default=DEFAULT_RPS, help="chiamate/secondo (pacing)")
ap.add_argument("--dry-run", action="store_true")
args = ap.parse_args()
now = datetime.now(UTC).replace(microsecond=0)
budget = Budget(rps=args.rps)
t0 = time.monotonic()
frames = []
for a in args.assets:
df = sweep(a, budget, now)
frames.append(df)
if df.empty:
print(f" {a}: NESSUNA riga — il giro e' fallito, non e' un mercato vuoto")
continue
st = df["quote_status"].value_counts().to_dict()
print(f" {a}: {len(df):4d} strumenti | ok {st.get('ok', 0)} | "
f"no_quote {st.get('no_quote', 0)} | error {st.get('error', 0)}")
out = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame(columns=["quote_status"])
dur = time.monotonic() - t0
print(f" {budget.calls} chiamate in {dur:.0f}s ({budget.calls/max(dur,1):.1f}/s) | "
f"429: {budget.rate_limited} | errori: {budget.errors}")
for e in budget.err_samples:
print(f" errore: {e}")
if not args.dry_run:
heartbeat(now, out, budget) # anche a giro fallito: il silenzio non e' un dato
if out.empty:
print(" NIENTE DA SCRIVERE — giro fallito")
return 1
bad = float((out["quote_status"] == "error").mean())
if bad > 0.20:
print(f" ATTENZIONE: {100*bad:.0f}% di quote in ERRORE — la riga c'e' ma il dato no.")
if args.dry_run:
print(" --dry-run: non scrivo")
return 0
p = write(out, now)
print(f" scritto {p} ({len(out)} righe)")
return 0
if __name__ == "__main__":
raise SystemExit(main())