"""Options-flow direction factors — pre-reg studies/2026-09-09-flow-prereg.md.""" import gzip, json, glob, os, numpy as np, pandas as pd, warnings from scipy import stats from sklearn.linear_model import LogisticRegression from sklearn.metrics import roc_auc_score warnings.filterwarnings("ignore") src=open("crown_policies.py").read().split("\ndf=run()")[0]; cp={}; exec(src, cp) FP=cp["FP"]; d=cp["load_nq"](); QUAD=cp["QUAD"] def load_day(p): j=json.load(gzip.open(p)); t=pd.DataFrame(j["qqq_ticks"]) if t.empty: return None t["et"]=pd.to_datetime(t.tape_time,utc=True).dt.tz_convert("America/New_York"); t["m"]=t.et.dt.hour*60+t.et.dt.minute t["np"]=pd.to_numeric(t.net_call_premium,errors="coerce").fillna(0)-pd.to_numeric(t.net_put_premium,errors="coerce").fillna(0) r=dict(date=pd.Timestamp(j["date"]).date(), q_1010=t[t.m<=610].np.sum(), q_day=t.np.sum(), q_last30=t[t.m>=945].np.sum()) td=pd.DataFrame(j["tide"]) if len(td): td["et"]=pd.to_datetime(td.timestamp).dt.tz_convert("America/New_York") if pd.to_datetime(td.timestamp).dt.tz is not None else pd.to_datetime(td.timestamp); td["m"]=td.et.dt.hour*60+td.et.dt.minute td["np"]=pd.to_numeric(td.net_call_premium,errors="coerce")-pd.to_numeric(td.net_put_premium,errors="coerce") a=td[td.m<=610]; r["t_1010"]=float(a.np.iloc[-1]) if len(a) else np.nan; r["t_day"]=float(td.np.iloc[-1]) else: r["t_1010"]=np.nan; r["t_day"]=np.nan return r FL=pd.DataFrame([x for x in (load_day(p) for p in sorted(glob.glob(f"{FP}/flow/*.json.gz"))) if x]).set_index("date").sort_index() print("flow days",len(FL),FL.index.min(),"→",FL.index.max()) # pre-open factors come from the PRIOR session F=pd.DataFrame(index=FL.index); F["F1"]=FL.q_1010; F["F2"]=FL.t_1010; F["F3"]=FL.q_day.shift(1); F["F4"]=FL.t_day.shift(1); F["F5"]=FL.q_last30.shift(1) for c in ["F1","F2","F3","F4","F5"]: F[c+"_z"]=F[c]/F[c].abs().rolling(60,min_periods=30).median().shift(1) # NQ session facts S={} for sday,g in d.groupby("sday"): rth=g[(g["min"]>=570)&(g["min"]<960)].reset_index(drop=True) if len(rth)<300 or sday in QUAD: continue b1030=rth[rth["min"]<=629]; b1010=rth[rth["min"]<=610] S[sday]=dict(rth=rth,o=rth.iloc[0]["open"],c=rth.iloc[-1]["close"],c1030=b1030.iloc[-1]["close"],c1010=b1010.iloc[-1]["close"]) rows=[] for sday,r in F.iterrows(): if sday not in S: continue s=S[sday]; rows.append(dict(date=sday,**{k:r[k] for k in F.columns},up30=float(s["c1030"]>s["o"]),up_close=float(s["c"]>s["o"]),up_1010_close=float(s["c"]>s["c1010"]),ret30=(s["c1030"]/s["o"]-1)*100,ret_close=(s["c"]/s["o"]-1)*100,ret_1010=(s["c"]/s["c1010"]-1)*100)) D=pd.DataFrame(rows).set_index("date"); print("sessions",len(D)) L=[f"# Options-flow direction factors — results · {pd.Timestamp.now():%Y-%m-%d %H:%M}\n",f"Pre-registration: `2026-09-09-flow-prereg.md`. Sessions {len(D)} ({D.index.min()} → {D.index.max()}), quad-witching excluded. Base rates: up30 {D.up30.mean()*100:.1f}%, up_close {D.up_close.mean()*100:.1f}%, up 10:10→close {D.up_1010_close.mean()*100:.1f}%.\n","## Descriptive — by factor tercile (z-scored by rolling 60-session median |x|)"] def wf_auc(feat,target): dd=D.dropna(subset=[feat,target]); yr=pd.Series(dd.index).apply(lambda x:x.year).values; ps=[];ys=[] for Y in (2025,2026): tr=dd[yr1 else np.nan for fac,targets in (("F1_z",[("up_1010_close","ret_1010")]),("F2_z",[("up_1010_close","ret_1010")]),("F3_z",[("up30","ret30"),("up_close","ret_close")]),("F4_z",[("up30","ret30"),("up_close","ret_close")]),("F5_z",[("up30","ret30"),("up_close","ret_close")])): dd=D.dropna(subset=[fac]); q=pd.qcut(dd[fac].rank(method="first"),3,labels=["neg","mid","pos"]) for up,ret in targets: g=dd.groupby(q); L.append(f"- **{fac[:-2]}** → {up}: "+" · ".join(f"{k} {g[up].mean()[k]*100:.1f}% ({g[ret].mean()[k]:+.3f}%)" for k in ["neg","mid","pos"])+f"; WF AUC {wf_auc(fac,up):.3f}; n={len(dd)}") L.append("\n## Policies (Part A fill rules)") mid=pd.Timestamp("2025-03-15").date(); res=[] for sday,r in D.iterrows(): s=S[sday]; rth=s["rth"]; f5=cp["five_min"](rth) for pol,fac,sig in (("FL1_qqq_1010","F1_z",611),("FL2_tide_1010","F2_z",611),("FL3_prior_day","F3_z",575)): z=r[fac]; side=int(np.sign(z)) if (np.isfinite(z) and abs(z)>=1) else 0 i5=int(np.searchsorted(f5.endmin.values,sig-1)); i5=min(i5,len(f5)-1) w=cp["walk"](rth,sig,side,f5,i5) if side else None res.append(dict(date=sday,pol=pol,side=side,z=z,trig=w is not None,**(w or {}))) R=pd.DataFrame(res); R.to_parquet("data/flow_study.parquet"); D.to_parquet("data/flow_features.parquet") for pol,g in R.groupby("pol"): t=g[g.trig] if len(t)<20: L.append(f"- **{pol}**: triggers {len(t)} — n too small."); continue h1=t[t.date=mid]; p=stats.ttest_1samp(t.r,0,alternative="greater").pvalue; ok=t.r.mean()>0 and p<0.0167 and h1.r.mean()>0 and h2.r.mean()>0 and len(t)>=100 L.append(f"- **{pol}**: triggers {len(t)} (long {(t.side>0).sum()}/short {(t.side<0).sum()}), W/L/T {(t.res=='win').sum()}/{(t.res=='loss').sum()}/{(t.res=='timeout').sum()}, R/trigger **{t.r.mean():+.3f}** (p={p:.3f}), halves {h1.r.mean():+.3f} (n={len(h1)}) / {h2.r.mean():+.3f} (n={len(h2)}), long {t[t.side>0].r.mean():+.3f} / short {t[t.side<0].r.mean():+.3f}, $ at 1 NQ {t.pts.sum()*20:+,.0f} → **{'VALIDATED' if ok else 'REFUTED'}**.") out="\n".join(L); open(f"{FP}/studies/2026-09-09-flow-results.md","w").write(out); print(out)