From 2c322419b69049a6ee5d59e4b59d4277ce86270e Mon Sep 17 00:00:00 2001 From: jaxperro Date: Wed, 22 Jul 2026 18:30:54 -0400 Subject: [PATCH] =?UTF-8?q?research:=20surgebot=20A2=20=E2=80=94=20measure?= =?UTF-8?q?ment=20arm=20relaunch=20(every-trigger=20$100=20FAKs=20+=20atte?= =?UTF-8?q?mpts=20stream=20+=20offline=20virtual-book=20replay)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit v1's cash-gated book halted at its pre-registered -50% line; post-mortem showed the ~2% cash-gated subsample was adversely selected (-$9/fill vs +$41/fill full-signal, same day). A2 samples every trigger and replays bankroll specs offline (surge_book_replay.py -> surge_book.json). Signal semantics verbatim; grades to surge_meas_ledger.jsonl. See #19. Co-Authored-By: Claude Fable 5 --- research/grade_surge.py | 63 +++-- research/nightly.sh | 6 +- research/surge_book_replay.py | 152 ++++++++++++ research/surgebot.py | 451 +++++++++++++++++++++++----------- 4 files changed, 504 insertions(+), 168 deletions(-) create mode 100644 research/surge_book_replay.py diff --git a/research/grade_surge.py b/research/grade_surge.py index 553fd0ed..cfecae92 100644 --- a/research/grade_surge.py +++ b/research/grade_surge.py @@ -1,12 +1,15 @@ #!/usr/bin/env python3 -"""Nightly chain-truth grading of the surge paper harness (wwf-surgebot). +"""Nightly chain-truth grading of the surge measurement harness (A2). -Pulls /data/surge_state.json off the box (read-only — never writes back), -re-grades every settled entry with CTF payout vectors (payouts.truth; the -harness's provisional CLOB winner flags lie on operator-resolved markets, -[[polymarket-resolution-truth]]), and appends one row per settle to -research/surge_paper_ledger.jsonl (keyed asset+ts, idempotent). Also emits -the running paper-book summary the #16 sprint decision reads.""" +Pulls /data/surge2_state.json AND /data/surge_attempts.jsonl off the box +(read-only — never writes back), re-grades every settled entry with CTF +payout vectors (payouts.truth; provisional CLOB winner flags lie on +operator-resolved markets), and appends one row per settle to +research/surge_meas_ledger.jsonl (keyed asset+ts, idempotent). The pulled +copies are KEPT on disk for surge_book_replay.py (virtual bankroll specs). + +v1 history is closed: surge_paper_ledger.jsonl and /data/surge_state.json +are frozen audit artifacts of the halted cash-gated book (2026-07-22).""" import json import os import shutil @@ -14,33 +17,53 @@ import subprocess import sys HERE = os.path.dirname(os.path.abspath(__file__)) -LEDGER = os.path.join(HERE, "surge_paper_ledger.jsonl") -TMP = os.path.join(HERE, ".surge_state.pull.json") +LEDGER = os.path.join(HERE, "surge_meas_ledger.jsonl") +STATE_PULL = os.path.join(HERE, ".surge2_state.pull.json") +ATT_PULL = os.path.join(HERE, ".surge_attempts.pull.jsonl") FLYCTL = shutil.which("flyctl") or "/opt/homebrew/bin/flyctl" +def sftp(remote, local): + subprocess.run([FLYCTL, "ssh", "sftp", "get", remote, local + ".new", + "-a", "wwf-surgebot"], capture_output=True, + timeout=600, stdin=subprocess.DEVNULL) + if os.path.exists(local + ".new") and os.path.getsize(local + ".new") > 0: + os.replace(local + ".new", local) # keep last good pull on failure + return True + try: + os.remove(local + ".new") + except FileNotFoundError: + pass + return False + + def main(): - r = subprocess.run([FLYCTL, "ssh", "sftp", "get", "/data/surge_state.json", - TMP, "-a", "wwf-surgebot"], capture_output=True, - timeout=300, stdin=subprocess.DEVNULL) - if not os.path.exists(TMP) or os.path.getsize(TMP) == 0: + ok = sftp("/data/surge2_state.json", STATE_PULL) + sftp("/data/surge_attempts.jsonl", ATT_PULL) + if not ok: print("[grade_surge] box unreachable or no state yet — skip") return 0 - st = json.load(open(TMP)) - os.remove(TMP) + st = json.load(open(STATE_PULL)) sys.path.insert(0, os.path.join(HERE, "..", "live")) import payouts have = set() + ledger_rows = 0 try: for ln in open(LEDGER): try: d = json.loads(ln) have.add((d["asset"], d["ts"])) + ledger_rows += 1 except Exception: pass except FileNotFoundError: pass settled = st.get("settled", []) + total_settles = st.get("counters", {}).get("settled_total", 0) + if total_settles > ledger_rows + len(settled): + print(f"[grade_surge] ⚠ TAPE LOSS: {total_settles} lifetime settles " + f"but only {ledger_rows} graded + {len(settled)} in state — " + f"raise SETTLED_TRIM or grade more often") new = [s for s in settled if (s["asset"], s["ts"]) not in have] if new: payouts.ensure(sorted({s["cond"] for s in new if s.get("cond")})) @@ -67,11 +90,15 @@ def main(): except FileNotFoundError: pass c = st.get("counters", {}) + ev = pnl_sum / total if total else None print(f"[grade_surge] +{graded} settles ({flips} provisional flips) · " - f"book: {total} settled {wins}W · chain P&L ${pnl_sum:+.2f} · " - f"cash ${st.get('cash', 0):.2f} · open {len(st.get('open', {}))} · " - f"lifetime trig {c.get('triggers')} fill {c.get('fills')} " + f"measurement: {total} settled {wins}W · chain P&L ${pnl_sum:+.2f}" + f"{f' (${ev:+.2f}/fill)' if ev is not None else ''} · " + f"open {len(st.get('open', {}))} · lifetime trig {c.get('triggers')} " + f"att {c.get('attempts')} fill {c.get('fills')} " f"crater {c.get('craters')}") + print("[grade_surge] capturability read — NOT the #16 verdict " + "(forward_ledger is); bankroll reads come from surge_book_replay") return 0 diff --git a/research/nightly.sh b/research/nightly.sh index fa1110d0..caea829e 100755 --- a/research/nightly.sh +++ b/research/nightly.sh @@ -52,12 +52,14 @@ done "$PY" forward.py >> forward.log 2>&1 "$PY" informed_set.py >> forward.log 2>&1 # surge harness reads this daily -"$PY" grade_surge.py >> forward.log 2>&1 || true # paper book -> chain truth +"$PY" grade_surge.py >> forward.log 2>&1 || true # A2 measurement -> chain truth +"$PY" surge_book_replay.py >> forward.log 2>&1 || true # virtual $100/5% book "$PY" grade_oracle.py >> forward.log 2>&1 || true # oracle paper -> chain truth cd .. git add research/forward_ledger.jsonl research/params/informed_set.json \ - research/surge_paper_ledger.jsonl research/oracle_paper_ledger.jsonl 2>/dev/null + research/surge_paper_ledger.jsonl research/oracle_paper_ledger.jsonl \ + research/surge_meas_ledger.jsonl research/surge_book.json 2>/dev/null if ! git diff --cached --quiet; then git commit -q -m "research: forward ledger $(date -u +%F) [skip ci]" git pull --rebase --autostash -q && git push -q diff --git a/research/surge_book_replay.py b/research/surge_book_replay.py new file mode 100644 index 00000000..be329af2 --- /dev/null +++ b/research/surge_book_replay.py @@ -0,0 +1,152 @@ +#!/usr/bin/env python3 +"""Virtual bankroll replay over the surge A2 attempts stream. + +Re-runs the v1 deployment spec ($100 bank · 5%-of-equity daily stakes with +$1 floor · cash-gated · max 2 open per event · skip-if-open-same-asset · +fill at best ask inside p_ref*1.05) against the FULL attempt record +(.surge_attempts.pull.jsonl), settling with chain truth where graded +(surge_meas_ledger.jsonl) and provisional payouts otherwise +(.surge2_state.pull.json). Writes research/surge_book.json for the /test +dashboard and the #19 Friday read. + +This is the split that fixes v1's flaw: the physical harness samples EVERY +trigger; bankroll specs are simulated here, where cash-gating can no longer +corrupt the sample. Change SPEC below (or add variants) freely — this file +is analysis, not signal; the frozen signal lives in the harness.""" +import heapq +import json +import os +import time + +HERE = os.path.dirname(os.path.abspath(__file__)) +ATT = os.path.join(HERE, ".surge_attempts.pull.jsonl") +STATE = os.path.join(HERE, ".surge2_state.pull.json") +LEDGER = os.path.join(HERE, "surge_meas_ledger.jsonl") +OUT = os.path.join(HERE, "surge_book.json") + +SPEC = {"bank": 100.0, "stake_pct": 0.05, "stake_floor": 1.0, + "event_cap": 2, "fee_rate": 0.03, "slip_cap": 0.05} + + +def load_payouts(): + """asset:ts -> (payout, settled_ts, chain?) — ledger beats state.""" + pay = {} + try: + st = json.load(open(STATE)) + for s in st.get("settled", []): + pay[f"{s['asset']}:{s['ts']}"] = (s["payout"], s["settled_ts"], False) + except FileNotFoundError: + pass + try: + for ln in open(LEDGER): + d = json.loads(ln) + pay[f"{d['asset']}:{d['ts']}"] = (d["chain_payout"], + d["settled_ts"], True) + except FileNotFoundError: + pass + return pay + + +def main(): + if not os.path.exists(ATT): + print("[book_replay] no attempts pull yet — skip") + return 0 + pay = load_payouts() + atts = [] + for ln in open(ATT): + try: + atts.append(json.loads(ln)) + except Exception: + pass + atts.sort(key=lambda a: a["ts"]) + + cash = SPEC["bank"] + day = "" + stake = SPEC["stake_floor"] + open_lots = {} # asset -> lot (v1: one per asset) + due = [] # heap of (settle_ts, asset) + curve = [] + settled = wins = taken = 0 + cash_skip = event_skip = open_skip = crater = unresolved_cap = 0 + pnl_real = 0.0 + graded_n = 0 + + def equity(): + return cash + sum(l["cost"] for l in open_lots.values()) + + def settle_due(now): + nonlocal cash, settled, wins, pnl_real, graded_n + while due and due[0][0] <= now: + _, asset = heapq.heappop(due) + lot = open_lots.pop(asset, None) + if lot is None: + continue + p, _, chain = pay[lot["key"]] + cash += lot["shares"] * p + pnl_real += lot["shares"] * p - lot["cost"] - lot["fee"] + settled += 1 + wins += p == 1.0 + graded_n += chain + + for a in atts: + now = a["ts"] + settle_due(now) + d = time.strftime("%Y-%m-%d", time.gmtime(now)) + if d != day: + day = d + stake = max(SPEC["stake_floor"], round(SPEC["stake_pct"] * equity(), 2)) + if len(curve) == 0 or now - curve[-1][0] >= 1800: + curve.append([int(now), round(equity(), 2)]) + asset = a["asset"] + if asset in open_lots: + open_skip += 1 + continue + ev = a.get("event") + if ev and sum(1 for l in open_lots.values() + if l["event"] == ev) >= SPEC["event_cap"]: + event_skip += 1 + continue + if cash < stake: + cash_skip += 1 + continue + ba = a.get("best_ask") + cap = min(a["p_ref"] * (1 + SPEC["slip_cap"]), 0.99) + if not a.get("filled") or ba is None or ba > cap: + crater += 1 + continue + key = f"{asset}:{a['ts']}" + info = pay.get(key) + shares = stake / ba + fee = SPEC["fee_rate"] * shares * min(ba, 1 - ba) + cash -= stake + fee + open_lots[asset] = {"key": key, "event": ev, "cost": stake, + "fee": fee, "shares": shares} + taken += 1 + if info is not None: + heapq.heappush(due, (info[1], asset)) + else: + unresolved_cap += 1 # stays open until a later run grades it + settle_due(float("inf") if not open_lots else time.time()) + curve.append([int(time.time()), round(equity(), 2)]) + + out = {"computed_at": time.strftime("%Y-%m-%d %H:%M UTC", time.gmtime()), + "spec": SPEC, "sem_ver_source": "a2 attempts stream", + "equity": round(equity(), 2), "cash": round(cash, 2), + "open_n": len(open_lots), "settled": settled, "wins": wins, + "losses": settled - wins, "pnl_realized": round(pnl_real, 2), + "chain_graded_settles": graded_n, + "counters": {"attempts_seen": len(atts), "taken": taken, + "cash_skip": cash_skip, "event_skip": event_skip, + "open_skip": open_skip, "crater": crater, + "open_unresolved": unresolved_cap}, + "curve": curve[-336:]} + json.dump(out, open(OUT, "w"), indent=1) + print(f"[book_replay] virtual ${out['equity']:.2f} " + f"(cash ${out['cash']:.2f}) · {settled} settled {wins}W · " + f"taken {taken}/{len(atts)} attempts " + f"(skips c{cash_skip}/e{event_skip}/o{open_skip}, crater {crater})") + return 0 + + +if __name__ == "__main__": + main() diff --git a/research/surgebot.py b/research/surgebot.py index 69ef2991..faa58f7c 100644 --- a/research/surgebot.py +++ b/research/surgebot.py @@ -1,29 +1,43 @@ #!/usr/bin/env python3 -"""surgebot — Study A's surge-momentum signal as a real-time PAPER harness. +"""surgebot A2 — Study A's surge signal as a real-time MEASUREMENT harness. -PAPER ONLY. No keys, no orders, no bot imports — this file is self-contained -(recorder-style: baked into the wwf-surgebot image, must not depend on the -repo to boot). It exists to answer the two questions the tape sim cannot: -does the signal compute in real time on the live stream, and what does the -paper book do at the $100/5%-stake deployment spec (#16 sprint plan). Its -entries are graded nightly against chain truth on the Mac (grade_surge.py); -believing ANY of it is gated on the pre-registered #16 forward verdict. +PAPER ONLY. No keys, no orders, no bot imports — self-contained (baked into +the wwf-surgebot image, must not depend on the repo to boot). -FROZEN signal (params/study_flow.json, 2026-07-20 — do not tune here): +WHY A2 (2026-07-22 post-mortem of the v1 paper book): v1 rehearsed the +$100/5% deployment spec physically — cash-gating meant it attempted ~2% of +triggers, and that subsample was ADVERSELY SELECTED (its 102 chain-graded +fills: −$9/fill at $100-scale, while the sim scored ALL of the same day's +triggers at +$41/fill; the sim run on v1's own triggers agreed with v1's +losses within ~$2 — execution physics validated, sampling condemned). +A cash-gated book cannot measure the signal AND rehearse the spec at once. + +A2 splits them: this box attempts EVERY cooldown-passed trigger (flat $100 +paper FAK walking the live asks inside p_ref*1.05, partial fills kept) and +appends one line per attempt — including top-5 raw ask levels — to +/data/surge_attempts.jsonl. Any bankroll spec (the $100/5% book, other +banks, caps, tighter slip) is then replayed OFFLINE from that stream +(research/surge_book_replay.py -> research/surge_book.json, nightly). + +FROZEN signal semantics (params/study_flow.json — do not tune here): informed set top-150 (fetched from the repo, regenerated nightly) trigger net informed flow >= $300 in 60s, one market band entry price 0.10-0.90 · niches sports+esports - cooldown 900s per token -Deployment spec (2026-07-21 sizing discussion): - bank $100 paper · stake 5% of equity, set once per UTC day, $1 floor - cash-gated all-or-nothing · max 2 open positions per real-world event -Paper fill = live CLOB best ask inside p_ref*1.05, else crater (the same -FAK model the copybot's paper mode uses); venue fee 3% * shares * min(p,1-p). -Settles provisionally from the CLOB closed/winner flags; the nightly -re-grades with CTF payout vectors (operator-resolved markets lie here). + cooldown 900s per token, set at trigger time +Unchanged from v1: signal code path is verbatim. Changed in A2 (execution/ +accounting only): every-trigger attempts (no cash gate, no event cap, no +skip-if-open — those live in the offline replay), $100 book-walk fills +(ledger-comparable; v1 was best-ask at ~$5), attempts log, multi-lot per +asset. THE #16 VERDICT BINDS TO forward_ledger.jsonl ONLY; this harness +measures capturability. Settles are provisional (CLOB closed/winner) until +the nightly CTF payout re-grade (grade_surge.py -> surge_meas_ledger.jsonl; +v1's surge_paper_ledger.jsonl is closed and untouched, as is its state at +/data/surge_state.json). Any semantics-altering fix bumps SEM_VER and +restarts the shakedown clock. """ import json import os +import queue import re import ssl import threading @@ -36,18 +50,27 @@ WS_URL = "wss://ws-live-data.polymarket.com" CLOB = "https://clob.polymarket.com" SET_URL = ("https://raw.githubusercontent.com/jaxperro/winning-wallet-finder/" "main/research/params/informed_set.json") -STATE = os.environ.get("SURGE_STATE", "/data/surge_state.json") +STATE = os.environ.get("SURGE_STATE", "/data/surge2_state.json") +ATTEMPTS = os.environ.get("SURGE_ATTEMPTS", "/data/surge_attempts.jsonl") +SEM_VER = "a2" # bump on ANY semantics-altering change -BANK = 100.0 -STAKE_PCT = 0.05 -EVENT_CAP = 2 +# ── FROZEN — must equal params/study_flow.json / forward.py ──────────────── FLOW_USD = 300.0 WINDOW_S = 60 BAND = (0.10, 0.90) COOLDOWN_S = 900 +NICHES = {"sports", "esports"} +STAKE = 100.0 # ledger stake (sf.STAKE) — walk the book FEE_RATE = 0.03 SLIP_CAP = 0.05 -NICHES = {"sports", "esports"} +# ── LIVE-ONLY (execution/plumbing, never signal) ─────────────────────────── +BOOK_WORKERS = 4 +BOOK_QUEUE = 16 # overflow => throttle (counted) +CLOB_PASS_CAP = 40 # settle GETs per pass, round-robin +SETTLED_TRIM = 10000 # nightly ledger is the durable record +SKIPS_TRIM = 500 +TOP_LEVELS = 5 # raw ask levels recorded per attempt + NICHE_PATTERNS = [ ("esports", ["lol:", "dota", "cs2", "csgo", "valorant", "esports", "bilibili", "map ", "game 1", "game 2", "game 3"]), @@ -91,17 +114,58 @@ def iso_ts(s): return None +def walk_asks(asks, cap, stake_usd=STAKE): + """FAK against a CLOB ask list: consume ascending price levels <= cap + until stake_usd is spent; partial fills kept (that IS what a real FAK + gets). Returns top raw levels too so ANY smaller stake / tighter cap can + be replayed offline. asks arrive UNSORTED from /book.""" + lv = sorted(((float(a["price"]), float(a["size"])) for a in asks or [])) + top = [[p, s] for p, s in lv[:TOP_LEVELS]] + if not lv: + return {"filled": False, "best_ask": None, "top": top} + best = lv[0][0] + spent = shares = 0.0 + levels = 0 + for pxx, sz in lv: + if pxx > cap or spent >= stake_usd - 1e-9: + break + take = min(sz, (stake_usd - spent) / pxx) + if take <= 0: + break + shares += take + spent += take * pxx + levels += 1 + if not shares: + return {"filled": False, "best_ask": best, "top": top} + return {"filled": True, "vwap": spent / shares, "shares": shares, + "usd": spent, "levels": levels, "best_ask": best, "top": top, + "partial": spent < stake_usd - 1e-9} + + +def fresh_state(): + return { + "sem_ver": SEM_VER, "boot_ts": int(time.time()), + "open": {}, "settled": [], "skips": [], + "pnl_realized": 0.0, + "counters": {"triggers": 0, "attempts": 0, "fills": 0, "craters": 0, + "throttle_skip": 0, "book_fail": 0, + "settled_total": 0, "settle_clob": 0}} + + class Surge: def __init__(self): - self.state = {"cash": BANK, "day": "", "day_stake": 5.0, - "open": {}, "settled": [], "skips": [], - "counters": {"triggers": 0, "fills": 0, "craters": 0, - "cash_skip": 0, "event_skip": 0}} + self.state = fresh_state() if os.path.exists(STATE): try: - self.state = json.load(open(STATE)) - log(f"resumed: cash ${self.state['cash']:.2f} · " - f"{len(self.state['open'])} open") + loaded = json.load(open(STATE)) + if loaded.get("sem_ver") != SEM_VER: + log(f"⚠ SEM_VER {loaded.get('sem_ver')} -> {SEM_VER} — " + f"semantics changed; shakedown clock restarts here") + loaded["sem_ver"] = SEM_VER + self.state = loaded + log(f"resumed: pnl {self.state['pnl_realized']:+.2f} · " + f"{len(self.state['open'])} open · " + f"{self.state['counters']['fills']} lifetime fills") except Exception as e: log(f"⚠ state load failed ({e}) — fresh book") self.informed = set() @@ -109,15 +173,15 @@ class Surge: self.win = {} # asset -> [(ts, ±usd)] self.last_trig = {} self.lock = threading.Lock() - self._roll_day(force=True) + self.book_q = queue.Queue(BOOK_QUEUE) - # ── config / sizing ──────────────────────────────────────────────────── + # ── config ───────────────────────────────────────────────────────────── def load_set(self): try: d = get_json(SET_URL, timeout=15) self.informed = {w.lower() for w in d["wallets"]} self.set_meta = {"generated_at": d.get("generated_at"), - "n": len(self.informed)} + "n": len(self.informed)} age_h = (time.time() - (d.get("generated_at") or 0)) / 3600 log(f"informed set: {len(self.informed)} wallets " f"({age_h:.0f}h old)" + (" ⚠ STALE >48h" if age_h > 48 else "")) @@ -125,24 +189,7 @@ class Surge: log(f"⚠ informed set fetch failed ({e}) — " f"keeping {len(self.informed)} cached") - def equity(self): - return self.state["cash"] + sum(p["cost"] for p in - self.state["open"].values()) - - def _roll_day(self, force=False): - day = time.strftime("%Y-%m-%d", time.gmtime()) - if force or day != self.state["day"]: - self.state["day"] = day - self.state["day_stake"] = max(1.0, round(STAKE_PCT * self.equity(), 2)) - log(f"day {day}: stake ${self.state['day_stake']:.2f} " - f"(5% of ${self.equity():.2f} equity)") - - def persist(self): - tmp = STATE + ".tmp" - json.dump(self.state, open(tmp, "w")) - os.replace(tmp, STATE) - - # ── signal ───────────────────────────────────────────────────────────── + # ── signal — FROZEN, verbatim v1 path ────────────────────────────────── def on_trade(self, p): w = (p.get("proxyWallet") or "").lower() if w not in self.informed: @@ -167,92 +214,182 @@ class Surge: if now - self.last_trig.get(asset, 0) < COOLDOWN_S: return self.last_trig[asset] = now - self.state["counters"]["triggers"] += 1 - self.execute(p, asset, px, flow, title) + c = self.state["counters"] + c["triggers"] += 1 + # A2: EVERY trigger becomes an attempt (no cash/event/open gate — + # bankroll specs are replayed offline from the attempts log) + job = {"asset": asset, "cond": p.get("conditionId"), + "event": event_key(p.get("eventSlug") or p.get("slug") + or ""), + "title": title[:60], "outcome": p.get("outcome"), + "niche": niche(title), "ts": round(now, 3), + "p_ref": px, "flow": round(flow)} + try: + self.book_q.put_nowait(job) + c["attempts"] += 1 + except queue.Full: + c["throttle_skip"] += 1 + self._skip(job, "throttle", None) + self.log_attempt({**job, "filled": False, "why": "throttle"}) - # ── paper execution (called under lock) ──────────────────────────────── - def execute(self, p, asset, p_ref, flow, title): - self._roll_day() - c = self.state["counters"] - ev = event_key(p.get("eventSlug") or p.get("slug") or "") - if asset in self.state["open"]: - return - n_ev = sum(1 for o in self.state["open"].values() if o["event"] == ev) - if ev and n_ev >= EVENT_CAP: - c["event_skip"] += 1 - log(f"skip (event cap) {title[:40]}") - return - stake = self.state["day_stake"] - if self.state["cash"] < stake: - c["cash_skip"] += 1 - log(f"skip (no cash: ${self.state['cash']:.2f}) {title[:40]}") - return + # ── paper execution (workers; network OUT of the lock) ───────────────── + def book_worker(self): + while True: + job = self.book_q.get() + try: + self.attempt(job) + except Exception as e: + log(f"⚠ attempt error: {e}") + + def attempt(self, job): + t_req = time.time() try: - book = get_json(f"{CLOB}/book?token_id={asset}") - asks = [float(a["price"]) for a in (book.get("asks") or [])] - ba = min(asks) if asks else None - except Exception as e: - log(f"skip (book fetch failed: {e}) {title[:40]}") - return - cap = min(p_ref * (1 + SLIP_CAP), 0.99) - if ba is None or ba > cap: - c["craters"] += 1 - self.state["skips"].append({"ts": int(time.time()), "asset": asset, - "p_ref": p_ref, "ba": ba, "flow": flow, - "title": title[:60], "why": "crater"}) - del self.state["skips"][:-500] - log(f"CRATER {title[:40]} ref {p_ref:.2f} ask " - f"{('%.2f' % ba) if ba else '—'}") - self.persist() - return - shares = stake / ba - fee = FEE_RATE * shares * min(ba, 1 - ba) - end_ts = None # expected resolution (dashboard ETA) - try: - m = get_json(f"{CLOB}/markets/{p.get('conditionId')}", timeout=5) - end_ts = iso_ts(m.get("end_date_iso") or "") + book = get_json(f"{CLOB}/book?token_id={job['asset']}") except Exception: - pass - self.state["cash"] -= stake + fee - self.state["open"][asset] = { - "ts": int(time.time()), "cond": p.get("conditionId"), - "event": ev, "title": title[:60], "end_ts": end_ts, - "outcome": p.get("outcome"), "p_ref": p_ref, "price": ba, - "shares": round(shares, 4), "cost": stake, "fee": round(fee, 4), - "flow": round(flow)} - c["fills"] += 1 - log(f"FILL {p.get('outcome')} · {title[:40]} @ {ba:.3f} " - f"(${stake:.2f}, flow ${flow:.0f}) · cash ${self.state['cash']:.2f}") - self.persist() + with self.lock: + self.state["counters"]["book_fail"] += 1 + self._skip(job, "book_fail", None) + self.log_attempt({**job, "filled": False, "why": "book_fail"}) + return + lat_ms = int((time.time() - t_req) * 1000) + cap = min(job["p_ref"] * (1 + SLIP_CAP), 0.99) + r = walk_asks(book.get("asks"), cap) + end_ts = None # expected resolution (dashboard ETA) + if r["filled"]: + try: + m = get_json(f"{CLOB}/markets/{job.get('cond')}", timeout=5) + end_ts = iso_ts(m.get("end_date_iso") or "") + except Exception: + pass + rec = {**job, "cap": round(cap, 4), "latency_ms": lat_ms, + "top": r["top"], "best_ask": r["best_ask"], + "filled": r["filled"]} + with self.lock: + c = self.state["counters"] + if not r["filled"]: + c["craters"] += 1 + self._skip(job, "crater", r["best_ask"]) + ba = "—" if r["best_ask"] is None else f"{r['best_ask']:.3f}" + log(f"CRATER {job['title'][:38]} ref {job['p_ref']:.3f} " + f"ask {ba} flow ${job['flow']} ({lat_ms}ms)") + else: + fee = FEE_RATE * r["shares"] * min(r["vwap"], 1 - r["vwap"]) + lot = {**job, "end_ts": end_ts, "price": round(r["vwap"], 5), + "best_ask": r["best_ask"], + "shares": round(r["shares"], 4), + "cost": round(r["usd"], 2), "fee": round(fee, 4), + "levels": r["levels"], "partial": r["partial"], + "latency_ms": lat_ms} + self.state["open"][f"{job['asset']}:{job['ts']}"] = lot + c["fills"] += 1 + rec.update({"price": lot["price"], "shares": lot["shares"], + "usd": lot["cost"], "fee": lot["fee"], + "partial": lot["partial"], "end_ts": end_ts}) + log(f"FILL {job['outcome']} · {job['title'][:38]} @ " + f"{lot['price']:.3f} (${lot['cost']:.2f}" + f"{' partial' if lot['partial'] else ''}, " + f"flow ${job['flow']}, {lat_ms}ms)") + self.persist() + self.log_attempt(rec) + + def log_attempt(self, rec): + """Append-only full attempt stream — the offline-replay dataset. + One JSON line per attempt (fills, craters, fails, throttles).""" + try: + with open(ATTEMPTS, "a") as fh: + fh.write(json.dumps(rec) + "\n") + except Exception as e: + log(f"⚠ attempts log write failed: {e}") + + def _skip(self, job, why, ba): # under lock + self.state["skips"].append( + {"ts": int(job["ts"]), "asset": job["asset"], "why": why, + "p_ref": job["p_ref"], "flow": job["flow"], "ba": ba, + "title": job["title"]}) + del self.state["skips"][:-SKIPS_TRIM] # ── provisional settles (nightly re-grades with chain truth) ─────────── - def settle_pass(self): - for asset, pos in list(self.state["open"].items()): + def settle_clob_pass(self): + with self.lock: + now = time.time() + due = list(self.state["open"].keys()) + due.sort(key=lambda lid: self.state["open"][lid].get("clob_ck", 0)) + due = due[:CLOB_PASS_CAP] + snap = {} + for lid in due: + lot = self.state["open"][lid] + lot["clob_ck"] = now # round-robin stamp + snap[lid] = (lot["cond"], lot["asset"], lot.get("end_ts")) + verdicts = {} + ends = {} + for lid, (cond, asset, end_ts) in snap.items(): try: - m = get_json(f"{CLOB}/markets/{pos['cond']}") + m = get_json(f"{CLOB}/markets/{cond}") except Exception: continue - if not pos.get("end_ts"): # backfill ETAs for pre-ETA fills - pos["end_ts"] = iso_ts(m.get("end_date_iso") or "") + if end_ts is None: # backfill ETAs for pre-ETA fills + ends[lid] = iso_ts(m.get("end_date_iso") or "") if not m.get("closed"): continue - pay = None for t in m.get("tokens") or []: if str(t.get("token_id")) == str(asset): - pay = 1.0 if t.get("winner") else 0.0 - if pay is None: - continue - self.state["cash"] += pos["shares"] * pay - self.state["settled"].append({**pos, "asset": asset, "payout": pay, - "provisional": True, - "settled_ts": int(time.time()), - "pnl": round(pos["shares"] * pay - - pos["cost"] - pos["fee"], 2)}) - del self.state["open"][asset] - log(f"SETTLE {'WON' if pay else 'lost'} {pos['title'][:38]} " - f"{'+' if pay else ''}{pos['shares'] * pay - pos['cost'] - pos['fee']:.2f} " - f"· cash ${self.state['cash']:.2f}") - self.persist() + verdicts[lid] = 1.0 if t.get("winner") else 0.0 + with self.lock: + for lid, e in ends.items(): + if lid in self.state["open"]: + self.state["open"][lid]["end_ts"] = e + n = 0 + for lid, pay in verdicts.items(): + lot = self.state["open"].get(lid) + if lot: + self._settle(lid, lot, pay) + n += 1 + if n or ends: + self.persist() + + def _settle(self, lid, lot, pay): # under lock + pnl = round(lot["shares"] * pay - lot["cost"] - lot["fee"], 2) + self.state["pnl_realized"] = round(self.state["pnl_realized"] + pnl, 2) + c = self.state["counters"] + c["settled_total"] += 1 + c["settle_clob"] += 1 + self.state["settled"].append( + {**lot, "payout": pay, "provisional": True, + "settled_ts": int(time.time()), "pnl": pnl}) + del self.state["settled"][:-SETTLED_TRIM] + del self.state["open"][lid] + word = "WON" if pay == 1.0 else "refund" if pay == 0.5 else "lost" + log(f"SETTLE {word} {lot['title'][:36]} {pnl:+.2f} · " + f"realized {self.state['pnl_realized']:+.2f}") + + # ── plumbing ─────────────────────────────────────────────────────────── + def persist(self): + tmp = STATE + ".tmp" + json.dump(self.state, open(tmp, "w")) + os.replace(tmp, STATE) + + def feed_dict(self): # under lock + stt = self.state + return { + "mode": "paper-measurement", "study": "surge-A (#16)", + "sem_ver": SEM_VER, "updated": int(time.time()), + "boot_ts": stt["boot_ts"], + "frozen": {"flow_usd": FLOW_USD, "window_s": WINDOW_S, + "band": BAND, "cooldown_s": COOLDOWN_S, + "stake": STAKE, "fee_rate": FEE_RATE, + "slip_cap": SLIP_CAP}, + "counters": stt["counters"], "pnl_realized": stt["pnl_realized"], + "informed": self.set_meta, "open_n": len(stt["open"]), + "open": [{**l, "id": lid} + for lid, l in list(stt["open"].items())[-200:]], + "settled": stt["settled"][-200:], "skips": stt["skips"][-100:], + "note": "verdict binds to forward_ledger.jsonl (#16); this feed " + "measures capturability — bankroll specs are replayed " + "offline from the attempts log (surge_book.json)"} + + +SUB = json.dumps({"action": "subscribe", "subscriptions": [ + {"topic": "activity", "type": "trades", "filters": ""}]}) def run_conn(tag, bot): @@ -261,9 +398,7 @@ def run_conn(tag, bot): state = {"fresh": time.time()} def on_open(ws): - ws.send(json.dumps({"action": "subscribe", "subscriptions": - [{"topic": "activity", "type": "trades", - "filters": ""}]})) + ws.send(SUB) state["fresh"] = time.time() log(f"ws[{tag}]: connected") @@ -307,7 +442,7 @@ def run_conn(tag, bot): def serve_feed(bot, port=8080): - """Read-only public feed for the /surge dashboard (paper data only — + """Read-only public feed for the /test dashboard (paper data only — nothing here can place an order or mutate state). CORS-open.""" from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer @@ -321,19 +456,7 @@ def serve_feed(bot, port=8080): self.end_headers() return with bot.lock: - st = bot.state - body = json.dumps({ - "mode": "paper", "updated": int(time.time()), - "bank": BANK, "equity": round(bot.equity(), 2), - "cash": round(st["cash"], 2), - "day": st["day"], "day_stake": st["day_stake"], - "stake_pct": STAKE_PCT, "event_cap": EVENT_CAP, - "flow_usd": FLOW_USD, "window_s": WINDOW_S, - "band": BAND, "counters": st["counters"], - "informed": bot.set_meta, - "open": [{**p, "asset": a} for a, p in st["open"].items()], - "settled": st["settled"][-200:], - "skips": st["skips"][-100:]}).encode() + body = json.dumps(bot.feed_dict()).encode() self.send_response(200) self.send_header("Content-Type", "application/json") self.send_header("Access-Control-Allow-Origin", "*") @@ -345,10 +468,32 @@ def serve_feed(bot, port=8080): ThreadingHTTPServer(("0.0.0.0", port), H).serve_forever() +def selftest(): + """Boot gate — hard exit on failure.""" + assert niche("LoL: T1 vs GenG map 2") == "esports" + assert niche("Pittsburgh Pirates vs. New York Yankees: O/U 9.5") == "sports" + assert niche("Will it rain in NYC") == "other" + assert event_key("mlb-pit-nyy-2026-07-22-game") == "mlb-pit-nyy-2026-07-22" + r = walk_asks([{"price": "0.60", "size": "200"}, + {"price": "0.55", "size": "100"}], cap=0.62) + assert r["filled"] and abs(r["usd"] - 100.0) < 1e-6, r + assert r["levels"] == 2 and r["best_ask"] == 0.55 and not r["partial"], r + assert r["top"][0] == [0.55, 100.0] and len(r["top"]) == 2, r + r = walk_asks([{"price": "0.55", "size": "10"}], cap=0.60) + assert r["filled"] and r["partial"] and abs(r["usd"] - 5.5) < 1e-9, r + r = walk_asks([{"price": "0.70", "size": "10"}], cap=0.60) + assert not r["filled"] and r["best_ask"] == 0.70, r + assert not walk_asks([], 0.5)["filled"] + log("selftest OK") + + def main(): + selftest() bot = Surge() bot.load_set() threading.Thread(target=serve_feed, args=(bot,), daemon=True).start() + for _ in range(BOOK_WORKERS): + threading.Thread(target=bot.book_worker, daemon=True).start() for tag in ("a", "b"): threading.Thread(target=run_conn, args=(tag, bot), daemon=True).start() last_set = time.time() @@ -356,19 +501,29 @@ def main(): while True: time.sleep(60) n += 1 + now = time.time() with bot.lock: - bot._roll_day() + bot.last_trig = {a: t for a, t in bot.last_trig.items() + if now - t < COOLDOWN_S} + # trim flow windows for tokens that went quiet + bot.win = {a: b for a, b in bot.win.items() + if b and b[-1][0] > now - 2 * WINDOW_S} if n % 2 == 0: + bot.settle_clob_pass() + if n % 5 == 0: with bot.lock: - bot.settle_pass() + bot.persist() if time.time() - last_set > 6 * 3600: bot.load_set() last_set = time.time() - c = bot.state["counters"] - log(f"paper ${bot.equity():.2f} (cash ${bot.state['cash']:.2f}) · " - f"open {len(bot.state['open'])} · trig {c['triggers']} " - f"fill {c['fills']} crater {c['craters']} " - f"skip {c['cash_skip']}c/{c['event_skip']}e") + with bot.lock: + c = bot.state["counters"] + msg = (f"paper pnl {bot.state['pnl_realized']:+.2f} · " + f"open {len(bot.state['open'])} · trig {c['triggers']} " + f"att {c['attempts']} fill {c['fills']} " + f"crat {c['craters']} thr {c['throttle_skip']} · " + f"settled {c['settled_total']}") + log(msg) if __name__ == "__main__":