#!/usr/bin/env python3 # ============================================================================= # Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai # ============================================================================= """Step 05 — RQ4: Greeks information decay by DTE + OI-concentration price magnets. Part A (raw-dependent): aggregate Greeks by days-to-expiry bucket and test their predictive content for next-day returns and realized variance. Part B: does the underlying price gravitate toward the maximum-open-interest strike? Built from the raw stores when available; otherwise the shipped ``price_magnet_data.parquet`` reproduces the summary table (the auxiliary LPM regression needs raw-only columns and is skipped in fallback mode). Inputs : options.duckdb + stock_5min.duckdb (parts A and B, optional), data/processed/realized_vol.parquet (part A), data/processed/price_magnet_data.parquet (part B fallback) Outputs: results/rq4_greeks_decay.csv, results/rq4_price_magnet.csv, data/processed/price_magnet_data.parquet (part B, raw mode) """ import warnings import numpy as np import pandas as pd import _bootstrap # noqa: F401 from wp7 import config from wp7.data_io import (RawDataUnavailableError, load_price_magnet, load_realized_vol, open_raw_db) from wp7.econometrics import add_constant, hc1_tstats, ols, r_squared, standardize warnings.filterwarnings('ignore') TICKERS = config.RQ4_TICKERS TICKER_STR = ",".join([f"'{t}'" for t in TICKERS]) # -------------------------------------------------------------------------- # Part A — Greeks information decay by DTE bucket (raw-dependent) # -------------------------------------------------------------------------- def greeks_decay() -> None: print("\n--- PART A: GREEKS INFORMATION DECAY ---") opt_con = open_raw_db("options") print(" Extracting Greeks by DTE buckets...") greeks_by_dte = opt_con.execute(f""" WITH base AS ( SELECT ticker, trade_date, expiry_date, call_put, (expiry_date - trade_date) AS dte, delta, gamma, vega, theta, open_interest, volume, CASE WHEN bid_iv > 0 AND ask_iv > 0 THEN (bid_iv + ask_iv)/2.0 WHEN ask_iv > 0 THEN ask_iv ELSE bid_iv END AS mid_iv FROM option_chain WHERE ticker IN ({TICKER_STR}) AND ask_price > 0 AND (expiry_date - trade_date) BETWEEN 1 AND 180 ) SELECT ticker, trade_date, CASE WHEN dte BETWEEN 1 AND 7 THEN '01_1w' WHEN dte BETWEEN 8 AND 14 THEN '02_2w' WHEN dte BETWEEN 15 AND 30 THEN '03_1m' WHEN dte BETWEEN 31 AND 60 THEN '04_2m' WHEN dte BETWEEN 61 AND 90 THEN '05_3m' WHEN dte BETWEEN 91 AND 180 THEN '06_6m' END AS dte_bucket, AVG(ABS(delta)) AS avg_abs_delta, AVG(gamma) AS avg_gamma, AVG(vega) AS avg_vega, AVG(ABS(theta)) AS avg_abs_theta, AVG(mid_iv) AS avg_iv, SUM(volume) AS total_vol, SUM(open_interest) AS total_oi, SUM(CASE WHEN call_put='c' THEN gamma * open_interest ELSE 0 END) AS call_gamma_oi, SUM(CASE WHEN call_put='p' THEN gamma * open_interest ELSE 0 END) AS put_gamma_oi, SUM(CASE WHEN call_put='c' THEN delta * open_interest ELSE 0 END) AS call_delta_oi, SUM(CASE WHEN call_put='p' THEN delta * open_interest ELSE 0 END) AS put_delta_oi FROM base GROUP BY ticker, trade_date, dte_bucket HAVING dte_bucket IS NOT NULL AND COUNT(*) >= 5 ORDER BY ticker, trade_date, dte_bucket """).fetchdf() opt_con.close() print(f" Greeks by DTE: {len(greeks_by_dte):,} rows") rv = load_realized_vol() rv = rv[rv['ticker'].isin(TICKERS)] greeks_by_dte['trade_date'] = pd.to_datetime(greeks_by_dte['trade_date']) gm = pd.merge(greeks_by_dte, rv[['ticker', 'trade_date', 'ret_1d', 'rv_fwd_1d', 'daily_return']], on=['ticker', 'trade_date'], how='inner') features = ['avg_gamma', 'avg_vega', 'avg_abs_theta', 'avg_iv', 'call_gamma_oi', 'put_gamma_oi'] decay_results = [] for dte_bucket in sorted(gm['dte_bucket'].dropna().unique()): sub = gm[gm['dte_bucket'] == dte_bucket].copy() for target, target_name in [('ret_1d', 'Return'), ('rv_fwd_1d', 'RV')]: data = sub[features + [target]].dropna() if len(data) < 100: continue X = add_constant(standardize(data[features].values)) y = data[target].values coefs = ols(X, y) r2 = r_squared(y, X @ coefs) _, t_stats = hc1_tstats(X, y, coefs, len(features)) decay_results.append({ 'dte_bucket': dte_bucket, 'target': target_name, 'r_squared': r2, 'n_obs': len(data), 'n_significant': np.sum(np.abs(t_stats[1:]) > 1.96), 'gamma_t': t_stats[1], 'vega_t': t_stats[2], 'theta_t': t_stats[3], 'iv_t': t_stats[4], }) decay_df = pd.DataFrame(decay_results) decay_df.to_csv(config.RESULTS_DIR / "rq4_greeks_decay.csv", index=False) print("\n GREEKS PREDICTIVE POWER BY DTE BUCKET:") print(decay_df[['dte_bucket', 'target', 'r_squared', 'n_significant', 'n_obs']] .to_string(index=False)) # -------------------------------------------------------------------------- # Part B — Price magnets # -------------------------------------------------------------------------- def build_magnet_data() -> pd.DataFrame: """Assemble the filtered price-magnet panel from the raw stores.""" opt_con = open_raw_db("options") print(" Extracting OI concentration data...") oi_data = opt_con.execute(f""" WITH daily_oi AS ( SELECT ticker, trade_date, strike, call_put, SUM(open_interest) AS oi, SUM(volume) AS vol FROM option_chain WHERE ticker IN ({TICKER_STR}) AND (expiry_date - trade_date) BETWEEN 1 AND 30 AND open_interest > 0 GROUP BY ticker, trade_date, strike, call_put ), max_oi AS ( SELECT ticker, trade_date, MAX(oi) AS max_oi_val FROM daily_oi GROUP BY ticker, trade_date ), top_strikes AS ( SELECT d.ticker, d.trade_date, d.strike, d.call_put, d.oi, m.max_oi_val, ROW_NUMBER() OVER (PARTITION BY d.ticker, d.trade_date ORDER BY d.oi DESC) AS rn FROM daily_oi d JOIN max_oi m ON d.ticker = m.ticker AND d.trade_date = m.trade_date ) SELECT ticker, trade_date, MAX(CASE WHEN rn = 1 THEN strike END) AS top1_strike, MAX(CASE WHEN rn = 1 THEN oi END) AS top1_oi, MAX(CASE WHEN rn = 2 THEN strike END) AS top2_strike, MAX(CASE WHEN rn = 2 THEN oi END) AS top2_oi, SUM(CASE WHEN call_put='c' THEN oi ELSE 0 END) AS total_call_oi, SUM(CASE WHEN call_put='p' THEN oi ELSE 0 END) AS total_put_oi, SUM(oi) AS total_oi FROM top_strikes WHERE rn <= 5 GROUP BY ticker, trade_date ORDER BY ticker, trade_date """).fetchdf() print(f" OI concentration data: {len(oi_data):,} rows") opt_con.close() stk_con = open_raw_db("stocks_5min") print(" Loading intraday close prices...") intraday_close = stk_con.execute(f""" SELECT symbol AS ticker, CAST(datetime AS DATE) AS trade_date, FIRST(open) AS day_open, LAST(close) AS day_close, MIN(low) AS day_low, MAX(high) AS day_high FROM ohlcv WHERE symbol IN ({TICKER_STR}) AND datetime >= '2010-01-01' GROUP BY symbol, CAST(datetime AS DATE) ORDER BY symbol, trade_date """).fetchdf() stk_con.close() print(f" Intraday aggregated: {len(intraday_close):,} rows") oi_data['trade_date'] = pd.to_datetime(oi_data['trade_date']) intraday_close['trade_date'] = pd.to_datetime(intraday_close['trade_date']) magnet = pd.merge(oi_data, intraday_close, on=['ticker', 'trade_date'], how='inner') magnet['dist_open_to_strike'] = np.abs(magnet['day_open'] - magnet['top1_strike']) / magnet['day_open'] magnet['dist_close_to_strike'] = np.abs(magnet['day_close'] - magnet['top1_strike']) / magnet['day_close'] magnet['moved_toward_strike'] = (magnet['dist_close_to_strike'] < magnet['dist_open_to_strike']).astype(int) magnet['oi_concentration'] = magnet['top1_oi'] / magnet['total_oi'] magnet_clean = magnet[(magnet['dist_open_to_strike'] < 0.10) & (magnet['dist_open_to_strike'] > 0.001)].copy() magnet_clean[['ticker', 'trade_date', 'oi_concentration', 'dist_open_to_strike', 'dist_close_to_strike', 'moved_toward_strike']].to_parquet( config.PRICE_MAGNET_PARQUET, index=False) return magnet_clean def magnet_analysis(magnet_clean: pd.DataFrame, has_raw_columns: bool) -> None: """Summary by OI-concentration quintile (+ LPM regression in raw mode).""" print(f"\n Filtered magnet data: {len(magnet_clean):,} rows") print(f" Fraction moved toward max-OI strike: " f"{magnet_clean['moved_toward_strike'].mean():.4f}") magnet_clean['oi_conc_quintile'] = pd.qcut( magnet_clean['oi_concentration'], 5, labels=['Q1_Low', 'Q2', 'Q3', 'Q4', 'Q5_High'], duplicates='drop') magnet_summary = magnet_clean.groupby('oi_conc_quintile', observed=True).agg( pct_moved_toward=('moved_toward_strike', 'mean'), avg_dist_open=('dist_open_to_strike', 'mean'), avg_dist_close=('dist_close_to_strike', 'mean'), n_obs=('moved_toward_strike', 'count'), ).round(4) print("\n PRICE MAGNET EFFECT BY OI CONCENTRATION:") print(magnet_summary.to_string()) magnet_summary.to_csv(config.RESULTS_DIR / "rq4_price_magnet.csv") if not has_raw_columns: print("\n [LPM regression skipped — needs 'total_oi' from the raw build]") return features = ['oi_concentration', 'dist_open_to_strike', 'total_oi'] sub = magnet_clean[features + ['moved_toward_strike']].dropna() X = add_constant(standardize(sub[features].values)) y = sub['moved_toward_strike'].values coefs = ols(X, y) _, t = hc1_tstats(X, y, coefs, len(features)) print("\n LOGISTIC (LPM) REGRESSION: moved_toward_strike ~") for i, name in enumerate(['const'] + features): sig = "**" if abs(t[i]) > 2.576 else ("*" if abs(t[i]) > 1.96 else "") print(f" {name:25s}: β={coefs[i]:8.5f}, t={t[i]:7.3f} {sig}") def main(): print("=" * 70) print("RQ4: GREEKS INFORMATION DECAY & PRICE MAGNETS") print("=" * 70) config.ensure_output_dirs() try: greeks_decay() except RawDataUnavailableError as exc: print(f"\n[Part A skipped — raw stores unavailable]\n{exc}") print("\n--- PART B: OI CONCENTRATION AS PRICE MAGNETS ---") try: magnet_clean = build_magnet_data() magnet_analysis(magnet_clean, has_raw_columns=True) except RawDataUnavailableError as exc: print(f"\n[Raw build skipped — reproducing summary from shipped parquet]\n{exc}\n") magnet_analysis(load_price_magnet(), has_raw_columns=False) print("\nRQ4 COMPLETE.") if __name__ == "__main__": main()