Files
my_polymarket_m5_dasboad/auto.py
T
2026-07-26 01:51:18 +08:00

1334 lines
66 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""
auto.py - BTC 5m 差价交易机器人 (趋势跟随策略)
策略核心: 赚取差价 (非持仓到结算)
跟随 BTC 趋势入场, 趋势延续时止盈 (5m 市场动量效应 > 均值回归)
入场信号 (加权评分, 正=偏UP, 负=偏DN):
1. BTC 位移 (权重 4, 核心): BTC 偏离 strike 的方向 = 趋势方向
2. BTC 加速度 (权重 3, 核心): 动量启动/衰减, 加速时跟随价值最高
3. 盘口压力失衡 (权重 2): order book bids/asks 失衡确认短期方向
4. 鲸鱼跟随 (权重 1): 大户持仓方向 = 趋势方向 (跟单, 不反向)
退出条件:
持仓盈利 > 15% -> 止盈 (趋势跟随给趋势发展空间)
持仓亏损 > 15% -> 止损 (严格止损)
持仓超过 120s -> 强制平仓 (避免结算尾部风险)
剩余时间 < 30s -> 强制平仓
前置条件: onchain_leaderboard.py 必须在运行 (提供盘口 + 鲸鱼持仓数据)
用法:
python auto.py # dry run (默认)
python auto.py --live # 实盘
python auto.py --amount 1 # 每笔 $1
python auto.py --market btc_5m # 指定市场
python auto.py --weight-btc-dev 4 --weight-btc-accel 3 # 调权重
"""
import argparse
import json
import sys
import time
import threading
import traceback
from datetime import datetime
# 屏蔽 SSL 警告 (代理自签证书)
import urllib3
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
# ── 复用 quick_trade 的基础设施 ──────────────────────────────────────────────
import quick_trade as qt
PROXY = qt._PROXY
PROXIES = qt._PROXIES
# ════════════════════════════════════════════════════════════════════════════
# 配置 (差价交易: 趋势跟随策略)
# ════════════════════════════════════════════════════════════════════════════
# 参数分两类:
# 1. 交易参数: amount, entry_delay, exit_deadline, take_profit, stop_loss, max_price
# 2. 信号评分参数: 4 个权重 + min_score_ratio
#
# 策略目标: 赚取差价 (非持仓到结算)
# - 跟随 BTC 趋势入场, 趋势延续时止盈 (动量效应)
# - BTC 位移 + 加速度为核心信号, 盘口/鲸鱼为确认信号
# - 严格止损 + 中止盈, 高频次小利润
# - 持仓时间短 (max_hold_time), 避免结算尾部风险
#
# 权重的含义: 每个信号的最大贡献值, 不是"所需分数"
# 4 个信号权重总和 10 = 理论最大 ±10
# 入场门槛 = min_score_ratio × 权重总和 (默认 0.5 × 10 = ±5)
# 调 min_score_ratio = 改变入场严格度 (0.3=宽松, 0.7=严格)
class Config:
dry_run = True # True=模拟, False=实盘
market = "btc_5m" # 市场: btc_5m/btc_15m/btc_1h/eth_5m...
amount = 1.0 # 每笔下单金额 (USD)
entry_delay = 15 # 回合开始后等待秒数 (等盘口形成, 太早无流动性)
exit_deadline = 60 # 强制平仓倒计时 (-1=禁用, 60=剩60s平仓, 避免回合末流动性枯竭)
min_hold_time = 90 # 入场时最小剩余时间 (秒), 不够则跳过
max_hold_time = 180 # 持仓最大时长 (秒), 超过强制平仓 (趋势需要时间发展)
reentry_cooldown = 30 # 平仓后冷却秒数, 期间不再入场 (避免过度交易, 0=立即再入场)
stop_loss_grace = 5 # 入场后宽限秒数, 期间不检查止损 (避免瞬时滑点误触发)
max_spread = 0.05 # 入场前盘口价差上限 (bid-ask > 5¢ 不入场, 流动性差)
max_retries = 3 # 下单失败最大重试次数
retry_delay = 0.1 # 重试间隔基数 (秒, 指数退避)
virtual_balance = 10.0 # dry-run 初始模拟余额
take_profit_pct = 0.25 # 持仓盈利止盈阈值 (backtest TP×SL grid: 25% × 30% 最优)
stop_loss_pct = 0.30 # 持仓亏损止损阈值 (backtest: 25% TP + 30% SL 总PnL +$71.69)
max_price = 0.75 # token 价格上限 (趋势可能推高, 允许 0.75 入场)
buy_slippage = 0.02 # 买入滑点 (ask + 此值 = 成交价, 0.02=2¢, 摩擦最小化)
sell_slippage = 0.02 # 卖出滑点 (bid - 此值 = 成交价, 0.02=2¢)
# 信号权重: 每个信号的最大贡献值 (实际打分 ≤ 此值)
# 策略核心: 跟随当前 3 连胜 (🔥) 鲸鱼的动态方向, 而不是 curated 列表
weight_btc_dev = 1 # BTC 位移 (弱确认): 偏离 strike 仅作辅助
weight_btc_accel = 2 # BTC 加速度 (弱确认): 动量辅助
weight_book = 1 # 盘口压力 (弱确认): 失衡辅助
weight_whale_follow = 0 # 鲸鱼加权投票: 禁用 (curated 列表噪声大, 已替换为 streak_leader)
weight_winners = 0 # 上轮赢家方向: 已合并入 streak_leader (combined 投票)
weight_streak_leader = 4 # 连赢鲸鱼 🔥 + 上轮赢家 综合信号 (核心): 加权投票 2:1
streak_combined_threshold = 0.30 # combined 投票入场门槛 (历史回测最优)
# 入场门槛 = 这个比例 × (权重总和)
# 0.5 = 50% 满分才入场 (默认), 0.3 = 宽松, 0.7 = 严格
min_score_ratio = 0.5
# ════════════════════════════════════════════════════════════════════════════
# 日志
# ════════════════════════════════════════════════════════════════════════════
_LOG_LOCK = threading.Lock()
def log(msg, level="INFO"):
ts = datetime.now().strftime("%H:%M:%S")
tag = f"[{level}]"
prefix = " " if level == "DATA" else ""
with _LOG_LOCK:
print(f"{ts} {tag} {prefix}{msg}", flush=True)
def log_trade(msg, success=True):
ts = datetime.now().strftime("%H:%M:%S")
icon = "✅" if success else "❌"
with _LOG_LOCK:
print(f"{ts} {icon} {msg}", flush=True)
# ════════════════════════════════════════════════════════════════════════════
# 回测数据记录器 (JSONL 格式, 只在 dry_run 模式写)
# ════════════════════════════════════════════════════════════════════════════
_recorder = None # 延迟到 main() 初始化
class Recorder:
def __init__(self, path):
self.path = path
self.fh = open(path, "a", encoding="utf-8")
self._lock = threading.Lock()
def write(self, event_type: str, data: dict):
row = {"ts": time.time(), "type": event_type, **data}
with self._lock:
self.fh.write(json.dumps(row, ensure_ascii=False) + "\n")
self.fh.flush()
def close(self):
self.fh.close()
def rec(event_type: str, **kw):
"""安全写记录 (recorder 不存在时静默)"""
global _recorder
if _recorder:
_recorder.write(event_type, kw)
# ════════════════════════════════════════════════════════════════════════════
# 从 onchain_leaderboard 的 WebSocket 读取全部状态 (100ms 推送)
# ════════════════════════════════════════════════════════════════════════════
# 单个 dict, 无锁读取 (GIL 保证原子引用赋值, 读取端只做 dict.get)
_ws_state = {
"leaderboard": {},
"strike": 0.0,
"round_ts": 0,
"round_end": 0,
"token_ids": {},
"up_price": 0.0, # best bid
"down_price": 0.0,
"up_ask": 0.0,
"down_ask": 0.0,
"crypto_current": 0.0,
"up_book": {"bids": [], "asks": []},
"down_book": {"bids": [], "asks": []},
"ts": 0.0, # 最后更新时间
}
def _ws_client():
"""后台线程: 连接 onchain_leaderboard WebSocket, 持续接收全部状态"""
import websocket
ws_url = f"ws://127.0.0.1:{qt.PORT}/ws"
while True:
try:
ws = websocket.create_connection(ws_url, timeout=10)
log(f"已连接排行榜 WebSocket: {ws_url}")
while True:
try:
raw = ws.recv()
if not raw:
continue
# 解析在锁外完成 (CPU 密集, 不需要持锁)
data = json.loads(raw)
# 快速赋值: 只持锁做字段拷贝
with _ws_lock:
for k in ("leaderboard", "token_ids", "up_book", "down_book"):
v = data.get(k)
if v:
_ws_state[k] = v
for k in ("strike", "round_ts", "round_end", "up_price",
"down_price", "up_ask", "down_ask", "crypto_current"):
v = data.get(k)
if v:
_ws_state[k] = v
_ws_state["ts"] = time.time()
# 同步到 qt.state 让 place_order 读到 token_ids
if _ws_state.get("token_ids"):
qt.state["token_ids"] = dict(_ws_state["token_ids"])
qt.state["strike"] = _ws_state.get("strike", 0)
qt.state["round_ts"] = _ws_state.get("round_ts", 0)
qt.state["round_end"] = _ws_state.get("round_end", 0)
except Exception:
break
ws.close()
except Exception:
pass
log("排行榜 WS 断开, 2s 后重连...", "WARN")
time.sleep(2)
_ws_lock = threading.Lock()
def ws_get(key, default=None):
"""无锁快读 (GIL 保证 dict.get 原子性)"""
return _ws_state.get(key, default)
def read_leaderboard() -> dict:
"""读取排行榜副本 (深拷贝 up/down 列表, 防止 WS 写入时 RuntimeError)"""
lb = _ws_state.get("leaderboard")
if not lb:
return {}
return {
"up": list(lb.get("up", [])),
"down": list(lb.get("down", [])),
"up_total": lb.get("up_total", 0),
"dn_total": lb.get("dn_total", 0),
"up_count": lb.get("up_count", 0),
"dn_count": lb.get("dn_count", 0),
"hedged_count": lb.get("hedged_count", 0),
"streak_up": lb.get("streak_up", 0),
"streak_dn": lb.get("streak_dn", 0),
"winners": dict(lb.get("winners", {})),
"streak4": dict(lb.get("streak4", {})),
}
def get_cached_price(side: str) -> tuple:
"""从 WS 缓存读取 best_bid / best_ask (优先用盘口数据, 再用 up_ask/down_ask 字段)"""
# 从 order book 读: 最优卖 = asks[0], 最优买 = bids[0]
# 盘口数据格式: {"p": price, "s": size}
book = _ws_state.get(f"{side}_book", {})
asks = book.get("asks", [])
bids = book.get("bids", [])
best_ask = min((a["p"] for a in asks), default=0) if asks else 0
best_bid = max((b["p"] for b in bids), default=0) if bids else 0
# 兜底: 用旧的 up_ask/down_ask 字段
if best_ask <= 0:
best_ask = _ws_state.get(f"{side}_ask", 0)
if best_bid <= 0:
best_bid = _ws_state.get(f"{side}_price", 0)
return best_bid, best_ask
def start_ws_client():
try:
import websocket
except ImportError:
log("缺少 websocket-client, 请 pip install websocket-client", "ERROR")
return False
threading.Thread(target=_ws_client, daemon=True).start()
return True
# ════════════════════════════════════════════════════════════════════════════
# BTC 实时价格 (从 WS 缓存读取, 100ms 精度)
# ════════════════════════════════════════════════════════════════════════════
def get_btc_price() -> float:
return _ws_state.get("crypto_current", 0.0)
# BTC 价格历史 (用于计算加速度/动量)
_btc_history = [] # [(ts, price), ...] 最近 N 个采样
_BTC_HIST_MAX = 30 # 保留最近 30 个采样 (约 30 秒, 每次 run_bot 循环采一个)
_btc_hist_lock = threading.Lock()
def sample_btc_price(btc: float):
"""在 run_bot 主循环里采样 BTC 价格, 用于计算动量/加速度"""
if btc <= 0:
return
now = time.time()
with _btc_hist_lock:
_btc_history.append((now, btc))
# 保留最近 _BTC_HIST_MAX 个 + 30 秒内的
cutoff = now - 30
while _btc_history and (_btc_history[0][0] < cutoff or len(_btc_history) > _BTC_HIST_MAX):
_btc_history.pop(0)
def get_btc_velocity() -> tuple:
"""计算 BTC 价格动量
返回 (velocity_per_sec, accel_per_sec2):
velocity = (current - past) / dt
accel = (v_now - v_past) / dt
无足够数据时返回 (0.0, 0.0)
"""
with _btc_hist_lock:
if len(_btc_history) < 4:
return 0.0, 0.0
now_ts, now_p = _btc_history[-1]
# 取 ~10s 前的价格算 velocity
past_ts = now_ts - 10
past_p = None
for ts, p in _btc_history:
if ts >= past_ts:
past_p = p
past_ts_actual = ts
break
if past_p is None:
return 0.0, 0.0
dt1 = now_ts - past_ts_actual
if dt1 <= 0:
return 0.0, 0.0
v_now = (now_p - past_p) / dt1 # $/s
# 取 ~20s 前的价格算过去的 velocity, 与现在比得到加速度
old_ts = now_ts - 20
old_p = None
for ts, p in _btc_history:
if ts >= old_ts:
old_p = p
old_ts_actual = ts
break
if old_p is None or len(_btc_history) < 6:
return v_now, 0.0
# 过去的 velocity: 用 20s 前 到 10s 前之间
mid_ts = now_ts - 10
mid_p = None
for ts, p in _btc_history:
if ts >= mid_ts:
mid_p = p
break
if mid_p is None:
return v_now, 0.0
dt2 = mid_ts - old_ts_actual
if dt2 <= 0:
return v_now, 0.0
v_past = (mid_p - old_p) / dt2
accel = (v_now - v_past) / max(now_ts - mid_ts, 0.1)
return v_now, accel
# ════════════════════════════════════════════════════════════════════════════
# 虚拟持仓 (dry-run 模式)
# ════════════════════════════════════════════════════════════════════════════
class VirtualPortfolio:
def __init__(self, balance):
self.balance = balance
self.position = None
self.trades = []
self.total_pnl = 0.0
def buy(self, side, price, shares, amount):
if self.position:
return False, "已有持仓"
cost = shares * price
if cost > self.balance:
return False, f"余额不足 (${self.balance:.2f} < ${cost:.2f})"
self.balance -= cost
self.position = {"side": side, "shares": shares, "entry_price": price, "entry_ts": time.time()}
self.trades.append({"action": "buy", "side": side, "price": price, "shares": shares, "ts": time.time()})
return True, f"买入 {side.upper()} {shares:.1f}股 @ {price:.2f}"
def sell(self, exit_price):
if not self.position:
return False, "无持仓", 0
pos = self.position
revenue = pos["shares"] * exit_price
cost = pos["shares"] * pos["entry_price"]
pnl = revenue - cost
self.balance += revenue
self.total_pnl += pnl
self.trades.append({"action": "sell", "side": pos["side"], "price": exit_price,
"shares": pos["shares"], "pnl": pnl, "ts": time.time()})
self.position = None
return True, f"卖出 {pos['side'].upper()} {pos['shares']:.1f}股 @ {exit_price:.2f}", pnl
def status(self):
return {
"balance": round(self.balance, 2),
"position": self.position,
"total_pnl": round(self.total_pnl, 2),
"trades": len(self.trades),
}
# ════════════════════════════════════════════════════════════════════════════
# 下单 (带重试 + 执行验证)
# ════════════════════════════════════════════════════════════════════════════
_virtual = None # 延迟到 main() 初始化 (Config.virtual_balance 可能被 CLI 覆盖)
# 实盘模式 P&L 追踪
_live_stats = {"trades": 0, "wins": 0, "losses": 0, "total_pnl": 0.0, "total_fees": 0.0}
def record_live_trade(pnl: float, fees: float = 0.0):
"""实盘: 记录单笔 P&L (从 pUSD 余额变化推算)"""
_live_stats["trades"] += 1
_live_stats["total_pnl"] += pnl
_live_stats["total_fees"] += fees
if pnl > 0:
_live_stats["wins"] += 1
elif pnl < 0:
_live_stats["losses"] += 1
def place_order(side: str, amount_usd: float, is_sell_all: bool = False) -> dict:
"""
下单统一入口 (买入 / SELL ALL)
返回 {"success": bool, "message": str, "price": float, "shares": float}
带: 滑点计算, 网络重试, 执行验证
价格来源: WS 缓存 (无 HTTP 延迟)
"""
for attempt in range(1, Config.max_retries + 1):
try:
token_id = _ws_state.get("token_ids", {}).get(side)
if not token_id:
with qt.state_lock:
token_id = qt.state["token_ids"].get(side)
if not token_id:
return {"success": False, "message": "回合未就绪, 无 token_id"}
# 从 WS 缓存读价格 (无 HTTP, 0ms 延迟)
bid, ask = get_cached_price(side)
if is_sell_all:
price = bid if bid > 0 else 0.5
order_price = round(max(price - Config.sell_slippage, 0.01), 2)
else:
price = ask if ask > 0 else bid
if price <= 0:
if attempt < Config.max_retries:
time.sleep(0.3 * attempt)
continue
return {"success": False, "message": "无可用盘口价格"}
order_price = round(min(price + Config.buy_slippage, 0.99), 2)
# ── Dry Run: 模拟成交 (与实盘对齐: 用 order_price = ask + 滑点) ──
if Config.dry_run:
if is_sell_all:
if not _virtual or not _virtual.position or _virtual.position["side"] != side:
return {"success": False, "message": "[DRY] 无持仓可卖"}
shares = _virtual.position["shares"]
ok, msg = _virtual.sell(order_price)[:2]
if ok:
pnl = _virtual.total_pnl
log_trade(f"[DRY] SELL ALL {side.upper()} {shares:.1f}股 @ {order_price:.2f} | 累计PnL: ${pnl:+.2f}")
return {"success": ok, "message": f"[DRY] {msg}", "price": order_price, "shares": shares}
else:
# shares = amount / order_price (含滑点), 余额检查也用含滑点价
shares = round(amount_usd / order_price, 2)
ok, msg = _virtual.buy(side, order_price, shares, amount_usd)
if ok:
log_trade(f"[DRY] BUY {side.upper()} {shares:.1f}股 @ {order_price:.2f} (${amount_usd})")
return {"success": ok, "message": f"[DRY] {msg}", "price": order_price, "shares": shares}
# ── 实盘下单 ──
if is_sell_all:
result = qt.execute_sell(side, "all")
else:
result = qt.execute_buy(side, amount_usd)
if result.get("success"):
# ── 执行验证: 轮询链上余额, 成交即返回 (50ms × 10 次, 共 500ms 上限) ──
# 替代原 time.sleep(0.5) + 查 1 次, 平均省 200-400ms, 最坏相同
w3 = qt.get_w3()
actual_bal = 0.0
verified = False
pre_bal = qt.get_token_balance(w3, token_id) # 下单前余额 (用于卖出验证)
for _ in range(10):
time.sleep(0.05)
actual_bal = qt.get_token_balance(w3, token_id)
if is_sell_all:
# 卖出验证: 余额相比下单前减少 (即 token 被消耗)
if actual_bal < pre_bal - 0.5:
verified = True
break
else:
# 买入验证: 余额 > 0 即成交
if actual_bal > 0:
verified = True
break
if not verified:
if not is_sell_all:
# 买入未成交: 可重试
if attempt < Config.max_retries:
log(f"买入后链上余额仍为 0, 第{attempt}次重试...", "WARN")
time.sleep(Config.retry_delay * attempt)
continue
log_trade(f"⚠️ 订单返回成功但链上余额为 0, 可能未成交", success=False)
else:
# 卖出未成交: 余额没变化, 警告但不重试 (避免重复下单把已成交的卖单当失败)
log_trade(f"⚠️ 卖出订单返回成功但链上余额未变 (pre={pre_bal:.1f} post={actual_bal:.1f}), 可能未成交", success=False)
shares = result.get("shares", amount_usd / price)
log_trade(f"{'SELL ALL' if is_sell_all else 'BUY'} {side.upper()} | {result.get('message', '')} | {result.get('elapsed_ms', 0)}ms")
return {"success": True, "message": result.get("message", ""),
"price": price, "shares": shares, "order_id": result.get("order_id", "")}
# ── 下单失败: 判断是否可重试 ──
msg = result.get("message", "")
retryable = any(kw in msg.lower() for kw in ["timeout", "network", "connection", "reset", "refused", "5xx", "502", "503"])
if retryable and attempt < Config.max_retries:
delay = Config.retry_delay * (2 ** (attempt - 1))
log(f"下单失败 (可重试): {msg} | 第{attempt}/{Config.max_retries}次, {delay:.1f}s 后重试", "WARN")
time.sleep(delay)
continue
log_trade(f"下单失败: {msg}", success=False)
return {"success": False, "message": msg}
except Exception as e:
if attempt < Config.max_retries:
delay = Config.retry_delay * (2 ** (attempt - 1))
log(f"下单异常: {e} | 第{attempt}/{Config.max_retries}次, {delay:.1f}s 后重试", "WARN")
time.sleep(delay)
continue
log_trade(f"下单异常 (放弃): {e}", success=False)
return {"success": False, "message": str(e)}
return {"success": False, "message": f"重试 {Config.max_retries} 次后仍失败"}
# ════════════════════════════════════════════════════════════════════════════
# 差价交易信号分析 (趋势跟随策略)
# ════════════════════════════════════════════════════════════════════════════
def analyze_signals(lb: dict, btc: float, strike: float) -> dict:
"""
差价交易信号分析 (趋势跟随 + 动量确认)
核心逻辑: BTC 5m 市场短期呈动量效应, 跟随趋势入场, 趋势延续时止盈
信号 1 (核心): BTC 位移 - BTC 偏离 strike 的方向 = 趋势方向
信号 2 (核心): BTC 加速度 - 价格变化率的变化, 加速时趋势更强
信号 3 (确认): 盘口压力 - orderbook 失衡方向与趋势一致
信号 4 (确认): 鲸鱼跟随 - 大户持仓方向 = 趋势方向 (跟随, 不反向)
正分 = 偏 UP, 负分 = 偏 DOWN
返回: {"score": int, "side": "up"/"down"/None, "details": [str], "signals": dict, "min_abs": int}
"""
score = 0
details = []
signals = {}
# ── 信号 1: BTC 位移 (核心, 权重最大) ──
# BTC 偏离 strike 的方向就是趋势方向, 偏离越大趋势越强
w_dev = Config.weight_btc_dev
if btc > 0 and strike > 0:
dev = btc - strike # 正=BTC上涨, 负=BTC下跌
if abs(dev) >= 10: # 至少偏离 $10 才算有趋势 (干跑日志显示锚定偏差 $10~16 几乎必有, 提高阈值过滤噪声)
# 映射: $10->1, $20->2, $30->3, $40+->4 (权重4封顶)
s = min(w_dev, int(abs(dev) / 10))
s = s if dev > 0 else -s
score += s
direction = "UP" if s > 0 else "DN"
sign = "+" if s > 0 else ""
signals["btc_dev"] = f"{direction} {sign}{s} (BTC ${dev:+.2f} vs 锚定)"
else:
signals["btc_dev"] = f"中性 (偏离 ${dev:+.2f} < $10, 无趋势)"
else:
signals["btc_dev"] = "无 BTC/锚定数据"
# ── 信号 2: BTC 加速度 (核心, 捕捉动量启动) ──
# velocity = 10s 内 BTC 价格变化率
# accel = velocity 的变化率, 正=加速上涨, 负=加速下跌
# 加速时趋势刚开始, 跟随价值最高
w_acc = Config.weight_btc_accel
velocity, accel = get_btc_velocity()
if abs(velocity) > 0.1: # 至少 $0.1/s 才算有动量
# velocity 映射: $0.1/s=1, $0.2/s=2, $0.3/s+=3 (权重3封顶)
s_v = min(w_acc, int(abs(velocity) / 0.1))
# 加速度同方向加成 (加速时给满分, 减速时减半)
if (velocity > 0 and accel >= 0) or (velocity < 0 and accel <= 0):
s = s_v # 加速, 满分
accel_desc = "加速"
else:
s = max(1, s_v // 2) # 减速, 半分
accel_desc = "减速"
s = s if velocity > 0 else -s
score += s
direction = "UP" if s > 0 else "DN"
sign = "+" if s > 0 else ""
signals["btc_accel"] = (f"{direction} {sign}{s} (v={velocity:+.3f}/s, "
f"a={accel:+.3f}/s², {accel_desc})")
else:
signals["btc_accel"] = f"中性 (v={velocity:+.3f}/s < 0.1, 无动量)"
# ── 信号 3: 盘口压力确认 (权重中等) ──
# orderbook bids 厚 = 买盘强, 跟随买盘方向
wb = Config.weight_book
up_book = _ws_state.get("up_book") or {"bids": [], "asks": []}
dn_book = _ws_state.get("down_book") or {"bids": [], "asks": []}
up_bid_d = sum(b.get("s", 0) for b in up_book.get("bids", []))
up_ask_d = sum(a.get("s", 0) for a in up_book.get("asks", []))
dn_bid_d = sum(b.get("s", 0) for b in dn_book.get("bids", []))
dn_ask_d = sum(a.get("s", 0) for a in dn_book.get("asks", []))
if (up_bid_d + up_ask_d) > 0 and (dn_bid_d + dn_ask_d) > 0:
# 归一化压力: 正 = UP 涨, 负 = DN 涨
up_pressure = (up_bid_d - up_ask_d) / max(up_bid_d + up_ask_d, 1)
dn_pressure = (dn_bid_d - dn_ask_d) / max(dn_bid_d + dn_ask_d, 1)
# 综合压力: up 买盘强 + dn 抛压强 = UP 涨
net_pressure = up_pressure - dn_pressure
if abs(net_pressure) > 0.2:
s = min(wb, int(abs(net_pressure) * 5))
s = s if net_pressure > 0 else -s
score += s
direction = "UP" if s > 0 else "DN"
sign = "+" if s > 0 else ""
signals["book"] = (f"{direction} {sign}{s} "
f"(UP bid/ask={up_bid_d:.0f}/{up_ask_d:.0f}, "
f"DN bid/ask={dn_bid_d:.0f}/{dn_ask_d:.0f})")
else:
signals["book"] = (f"中性 (UP bid/ask={up_bid_d:.0f}/{up_ask_d:.0f}, "
f"DN bid/ask={dn_bid_d:.0f}/{dn_ask_d:.0f})")
else:
signals["book"] = "无盘口深度数据"
# ── 信号 4: 连赢鲸鱼 (🔥 streak=True) 动态信号 (核心 alpha) ──
# 只追踪当前 3 连胜的鲸鱼, 忽略 curated 列表
# 两个维度: (a) 当前集中度 (b) 跨回合翻向 (预言信号)
ws = Config.weight_streak_leader
cur_sides = {}
for e in lb.get("up", []):
if e.get("streak") and not e.get("hedged"):
cur_sides[e["addr"]] = "up"
for e in lb.get("down", []):
if e.get("streak") and not e.get("hedged"):
cur_sides[e["addr"]] = "down"
cur_up_n = sum(1 for v in cur_sides.values() if v == "up")
cur_dn_n = sum(1 for v in cur_sides.values() if v == "down")
cur_total = cur_up_n + cur_dn_n
if cur_total >= 2:
# (a) 当前集中度 (核心)
concentration = (cur_up_n - cur_dn_n) / cur_total
# (b) 上轮赢家方向 (新: 加入 combined 投票, 历史回测显示 winners_concentration 单独也有 86% 胜率)
winners_data = lb.get("winners") or {}
w_up_n = winners_data.get("up_n", 0)
w_dn_n = winners_data.get("dn_n", 0)
w_total = w_up_n + w_dn_n
if w_total >= 2:
winners_concentration = (w_up_n - w_dn_n) / w_total
else:
winners_concentration = 0.0
# (c) 集合变化: 新加 streak whales 的方向 vs 离开的
if _streak_snapshot.get("sides"):
prev_sides = _streak_snapshot["sides"]
left_addrs = set(prev_sides.keys()) - set(cur_sides.keys())
left_up = sum(1 for a in left_addrs if prev_sides.get(a) == "up")
left_dn = sum(1 for a in left_addrs if prev_sides.get(a) == "down")
new_addrs = set(cur_sides.keys()) - set(prev_sides.keys())
new_up = sum(1 for a in new_addrs if cur_sides.get(a) == "up")
new_dn = sum(1 for a in new_addrs if cur_sides.get(a) == "down")
net_set_change = (new_up - new_dn) - (left_up - left_dn)
set_change_score = 0.0
if abs(net_set_change) > 0:
set_change_score = net_set_change / max(new_up + new_dn + left_up + left_dn, 1)
else:
set_change_score = 0.0
new_up = new_dn = left_up = left_dn = 0
# 综合: 加权投票 streak(2) + winners(1), 各按样本量比例加权
# 历史回测: combined 模式 (threshold=0.15) 胜率 88.1%, avg/trade +$0.14
# 公式对齐 backtest: 2*streak_conc*share + 1*winners_conc*share
# streak 仅有时: 2*streak_conc (max 2)
# winners 仅有时: 1*winners_conc (max 1)
total_sample = cur_total + w_total
if total_sample > 0:
vote_score = (2.0 * concentration * (cur_total / total_sample)
+ 1.0 * winners_concentration * (w_total / total_sample))
else:
vote_score = 0.0
# 直接用 vote_score 作为 combined (历史回测对齐)
# backtest: combined threshold=0.30 → 91.4% 胜率, 420 trades, +$66 PnL
combined = vote_score
if abs(combined) >= Config.streak_combined_threshold:
s = min(ws, max(1, int(abs(combined) * ws * 1.4)))
s = s if combined > 0 else -s
score += s
direction = "UP" if s > 0 else "DN"
sign = "+" if s > 0 else ""
signals["streak_leader"] = (
f"{direction} {sign}{s} (🔥 {cur_up_n}U/{cur_dn_n}D, "
f"上轮赢家 {w_up_n}U/{w_dn_n}D, 综合 {combined:+.0%})"
)
else:
signals["streak_leader"] = (
f"中性 (🔥 {cur_up_n}U/{cur_dn_n}D, "
f"上轮赢家 {w_up_n}U/{w_dn_n}D, 综合 {combined:+.0%})"
)
else:
signals["streak_leader"] = f"无连赢鲸鱼 ({cur_total} 🔥)"
# ── 汇总 ──
max_score = (Config.weight_btc_dev + Config.weight_btc_accel
+ Config.weight_book + Config.weight_whale_follow + Config.weight_winners
+ Config.weight_streak_leader)
min_abs = max(1, round(max_score * Config.min_score_ratio))
side = "up" if score >= min_abs else "down" if score <= -min_abs else None
for k, v in signals.items():
details.append(f" [{k}] {v}")
return {"score": score, "side": side, "details": details, "signals": signals,
"min_abs": min_abs, "max_score": max_score,
"streak_up": cur_up_n, "streak_dn": cur_dn_n,
"streak_concentration": concentration if cur_total >= 2 else 0.0}
# ════════════════════════════════════════════════════════════════════════════
# 策略状态机
# ════════════════════════════════════════════════════════════════════════════
class BotState:
IDLE = "IDLE"
WAITING_ENTRY = "WAITING_ENTRY"
IN_POSITION = "IN_POSITION"
_bot_state = BotState.IDLE
_position_info = None # {"side", "entry_score", "entry_ts", "shares"}
_last_exit_ts = 0.0 # 上次平仓时间戳, 用于 reentry_cooldown 判断
# 连赢鲸鱼 (streak=True) 的跨回合分布追踪 — 用于"动态变化"信号
# 每回合开局时快照 {addr: "up"/"down"}; 信号基于"翻了 vs 上回合"
_streak_snapshot: dict = {} # {"round_ts": int, "sides": {addr: "up"/"down"}}
_streak_snapshot_round_ts: int = 0 # 已快照的回合 ts
# ════════════════════════════════════════════════════════════════════════════
# 获取回合信息
# ════════════════════════════════════════════════════════════════════════════
def get_round_info(mkt_key: str) -> dict:
"""从 WS 缓存读回合信息 (无 HTTP, 0ms)"""
mkt = qt.MARKETS[mkt_key]
duration = mkt["duration"]
now = int(time.time())
round_ts = now - (now % duration)
round_end = round_ts + duration
token_ids = ws_get("token_ids", {})
strike = ws_get("strike", 0)
ws_round = ws_get("round_ts", 0)
# WS 数据落后: 自己拉
if ws_round != round_ts or not token_ids:
token_ids = qt.get_round_token_ids(round_ts, mkt["slug_prefix"], mkt.get("slug_type", ""))
if strike <= 0:
strike, _ = qt.get_strike_price(round_ts, duration, mkt["symbol"], mkt["variant"])
return {"round_ts": round_ts, "round_end": round_end, "token_ids": token_ids, "strike": strike}
# ════════════════════════════════════════════════════════════════════════════
# 综合市场状态捕获 (entry/exit/snapshot 时调用, 用于回测)
# ════════════════════════════════════════════════════════════════════════════
def capture_market_state(round_ts: int = 0, strike: float = 0.0,
remaining: float = 0.0, side: str = "") -> dict:
"""一次性捕获所有维度, 供 entry/exit/snapshot 写入 jsonl"""
btc = get_btc_price()
velocity, accel = get_btc_velocity()
up_bid, up_ask = get_cached_price("up")
dn_bid, dn_ask = get_cached_price("down")
# 盘口深度 (bids/asks 总量)
up_book = _ws_state.get("up_book") or {"bids": [], "asks": []}
dn_book = _ws_state.get("dn_book") or {"bids": [], "asks": []}
up_bid_depth = sum(b.get("s", 0) for b in up_book.get("bids", []))
up_ask_depth = sum(a.get("s", 0) for a in up_book.get("asks", []))
dn_bid_depth = sum(b.get("s", 0) for b in dn_book.get("bids", []))
dn_ask_depth = sum(a.get("s", 0) for a in dn_book.get("asks", []))
total_depth = up_bid_depth + up_ask_depth + dn_bid_depth + dn_ask_depth
up_depth_share = (up_bid_depth + up_ask_depth) / total_depth if total_depth else 0.5
# 🔥 streak whales (核心 alpha)
lb = _ws_state.get("leaderboard") or {}
up_entries = [e for e in lb.get("up", []) if e.get("streak") and not e.get("hedged")]
dn_entries = [e for e in lb.get("down", []) if e.get("streak") and not e.get("hedged")]
streak_up_n = len(up_entries)
streak_dn_n = len(dn_entries)
streak_up_shares = sum(e.get("shares", 0) for e in up_entries)
streak_dn_shares = sum(e.get("shares", 0) for e in dn_entries)
# 上轮赢家
winners = lb.get("winners") or {}
# 整体鲸鱼分布 (非 streak)
all_whale_up = sum(e.get("shares", 0) for e in lb.get("up", []) if e.get("whale"))
all_whale_dn = sum(e.get("shares", 0) for e in lb.get("down", []) if e.get("whale"))
# round timing
now = time.time()
elapsed = (now - round_ts) if round_ts else 0
state = {
"ts": now,
"btc": round(btc, 2),
"strike": round(strike, 2),
"btc_dev": round(btc - strike, 2) if strike > 0 else 0,
"btc_dev_pct": round((btc - strike) / strike * 100, 4) if strike > 0 else 0,
"btc_velocity": round(velocity, 4),
"btc_accel": round(accel, 4),
"up_bid": up_bid, "up_ask": up_ask,
"dn_bid": dn_bid, "dn_ask": dn_ask,
"up_spread": round(up_ask - up_bid, 4) if up_ask and up_bid else 0,
"dn_spread": round(dn_ask - dn_bid, 4) if dn_ask and dn_bid else 0,
"up_bid_depth": round(up_bid_depth, 0),
"up_ask_depth": round(up_ask_depth, 0),
"dn_bid_depth": round(dn_bid_depth, 0),
"dn_ask_depth": round(dn_ask_depth, 0),
"book_imbalance": round((up_bid_depth - up_ask_depth) - (dn_bid_depth - dn_ask_depth), 1),
"up_depth_share": round(up_depth_share, 3),
"streak_up_n": streak_up_n,
"streak_dn_n": streak_dn_n,
"streak_up_shares": round(streak_up_shares, 0),
"streak_dn_shares": round(streak_dn_shares, 0),
"streak_concentration": round(
(streak_up_n - streak_dn_n) / max(streak_up_n + streak_dn_n, 1), 3),
"streak4_total": (_ws_state.get("streak4") or {}).get("total", 0),
"streak4_up": (_ws_state.get("streak4") or {}).get("up", 0),
"streak4_dn": (_ws_state.get("streak4") or {}).get("dn", 0),
"streak_up_total_field": lb.get("streak_up", 0),
"streak_dn_total_field": lb.get("streak_dn", 0),
"all_whale_up_shares": round(all_whale_up, 0),
"all_whale_dn_shares": round(all_whale_dn, 0),
"winners_total": winners.get("total", 0),
"winners_up_n": winners.get("up_n", 0),
"winners_dn_n": winners.get("dn_n", 0),
"winners_up_s": round(winners.get("up_s", 0), 0),
"winners_dn_s": round(winners.get("dn_s", 0), 0),
"winners_side": winners.get("side", ""),
"hedged_count": lb.get("hedged_count", 0),
"up_total_shares": lb.get("up_total", 0),
"dn_total_shares": lb.get("dn_total", 0),
"up_count": lb.get("up_count", 0),
"dn_count": lb.get("dn_count", 0),
"round_ts": round_ts,
"round_elapsed": round(elapsed, 1),
"round_remaining": round(remaining, 1),
"side": side,
}
return state
# ════════════════════════════════════════════════════════════════════════════
# 主策略循环
# ════════════════════════════════════════════════════════════════════════════
def run_bot():
global _bot_state, _position_info, _last_exit_ts, _streak_snapshot, _streak_snapshot_round_ts
mkt = qt.MARKETS[Config.market]
log(f"🤖 差价交易机器人启动 | 市场: {mkt['label']} | 模式: {'DRY RUN' if Config.dry_run else '实盘'}")
max_score = (Config.weight_btc_dev + Config.weight_btc_accel
+ Config.weight_book + Config.weight_whale_follow)
min_abs = max(1, round(max_score * Config.min_score_ratio))
log(f" 策略: 跟随 🔥 连赢鲸鱼 + 上轮赢家 综合信号 (backtest 验证) | 金额: ${Config.amount}")
log(f" 权重: BTC位移{Config.weight_btc_dev}/BTC加速{Config.weight_btc_accel}/盘口{Config.weight_book}/鲸鱼跟{Config.weight_whale_follow}/上轮赢家{Config.weight_winners}/连赢鲸鱼{Config.weight_streak_leader} | 最低入场: {min_abs}/{max_score} ({Config.min_score_ratio*100:.0f}%)")
ed = f"{Config.exit_deadline}s" if Config.exit_deadline > 0 else "禁用"
mhd = f"{Config.max_hold_time}s" if Config.max_hold_time > 0 else "禁用"
cd = f"{Config.reentry_cooldown}s" if Config.reentry_cooldown > 0 else "禁用"
log(f" 入场延迟: {Config.entry_delay}s | 持仓上限: {mhd} | 强制平仓: {ed} | 平仓冷却: {cd}")
log(f" 止盈: {Config.take_profit_pct*100:.0f}pct | 止损: {Config.stop_loss_pct*100:.0f}pct (宽限 {Config.stop_loss_grace}s)")
log(f" 价格上限: {Config.max_price:.2f} | 买入滑点: {Config.buy_slippage:.2f} | 卖出滑点: {Config.sell_slippage:.2f} | 价差上限: {Config.max_spread:.3f}")
log(f" ⚠️ 需 onchain_leaderboard.py 在运行 (提供盘口+鲸鱼持仓数据)")
if Config.dry_run:
log(f" 虚拟初始余额: ${Config.virtual_balance:.2f}")
else:
if not qt.ACCOUNTS:
log("❌ 未配置账号, 请在 .env 中设置 QUICK_PRIVATE_KEY 和 QUICK_FUNDER", "ERROR")
return
log(f" 交易账号: {qt.ACCOUNTS[0]['label']}")
last_round_ts = 0
lb_warned = False
_last_signal_log = 0
_last_pos_log = 0
_last_wait_log = 0
_last_price_skip_log = 0
while True:
try:
info = get_round_info(Config.market)
round_ts = info["round_ts"]
round_end = info["round_end"]
token_ids = info["token_ids"]
strike = info["strike"]
now = time.time()
remaining = round_end - now
btc = get_btc_price()
sample_btc_price(btc) # 采样 BTC 价格用于计算加速度
# ── 新回合检测 ──
if round_ts != last_round_ts:
if token_ids and strike > 0:
last_round_ts = round_ts
_bot_state = BotState.WAITING_ENTRY
_position_info = None
_last_exit_ts = 0.0 # 新回合重置冷却状态 (本回合还没平仓过)
lb_warned = False
rt = datetime.fromtimestamp(round_ts).strftime("%H:%M:%S")
log(f"{'─'*60}")
log(f"🔄 新回合 {rt} | 锚定: ${strike:,.2f} | BTC: ${btc:,.2f} | 剩余: {remaining:.0f}s")
else:
time.sleep(2)
continue
# ── 持仓中: 监控退出条件 (200ms 高频, 无 HTTP) ──
if _bot_state == BotState.IN_POSITION and _position_info:
pos = _position_info
side = pos["side"]
entry_price = pos["entry_price"]
# 从 WS 缓存读当前 bid (无 HTTP, 0ms)
current_bid, _ = get_cached_price(side)
# WS 数据陈旧检测 (>5s 无更新, 警告)
ws_age = time.time() - _ws_state.get("ts", 0)
if ws_age > 5 and current_bid > 0:
log(f"⚠️ WS 数据陈旧 {ws_age:.0f}s, 价格可能不准", "WARN")
# 退出条件 1: 止盈 (考虑卖出滑点, 确保实际成交价仍有盈利)
if Config.take_profit_pct > 0 and current_bid > 0 and entry_price > 0:
actual_exit = current_bid - Config.sell_slippage
pnl_pct = (actual_exit - entry_price) / entry_price
if pnl_pct >= Config.take_profit_pct:
log(f"💰 止盈 (token {entry_price:.2f} -> {current_bid:.2f}, 实际成交 {actual_exit:.2f}, {pnl_pct*100:+.0f}pct)")
_do_exit(side, "take_profit")
continue
# 退出条件 2: 止损 (实际成交价 = bid - SELL_SLIPPAGE)
# 入场后宽限期内不检查止损 (避免瞬时滑点/盘口抖动误触发)
# 注意: 止盈不受宽限期影响 (有利润就该锁)
held_seconds = time.time() - pos.get("entry_ts", time.time())
if (Config.stop_loss_pct > 0 and current_bid > 0 and entry_price > 0
and held_seconds >= Config.stop_loss_grace):
actual_exit = current_bid - Config.sell_slippage
pnl_pct = (actual_exit - entry_price) / entry_price
if pnl_pct <= -Config.stop_loss_pct:
log(f"🛑 止损 (token {entry_price:.2f} -> {current_bid:.2f}, 实际成交 {actual_exit:.2f}, {pnl_pct*100:.0f}pct, 持仓{held_seconds:.0f}s)")
_do_exit(side, "stop_loss")
continue
# 退出条件 3: 持仓超过最大时长 (防止一直拿到结算)
if Config.max_hold_time > 0:
held_seconds = time.time() - pos.get("entry_ts", time.time())
if held_seconds >= Config.max_hold_time:
log(f"⏳ 持仓超时 ({held_seconds:.0f}s >= {Config.max_hold_time}s), 强制平仓")
_do_exit(side, "max_hold_time")
continue
# 退出条件 4: 时间到
if Config.exit_deadline > 0 and remaining <= Config.exit_deadline:
log(f"⏰ 时间到, 强制平仓 | 剩余 {remaining:.0f}s")
_do_exit(side, "exit_deadline")
continue
# 持仓中追踪极值 (每 200ms, 不阻塞退出检查)
if current_bid > 0 and entry_price > 0:
cur_pnl_pct = (current_bid - entry_price) / entry_price * 100
if cur_pnl_pct > pos.get("peak_pnl_pct", float("-inf")):
pos["peak_pnl_pct"] = cur_pnl_pct
if cur_pnl_pct < pos.get("trough_pnl_pct", float("inf")):
pos["trough_pnl_pct"] = cur_pnl_pct
if current_bid > pos.get("max_bid", 0):
pos["max_bid"] = current_bid
if current_bid < pos.get("min_bid", 999) or pos.get("min_bid", 999) == 0:
pos["min_bid"] = current_bid
# 状态打印 (每 2s 一次, 不阻塞退出检查)
if int(now) >= _last_pos_log + 2:
_last_pos_log = int(now)
pnl_str = ""
if current_bid > 0 and entry_price > 0:
pnl_pct = (current_bid - entry_price) / entry_price * 100
pnl_str = f" | token {entry_price:.2f}->{current_bid:.2f} ({pnl_pct:+.0f}pct)"
log(f"📊 持仓 {side.upper()} | BTC ${btc:,.2f} vs 锚定 ${strike:,.2f} | 剩余 {remaining:.0f}s{pnl_str}", "DATA")
# 每 5 秒记录一次快照 (用于回测)
if int(now) % 5 == 0:
rec("snapshot", round_ts=round_ts, remaining=remaining, strike=strike,
btc=btc, bid=current_bid, side=side, entry_price=entry_price,
state=capture_market_state(round_ts, strike, remaining, side),
peak_pnl_pct=pos.get("peak_pnl_pct", 0.0),
trough_pnl_pct=pos.get("trough_pnl_pct", 0.0))
time.sleep(0.2) # 200ms 高频监控
continue
# ── 等待入场 (高频, 无 HTTP) ──
if _bot_state == BotState.WAITING_ENTRY:
if not token_ids or strike <= 0:
time.sleep(0.3)
continue
elapsed = now - round_ts
# 等待入场延迟 (仅回合开头生效: 若已平仓过, _last_exit_ts > 0, 不再等 entry_delay)
if _last_exit_ts == 0 and elapsed < Config.entry_delay:
if int(now) >= _last_wait_log + 5:
_last_wait_log = int(now)
log(f"⏳ 等待大户建仓 ({Config.entry_delay-elapsed:.0f}s) | BTC ${btc:,.2f}", "DATA")
time.sleep(0.3)
continue
# 平仓后冷却 (避免同回合过度交易)
if _last_exit_ts > 0:
cooldown_remaining = Config.reentry_cooldown - (now - _last_exit_ts)
if cooldown_remaining > 0:
if int(now) >= _last_wait_log + 5:
_last_wait_log = int(now)
log(f"🧊 平仓冷却中 ({cooldown_remaining:.0f}s) | BTC ${btc:,.2f}", "DATA")
time.sleep(0.3)
continue
# 剩余时间不够持有 → 跳过 (避免回合末成交, 止盈/止损无法触发)
if remaining < Config.min_hold_time:
if int(now) >= _last_wait_log + 5:
_last_wait_log = int(now)
log(f"⏭️ 剩余时间 {remaining:.0f}s < {Config.min_hold_time}s, 跳过本回合", "DATA")
if Config.exit_deadline > 0:
_bot_state = BotState.IDLE
time.sleep(2)
continue
# 读取链上情报
lb = read_leaderboard()
if not lb:
if not lb_warned:
log("⚠️ 无排行榜数据 (onchain_leaderboard 未运行?)", "WARN")
lb_warned = True
if Config.exit_deadline > 0 and remaining <= Config.exit_deadline:
log("⏭️ 无信号数据, 跳过本回合")
_bot_state = BotState.IDLE
time.sleep(2)
continue
lb_warned = False
# 快照本回合的 🔥 鲸鱼分布 (用于下一回合的"翻向"信号)
# 一次性捕获 (entry_delay 后, 此时鲸鱼已建仓)
if _streak_snapshot_round_ts != round_ts:
snap_sides = {}
for e in lb.get("up", []):
if e.get("streak") and not e.get("hedged"):
snap_sides[e["addr"]] = "up"
for e in lb.get("down", []):
if e.get("streak") and not e.get("hedged"):
snap_sides[e["addr"]] = "down"
_streak_snapshot = {"round_ts": round_ts, "sides": snap_sides}
_streak_snapshot_round_ts = round_ts
if snap_sides:
log(f"📸 🔥快照 r={datetime.fromtimestamp(round_ts).strftime('%H:%M:%S')}: "
f"{sum(1 for v in snap_sides.values() if v == 'up')}UP/"
f"{sum(1 for v in snap_sides.values() if v == 'down')}DN", "DATA")
analysis = analyze_signals(lb, btc, strike)
score = analysis["score"]
side = analysis["side"]
# 已有持仓时不再分析新信号 (避免日志刷屏)
if _position_info:
time.sleep(2)
continue
if side:
# 价格过滤: 入场价(含滑点)高于阈值则不下单
# 持续监控: 价格降下来后会自动下单
bid, ask = get_cached_price(side)
order_price = round(min(ask + Config.buy_slippage, 0.99), 2) if ask > 0 else 0
if Config.max_price < 1.0 and order_price > Config.max_price:
if int(now) >= _last_price_skip_log + 5:
_last_price_skip_log = int(now)
log(f"⏸️ 价格过高跳过: {side.upper()} 入场价{order_price:.2f} > 阈值 {Config.max_price:.2f}", "DATA")
time.sleep(2)
continue
# 盘口价差过滤: bid-ask 价差过大 = 流动性差, 入场即被套
# 双边滑点 + 大价差 = 入场就亏损, 必须过滤
if bid > 0 and ask > 0:
spread = ask - bid
if spread > Config.max_spread:
if int(now) >= _last_price_skip_log + 5:
_last_price_skip_log = int(now)
log(f"⏸️ 价差过大跳过: {side.upper()} spread={spread:.3f} > 阈值 {Config.max_spread:.3f} (bid={bid:.2f} ask={ask:.2f})", "DATA")
time.sleep(2)
continue
log(f"📡 入场信号: {side.upper()} (评分 {score:+d}/{analysis['max_score']}, 需 ±{analysis['min_abs']})")
for d in analysis["details"]:
log(d, "DATA")
result = place_order(side, Config.amount, is_sell_all=False)
if result.get("success"):
_bot_state = BotState.IN_POSITION
entry_price = result.get("price", 0)
entry_shares = result.get("shares", 0)
_position_info = {
"side": side,
"entry_price": entry_price,
"entry_score": score,
"entry_ts": now,
"shares": entry_shares,
"entry_signal": side.upper(),
"entry_signals": analysis["signals"],
"peak_pnl_pct": 0.0,
"trough_pnl_pct": 0.0,
"max_bid": entry_price,
"min_bid": entry_price,
"entry_streak_concentration": analysis.get("streak_concentration", 0.0),
}
rec("entry", round_ts=round_ts, strike=strike, btc=btc,
side=side, score=score, signals=analysis["signals"],
price=entry_price, shares=entry_shares, amount=Config.amount,
streak_up=analysis.get("streak_up", 0),
streak_dn=analysis.get("streak_dn", 0),
streak_concentration=analysis.get("streak_concentration", 0.0),
state=capture_market_state(round_ts, strike, remaining, side),
virtual_balance=round(_virtual.balance, 2) if _virtual else 0)
log(f"✅ 已入场 {side.upper()} | {entry_shares:.1f}股 @ {entry_price:.2f}")
else:
log(f"❌ 入场失败: {result.get('message', '')}", "WARN")
if Config.exit_deadline <= 0 or remaining > Config.exit_deadline + 10:
time.sleep(2)
else:
_bot_state = BotState.IDLE
else:
if int(now) >= _last_signal_log + 10:
_last_signal_log = int(now)
log(f"🔍 信号不足 (评分 {score:+d}, 需 ±{analysis['min_abs']}) | 剩余 {remaining:.0f}s", "DATA")
if Config.exit_deadline > 0 and remaining <= Config.exit_deadline:
log(f"⏭️ 信号不足, 跳过本回合 (评分 {score:+d})")
_bot_state = BotState.IDLE
time.sleep(0.5)
continue
# ── IDLE ──
time.sleep(2)
except KeyboardInterrupt:
raise
except Exception as e:
log(f"主循环异常: {e}", "ERROR")
traceback.print_exc()
time.sleep(5)
def _do_exit(side: str, reason: str = ""):
"""执行平仓 + 重置状态 (带即时重试, 无 sleep)
平仓后回到 WAITING_ENTRY (而非 IDLE), 允许同回合再次入场, 但需等 reentry_cooldown 秒
"""
global _bot_state, _position_info, _last_exit_ts
entry = _position_info or {}
for i in range(Config.max_retries + 1):
result = place_order(side, 0, is_sell_all=True)
if result.get("success"):
if i > 0:
log_trade(f"重试平仓成功 {side.upper()} (第{i}次)")
else:
log_trade(f"已平仓 {side.upper()}")
break
if i < Config.max_retries:
log_trade(f"平仓失败: {result.get('message', '')} (重试 {i+1}/{Config.max_retries})", success=False)
_bot_state = BotState.WAITING_ENTRY
_position_info = None
_last_exit_ts = time.time()
# 记录 (dry-run 回测数据)
entry_price = entry.get("entry_price", 0)
entry_shares = entry.get("shares", 0)
exit_price = result.get("price", 0) if result and result.get("success") else 0
if entry_price > 0 and entry_shares > 0 and exit_price > 0:
cost = entry_shares * entry_price
proceeds = entry_shares * exit_price
pnl = proceeds - cost
pct = (exit_price - entry_price) / entry_price * 100
rec("exit", reason=reason, entry_price=entry_price, exit_price=exit_price,
shares=entry_shares, pnl=round(pnl, 2), pnl_pct=round(pct, 1),
hold_seconds=round(time.time() - entry.get("entry_ts", time.time()), 1),
btc=get_btc_price(),
strike=ws_get("strike", 0),
peak_pnl_pct=entry.get("peak_pnl_pct", 0.0),
trough_pnl_pct=entry.get("trough_pnl_pct", 0.0),
max_bid=entry.get("max_bid", entry_price),
min_bid=entry.get("min_bid", exit_price),
state=capture_market_state(ws_get("round_ts", 0), ws_get("strike", 0), 0, side),
entry_signal=entry.get("entry_signal", ""),
entry_signals=entry.get("entry_signals", {}))
# 计算单笔 P&L (实盘)
if not Config.dry_run:
entry_price = entry.get("entry_price", 0)
entry_shares = entry.get("shares", 0)
exit_price = result.get("price", 0) if result and result.get("success") else 0
if entry_price > 0 and entry_shares > 0 and exit_price > 0:
cost = entry_shares * entry_price
proceeds = entry_shares * exit_price
pnl = proceeds - cost
record_live_trade(pnl)
pct = (exit_price - entry_price) / entry_price * 100
log_trade(f"本笔 PnL: ${pnl:+.2f} ({pct:+.1f}%) | 累计: ${_live_stats['total_pnl']:+.2f} | {_live_stats['wins']}{_live_stats['losses']}负")
if Config.dry_run and _virtual:
s = _virtual.status()
log(f"💼 虚拟账户 | 余额: ${s['balance']:.2f} | 累计PnL: ${s['total_pnl']:+.2f} | 交易: {s['trades']}笔")
# ════════════════════════════════════════════════════════════════════════════
# CLI 入口
# ════════════════════════════════════════════════════════════════════════════
def main():
parser = argparse.ArgumentParser(description="BTC 5m 差价交易机器人 (趋势跟随策略)")
parser.add_argument("--live", action="store_true", help="实盘交易 (默认 dry run)")
parser.add_argument("--amount", type=float, default=None, help="每笔金额 (USD)")
parser.add_argument("--market", type=str, default="btc_5m", help="市场 (btc_5m/btc_15m/eth_5m...)")
parser.add_argument("--min-score-ratio", type=float, default=None, help="最低入场信号分占满分比例 (0.5=50pct, 范围 0~1)")
parser.add_argument("--take-profit", type=float, default=None, help="止盈: token涨幅 (0.10=10pct, -1=禁用)")
parser.add_argument("--stop-loss", type=float, default=None, help="止损: token跌幅 (0.15=15pct, -1=禁用)")
parser.add_argument("--exit-deadline", type=int, default=None, help="强制平仓倒计时 (秒, -1=禁用)")
parser.add_argument("--reentry-cooldown", type=int, default=None, help="平仓后冷却秒数 (同回合再次入场前等待, 0=立即, 默认10)")
parser.add_argument("--stop-loss-grace", type=int, default=None, help="入场后止损宽限期 (秒, 默认5, 期间不检查止损)")
parser.add_argument("--max-spread", type=float, default=None, help="入场前盘口价差上限 (0.05=5¢, 默认0.05)")
parser.add_argument("--weight-btc-dev", type=int, default=None, help="BTC 位移权重 (核心), 默认4")
parser.add_argument("--weight-btc-accel", type=int, default=None, help="BTC 加速度权重 (核心), 默认3")
parser.add_argument("--weight-book", type=int, default=None, help="盘口压力失衡权重, 默认2")
parser.add_argument("--weight-whale-follow", type=int, default=None, help="鲸鱼跟随权重 (curated 列表, 已弃用, 留作兼容), 默认0")
parser.add_argument("--weight-winners", type=int, default=None, help="上轮赢家方向权重, 默认2")
parser.add_argument("--weight-streak-leader", type=int, default=None, help="连赢鲸鱼 🔥 动态信号权重 (核心), 默认4")
parser.add_argument("--streak-threshold", type=float, default=None, help="combined 投票入场门槛 (backtest: 0.30-0.45), 默认0.30")
parser.add_argument("--max-price", type=float, default=None, help="token价格上限 (0.60=60¢, 1.0=禁用)")
parser.add_argument("--buy-slippage", type=float, default=None, help="买入滑点 (0.03=3¢)")
parser.add_argument("--sell-slippage", type=float, default=None, help="卖出滑点 (0.03=3¢)")
parser.add_argument("--balance", type=float, default=None, help="dry-run 初始余额")
args = parser.parse_args()
Config.dry_run = not args.live
if args.market is not None:
Config.market = args.market
if args.amount is not None:
Config.amount = args.amount
if args.min_score_ratio is not None:
Config.min_score_ratio = max(0.0, min(1.0, args.min_score_ratio))
if args.take_profit is not None:
Config.take_profit_pct = args.take_profit
if args.stop_loss is not None:
Config.stop_loss_pct = args.stop_loss
if args.exit_deadline is not None:
Config.exit_deadline = args.exit_deadline
if args.reentry_cooldown is not None:
Config.reentry_cooldown = args.reentry_cooldown
if args.stop_loss_grace is not None:
Config.stop_loss_grace = args.stop_loss_grace
if args.max_spread is not None:
Config.max_spread = args.max_spread
if args.weight_btc_dev is not None:
Config.weight_btc_dev = args.weight_btc_dev
if args.weight_btc_accel is not None:
Config.weight_btc_accel = args.weight_btc_accel
if args.weight_book is not None:
Config.weight_book = args.weight_book
if args.weight_whale_follow is not None:
Config.weight_whale_follow = args.weight_whale_follow
if args.weight_winners is not None:
Config.weight_winners = args.weight_winners
if args.weight_streak_leader is not None:
Config.weight_streak_leader = args.weight_streak_leader
if args.streak_threshold is not None:
Config.streak_combined_threshold = max(0.05, min(1.0, args.streak_threshold))
if args.max_price is not None:
Config.max_price = args.max_price
if args.buy_slippage is not None:
Config.buy_slippage = args.buy_slippage
if args.sell_slippage is not None:
Config.sell_slippage = args.sell_slippage
if args.balance is not None:
Config.virtual_balance = args.balance
if Config.market not in qt.MARKETS:
log(f"未知市场: {Config.market}", "ERROR")
sys.exit(1)
# 连接 onchain_leaderboard 的 WebSocket
if not start_ws_client():
log("无法启动排行榜 WS 客户端, 退出", "ERROR")
sys.exit(1)
log("等待 WS 数据就绪 (BTC 价格)...")
for _ in range(30):
if get_btc_price() > 0:
break
time.sleep(1)
else:
log("无法获取 BTC 价格, 请检查 onchain_leaderboard 是否在运行", "ERROR")
sys.exit(1)
log(f"BTC 现价: ${get_btc_price():,.2f}")
if not Config.dry_run:
try:
qt.get_clob_client()
except Exception as e:
log(f"CLOB 初始化失败: {e}", "ERROR")
sys.exit(1)
else:
global _virtual, _recorder
_virtual = VirtualPortfolio(Config.virtual_balance)
_recorder = Recorder(f"backtest_{datetime.now().strftime('%m%d_%H%M')}.jsonl")
try:
run_bot()
except KeyboardInterrupt:
log("\n👋 退出")
if Config.dry_run and _virtual:
s = _virtual.status()
log(f"最终账户: ${s['balance']:.2f} | PnL: ${s['total_pnl']:+.2f} | {s['trades']}笔交易")
elif not Config.dry_run:
log(f"实盘总结: {_live_stats['trades']}笔 | {_live_stats['wins']}{_live_stats['losses']}负 | 累计PnL: ${_live_stats['total_pnl']:+.2f}")
finally:
if _recorder:
_recorder.close()
log(f"📁 回测数据已保存: {_recorder.path}")
if __name__ == "__main__":
main()