Files
winning-wallet-finder_github/recorder/ingest.py
T

127 lines
5.4 KiB
Python

#!/usr/bin/env python3
"""Nightly tape ingest (daily.sh): pull closed RTDS segments off the
wwf-recorder volume into live/rtds.duckdb, delete on the box ONLY after the
inserted row count matches. Own DB file on purpose — must never contend with
cache.duckdb's single-writer lock (2026-07-17 collision). Transport is
`flyctl ssh console` + base64 (no ingress on the recorder by design)."""
import base64
import gzip
import json
import os
import shutil
import ssl
import subprocess
import time
import urllib.request
HERE = os.path.dirname(os.path.abspath(__file__))
DB = os.path.join(HERE, "..", "live", "rtds.duckdb")
APP, SEG = "wwf-recorder", "/data/segments"
def ping(msg):
"""Best-effort Discord ping to the DAILY channel (user ask 2026-07-19:
the nightly pull announces start + finish). Never fatal."""
try:
hook = json.load(open(os.path.join(HERE, "..", "config.json"))).get("daily_webhook")
if not hook:
return
req = urllib.request.Request(hook, data=json.dumps({"content": msg}).encode(),
headers={"Content-Type": "application/json",
"User-Agent": "Mozilla/5.0"})
urllib.request.urlopen(req, timeout=10,
context=ssl._create_unverified_context()).read()
except Exception:
pass
# launchd runs daily.sh with a minimal PATH (no /opt/homebrew/bin) — the
# 2026-07-18 run died on FileNotFoundError: 'flyctl'. Resolve it explicitly.
FLYCTL = shutil.which("flyctl") or "/opt/homebrew/bin/flyctl"
def box(cmd):
r = subprocess.run([FLYCTL, "ssh", "console", "-a", APP, "-C",
f"bash -c '{cmd}'"], capture_output=True, text=True,
timeout=900) # busy-hour segments (~15MB gz) outgrow 300s over ssh
return r.stdout
def main():
import duckdb
con = duckdb.connect(DB)
con.execute("""CREATE TABLE IF NOT EXISTS trades(
ts DOUBLE, wallet VARCHAR, asset VARCHAR, cond VARCHAR, side VARCHAR,
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]
ping(f"📼 tape ingest started — {len(segs)} segment(s) queued")
total = 0
for s in sorted(segs):
# closes #11: sftp moves the gz BINARY in one session (base64-over-
# console was +33% bytes and needed 900s timeouts at busy-hour sizes);
# the old path stays as the fallback because sftp exits 0 even on
# some failures — the gunzip is the integrity check either way.
lines = None
local = os.path.join("/tmp", s)
try:
subprocess.run([FLYCTL, "ssh", "sftp", "get", f"{SEG}/{s}", local,
"-a", APP], capture_output=True, timeout=900)
with gzip.open(local, "rb") as fh:
lines = fh.read().decode().splitlines()
except Exception:
lines = None
finally:
try:
os.remove(local)
except OSError:
pass
if lines is None:
try:
raw = box(f"base64 {SEG}/{s}")
lines = gzip.decompress(base64.b64decode(raw)).decode().splitlines()
except Exception as e:
# BOTH transports failed — skip, never crash: a single monster
# segment (rtds_20260718_21, 2026-07-19) blocked the whole
# sorted backlog behind it when this timeout escaped the try
print(f"[ingest] {s}: both transports failed ({type(e).__name__}) "
"— left on box, continuing")
continue
rows, aux = [], []
for ln in lines:
try:
d = json.loads(ln)
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) + len(aux), int(time.time())])
con.execute("COMMIT")
box(f"rm {SEG}/{s}")
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")
ping(f"📼 tape ingest done: +{total:,} rows across {len(segs)} segment(s) "
f"· rtds.duckdb now {n:,} trades")
if __name__ == "__main__":
main()