research: surgebot A2 — measurement arm relaunch (every-trigger $100 FAKs + attempts stream + offline virtual-book replay)

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 <noreply@anthropic.com>
This commit is contained in:
jaxperro
2026-07-22 18:30:54 -04:00
parent 4a406ac1c5
commit 2c322419b6
4 changed files with 504 additions and 168 deletions
+45 -18
View File
@@ -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
+4 -2
View File
@@ -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
+152
View File
@@ -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()
+303 -148
View File
@@ -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__":