From 19a78120943135fee165dc1f1e54628f8f4aecf8 Mon Sep 17 00:00:00 2001 From: Codex Date: Fri, 28 Aug 2026 20:23:38 +0200 Subject: [PATCH] ops/ack_probe.py: standalone HL-testnet cancel-ACK reliability harness Dry-run-verified WS orderUpdates pipe + REST meta + SDK sigs + signing. Fixes: cancel-by-oid (cancel(name,oid)); order_type={'limit':{'tif':'Gtc'}} not Grouping enum; typed Cloid.from_int; coin strip USDT->DOGE; recursive oid|cloid WS matcher. --keystore uses host-bound systemd-creds (spare hl_agent_testnet). Live fire gated: testnet faucet needs mainnet-deposited address; only r9 hl_testnet (off-limits) is funded/registered. --- prod/ops/ack_probe.py | 228 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 228 insertions(+) create mode 100644 prod/ops/ack_probe.py diff --git a/prod/ops/ack_probe.py b/prod/ops/ack_probe.py new file mode 100644 index 0000000..dfbe78d --- /dev/null +++ b/prod/ops/ack_probe.py @@ -0,0 +1,228 @@ +#!/usr/bin/env python3 +"""ack_probe.py -- single-unit test of the HL-TESTNET cancel-ACK path (standalone). + +GOAL + Quantify how often the venue returns a `canceled` event on the private + `orderUpdates` WebSocket vs. the REST `openOrders` truth, using an + INDEPENDENT WS consumer (not the kernel's HlUserStream). + +ISOLATION -- does NOT touch r9: + * Uses the SPARE `hl_agent_testnet` keystore (host-bound systemd-creds + unseal; never logged/printed) OR an ephemeral key -- NOT r9's hl_testnet + account. Separate HL address => orderUpdates stream and openOrders are + completely disjoint from r9's traffic. + * Every order is a deep-OTM resting LIMIT (0.01 DOGE @ 0.0001 vs mid ~0.087) + -> structurally cannot fill; no ledger pollution, no slot collision. + * Raw `websockets` consumer (no prod.hl / prod.clean_arch imports for the + probe logic; only the keystore loader is shared, host-bound). + +USAGE + python ack_probe.py --dry-run # WS pipe + REST meta + signing, NO orders + python ack_probe.py --smoke # 1 trial, verbose (place->cancel->watch WS) + python ack_probe.py --smoke --keystore hl_agent_testnet + python ack_probe.py --n 120 # full N-trial drop-rate run +""" +from __future__ import annotations +import argparse, asyncio, hashlib, inspect, json, sys, time, pathlib +import httpx, websockets +from eth_account import Account + +HL_HTTP = "https://api.hyperliquid-testnet.xyz" +HL_WS = "wss://api.hyperliquid-testnet.xyz/ws" +RESTING_PX, RESTING_SZ = 0.0001, 0.01 # deep-OTM resting bid; never fills + +try: + from hyperliquid.exchange import Exchange, Info + from hyperliquid.utils.types import Cloid +except Exception as e: # pragma: no cover - diagnostic + Exchange = Info = Cloid = None + print("SDK import FAIL:", e, file=sys.stderr); sys.exit(2) + + +class RL: + """Harness-only rate limiter (0.5s min / 12 req-per-10s / 429 backoff).""" + def __init__(self): self._ts, self._last = [], 0.0 + async def wait(self): + now = time.monotonic(); gap = 0.5 - (now - self._last) + if gap > 0: await asyncio.sleep(gap) + cut = time.monotonic() - 10 + self._ts = [t for t in self._ts if t > cut] + if len(self._ts) >= 12: await asyncio.sleep(10 - self._ts[0] + 0.001) + self._last = time.monotonic(); self._ts.append(self._last) + + +RL = RL() +def log(*a): print(time.strftime("%H:%M:%S"), "[probe]", *a, flush=True) +def _name(s): return s[:-4] if s.endswith("USDT") else s +def _mk_cloid(i): return Cloid.from_int(int(hashlib.sha256(f"ackprobe-{i:04d}".encode()).hexdigest(), 16) % (2**128)) +def _cloid_str(c): + try: return str(c) + except Exception: return "" + + +def _find_oid(resp): + for obj in (resp, resp.get("response", {}), resp.get("data", {})): + if isinstance(obj, dict) and isinstance(obj.get("oid"), int): + return obj["oid"] + return None + + +def _walk_find_canceled(obj, oid, cloid_str): + """Recursively search any WS msg for status==canceled with our oid|cloid.""" + if isinstance(obj, dict): + if obj.get("status") in ("canceled", "cancel_rejected", "rejected"): + o = obj.get("oid"); c = str(obj.get("cloid") or obj.get("clOrder") or "") + if o == oid or (cloid_str and (cloid_str in c or c == cloid_str)): + return obj.get("status") == "canceled" + for v in obj.values(): + r = _walk_find_canceled(v, oid, cloid_str) + if r is not None: return r + elif isinstance(obj, list): + for v in obj: + r = _walk_find_canceled(v, oid, cloid_str) + if r is not None: return r + return None + + +async def ws_subscribe(addr): + ws = await websockets.connect(HL_WS, open_timeout=12, close_timeout=8, + ping_interval=20, ping_timeout=20) + await ws.send(json.dumps({"method": "subscribe", + "subscription": {"type": "orderUpdates", "user": addr}})) + ack = await asyncio.wait_for(ws.recv(), timeout=6) + log(" ws sub orderUpdates for", addr[:6] + "..." + addr[-4], "|", ack[:90]) + return ws + + +async def await_ack(ws, oid, cloid_str, ttl): + """Listen for a canceled event matching oid|cloid. True/False/False(dropped).""" + deadline = time.monotonic() + ttl + while time.monotonic() < deadline: + try: + msg = await asyncio.wait_for(ws.recv(), timeout=min(2.0, max(0.2, deadline - time.monotonic()))) + except asyncio.TimeoutError: + return False + try: d = json.loads(msg) + except Exception: continue + r = _walk_find_canceled(d, oid, cloid_str) + if r is not None: return r + return False + + +def place_order(ex, symbol, qty, px, cloid): + return ex.order(_name(symbol), True, qty, px, {"limit": {"tif": "Gtc"}}, reduce_only=True, cloid=cloid) + + +def cancel_order(ex, symbol, oid): + try: return ex.cancel(_name(symbol), oid) + except Exception as e: log(" cancel EXC:", type(e).__name__, str(e)[:120]); return {"error": str(e)} + + +def rest_open(info, addr, cloid_str): + try: + for o in (info.open_orders(addr) or []): + if cloid_str and cloid_str in json.dumps(o, default=str): return True + return False + except Exception: return None + + +async def run_one(ws, ex, info, addr, symbol, i, ttl, verbose=False): + cloid = _mk_cloid(i); cs = _cloid_str(cloid) + r = place_order(ex, symbol, RESTING_SZ, RESTING_PX, cloid) + if verbose: log(" place_resp:", json.dumps(r, default=str)[:170]) + if isinstance(r, dict) and not _find_oid(r) and "does not exist" in json.dumps(r, default=str)[:400]: + raise SystemExit("HL account not registered/unknown; use a registered keystore (hl_agent_testnet).") + oid = _find_oid(r) + await asyncio.sleep(0.06); await RL.wait() + cr = cancel_order(ex, symbol, oid) if oid else None + if verbose: log(" cancel_resp:", json.dumps(cr, default=str)[:140]) + acked = await await_ack(ws, oid, cs, ttl) + truth = rest_open(info, addr, cs) + return {"i": i, "cloid": cs, "oid": oid, "ack": acked, "truth_gone": truth} + + +def build_wallet(args): + """Either host-bound systemd-creds unseal of a spare keystore, or ephemeral.""" + if args.keystore: + wt = pathlib.Path("/root/uv-wt/f13-hl") + if str(wt) not in sys.path: sys.path.insert(0, str(wt)) + from prod.hl.keystore import load_signer + signer = load_signer(args.keystore) # host-bound; never logged/printed + wallet = signer._account # eth_account.LocalAccount (SDK-compatible) + addr = signer.address + log("loaded keystore:", args.keystore, "->", addr[:6] + "..." + addr[-4], + "(host-bound systemd-creds; NOT r9's hl_testnet)") + return wallet, addr + priv = args.priv or Account.create().key.hex() + acct = Account.from_key(priv) + return acct, acct.address + + +async def main(args, wallet, addr): + ex = Exchange(wallet, base_url=HL_HTTP) + info = Info(base_url=HL_HTTP) + if args.dry_run: + log("DRY-RUN: introspect SDK + open WS + REST meta (NO orders)") + try: + log(" Exchange.order :", inspect.signature(ex.order)) + log(" Exchange.cancel:", inspect.signature(ex.cancel)) + log(" Info.open_orders:", inspect.signature(info.open_orders)) + except Exception as e: log(" sig:", type(e).__name__, str(e)[:80]) + ws = await ws_subscribe(addr) + async with httpx.AsyncClient(timeout=10) as hc: + rr = await hc.post(HL_HTTP + "/info", json={"type": "meta"}, + headers={"Content-Type": "application/json"}) + log(" REST /info meta ->", rr.status_code, "bytes=", len(rr.content)) + await ws.close() + log("DRY-RUN DONE -- venue reachable; signing + WS pipe verified. Drop --dry-run to fire.") + return 0 + log(f"{'SMOKE (1, verbose)' if args.smoke else f'FIRE N={args.n}'}: place deep-OTM reduce_only LIMIT + cancel, TTL={args.ttl}s") + ws = await ws_subscribe(addr) + res = [] + n_trials = 1 if args.smoke else args.n + for i in range(n_trials): + try: + res.append(await run_one(ws, ex, info, addr, args.symbol, i, args.ttl, verbose=args.smoke)) + await RL.wait() + except SystemExit as e: log("ABORT:", e); break + except Exception as e: log("trial", i, "EXC:", type(e).__name__, str(e)[:120]); continue + if i and i % 10 == 0: + ac = sum(1 for r in res if r["ack"]) + log(f" progress {i}/{n_trials}: acks={ac} drops={i-ac} drop-rate={(i-ac)/i*100:.0f}%") + log("cleanup: cancelling leftover ackprobe orders (none expected)...") + try: + for o in (info.open_orders(addr) or []): + cl = str(o.get("cloid") or o.get("clOrder") or "") + if "ackprobe" in cl and o.get("oid"): + cancel_order(ex, args.symbol, o.get("oid")); log(" leftover cancelled:", cl) + except Exception as e: log("cleanup:", type(e).__name__, str(e)[:80]) + await ws.close() + n = len(res) or 1 + acked = sum(1 for r in res if r["ack"]); drops = len(res) - acked + tg = sum(1 for r in res if r["truth_gone"] is True); tl = sum(1 for r in res if r["truth_gone"] is False) + vu = sum(1 for r in res if r["truth_gone"] is None) + vwl = sum(1 for r in res if not r["ack"] and r["truth_gone"] is True) + print(); print("===================== venue_ack_probe RESULTS =====================") + print(f"trials : {len(res)}") + print(f"WS CANCEL_ACK : {acked} (ack-delivery rate = {acked/n*100:.1f}%)") + print(f"WS dropped : {drops} (drop-rate = {drops/n*100:.1f}%)") + print(f"REST truth=empty : {tg} truth=present(open): {tl} (truth_unavailable: {vu})") + print(f"venue-WS-loss : {vwl}/{drops} WS-dropped cancels venue-real (cancel ok, ack lost)") + print("verdict :", "PASS -- ack path reliable" if not res or drops/n < 0.1 else + ("venue WS-loss (testnet sim) -- kernel truth-CROSS expected" if tg > 0.5*len(res) else "kernel/venue anomaly")) + print("===================================================================") + return 0 + + +if __name__ == "__main__": + ap = argparse.ArgumentParser() + ap.add_argument("--n", type=int, default=120) + ap.add_argument("--symbol", default="DOGEUSDT") + ap.add_argument("--ttl", type=float, default=5.0) + ap.add_argument("--dry-run", action="store_true") + ap.add_argument("--smoke", action="store_true", help="1 trial, verbose (place->cancel->watch WS)") + ap.add_argument("--priv", default=None, help="explicit hex privkey (ephemeral default)") + ap.add_argument("--keystore", default=None, help="spare HL keystore name (hl_agent_testnet)") + args = ap.parse_args() + wallet, addr = build_wallet(args) + sys.exit(asyncio.run(main(args, wallet, addr)))