"""FETCH + CERTIFY universo azioni/ETF da IB (ADJUSTED_LAST) -> data/raw/eq__1d.parquet. Apre il fronte EQUITY (branch research/equities-ib). Disciplina v2.0.0: PRIMA il dato certificato, POI la strategia. IB dà storia daily aggiustata per dividendi+split (ADJUSTED_LAST), profonda (SPY dal 1996), sul conto paper. Namespace dedicato 'eq_' (NON tocca i parquet crypto). UNIVERSO (prima ricerca = momentum cross-sectional settoriale, l'edge robusto plausibile in equity): * 11 SPDR settoriali (XLK..XLC); * broad/macro SPY QQQ IWM TLT GLD HYG. NB: i 9 settori "classici" partono 1998; XLRE 2015, XLC 2018 -> lo start COMUNE a 11 e' 2018. Per backtest lunghi usare i 9 classici (1998+) o accettare lo start 2018 per gli 11. CERTIFICAZIONE (gemello equity di certify_feed.py): (1) integrità: barre, range, date monotone, duplicati, flat bars (close invariato); (2) gap: run di giorni-lavorativi mancanti > 5 (festivi normali, buchi lunghi = sospetti); (3) sanità ritorni: max |daily ret| (un >50% non-evento = errore di adjustment); (4) sanità adjustment: primo close aggiustato << ultimo (i dividendi abbassano lo storico); (5) SPLIT NON AGGIUSTATI (aggiunto 2026-07-25): ADJUSTED_LAST non sempre aggiusta gli split, e il check (3) NON li vede — uno split 2:1 fa esattamente -50%, sul filo della soglia. IWM (gamba di GTAA01 in produzione) ed EFA avevano uno split non aggiustato il 2005-06-09 e passavano come OK. Rilevatore in src/data/eq_splits.py -> status SPLIT-NON-AGG. PREREQUISITO: gateway IB paper su 127.0.0.1:4002 (docker compose up -d ib-gateway). uv run --with ib_async python scripts/research/fetch_ib_equities.py """ import sys, time from pathlib import Path import numpy as np, pandas as pd ROOT = Path(__file__).resolve().parents[2] sys.path.insert(0, str(ROOT)) from src.data.eq_splits import detect_unadjusted_splits # noqa: E402 RAW = ROOT / "data" / "raw" RAW.mkdir(parents=True, exist_ok=True) SECTORS = ["XLK", "XLF", "XLE", "XLV", "XLI", "XLP", "XLY", "XLU", "XLB", "XLRE", "XLC"] BROAD = ["SPY", "QQQ", "IWM", "TLT", "GLD", "HYG"] # espansione "diversi mercati" (intl / bond / credito / commodity / settori extra) per il lead-lag crypto BROAD2 = ["DIA", "EFA", "EEM", "FXI", "EWJ", "AGG", "LQD", "IEF", "USO", "SLV", "DBC", "VNQ"] UNIVERSE = SECTORS + BROAD + BROAD2 # Prima quotazione degli strumenti che contano (le sei gambe di GTAA01). RIFERIMENTO DICHIARATO, # fonte secondaria: serve solo a distinguere "questo ETF e' giovane" da "il feed ci da' meno storia # di quella che esiste". Trovato il 2026-08-07: TLT parte dal 2016-02-03 invece che dal 2002 — IB su # questo conto NON serve barre precedenti (richiesta retro esplicita -> 0 barre), quindi non e' un # fetch da rifare ma un limite da dichiarare. Nessuna certificazione se n'era accorta perche' tutte # guardavano DENTRO la serie (integrita', gap, spike, split) e nessuna la sua LUNGHEZZA. PRIMA_QUOTAZIONE = {"SPY": "1993-01-22", "QQQ": "1999-03-10", "IWM": "2000-05-22", "TLT": "2002-07-22", "GLD": "2004-11-18", "HYG": "2007-04-04"} DURATION = "30 Y" # tetto della richiesta: una serie che parte QUI non e' troncata, e' al cap def certify(sym: str, df: pd.DataFrame, prev: pd.DataFrame | None = None) -> dict: """`prev` = la versione GIA' SU DISCO. Senza, la certificazione non puo' vedere cio' che la serie ha PERSO: tutti i controlli esistenti guardano dentro la serie che hanno in mano.""" if df.empty: return {"sym": sym, "n": 0, "status": "VUOTO"} idx = df.index dup = int(idx.duplicated().sum()) mono = bool(idx.is_monotonic_increasing) c = df["close"].values.astype(float) ret = np.diff(c) / c[:-1] flat = int((ret == 0).sum()) maxret = float(np.max(np.abs(ret))) if len(ret) else 0.0 # gap: giorni lavorativi attesi vs presenti, run lunghi mancanti bdays = pd.bdate_range(idx[0], idx[-1]) missing = len(bdays) - len(idx.intersection(bdays)) gaps = bdays.difference(idx) longgap = 0 if len(gaps): g = pd.Series(1, index=gaps).resample("1D").sum().fillna(0) # conta run consecutivi di bday mancanti s = (gaps.to_series().diff().dt.days.fillna(1) > 3).cumsum() longgap = int((gaps.to_series().groupby(s).size() > 5).sum()) span_y = (idx[-1] - idx[0]).days / 365.25 adj_ratio = round(float(c[0] / c[-1]), 3) # primo/ultimo: <1 atteso (storico abbassato dai div) # SPLIT NON AGGIUSTATI: IB ADJUSTED_LAST non sempre li aggiusta. NB la guardia `maxret > 0.5` # NON li vede: uno split 2:1 non aggiustato fa ESATTAMENTE -50% e cade sul filo della soglia # (IWM 2005-06-09 passava a 49.5% con status OK, ed e' una gamba di GTAA01 in produzione). # Vedi src/data/eq_splits.py per il discriminante split-vs-crollo (range intraday). splits = detect_unadjusted_splits(df) # --- storia PERSA: due guardie, perche' i due difetti non si vedono nello stesso modo --- # (a) REGRESSIONE rispetto al disco: la serie nuova parte dopo la vecchia, o ha meno barre. # Prende una troncatura il giorno in cui compare — ma non una che c'era gia'. primo = idx[0].tz_localize(None) if idx[0].tzinfo else idx[0] ultimo = idx[-1].tz_localize(None) if idx[-1].tzinfo else idx[-1] perse = 0 if prev is not None and len(prev): vecchio = prev.index[0] vecchio = vecchio.tz_localize(None) if vecchio.tzinfo else vecchio perse = max(0, (primo - vecchio).days) tronca = perse > 10 or (prev is not None and len(prev) > 0 and len(df) < len(prev) - 5) # (b) DISTANZA DALLA QUOTAZIONE: prende anche una troncatura presente da sempre. Non e' un # difetto se la serie parte al tetto della richiesta (30 anni indietro): quello e' il cap. corta = 0.0 if sym in PRIMA_QUOTAZIONE: atteso = max(pd.Timestamp(PRIMA_QUOTAZIONE[sym]), ultimo - pd.DateOffset(years=int(DURATION.split()[0]))) corta = max(0, (primo - atteso).days) / 365.25 status = "OK" if dup or not mono: status = "INTEGRITA'" elif tronca: status = "TRONCATO" # il feed ha PERSO storia rispetto a ieri elif splits: status = "SPLIT-NON-AGG" elif maxret > 0.5: status = "SPIKE?" elif longgap > 0: status = "GAP-LUNGO" elif corta >= 1.0: status = "STORIA-CORTA" # meno storia di quella che lo strumento HA elif span_y < 1: status = "corto<1y" return {"sym": sym, "n": len(df), "primo": idx[0].date(), "ultimo": idx[-1].date(), "anni": round(span_y, 1), "dup": dup, "mono": mono, "flat": flat, "maxret%": round(maxret * 100, 1), "miss_bd": missing, "gap_lunghi": longgap, "adj_first/last": adj_ratio, "persi_g": perse, "manca_a": round(corta, 1), "status": status, "splits": [f"{s['date'].date()} 1:{s['factor']:g}" for s in splits]} def main(): try: from ib_async import IB, Stock except Exception: print("ib_async assente. Esegui con: uv run --with ib_async python scripts/research/fetch_ib_equities.py") sys.exit(2) ib = IB() try: # `readonly=True`: questo client SCARICA STORICO e non manda ordini mai. Non e' solo # igiene — in `ib_async.IB.connectAsync` le richieste "open orders" e "completed orders" # esistono SOLO se il client non e' readonly (`if not readonly: reqs[...]`), e sul gateway # paper non rispondono: ogni notte, 63 giri su 63, il log del cron si prendeva # open orders request timed out # completed orders request timed out # Due righe d'errore innocue ripetute per sempre sono il modo in cui un errore VERO # smette di farsi notare (P14). Qui si toglie la CAUSA, non si filtra il messaggio. # ⚠️ OSSERVATO il 2026-08-28: a meta' pomeriggio il gateway NON serve storico — # `reqHistoricalData` va in timeout e lo script stampa "0 barre (subscription?)" per ogni # simbolo. Quattro tentativi fra le 13:05 e le 15:35 UTC, tutti a zero, con e senza # `readonly`; il giro del cron delle 00:30 riesce ogni notte. La causa sta nel gateway, # non qui, e non e' stata diagnosticata. Non e' pericoloso: con 0 barre lo script NON # sovrascrive i parquet, quindi un giro fallito lascia il dato di ieri invece di romperlo. ib.connect("127.0.0.1", 4002, clientId=90, timeout=15, readonly=True) except Exception as e: print(f"[CONNESSIONE FALLITA] 127.0.0.1:4002 -> {repr(e)[:120]}\n Avvia: docker compose up -d ib-gateway") sys.exit(1) print("=" * 104) print(f" FETCH + CERTIFY azioni/ETF (ADJUSTED_LAST) -> data/raw/eq_* | acct {ib.managedAccounts()}") print("=" * 104) rep, ok = [], [] force = "--force" in sys.argv[1:] universe = UNIVERSE if "--only" in sys.argv[1:]: # refresh mirato (es. solo i 6 ETF GTAA per il cron) universe = sys.argv[sys.argv.index("--only") + 1].upper().split(",") force = True # --only implica refresh dei simboli indicati for sym in universe: out_path = RAW / f"eq_{sym.lower()}_1d.parquet" if out_path.exists() and not force: print(f" {sym:5} GIA' SU DISCO -> skip (usa --force per riscaricare)") ok.append(sym) continue con = Stock(sym, "SMART", "USD") try: bars = ib.reqHistoricalData(con, endDateTime="", durationStr=DURATION, barSizeSetting="1 day", whatToShow="ADJUSTED_LAST", useRTH=True, formatDate=1, timeout=60) except Exception as e: print(f" {sym:5} ERR {repr(e)[:70]}"); rep.append({"sym": sym, "status": "ERR"}); time.sleep(1.2); continue if not bars: print(f" {sym:5} 0 barre (subscription?)"); rep.append({"sym": sym, "n": 0, "status": "VUOTO"}); time.sleep(1.2); continue df = pd.DataFrame([(pd.Timestamp(str(b.date)), b.open, b.high, b.low, b.close, b.volume) for b in bars], columns=["ts", "open", "high", "low", "close", "volume"]).set_index("ts").sort_index() prev = None if out_path.exists(): try: pv = pd.read_parquet(out_path) pv.index = pd.to_datetime(pv["timestamp"], unit="ms") prev = pv.sort_index() except Exception: prev = None c = certify(sym, df, prev) rep.append(c) if c["status"] == "TRONCATO": # NON si sovrascrive e NON si fonde: ADJUSTED_LAST e' ri-aggiustato all'indietro a ogni # dividendo, quindi incollare una vintage vecchia a una nuova crea un salto sul giunto — # un difetto peggiore di quello che si voleva evitare. Si tiene il file coerente e si urla. print(f" {sym:5} ⚠️ TRONCATO: {c['primo']} contro {prev.index[0].date()} su disco " f"({c['persi_g']}g persi, {c['n']} barre contro {len(prev)}) -> FILE NON TOCCATO") time.sleep(1.2) continue if c.get("n", 0) > 0: out = df.copy() # ms epoch (come i parquet crypto), robusto alla risoluzione datetime64 (s/us/ns) out["timestamp"] = out.index.astype("datetime64[ms]").astype("int64") out.reset_index(drop=True).to_parquet(RAW / f"eq_{sym.lower()}_1d.parquet") if c["status"] == "OK": ok.append(sym) print(f" {sym:5} n={c.get('n',0):>5} {str(c.get('primo','')):>10}->{str(c.get('ultimo',''))} " f"{c.get('anni','?')}y flat={c.get('flat','?')} maxret={c.get('maxret%','?')}% " f"miss_bd={c.get('miss_bd','?')} gapL={c.get('gap_lunghi','?')} adj={c.get('adj_first/last','?')} [{c['status']}]" + (f" split={c['splits']}" if c.get("splits") else "") + (f" manca_a_quotazione={c['manca_a']}a" if c.get("manca_a", 0) >= 1.0 else "")) time.sleep(1.2) # pacing IB print("-" * 104) guasti = [r for r in rep if r.get("status") in ("TRONCATO", "STORIA-CORTA")] if guasti: print(" ⚠️ STORIA MANCANTE (le certificazioni locali non la vedono: guardano DENTRO la serie):") for r in guasti: print(f" {r['sym']:5} [{r['status']}] parte dal {r['primo']} — " f"{r.get('manca_a', 0)}a dopo la quotazione, {r.get('persi_g', 0)}g persi vs disco") print(f" CERTIFICATI OK ({len(ok)}/{len(UNIVERSE)}): {ok}") sec_ok = [s for s in SECTORS if s in ok] print(f" settori OK: {len(sec_ok)}/11 {sec_ok}") print(f" -> scritti in data/raw/eq__1d.parquet (ADJUSTED_LAST, namespace dedicato).") ib.disconnect() if __name__ == "__main__": main()