#!/usr/bin/env python3 """Bucket public-SNAT Google-family fail windows (H84 ≤15m vs H80 ≥2h). Reads ClickHouse flows_raw (nat_src × Google DIP), builds 5-minute ACK% series, finds contiguous fail runs, labels duration buckets. Candidate metrics only — does NOT claim TSPU attribution. python3 scripts/bucket_snat_fail_windows.py \\ --nats 203.0.113.55,203.0.113.22,203.0.113.30 --hours 12 python3 scripts/bucket_snat_fail_windows.py \\ --nats-file docs/evidence/reblock_risk_14snat_20260914/nats.txt --hours 18 # Offline: matrix/probe TSV with columns ts,src,result (OK_TCP|FAIL|…) python3 scripts/bucket_snat_fail_windows.py --from-tsv path/to/probes.tsv See: docs/evidence/tspu_forum_synthesis_20260914.md """ from __future__ import annotations import argparse import csv import json import sys from collections import defaultdict from datetime import datetime, timezone from pathlib import Path ROOT = Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT / "scripts")) from alert_thresholds import ch_query, load_dotenv # noqa: E402 # Duration buckets aligned with forum digs (H84 ~10 min; H80 multi-hour) H84_MAX_MIN = 15 H80_MIN_MIN = 120 def _parse_nats(raw: str) -> list[str]: parts = [] for p in raw.replace("\n", ",").split(","): p = p.strip() if p and not p.startswith("#"): parts.append(p) return parts def _label(duration_min: float) -> str: if duration_min <= H84_MAX_MIN: return "h84_like" if duration_min >= H80_MIN_MIN: return "h80_like" return "mid" def _runs_from_flags( series: list[tuple[datetime, bool]], bucket_min: int ) -> list[dict]: """series: sorted (bucket_start, is_fail). Gaps reset a run.""" runs: list[dict] = [] i = 0 while i < len(series): t0, fail = series[i] if not fail: i += 1 continue j = i while j + 1 < len(series): t_next, f_next = series[j + 1] gap = (t_next - series[j][0]).total_seconds() / 60.0 if not f_next or gap > bucket_min * 1.5: break j += 1 t1 = series[j][0] # Inclusive span of fail buckets duration_min = (t1 - t0).total_seconds() / 60.0 + bucket_min runs.append( { "start_utc": t0.strftime("%Y-%m-%dT%H:%M:%SZ"), "end_utc": t1.strftime("%Y-%m-%dT%H:%M:%SZ"), "duration_min": round(duration_min, 1), "fail_buckets": j - i + 1, "bucket": _label(duration_min), } ) i = j + 1 return runs def fetch_nat_series( nats: list[str], hours: int, bucket_min: int, ack_fail_pct: float, min_flows: int ) -> dict[str, list[tuple[datetime, bool, dict]]]: """Return nat → list of (t, is_fail, metrics). Prefer nat_src path (CGNAT).""" in_list = ", ".join(f"'{n}'" for n in nats) sql = f""" SELECT toStartOfInterval(time_received, INTERVAL {int(bucket_min)} MINUTE) AS t, replaceAll(IPv6NumToString(nat_src_addr), '::ffff:', '') AS nat_ip, count() AS flows, countIf(has_ack = 1) AS ack, countIf(has_syn_no_ack = 1) AS syn_na, round(100. * countIf(has_ack = 1) / nullIf(count(), 0), 1) AS ack_pct FROM netflow.flows_raw WHERE time_received >= now() - INTERVAL {int(hours)} HOUR AND proto = 6 AND has_nat_src = 1 AND nat_src_addr IS NOT NULL AND nat_src_addr != src_addr AND replaceAll(IPv6NumToString(nat_src_addr), '::ffff:', '') IN ({in_list}) AND ( IPv6NumToString(dst_addr) LIKE '::ffff:142.250.%' OR IPv6NumToString(dst_addr) LIKE '::ffff:142.251.%' OR IPv6NumToString(dst_addr) LIKE '::ffff:216.58.%' ) GROUP BY t, nat_ip ORDER BY nat_ip, t """ rows = ch_query(sql, timeout=180) by_nat: dict[str, list[tuple[datetime, bool, dict]]] = defaultdict(list) for r in rows: t_raw = r["t"] if isinstance(t_raw, str): # ClickHouse JSONEachRow: '2026-09-14 12:00:00' or ISO t_raw = t_raw.replace("Z", "").replace("T", " ") t = datetime.strptime(t_raw[:19], "%Y-%m-%d %H:%M:%S").replace( tzinfo=timezone.utc ) else: t = t_raw flows = int(r.get("flows") or 0) ack_pct = float(r.get("ack_pct") or 0) is_fail = flows >= min_flows and ack_pct <= ack_fail_pct by_nat[r["nat_ip"]].append( ( t, is_fail, { "flows": flows, "ack": int(r.get("ack") or 0), "syn_na": int(r.get("syn_na") or 0), "ack_pct": ack_pct, }, ) ) return by_nat def runs_from_tsv(path: Path, bucket_min: int) -> dict[str, list[dict]]: """Offline probe TSV: ts[,src|nat],result — FAIL* = fail, OK* = ok.""" # Flexible headers with path.open(newline="") as f: sample = f.read(4096) f.seek(0) dialect = csv.Sniffer().sniff(sample, delimiters="\t,") reader = csv.DictReader(f, dialect=dialect) fieldmap = {k.lower(): k for k in (reader.fieldnames or [])} def col(*names: str) -> str | None: for n in names: if n in fieldmap: return fieldmap[n] return None c_ts = col("ts", "ts_utc", "time", "t") c_src = col("src", "src_ip", "nat", "nat_ip", "egress") c_res = col("result", "status", "class", "outcome") if not c_ts or not c_src or not c_res: raise SystemExit( f"TSV needs ts/src/result columns; got {list(fieldmap)}" ) # Per src, 5m buckets: any FAIL in bucket → fail buckets: dict[str, dict[datetime, list[bool]]] = defaultdict( lambda: defaultdict(list) ) for row in reader: ts_s = (row.get(c_ts) or "").strip() src = (row.get(c_src) or "").strip() res = (row.get(c_res) or "").strip().upper() if not ts_s or not src: continue ts_s = ts_s.replace("Z", "+00:00") try: t = datetime.fromisoformat(ts_s) except ValueError: try: t = datetime.strptime(ts_s[:19], "%Y-%m-%d %H:%M:%S").replace( tzinfo=timezone.utc ) except ValueError: continue if t.tzinfo is None: t = t.replace(tzinfo=timezone.utc) # Floor to bucket epoch = int(t.timestamp()) floor = epoch - (epoch % (bucket_min * 60)) bt = datetime.fromtimestamp(floor, tz=timezone.utc) is_fail = res.startswith("FAIL") or res in ( "WRAP_FAIL", "TIMEOUT", "CONNECTING", ) is_ok = res.startswith("OK") or res in ("FINISHED", "DONE") if is_fail: buckets[src][bt].append(True) elif is_ok: buckets[src][bt].append(False) out: dict[str, list[dict]] = {} for src, bmap in buckets.items(): series = [] for bt in sorted(bmap): flags = bmap[bt] # Fail if any fail and no ok? Prefer majority fail, or any fail if mixed probe noise is_fail = any(flags) and (sum(flags) >= max(1, len(flags) // 2)) series.append((bt, is_fail)) out[src] = _runs_from_flags(series, bucket_min) return out def summarize(runs_by_nat: dict[str, list[dict]]) -> dict: totals = {"h84_like": 0, "mid": 0, "h80_like": 0, "nats": 0, "runs": 0} for nat, runs in runs_by_nat.items(): totals["nats"] += 1 for r in runs: totals["runs"] += 1 totals[r["bucket"]] = totals.get(r["bucket"], 0) + 1 return totals def main() -> int: load_dotenv(ROOT / ".env") ap = argparse.ArgumentParser(description=__doc__) ap.add_argument( "--nats", default="", help="comma-separated public SNAT list", ) ap.add_argument("--nats-file", type=Path, help="one IP per line") ap.add_argument("--hours", type=int, default=12, help="lookback hours (TTL often 1d)") ap.add_argument("--bucket-min", type=int, default=5) ap.add_argument("--ack-fail-pct", type=float, default=25.0) ap.add_argument("--min-flows", type=int, default=3) ap.add_argument( "--from-tsv", type=Path, help="offline probe TSV (ts,src,result) instead of ClickHouse", ) ap.add_argument("--json", action="store_true", help="machine-readable stdout") args = ap.parse_args() if args.from_tsv: runs_by_nat = runs_from_tsv(args.from_tsv, args.bucket_min) series_note = f"offline TSV {args.from_tsv}" else: nats = _parse_nats(args.nats) if args.nats_file and args.nats_file.is_file(): nats.extend(_parse_nats(args.nats_file.read_text())) # de-dupe preserve order seen: set[str] = set() nats = [n for n in nats if not (n in seen or seen.add(n))] if not nats: ap.error("provide --nats / --nats-file or --from-tsv") by_nat = fetch_nat_series( nats, args.hours, args.bucket_min, args.ack_fail_pct, args.min_flows ) runs_by_nat = {} for nat in nats: series = [(t, f) for t, f, _m in by_nat.get(nat, [])] runs_by_nat[nat] = _runs_from_flags(series, args.bucket_min) series_note = ( f"CH Google DIP fwd nat_src, hours={args.hours}, " f"fail if flows>={args.min_flows} and ack_pct<={args.ack_fail_pct}" ) totals = summarize(runs_by_nat) payload = { "hypothesis_test": "H84_vs_H80_duration", "disclaimer": "candidate buckets only — not TSPU attribution", "source": series_note, "thresholds_min": {"h84_like_max": H84_MAX_MIN, "h80_like_min": H80_MIN_MIN}, "totals": totals, "by_nat": runs_by_nat, } if args.json: print(json.dumps(payload, ensure_ascii=False, indent=2)) return 0 print(f"# SNAT fail-window buckets ({series_note})") print( f"# labels: h84_like≤{H84_MAX_MIN}m | mid | h80_like≥{H80_MIN_MIN}m " f"— NOT a TSPU verdict\n" ) for nat, runs in runs_by_nat.items(): if not runs: print(f"{nat}\t(no fail runs in window)") continue print(f"## {nat} ({len(runs)} run(s))") for r in runs: print( f" {r['bucket']:10} {r['duration_min']:7.1f}m " f"{r['start_utc']} → {r['end_utc']} " f"buckets={r['fail_buckets']}" ) print() print( f"TOTALS nats={totals['nats']} runs={totals['runs']} " f"h84_like={totals['h84_like']} mid={totals['mid']} " f"h80_like={totals['h80_like']}" ) print( "\nInterpret: many h84_like → cooldown-like (forum H84); " "dominant h80_like → sticky/reblock (H80). Empty on control SNAT = good." ) return 0 if __name__ == "__main__": raise SystemExit(main())