From 77b2457523689e6d3957596357a517fe071a160e Mon Sep 17 00:00:00 2001 From: jaxperro Date: Sun, 19 Jul 2026 15:55:10 -0400 Subject: [PATCH] =?UTF-8?q?recorder:=20FULL=20firehose=20=E2=80=94=20activ?= =?UTF-8?q?ity/*=20(trades=20+=20orders=5Fmatched),=20comments,=20crypto?= =?UTF-8?q?=5Fprices;=20aux=20segment=20family=20+=20ingest=20branch=20(pr?= =?UTF-8?q?obed:=20rfq/prices=20dead)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- recorder/ingest.py | 22 ++++++++---- recorder/recorder.py | 85 +++++++++++++++++++++++++++++++------------- 2 files changed, 75 insertions(+), 32 deletions(-) diff --git a/recorder/ingest.py b/recorder/ingest.py index 9adf93b1..6d01abd5 100644 --- a/recorder/ingest.py +++ b/recorder/ingest.py @@ -37,6 +37,8 @@ def main(): price DOUBLE, size DOUBLE, tx VARCHAR, title VARCHAR)""") con.execute("""CREATE TABLE IF NOT EXISTS ingested( segment VARCHAR PRIMARY KEY, rows BIGINT, ingested_at BIGINT)""") + con.execute("""CREATE TABLE IF NOT EXISTS aux( + ts DOUBLE, topic VARCHAR, type VARCHAR, payload VARCHAR)""") have = {r[0] for r in con.execute("SELECT segment FROM ingested").fetchall()} segs = [s for s in box(f"ls {SEG}").split() if s.endswith(".gz") and s not in have] @@ -48,24 +50,30 @@ def main(): except Exception as e: print(f"[ingest] {s}: fetch/decode failed ({e}) — left on box") continue - rows = [] + rows, aux = [], [] for ln in lines: try: d = json.loads(ln) - rows.append((d.get("ts"), d.get("wallet"), d.get("asset"), - d.get("cond"), d.get("side"), d.get("price"), - d.get("size"), d.get("tx"), d.get("title"))) + if s.startswith("aux_"): + aux.append((d.get("ts"), d.get("topic"), d.get("type"), + json.dumps(d.get("payload")))) + else: + rows.append((d.get("ts"), d.get("wallet"), d.get("asset"), + d.get("cond"), d.get("side"), d.get("price"), + d.get("size"), d.get("tx"), d.get("title"))) except Exception: pass con.execute("BEGIN") # audit 3.10: atomic pair — a crash if rows: # between the two inserts re-ingested con.executemany("INSERT INTO trades VALUES (?,?,?,?,?,?,?,?,?)", rows) + if aux: + con.executemany("INSERT INTO aux VALUES (?,?,?,?)", aux) con.execute("INSERT INTO ingested VALUES (?,?,?)", - [s, len(rows), int(time.time())]) + [s, len(rows) + len(aux), int(time.time())]) con.execute("COMMIT") box(f"rm {SEG}/{s}") - total += len(rows) - print(f"[ingest] {s}: {len(rows)} rows") + total += len(rows) + len(aux) + print(f"[ingest] {s}: {len(rows)} trades + {len(aux)} aux") n = con.execute("SELECT count(*) FROM trades").fetchone()[0] print(f"[ingest] +{total} rows · rtds.duckdb now {n:,} trades") diff --git a/recorder/recorder.py b/recorder/recorder.py index 0fe0c0be..ea0a6740 100644 --- a/recorder/recorder.py +++ b/recorder/recorder.py @@ -24,8 +24,15 @@ import websocket URL = "wss://ws-live-data.polymarket.com" DIR = os.environ.get("TAPE_DIR", "/data/segments") +# FULL capture (2026-07-19, probed): activity/* = trades + orders_matched +# (the maker side of every match, ~1.1k/min), comments (market chatter, +# tiny), crypto_prices (~300/min). rfq + prices probed dead. Trades keep the +# rtds_* segment family (ingest pipeline unchanged); everything else tapes +# raw into a parallel aux_* family. SUB = json.dumps({"action": "subscribe", "subscriptions": [ - {"topic": "activity", "type": "trades", "filters": ""}]}) + {"topic": "activity", "type": "*", "filters": ""}, + {"topic": "comments", "type": "*", "filters": ""}, + {"topic": "crypto_prices", "type": "*", "filters": ""}]}) def log(m): @@ -40,22 +47,22 @@ class Tape: def __init__(self): os.makedirs(DIR, exist_ok=True) - self.hour, self.fh, self.n = None, None, 0 - self.msgs = self.dupes = self.gaps = 0 + self.files = {} # family -> [hour, fh, rows] + self.msgs = self.aux = self.dupes = self.gaps = 0 self.last_msg = time.time() self._wlock = threading.Lock() self._recent = {} # dedupe key -> ts - def _rotate(self, hour): - if self.fh: - self.fh.close() - plain = os.path.join(DIR, f"rtds_{self.hour}.jsonl") + def _rotate(self, fam, hour): + cur = self.files.get(fam) + if cur and cur[1]: + cur[1].close() + plain = os.path.join(DIR, f"{fam}_{cur[0]}.jsonl") with open(plain, "rb") as i, gzip.open(plain + ".gz", "wb") as o: shutil.copyfileobj(i, o) os.remove(plain) - log(f"rotated {self.hour} ({self.n} rows)") - self.hour, self.n = hour, 0 - self.fh = open(os.path.join(DIR, f"rtds_{hour}.jsonl"), "a") + log(f"rotated {fam}_{cur[0]} ({cur[2]} rows)") + self.files[fam] = [hour, open(os.path.join(DIR, f"{fam}_{hour}.jsonl"), "a"), 0] # disk guard: drop oldest closed segments past 85% usage try: st = os.statvfs(DIR) @@ -69,25 +76,48 @@ class Tape: except Exception: pass + def _dedup(self, key, now): + if key in self._recent: + self.dupes += 1 + return True + self._recent[key] = now + if len(self._recent) > 80000: # ~2 busy minutes; prune old + cut = now - 120 + self._recent = {k: t for k, t in self._recent.items() if t > cut} + return False + + def _fh(self, fam): + hour = time.strftime("%Y%m%d_%H", time.gmtime()) + if fam not in self.files or self.files[fam][0] != hour: + self._rotate(fam, hour) + return self.files[fam] + + def write_aux(self, topic, ty, p, raw_len): + """Everything that isn't a trade: raw payload, schema-free.""" + now = time.time() + with self._wlock: + self.last_msg = now + if self._dedup((topic, ty, hash(json.dumps(p, sort_keys=True))), now): + return + cur = self._fh("aux") + cur[1].write(json.dumps({"ts": round(now, 3), "topic": topic, + "type": ty, "payload": p}) + "\n") + cur[2] += 1 + self.aux += 1 + def write(self, p): key = (p.get("transactionHash"), p.get("asset"), p.get("side"), str(p.get("size")), str(p.get("price"))) now = time.time() with self._wlock: self.last_msg = now # either socket delivering = stream alive - if key in self._recent: - self.dupes += 1 + if self._dedup(key, now): return - self._recent[key] = now - if len(self._recent) > 40000: # ~2 busy minutes; prune old - cut = now - 120 - self._recent = {k: t for k, t in self._recent.items() if t > cut} self._do_write(p, now) def _do_write(self, p, now): - hour = time.strftime("%Y%m%d_%H", time.gmtime()) - if hour != self.hour: - self._rotate(hour) + cur = self._fh("rtds") + self.fh, self.n = cur[1], cur[2] ts = p.get("timestamp") or 0 if ts > 1e12: ts /= 1000.0 @@ -98,7 +128,7 @@ class Tape: "side": p.get("side"), "price": p.get("price"), "size": p.get("size"), "tx": p.get("transactionHash"), "title": str(p.get("title") or "")[:60]}) + "\n") - self.n += 1 + cur[2] += 1 self.msgs += 1 @@ -138,8 +168,11 @@ def run_conn(tag, tape): m = json.loads(raw) except Exception: return - if m.get("topic") == "activity" and m.get("type") == "trades": + topic, ty = m.get("topic"), m.get("type") + if topic == "activity" and ty == "trades": tape.write(m.get("payload") or {}) + elif topic: # orders_matched / comments / prices + tape.write_aux(topic, ty, m.get("payload") or {}, len(raw)) try: app = websocket.WebSocketApp(URL, on_open=on_open, on_message=on_message) @@ -157,12 +190,14 @@ def main(): tape = Tape() for tag in ("a", "b"): # dual-connection capture threading.Thread(target=run_conn, args=(tag, tape), daemon=True).start() - last = lastd = 0 + last = lastd = lasta = 0 while True: time.sleep(60) - log(f"tape: {tape.msgs - last} msg/min · dupes {tape.dupes - lastd}/min " - f"· hour rows {tape.n} · reconnects {tape.gaps}") - last, lastd = tape.msgs, tape.dupes + tr = tape.files.get("rtds"); ax = tape.files.get("aux") + log(f"tape: {tape.msgs - last} trades/min · aux {tape.aux - lasta}/min " + f"· dupes {tape.dupes - lastd}/min · hour {tr[2] if tr else 0}+" + f"{ax[2] if ax else 0} · reconnects {tape.gaps}") + last, lastd, lasta = tape.msgs, tape.dupes, tape.aux if __name__ == "__main__":