Files
PythagorasGoal/scripts/research/fetch_ib_equities.py
T
Adriano Dal Pastro 2beb11b764 fetch IB: client in readonly — via due righe d'errore da 63 notti su 63
`fetch_ib_equities.py` si connetteva senza `readonly`, e ogni notte il log del
cron si prendeva

    open orders request timed out
    completed orders request timed out

126 righe in 63 giri. Due errori innocui ripetuti per sempre sono il modo in cui
un errore VERO smette di farsi notare (P14): su `cron_daily.log`, 126 dei 146
match di "error" erano questi.

CAUSA, letta nel sorgente di ib_async e non indovinata: in `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.
Questo client scarica storico e non manda ordini mai: `readonly=True` e' insieme
la cura del rumore e la dichiarazione corretta di cosa fa. Tolta la CAUSA, non
filtrato il messaggio — filtrarlo avrebbe nascosto anche il giorno in cui quel
timeout significasse qualcosa.

VERIFICATO con un A/B sul solo flag, contro il gateway vero:
  · readonly=False -> le due righe compaiono, e si apre sul gateway il dialogo
    modale "API client needs write access action confirmation" (visto nei log del
    container, resta su ~75s);
  · readonly=True  -> nessuna delle due righe, nessun dialogo.

⚠️ CIO' CHE NON E' STATO VERIFICATO, e va detto: in nessuna delle quattro prove
fra le 13:05 e le 15:35 UTC il gateway ha servito storico — 0 barre con ENTRAMBI
i flag, quindi la causa non e' questa modifica, ma non ho potuto confermare
end-to-end che il fetch continui a riportare barre. La conferma e' il log del
cron di stanotte: se SPY/QQQ/IWM/TLT/GLD/HYG tornano con le loro barre e senza
le due righe di timeout, e' a posto; se tornano tutti a 0, si revoca il flag.
Il fallimento e' comunque innocuo: con 0 barre lo script NON sovrascrive i
parquet (verificato: eq_spy/eq_qqq intatti dopo i tentativi falliti).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-28 13:31:22 +00:00

224 lines
13 KiB
Python

"""FETCH + CERTIFY universo azioni/ETF da IB (ADJUSTED_LAST) -> data/raw/eq_<sym>_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_<sym>_1d.parquet (ADJUSTED_LAST, namespace dedicato).")
ib.disconnect()
if __name__ == "__main__":
main()