#!/usr/bin/env python3 """ quick_trade.py — 多市场一键交易面板 启动: python quick_trade.py 打开: http://localhost:8890 快捷键: ↑ = 买 UP ↓ = 买 DOWN Q = 卖 UP ALL W = 卖 DOWN ALL 1-5 = 选金额 ($5/$10/$20/$50/$100) """ import asyncio import json import math import os import sys import time import threading import webbrowser from collections import deque from datetime import datetime, timezone import aiohttp from aiohttp import web import requests from dotenv import load_dotenv from web3 import Web3 load_dotenv() # ── 代理检测 (国内 Polymarket/Alchemy 被墙) ───────────────────────────────── def _detect_proxy(): import socket for key in ("PROXY_URL", "HTTPS_PROXY", "https_proxy", "ALL_PROXY", "all_proxy", "HTTP_PROXY", "http_proxy"): val = os.environ.get(key, "").strip() if val: return val for port in (9098, 9099, 7890, 7897, 1080, 8080, 8889, 10809): try: s = socket.create_connection(("127.0.0.1", port), timeout=0.3) s.close() return f"http://127.0.0.1:{port}" except (OSError, socket.timeout): continue return None _PROXY = _detect_proxy() _PROXIES: dict = {"http": _PROXY, "https": _PROXY} if _PROXY else {} if _PROXY: print(f"[Proxy] {_PROXY}") # 同步到标准 env 变量, 让 py_clob_client_v2 / requests / aiohttp 等三方库自动走代理 os.environ["HTTP_PROXY"] = _PROXY os.environ["HTTPS_PROXY"] = _PROXY # ── 多账号配置 ───────────────────────────────────────────────────────────────── CHAIN_ID = 137 SIGNATURE_TYPE = 3 # Polymarket Deposit Wallet (POLY_1271, V2 新钱包类型) CLOB_HTTP = "https://clob.polymarket.com" CTF_ADDR = "0x4D97DCd97eC945f40cF65F87097ACe5EA0476045" def _load_accounts(): """从 .env 加载所有 QUICK 账号: QUICK_PRIVATE_KEY / QUICK_FUNDER (第1个), QUICK_PRIVATE_KEY_2 / QUICK_FUNDER_2 (第2个), ...""" accts = [] # 第 1 个 (无编号) pk = os.getenv("QUICK_PRIVATE_KEY", "") fd = os.getenv("QUICK_FUNDER", "") rk = os.getenv("QUICK_RELAYER_API_KEY", "") ra = os.getenv("QUICK_RELAYER_API_KEY_ADDRESS", "") if pk and fd: accts.append({"private_key": pk, "funder": fd, "relayer_key": rk, "relayer_addr": ra, "label": fd[:6] + ".." + fd[-4:]}) # 第 2~N 个 (有编号, 兼容 _2 和 _02 格式) for i in range(2, 20): pk = os.getenv(f"QUICK_PRIVATE_KEY_{i}", "") or os.getenv(f"QUICK_PRIVATE_KEY_{i:02d}", "") fd = os.getenv(f"QUICK_FUNDER_{i}", "") or os.getenv(f"QUICK_FUNDER_{i:02d}", "") rk = os.getenv(f"QUICK_RELAYER_API_KEY_{i}", "") or os.getenv(f"QUICK_RELAYER_API_KEY_{i:02d}", "") ra = os.getenv(f"QUICK_RELAYER_API_KEY_ADDRESS_{i}", "") or os.getenv(f"QUICK_RELAYER_API_KEY_ADDRESS_{i:02d}", "") if pk and fd: accts.append({"private_key": pk, "funder": fd, "relayer_key": rk, "relayer_addr": ra, "label": fd[:6] + ".." + fd[-4:]}) else: break return accts ACCOUNTS = _load_accounts() _active_account_idx = 0 _account_lock = threading.Lock() # 便捷访问当前账号 def _current_account(): with _account_lock: idx = _active_account_idx if idx < len(ACCOUNTS): return ACCOUNTS[idx] return {"private_key": "", "funder": "", "relayer_key": "", "relayer_addr": "", "label": "N/A"} # 兼容旧代码 PRIVATE_KEY = ACCOUNTS[0]["private_key"] if ACCOUNTS else "" FUNDER = ACCOUNTS[0]["funder"] if ACCOUNTS else "" # ── 多市场配置 ───────────────────────────────────────────────────────────────── MARKETS = { "btc_5m": {"label": "BTC 5m", "slug_prefix": "btc-updown-5m", "duration": 300, "symbol": "BTC", "variant": "fiveminute", "rtds": "btc/usd"}, "btc_15m": {"label": "BTC 15m", "slug_prefix": "btc-updown-15m", "duration": 900, "symbol": "BTC", "variant": "fifteenminute", "rtds": "btc/usd"}, "btc_1h": {"label": "BTC 1h", "slug_prefix": "bitcoin", "duration": 3600, "symbol": "BTC", "variant": "onehour", "rtds": "btc/usd", "slug_type": "hourly"}, "eth_5m": {"label": "ETH 5m", "slug_prefix": "eth-updown-5m", "duration": 300, "symbol": "ETH", "variant": "fiveminute", "rtds": "eth/usd"}, "eth_15m": {"label": "ETH 15m", "slug_prefix": "eth-updown-15m", "duration": 900, "symbol": "ETH", "variant": "fifteenminute", "rtds": "eth/usd"}, "eth_1h": {"label": "ETH 1h", "slug_prefix": "ethereum", "duration": 3600, "symbol": "ETH", "variant": "onehour", "rtds": "eth/usd", "slug_type": "hourly"}, } active_market = "btc_5m" active_market_lock = threading.Lock() PORT = 8890 ALCHEMY_KEY = os.getenv("ALCHEMY_KEY", "") RPC_URL = f"https://polygon-mainnet.g.alchemy.com/v2/{ALCHEMY_KEY}" POLY_BUILDER_CODE = "0xe7bd311abb3a706580f8f5dea993360e4805582f0a7cceb3d1a3abf87bbf7697" # 滑点配置 BUY_SLIPPAGE = 0.10 # 买入滑点 (10 cents) SELL_SLIPPAGE = 0.10 # 卖出滑点 (10 cents) # ── 全局状态 ────────────────────────────────────────────────────────────────── state_lock = threading.Lock() state = { "round_ts": 0, "round_end": 0, "token_ids": {}, "strike": 0.0, "crypto_current": 0.0, # BTC 或 ETH 现价 "up_price": 0.0, "down_price": 0.0, "up_ask": 0.0, "down_ask": 0.0, "up_position": 0.0, "down_position": 0.0, "usdc_balance": 0.0, "keys_ok": bool(ACCOUNTS), "up_book": {"bids": [], "asks": []}, "down_book": {"bids": [], "asks": []}, } # RTDS 推送缓存: symbol -> price _rtds_prices = {"btc/usd": 0.0, "eth/usd": 0.0} trade_log = deque(maxlen=50) # ── Web3 连接 ───────────────────────────────────────────────────────────────── _w3 = None def get_w3(): global _w3 if _w3: return _w3 try: from web3.middleware import ExtraDataToPOAMiddleware as POA except ImportError: try: from web3.middleware import geth_poa_middleware as POA except ImportError: POA = None _rpc_session = requests.Session() if _PROXIES: _rpc_session.proxies.update(_PROXIES) _w3 = Web3(Web3.HTTPProvider( RPC_URL, request_kwargs={"timeout": 10}, session=_rpc_session, )) if POA: try: _w3.middleware_onion.inject(POA, layer=0) except Exception: pass return _w3 # ── CLOB Client (多账号, 延迟初始化, 凭证缓存) ───────────────────────────────── _clob = None _clob_lock = threading.Lock() _clob_account_idx = -1 # 当前 _clob 对应的账号索引 _CREDS_DIR = os.path.dirname(os.path.abspath(__file__)) def _creds_file_for(funder): return os.path.join(_CREDS_DIR, f".quick_trade_creds_{funder[:8].lower()}.json") def _load_cached_creds(funder): """从本地文件加载缓存的 API 凭证""" try: path = _creds_file_for(funder) if not os.path.exists(path): return None with open(path, "r") as f: data = json.load(f) if data.get("funder") != funder: return None from py_clob_client_v2.clob_types import ApiCreds return ApiCreds( api_key=data["api_key"], api_secret=data["api_secret"], api_passphrase=data["api_passphrase"], ) except Exception: return None def _save_creds(funder, creds): """缓存 API 凭证到本地文件""" try: data = { "funder": funder, "api_key": creds.api_key, "api_secret": creds.api_secret, "api_passphrase": creds.api_passphrase, } with open(_creds_file_for(funder), "w") as f: json.dump(data, f) except Exception: pass def get_clob_client(): global _clob, _clob_account_idx acct = _current_account() pk, funder = acct["private_key"], acct["funder"] with _account_lock: idx = _active_account_idx with _clob_lock: if _clob is not None and _clob_account_idx == idx: return _clob if not pk or not funder: raise RuntimeError("账号未配置") from py_clob_client_v2.client import ClobClient _clob = ClobClient( host=CLOB_HTTP, key=pk, chain_id=CHAIN_ID, signature_type=SIGNATURE_TYPE, funder=funder, ) cached = _load_cached_creds(funder) if cached: _clob.set_api_creds(cached) print(f"✅ CLOB client [{acct['label']}] (缓存凭证)") else: creds = _clob.create_or_derive_api_key() _clob.set_api_creds(creds) _save_creds(funder, creds) print(f"✅ CLOB client [{acct['label']}] (首次签名, 已缓存)") _clob_account_idx = idx return _clob # ── 回合信息 ────────────────────────────────────────────────────────────────── def _make_hourly_slug(crypto_name: str, round_ts: int) -> str: """生成 1h 市场的 slug, 如 bitcoin-up-or-down-april-25-2026-11am-et""" from datetime import timedelta # round_ts 是 UTC 整小时, 转换为 ET (UTC-4) et_dt = datetime.fromtimestamp(round_ts, tz=timezone.utc) - timedelta(hours=4) month = et_dt.strftime("%B").lower() # april day = et_dt.day year = et_dt.year hour = et_dt.hour ampm = "am" if hour < 12 else "pm" h12 = hour % 12 if h12 == 0: h12 = 12 return f"{crypto_name}-up-or-down-{month}-{day}-{year}-{h12}{ampm}-et" def get_round_token_ids(round_ts: int, slug_prefix: str = "btc-updown-5m", slug_type: str = "") -> dict: """获取回合的 UP/DOWN token_id""" if slug_type == "hourly": slug = _make_hourly_slug(slug_prefix, round_ts) else: slug = f"{slug_prefix}-{round_ts}" try: r = requests.get( "https://gamma-api.polymarket.com/events", params={"slug": slug}, timeout=10, proxies=_PROXIES, ) events = r.json() if not events or not isinstance(events, list): return {} mkt = events[0].get("markets", [{}])[0] tokens = mkt.get("clobTokenIds", "[]") outcomes = mkt.get("outcomes", "[]") if isinstance(tokens, str): tokens = json.loads(tokens) if isinstance(outcomes, str): outcomes = json.loads(outcomes) result = {} for i, o in enumerate(outcomes): if i < len(tokens): result[o.lower()] = tokens[i] return result except Exception: return {} def _robust_get_json(url, timeout=5, headers=None) -> dict: """ 发送 GET 请求并解析为 JSON。 首先尝试使用 requests.get,若失败或遭遇非 200 状态码, 则自动 fallback 使用系统 curl 命令抓取以绕过 Cloudflare WAF 拦截。 """ try: r = requests.get(url, timeout=timeout, verify=False, headers=headers, proxies=_PROXIES) if r.status_code == 200: return r.json() except Exception: pass # Fallback: 使用系统 curl import subprocess, json try: headers_list = [] if headers: for k, v in headers.items(): headers_list.extend(["-H", f"{k}: {v}"]) cmd = (["curl", "-s", "-x", _PROXY] + headers_list + [url]) if _PROXY else (["curl", "-s"] + headers_list + [url]) res = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout) if res.returncode == 0: return json.loads(res.stdout.strip()) except Exception: pass return None def get_strike_price(round_ts: int, duration: int = 300, symbol: str = "BTC", variant: str = "fiveminute") -> tuple: """获取本回合的锚定价。 返回 (price, confirmed): confirmed=True 表示数据已稳定可缓存。 注: 统一用 variant=fiveminute, 因为 fifteenminute/onehour variant 返回的是整点 openPrice 而非该回合真实起始价。 """ try: start = datetime.fromtimestamp(round_ts, tz=timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") end = datetime.fromtimestamp(round_ts + duration, tz=timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") url = (f"https://polymarket.com/api/crypto/crypto-price?symbol={symbol}" f"&eventStartTime={start}&variant=fiveminute&endDate={end}") data = _robust_get_json(url, timeout=5, headers={"User-Agent": "Mozilla/5.0"}) if data: op = data.get("openPrice") incomplete = data.get("incomplete", False) if op is not None and float(op) > 0: # incomplete=True 表示数据可能不准, 需要继续重试 return float(op), not incomplete except Exception: pass return 0, False # ── Token 余额查询 ─────────────────────────────────────────────────────────── _ctf_abi = json.loads( '[{"inputs":[{"name":"account","type":"address"},{"name":"id","type":"uint256"}],' '"name":"balanceOf","outputs":[{"name":"","type":"uint256"}],' '"type":"function","constant":true}]' ) def get_token_balance(w3, token_id: str) -> float: """查询当前账号 CTF token 持仓 (shares)""" try: funder = _current_account()["funder"] ctf = w3.eth.contract(address=w3.to_checksum_address(CTF_ADDR), abi=_ctf_abi) bal = ctf.functions.balanceOf(w3.to_checksum_address(funder), int(token_id)).call() return bal / 1e6 except Exception: return 0 # ══════════════════════════════════════════════════════════════════════════════ # 交易执行 # ══════════════════════════════════════════════════════════════════════════════ def execute_buy(side: str, amount_usd: float) -> dict: """市价买入 (FOK)""" ts = time.time() with state_lock: token_id = state["token_ids"].get(side) ask = state.get(f"{side}_ask", 0) bid = state.get(f"{side}_price", 0) if not token_id: return {"success": False, "message": "回合未就绪,无 token"} price = ask if ask > 0 else bid if price <= 0: return {"success": False, "message": "无可用盘口价格"} try: client = get_clob_client() from py_clob_client_v2.clob_types import MarketOrderArgs, OrderType from py_clob_client_v2.order_builder.constants import BUY buy_price = round(min(price + BUY_SLIPPAGE, 0.99), 2) args = MarketOrderArgs( token_id=token_id, amount=round(amount_usd, 2), price=buy_price, side=BUY, order_type=OrderType.FOK, builder_code=POLY_BUILDER_CODE, ) signed = client.create_market_order(args) resp = client.post_order(signed, OrderType.FOK) elapsed_ms = int((time.time() - ts) * 1000) order_id = resp.get("orderID", resp.get("id", "")) shares = round(amount_usd / price, 1) msg = f"BUY {side.upper()} {shares}sh @ {price:.2f} (${amount_usd})" _log_trade(True, msg, elapsed_ms) return {"success": True, "message": msg, "order_id": order_id, "elapsed_ms": elapsed_ms} except Exception as e: elapsed_ms = int((time.time() - ts) * 1000) msg = f"BUY {side.upper()} 失败: {str(e)[:100]}" _log_trade(False, msg, elapsed_ms) return {"success": False, "message": msg, "elapsed_ms": elapsed_ms} def execute_sell(side: str, amount) -> dict: """市价卖出 (FOK). amount='all' 或 USD 金额.""" ts = time.time() with state_lock: token_id = state["token_ids"].get(side) bid = state.get(f"{side}_price", 0) position = state.get(f"{side}_position", 0) if not token_id: return {"success": False, "message": "回合未就绪,无 token"} if position <= 0.01: return {"success": False, "message": "无持仓可卖"} # 确定卖出份额 if amount == "all": # 从链上读取最新余额,避免缓存不准 w3 = get_w3() actual = get_token_balance(w3, token_id) shares = math.floor(actual * 100) / 100 else: usd = float(amount) price = bid if bid > 0 else 0.5 shares = usd / price shares = min(shares, position) shares = math.floor(shares * 100) / 100 if shares <= 0: return {"success": False, "message": "可卖份额为 0"} try: client = get_clob_client() from py_clob_client_v2.clob_types import MarketOrderArgs, OrderType from py_clob_client_v2.order_builder.constants import SELL sell_price = round(max((bid if bid > 0 else 0.5) - SELL_SLIPPAGE, 0.01), 2) args = MarketOrderArgs( token_id=token_id, amount=round(shares, 2), price=sell_price, side=SELL, order_type=OrderType.FOK, builder_code=POLY_BUILDER_CODE, ) signed = client.create_market_order(args) resp = client.post_order(signed, OrderType.FOK) elapsed_ms = int((time.time() - ts) * 1000) order_id = resp.get("orderID", resp.get("id", "")) label = "ALL " if amount == "all" else "" msg = f"SELL {label}{side.upper()} {shares:.1f}sh @ {sell_price:.2f}" _log_trade(True, msg, elapsed_ms) return {"success": True, "message": msg, "order_id": order_id, "elapsed_ms": elapsed_ms} except Exception as e: elapsed_ms = int((time.time() - ts) * 1000) msg = f"SELL {side.upper()} 失败: {str(e)[:100]}" _log_trade(False, msg, elapsed_ms) return {"success": False, "message": msg, "elapsed_ms": elapsed_ms} def _log_trade(success: bool, message: str, elapsed_ms: int): """添加交易记录""" trade_log.appendleft({ "time": datetime.now().strftime("%H:%M:%S"), "success": success, "message": message, "elapsed_ms": elapsed_ms, }) def handle_trade(data: dict) -> dict: """统一交易入口""" action = data.get("action", "") side = data.get("side", "") amount = data.get("amount", 20) # cancel_all 不需要 side if action == "cancel_all_orders": return execute_cancel_all() if side not in ("up", "down"): return {"success": False, "message": "无效方向"} if action == "buy": return execute_buy(side, float(amount)) elif action == "sell": return execute_sell(side, amount) elif action == "sell_all": return execute_sell(side, "all") elif action in ("limit_buy", "limit_sell"): limit_price = data.get("price", 0) if not limit_price or float(limit_price) <= 0: return {"success": False, "message": "请输入限价"} return execute_limit_order( side=side, buy_or_sell="buy" if action == "limit_buy" else "sell", limit_price=float(limit_price), amount_usd=float(amount), ) else: return {"success": False, "message": f"未知操作: {action}"} def execute_cancel_all() -> dict: """取消所有挂单""" ts = time.time() try: client = get_clob_client() resp = client.cancel_all() elapsed_ms = int((time.time() - ts) * 1000) # resp 通常返回 {"canceled": [...], "not_canceled": [...]} canceled = resp.get("canceled", []) if isinstance(resp, dict) else [] count = len(canceled) if canceled else 0 msg = f"已取消 {count} 笔挂单" _log_trade(True, msg, elapsed_ms) return {"success": True, "message": msg, "elapsed_ms": elapsed_ms} except Exception as e: elapsed_ms = int((time.time() - ts) * 1000) msg = f"取消挂单失败: {str(e)[:100]}" _log_trade(False, msg, elapsed_ms) return {"success": False, "message": msg, "elapsed_ms": elapsed_ms} def execute_limit_order(side: str, buy_or_sell: str, limit_price: float, amount_usd: float) -> dict: """GTC 限价单""" ts = time.time() with state_lock: token_id = state["token_ids"].get(side) if not token_id: return {"success": False, "message": "回合未就绪,无 token"} limit_price = round(limit_price, 2) if limit_price <= 0 or limit_price >= 1: return {"success": False, "message": f"价格 {limit_price} 无效 (需 0.01-0.99)"} # USD → shares shares = round(amount_usd / limit_price, 2) if shares <= 0: return {"success": False, "message": "份额为 0"} try: client = get_clob_client() from py_clob_client_v2.clob_types import OrderArgs, OrderType from py_clob_client_v2.order_builder.constants import BUY, SELL args = OrderArgs( token_id=token_id, price=limit_price, size=shares, side=BUY if buy_or_sell == "buy" else SELL, builder_code=POLY_BUILDER_CODE, ) signed = client.create_order(args) resp = client.post_order(signed, OrderType.GTC) elapsed_ms = int((time.time() - ts) * 1000) order_id = resp.get("orderID", resp.get("id", "")) label = "BUY" if buy_or_sell == "buy" else "SELL" msg = f"LIMIT {label} {side.upper()} {shares:.1f}sh @ {limit_price:.2f} (${amount_usd})" _log_trade(True, msg, elapsed_ms) return {"success": True, "message": msg, "order_id": order_id, "elapsed_ms": elapsed_ms} except Exception as e: elapsed_ms = int((time.time() - ts) * 1000) label = "BUY" if buy_or_sell == "buy" else "SELL" msg = f"LIMIT {label} {side.upper()} 失败: {str(e)[:100]}" _log_trade(False, msg, elapsed_ms) return {"success": False, "message": msg, "elapsed_ms": elapsed_ms} # ══════════════════════════════════════════════════════════════════════════════ # 后台工作线程 # ══════════════════════════════════════════════════════════════════════════════ def state_worker(): """后台线程: 回合管理、锚定价、持仓、余额 (每 3 秒, 慢速任务)""" import urllib3 urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) last_round = 0 strike_confirmed_round = 0 last_market_key = "" while True: try: with active_market_lock: mkt_key = active_market mkt = MARKETS[mkt_key] duration = mkt["duration"] # 市场切换 → 强制刷新 if mkt_key != last_market_key: last_round = 0 strike_confirmed_round = 0 last_market_key = mkt_key # 更新 RTDS → state.crypto_current rtds_key = mkt["rtds"] with state_lock: state["crypto_current"] = _rtds_prices.get(rtds_key, 0.0) now = int(time.time()) round_ts = now - (now % duration) # ── 新回合检测 ── if round_ts != last_round: token_ids = get_round_token_ids(round_ts, mkt["slug_prefix"], mkt.get("slug_type", "")) if token_ids: with state_lock: state["round_ts"] = round_ts state["round_end"] = round_ts + duration state["token_ids"] = token_ids state["up_position"] = 0 state["down_position"] = 0 state["strike"] = 0 last_round = round_ts rt = datetime.fromtimestamp(round_ts).strftime("%H:%M:%S") print(f"🔄 [{mkt['label']}] 新回合 {rt} | 等待锚定价...") # ── 锚定价 (持续重试) ── with state_lock: current_round = state["round_ts"] strike_pending = False if current_round > 0 and strike_confirmed_round != current_round: strike_pending = True strike, confirmed = get_strike_price(current_round, duration, mkt["symbol"], mkt["variant"]) if strike > 0: with state_lock: state["strike"] = strike if confirmed: strike_confirmed_round = current_round rt = datetime.fromtimestamp(current_round).strftime("%H:%M:%S") print(f"📌 [{mkt['label']}] 回合 {rt} 锚定价: ${strike:,.2f}") strike_pending = False elif not hasattr(state_worker, '_strike_notified') or state_worker._strike_notified != current_round: state_worker._strike_notified = current_round rt = datetime.fromtimestamp(current_round).strftime("%H:%M:%S") print(f"⏳ [{mkt['label']}] 回合 {rt} 锚定价(未确认): ${strike:,.2f}") # ── 资产查询 (pUSD 和持仓, 走 CLOB API 不消耗 Alchemy CU, 每5秒刷新) ── with state_lock: tids = dict(state["token_ids"]) acct = _current_account() _now = time.time() if not hasattr(state_worker, '_last_bal_t'): state_worker._last_bal_t = 0 state_worker._last_bal_idx = -1 # 切换账号后立即刷新余额 with _account_lock: cur_idx = _active_account_idx acct_changed = (cur_idx != getattr(state_worker, '_last_bal_idx', -1)) if acct["private_key"] and acct["funder"] and (acct_changed or _now - state_worker._last_bal_t >= 5): state_worker._last_bal_t = _now state_worker._last_bal_idx = cur_idx try: from py_clob_client_v2.clob_types import BalanceAllowanceParams, AssetType client = get_clob_client() # 1. 查询 pUSD 余额 data = client.get_balance_allowance( BalanceAllowanceParams(asset_type=AssetType.COLLATERAL) ) if isinstance(data, dict): b = data.get("balance", 0) if not b: b = data.get("collateral", {}).get("balance", 0) with state_lock: state["usdc_balance"] = float(b) / 1e6 if b else 0 # 2. 查询持仓 (Token ID) for side in ("up", "down"): tid = tids.get(side) if tid: c_data = client.get_balance_allowance( BalanceAllowanceParams(asset_type=AssetType.CONDITIONAL, token_id=str(tid)) ) if isinstance(c_data, dict): cb = c_data.get("balance", 0) with state_lock: state[f"{side}_position"] = float(cb) / 1e6 if cb else 0 except Exception: pass except Exception as e: print(f"[state_worker] 错误: {e}") # 锚定价待确认时加快轮询 time.sleep(1 if strike_pending else 3) def clob_price_worker(): """后台线程: 高频 CLOB 盘口价轮询 (每秒, 独立线程)""" session = requests.Session() session.headers.update({"Accept": "application/json"}) if _PROXIES: session.proxies.update(_PROXIES) while True: try: with state_lock: tids = dict(state["token_ids"]) for side in ("up", "down"): tid = tids.get(side) if not tid: continue try: r = session.get( f"{CLOB_HTTP}/book", params={"token_id": tid}, timeout=2, ) if r.status_code == 200: book = r.json() bids = book.get("bids", []) asks = book.get("asks", []) best_bid = max((float(b["price"]) for b in bids), default=0) if bids else 0 best_ask = min((float(a["price"]) for a in asks), default=0) if asks else 0 # Top-5 orderbook sorted_bids = sorted(bids, key=lambda b: float(b["price"]), reverse=True)[:10] sorted_asks = sorted(asks, key=lambda a: float(a["price"]))[:10] top_bids = [{"p": float(b["price"]), "s": float(b["size"])} for b in sorted_bids] top_asks = [{"p": float(a["price"]), "s": float(a["size"])} for a in sorted_asks] with state_lock: state[f"{side}_price"] = best_bid state[f"{side}_ask"] = best_ask state[f"{side}_book"] = {"bids": top_bids, "asks": top_asks} except Exception: pass except Exception: pass time.sleep(0.1) # 嵌入模式标志 (由 onchain_leaderboard 设置) _embedded = False # ══════════════════════════════════════════════════════════════════════════════ # aiohttp Web 应用 # ══════════════════════════════════════════════════════════════════════════════ app = web.Application() async def index(request): return web.Response(text=HTML_TEMPLATE, content_type="text/html", charset="utf-8") async def trade_api(request): data = await request.json() loop = asyncio.get_event_loop() result = await loop.run_in_executor(None, handle_trade, data) return web.json_response(result) async def switch_market_api(request): global active_market data = await request.json() mkt_key = data.get("market", "") if mkt_key not in MARKETS: return web.json_response({"success": False, "message": f"未知市场: {mkt_key}"}) with active_market_lock: active_market = mkt_key # 清空 state 让 worker 立即刷新 with state_lock: state["round_ts"] = 0 state["round_end"] = 0 state["token_ids"] = {} state["strike"] = 0 state["up_price"] = 0 state["down_price"] = 0 state["up_ask"] = 0 state["down_ask"] = 0 state["up_position"] = 0 state["down_position"] = 0 label = MARKETS[mkt_key]["label"] print(f"🔀 切换市场 → {label}") return web.json_response({"success": True, "message": f"已切换到 {label}"}) async def switch_account_api(request): global _active_account_idx data = await request.json() idx = data.get("index", 0) if idx < 0 or idx >= len(ACCOUNTS): return web.json_response({"success": False, "message": f"无效账号索引: {idx}"}) with _account_lock: _active_account_idx = idx acct = ACCOUNTS[idx] # 同步 relayer 环境变量 if acct.get("relayer_key"): os.environ["RELAYER_API_KEY"] = acct["relayer_key"] if acct.get("relayer_addr"): os.environ["RELAYER_API_KEY_ADDRESS"] = acct["relayer_addr"] # 清空持仓 & 余额 (让 worker 重新查询) with state_lock: state["up_position"] = 0 state["down_position"] = 0 state["usdc_balance"] = 0 print(f"👛 切换账号 → [{acct['label']}]") # 后台初始化新账号的 CLOB client def _init(): try: get_clob_client() except Exception as e: print(f" ❌ CLOB 初始化失败: {e}") threading.Thread(target=_init, daemon=True).start() return web.json_response({"success": True, "message": f"已切换到 {acct['label']}"}) async def ws_handler(request): ws = web.WebSocketResponse() await ws.prepare(request) try: while not ws.closed: with state_lock: data = dict(state) with active_market_lock: data["active_market"] = active_market data["market_label"] = MARKETS.get(active_market, {}).get("label", "") data["trades"] = list(trade_log) data["server_time"] = time.time() # 账号信息 with _account_lock: data["active_account"] = _active_account_idx data["accounts"] = [a["label"] for a in ACCOUNTS] await ws.send_json(data) await asyncio.sleep(0.1) except Exception: pass return ws app.router.add_get("/", index) app.router.add_post("/api/trade", trade_api) app.router.add_post("/api/switch_market", switch_market_api) app.router.add_post("/api/switch_account", switch_account_api) app.router.add_get("/ws", ws_handler) # ══════════════════════════════════════════════════════════════════════════════ # HTML 模板 (嵌入式, 单文件部署) # ══════════════════════════════════════════════════════════════════════════════ HTML_TEMPLATE = r"""