From ff84b3d4b0ca8f9314d692c616be1524df5cb411 Mon Sep 17 00:00:00 2001 From: claude_dev Date: Sat, 11 Jul 2026 00:03:46 +0800 Subject: [PATCH] =?UTF-8?q?feat(live):=20D-3=20sanguo=E5=AE=9E=E7=9B=98?= =?UTF-8?q?=E5=88=86=E6=94=AF(=E5=BD=B1=E5=AD=90=E4=B8=8B=E5=8D=95)+D?= =?UTF-8?q?=E6=9C=9F=E8=AE=BE=E8=AE=A1=E6=96=87=E6=A1=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit D-3 模式A影子下单(spec §5): - bridge_client.py: QMT bridge HTTP客户端(urllib, X-Bridge-Token, 失败不抛返回None) - live_orchestrator: _shadow_trades_to_bridge 当日成交POST bridge(默认enabled=false) - persistence: paper_shadow_orders幂等表+save_shadow_order/is_trade_shadowed - config: data_platform.yaml加live段, token走env(BRIDGE_TOKEN) - to_bridge_code symbol转换与guess_exchange一致(2位前缀) 安全: enabled=false默认关+token走env+幂等防重复+影子失败不阻断live_step docs: phase3d-live-trading-design.md(D期完整设计) --- config/data_platform.yaml | 7 + .../2026-07-10-phase3d-live-trading-design.md | 173 ++++++++++++++++++ sanguo_data/config.py | 2 + sanguo_trader/bridge_client.py | 100 ++++++++++ sanguo_trader/live_orchestrator.py | 63 +++++++ sanguo_trader/persistence.py | 35 ++++ 6 files changed, 380 insertions(+) create mode 100644 docs/superpowers/specs/2026-07-10-phase3d-live-trading-design.md create mode 100644 sanguo_trader/bridge_client.py diff --git a/config/data_platform.yaml b/config/data_platform.yaml index 9b03f97..a31cab2 100644 --- a/config/data_platform.yaml +++ b/config/data_platform.yaml @@ -39,3 +39,10 @@ performance: # 资金占用成本归因(spec §195):年化无风险利率,每策略占用资金按此日扣归因到 PnL risk_free_rate: 0.02 + +# 实盘集成(D期,spec §5)— 默认关闭,D-4a 联调再开 +# bridge_token 走环境变量 BRIDGE_TOKEN,不写入 yaml(不进 git) +live: + enabled: false # 总开关(false=影子分支整个跳过,live_step 行为不变) + bridge_url: https://bridge.mysanguo.top + shadow: true # 模式A影子下单(模拟撮合为准,信号同步POST bridge影子) diff --git a/docs/superpowers/specs/2026-07-10-phase3d-live-trading-design.md b/docs/superpowers/specs/2026-07-10-phase3d-live-trading-design.md new file mode 100644 index 0000000..e74b752 --- /dev/null +++ b/docs/superpowers/specs/2026-07-10-phase3d-live-trading-design.md @@ -0,0 +1,173 @@ +# Phase 3D 实盘交易集成设计(miniQMT bridge 架构) + +> 日期:2026-07-10 | 状态:设计中(网络层已完成,bridge/编排待开发) +> 前序:Phase 3C 模拟盘(PaperEngine 双源+双层记账,commit 1646903/0656108 已完成) +> 关联:[[khquant-analysis]](xtquant 参考)、vps-access skill、`docs/data-platform/daily-update-design.md` + +--- + +## 1. 背景与目标 + +C 期模拟盘已端到端跑通(PaperEngine:raw/qfq 双源 + 总账/分户双层记账 + 软限额 + 占用成本 + 分红送股)。D 期目标:**接入实盘**,实现「模拟→实盘」同引擎切换。 + +核心约束(决定架构): +- **miniQMT 仅 Windows 桌面**,必须登录常驻,提供 `xtquant` Python API +- **sanguo 跑 NAS Linux Docker 容器**(`sanguo_vnpy_v2`) +- **Windows 与 NAS 在不同网络**(异地),需跨网打通 +- 实走为**日线级**(每日 20:30 单根 bar 推进),对延迟不敏感 + +→ 结论:唯一可行路径是 **miniQMT(xtquant)**,PTrade/QMT 完整版封闭不可集成(详见 [[khquant-analysis]] 同源调研)。 + +--- + +## 2. 总体架构 + +``` +[网络A · 局域网] [公网 VPS] [网络B · 异地] + 43.133.235.218 +NAS Docker ┌─ frps(:7000) Windows 机器 + sanguo_vnpy_v2 ──HTTPS────────►│ Caddy(:443) ├ miniQMT 客户端(登录常驻) + live_orchestrator 下单/查询 │ bridge.mysanguo.top ├ bridge 服务(:8765) ← D期写 + │ → 127.0.0.1:18765 │ ↓ xtquant 本地 IPC + └─ frps(:18765) ◄─frp隧道──── └ frpc(连 frps:7000) +``` + +下单链路(6 跳): +`sanguo(NAS) → 公网 → Caddy(VPS:443) → frps(18765) → frp隧道 → Windows frpc → bridge(:8765) → xtquant → miniQMT` + +日线级单 bar 推进,延迟完全无感;高频不适用(非本项目场景)。 + +--- + +## 3. 网络层(已完成 ✅) + +| 组件 | 配置 | 状态 | +|------|------|------| +| Windows frpc | `frp v0.69.1`,`frpc.toml`(serverAddr 43.133.235.218:7000 + token + qmt-bridge proxy 8765→18765)| ✅ 已连(VPS 18765 监听确认)| +| frps(VPS)| 无 allowPorts 限制 | ✅ | +| Caddy(VPS)| `bridge.mysanguo.top { reverse_proxy 127.0.0.1:18765 }` | ✅ validate + reload | +| DNS(NameSilo)| `bridge A 43.133.235.218` TTL 3600 | ✅ 生效 | +| HTTPS 证书 | Caddy 自动 ACME(TLS1.3)| ✅ | +| 全链路验证 | `curl https://bridge.mysanguo.top` → `502 server: Caddy` | ✅ 502=隧道通到 Windows(8765 待写)| + +Windows frpc 开机自启:任务计划程序(`sanguo-frpc`,onstart + 失败重启)或启动文件夹,见 NAS `/volume1/stock/frp_windows/README.md`。 + +--- + +## 4. bridge 服务设计(Windows 端,D-1 待开发) + +### 4.1 形态 +- **FastAPI** HTTP 服务,监听 `127.0.0.1:8765`(仅本地,frpc 转发外部流量) +- 启动时初始化 `xtquant.XtQuantTrader`(连本地 miniQMT 客户端)+ `xtdata`(行情) +- 自启:任务计划程序(`sanguo-bridge`,onstart,依赖 miniQMT 客户端已登录) + +### 4.2 接口(最小集,YAGNI) + +| 方法 | 路径 | 入参 | 返回 | 说明 | +|------|------|------|------|------| +| GET | `/health` | — | `{status, miniqmt_connected}` | 健康检查(frpc/Caddy 探活)| +| POST | `/order` | `{code, action:buy/sell, price, volume, reason}` | `{order_id, ok}` | 下单(`xt_trader.order_stock`)| +| POST | `/cancel` | `{order_id}` | `{ok}` | 撤单 | +| GET | `/account` | — | `{cash, frozen, market_value, total}` | 资金(`xt_trader.query_stock_asset`)| +| GET | `/positions` | — | `[{code, volume, can_use, avg_price, ...}]` | 持仓(`xt_trader.query_stock_positions`)| +| GET | `/trades` | `?since=` | `[{code, action, price, volume, time}]` | 成交回报(账本同步用)| + +> 股票代码格式:xtquant 用 `600000.SH` / `000001.SZ`(sanguo 内部 `sh600000`,bridge 做转换)。 + +### 4.3 xtquant 调用要点(参考 OSkhQuant 架构,不抄代码) +- `xt_trader = XtQuantTrader(path, session_id)`;`xt_trader.start()`;`connect()` 后 `subscribe(account)` +- 下单:`xt_trader.order_stock(account, code, order_type, volume, price_type, price, strategy_name, order_remark)` + - `price_type`:限价 `XT_PRICE_LIMITED` / 市价 `XT_PRICE_LATEST_PRICE` 等 + - A 股 T+1:`query_stock_positions` 的 `can_use_volume` 即可卖量(buy 当日计 frozen) +- 行情:`xtdata.download_history_data` + `get_market_data_ex`(实盘 step 用实时,非历史) +- 回调:`xt_trader.register_callback` 异步接收成交通知 + +--- + +## 5. sanguo 端改造(D-3 待开发) + +### 5.1 配置(config/data_platform.yaml 加) +```yaml +live: + bridge_url: https://bridge.mysanguo.top + bridge_token: ${BRIDGE_TOKEN} # 环境变量,不进 git + enabled: false # 总开关,模拟联调时再开 +``` + +### 5.2 live_orchestrator 改造(`sanguo_trader/live_orchestrator.py`) +现状:`live_step` 恢复状态 → warmup → 当日 bar → `engine.step` → 存状态。`step` 返回 `(pending_new, closes)`。 + +D 期加「实盘执行分支」:当 `account.mode == 'live'` 且 `live.enabled`: +1. `step` 产生的**当日成交**(`closes`)→ 同步 POST `/order` 到 bridge(真实下单到 miniQMT) +2. 次日开盘前,从 bridge `GET /positions` `/account` 拉真实持仓/资金,**校正** account 账本(真实回报为准,纠模拟撮合漂移) +3. 鉴权:每个请求带 `X-Bridge-Token` header + +### 5.3 模拟撮合 vs 实盘下单的关系(关键设计决策) + +| 模式 | 说明 | D 期采用 | +|------|------|---------| +| **A 影子下单**(推荐先)| PaperEngine 照常模拟撮合(账本准),同时把信号 POST bridge「影子」下单到 miniQMT 模拟环境,**对比两者**验证一致性 | ✅ 联调期 | +| **B 实盘驱动** | 真实下单 + 成交回报驱动账本,PaperEngine 退化为信号生成器 | 切实盘后 | + +→ 联调先用 A(模拟盘端到端,零资金风险),一致性验证后切实盘切 B。 + +--- + +## 6. 鉴权与安全(⚠️ bridge.mysanguo.top 已公网暴露) + +**实测**:域名一上线即被扫描器(`81.171.74.60` 等)打 `/dump.sql` `/wp-config.php` `/secrets.json`。必须: + +1. **共享密钥**:每个请求 header `X-Bridge-Token: `,bridge 校验,不符 401。token 走环境变量(sanguo + Windows bridge 两端同值),**不进 git** +2. **最小接口**:只放 §4.2 的接口,不暴露 xtquant 全能力 +3. **限速**:FastAPI middleware 限流(防爆破) +4. **可选 IP 白名单**:sanguo 经 VPS 反代,bridge 看到的源 IP 是 VPS(43.133.235.218)→ bridge 可加白名单只接受 frps 来源 +5. **审计日志**:bridge 记录每笔下单(code/action/volume/price/来源 IP/时间),便于复盘异常 + +--- + +## 7. 端到端联调方案(模拟盘先行) + +| 阶段 | 环境 | 资金风险 | 目标 | +|------|------|---------|------| +| D-4a | miniQMT **模拟客户端**(现已在跑)| 零 | bridge 端到端打通:sanguo 信号 → bridge → miniQMT 模拟下单 → 回报 | +| D-4b | 模拟客户端 + **影子对比** | 零 | PaperEngine 模拟撮合 vs bridge 真实下单,验证一致性(成交价/持仓/资金)| +| D-4c | **小资金实盘**(切实盘账户)| 低 | 真金白银小单验证,切换模式 B | +| D-4d | 正式实盘 | 正常 | 纳入每日 20:30 scheduler | + +--- + +## 8. 任务拆分(D 期工作清单) + +| 编号 | 任务 | 端 | 依赖 | +|------|------|-----|------| +| D-1 | bridge 服务(FastAPI + xtquant + 鉴权 + 6 接口)| Windows | — | +| D-2 | bridge 自启(任务计划程序)+ frpc 自启确认 | Windows | D-1 | +| D-3 | sanguo live_orchestrator 实盘分支 + config + 鉴权客户端 | NAS | D-1 接口定义 | +| D-4a | 模拟盘端到端联调(影子下单打通)| 两端 | D-1/D-3 | +| D-4b | 模拟撮合 vs 实盘下单一致性验证 | 两端 | D-4a | +| D-4c | 小资金实盘切模式 B | 两端 | D-4b | +| D-5 | 文档/验收/部署 | — | D-4 | + +--- + +## 9. 风险与兜底 + +| 风险 | 影响 | 兜底 | +|------|------|------| +| Windows/miniQMT/frpc/bridge 四常驻,任一断 | 下单链路断 | 日线级 → 「断线次日补」+ bridge `/health` 探活 + scheduler 重试 | +| VPS 单点 | 全链路断 | 接受(日线级);备选 Tailscale 直连绕 VPS | +| bridge token 泄露 | 任意人可下单 | 环境变量 + 不 commit + 审计日志 + 限速 | +| miniQMT 停新申请(2026/7/6)| 新账户无法开 | 老账户可用;新账户换其他提供 miniQMT 券商(华泰/中泰/国信)| +| 模拟撮合与实盘成交价漂移 | 账本不准 | 模式 B 以 bridge 回报为准校正 | +| 公网扫描/攻击 | bridge 被打 | §6 鉴权 + 最小接口 + 限速 | + +--- + +## 10. 相关 + +- 源码参考:`/volume1/KnowledgeBase/github-repos/OSkhQuant`(xtquant 调用链路,CC BY-NC 仅参考架构) +- [[khquant-analysis]] wiki +- [[vps-deployment]] wiki(FRP/Caddy 基建) +- vps-access skill +- 前序:`docs/superpowers/specs/2026-07-07-phase3c-paper-trading-design.md` +- 网络层落地:NAS `/volume1/stock/frp_windows/`(frpc + README) diff --git a/sanguo_data/config.py b/sanguo_data/config.py index cebc15e..8f2b80c 100644 --- a/sanguo_data/config.py +++ b/sanguo_data/config.py @@ -10,6 +10,7 @@ class DataConfig: validation: dict performance: dict risk_free_rate: float = 0.02 # 年化无风险利率(spec §195 资金占用成本归因) + live: dict | None = None # 实盘集成(D期 spec §5),None/缺省=disabled def load_config(path: str) -> DataConfig: try: @@ -29,6 +30,7 @@ def load_config(path: str) -> DataConfig: validation=raw.get("validation", {}), performance=raw.get("performance", {}), risk_free_rate=float(raw.get("risk_free_rate", 0.02)), + live=raw.get("live"), ) diff --git a/sanguo_trader/bridge_client.py b/sanguo_trader/bridge_client.py new file mode 100644 index 0000000..260d4a5 --- /dev/null +++ b/sanguo_trader/bridge_client.py @@ -0,0 +1,100 @@ +"""sanguo 端 QMT bridge HTTP 客户端(D-3 影子下单)。 + +封装跨网调 Windows bridge 的 HTTP 调用(POST /order、GET /account、GET /positions)。 +设计原则:失败不抛异常,记 logger.warning 返回 None —— 影子下单是旁路, +绝不能阻断 PaperEngine 的模拟盘撮合主流程(spec §5 模式 A)。 + +依赖:仅标准库 urllib(不引入 requests/httpx),避免新增第三方依赖。 +接口契约见 sanguo_qmt_bridge/README.md。 +""" +import json +import logging +import urllib.error +import urllib.request +from typing import Any + +logger = logging.getLogger(__name__) + +_TIMEOUT = 10 # 秒 + + +def to_bridge_code(symbol: str) -> str: + """sanguo symbol → bridge code(sh600000 / sz000001)。 + + sanguo 内部 symbol 为纯数字码(如 '600000'、'000001'), + paper_trades.symbol 即此格式(见 engine._match 的 save_trade 入参)。 + bridge 接受 sh/sz 前缀码(见 sanguo_qmt_bridge/README.md「代码格式」)。 + + 规则(与 sanguo_data.datareader.guess_exchange 一致,避免在客户端 import vnpy): + 沪市 SSE:60/68/51/56/58 开头 → sh + 深市 SZSE:00/30/15 开头 → sz + 默认 → sh(guess_exchange 兜底亦为 SSE) + + 若 symbol 已含 sh/sz 前缀,直接小写返回。 + """ + if symbol[:2].lower() in ("sh", "sz"): + return symbol.lower() + if symbol.startswith(("60", "68", "51", "56", "58")): + return f"sh{symbol}" + if symbol.startswith(("00", "30", "15")): + return f"sz{symbol}" + return f"sh{symbol}" + + +class BridgeClient: + """QMT bridge HTTP 客户端(影子下单旁路)。 + + 每个请求带 header X-Bridge-Token。任何网络/解析失败均返回 None 并记 warning, + 不向调用方抛异常 —— 影子下单失败不得阻断 live_step 主流程。 + """ + + def __init__(self, url: str, token: str) -> None: + self._url = url.rstrip("/") + self._token = token + + def _headers(self) -> dict[str, str]: + return {"X-Bridge-Token": self._token, "Content-Type": "application/json"} + + def _post(self, path: str, payload: dict[str, Any]) -> dict | None: + """POST JSON,返回解析后 dict;失败记 warning 返回 None。""" + data = json.dumps(payload).encode("utf-8") + req = urllib.request.Request( + f"{self._url}{path}", data=data, headers=self._headers(), method="POST" + ) + try: + with urllib.request.urlopen(req, timeout=_TIMEOUT) as resp: + return json.loads(resp.read().decode("utf-8")) + except (urllib.error.URLError, TimeoutError, json.JSONDecodeError, OSError) as e: + logger.warning("bridge POST %s 失败: %s", path, e) + return None + + def _get(self, path: str) -> dict | None: + """GET JSON,返回解析后 dict;失败记 warning 返回 None。""" + req = urllib.request.Request( + f"{self._url}{path}", headers=self._headers(), method="GET" + ) + try: + with urllib.request.urlopen(req, timeout=_TIMEOUT) as resp: + return json.loads(resp.read().decode("utf-8")) + except (urllib.error.URLError, TimeoutError, json.JSONDecodeError, OSError) as e: + logger.warning("bridge GET %s 失败: %s", path, e) + return None + + def place_order(self, code: str, action: str, price: float, volume: int, + price_type: str = "limit", reason: str = "") -> dict | None: + """POST /order → {ok: bool, order_id: int} 或 {ok: false, error: str};失败返回 None。""" + return self._post("/order", { + "code": code, "action": action, "price": price, + "volume": volume, "price_type": price_type, "reason": reason, + }) + + def get_account(self) -> dict | None: + """GET /account → {ok, cash, frozen, market_value, total};失败返回 None。""" + return self._get("/account") + + def get_positions(self) -> list | None: + """GET /positions → positions 列表;失败返回 None。""" + resp = self._get("/positions") + if resp is None: + return None + return resp.get("positions") diff --git a/sanguo_trader/live_orchestrator.py b/sanguo_trader/live_orchestrator.py index 03a7f76..9136d27 100644 --- a/sanguo_trader/live_orchestrator.py +++ b/sanguo_trader/live_orchestrator.py @@ -11,6 +11,7 @@ listing_days 从 IPO 日算(首版 stub 0);realized_pnl 未恢复(归因 """ import json import logging +import os import sqlite3 from datetime import datetime, timedelta @@ -23,6 +24,7 @@ from .position_ledger import PositionLedger from .persistence import ( load_last_balance, load_positions, save_positions, load_pending_orders, save_pending_orders, + save_shadow_order, is_trade_shadowed, ) from .strategy_runner import StrategyRunner @@ -164,6 +166,67 @@ def live_step(db_path: str, account_id: int, data_source, cfg, today: str | None logger.info("live_step %s 完成 @%s,pending=%d positions=%d", account_id, today, len(pending_new), len(account.positions)) + # 7. 影子下单(D-3,spec §5 模式 A):当日成交 POST bridge,默认关闭 + # enabled=false 时 _shadow_trades_to_bridge 立即 return,live_step 行为完全不变 + _shadow_trades_to_bridge(db_path, account_id, today, cfg) + + +def _shadow_trades_to_bridge(db_path: str, account_id: int, today: str, cfg) -> None: + """当日成交影子下单到 bridge(D-3,spec §5 模式 A)。 + + PaperEngine 模拟撮合为准,当日成交信号同步 POST 到 bridge 影子下单到 miniQMT。 + 默认关闭(cfg.live.enabled=false);任何失败仅记日志,不阻断 live_step。 + 幂等:paper_shadow_orders UNIQUE(account_id, trade_id) 保证 scheduler 重跑不重复下单。 + """ + try: + live_cfg = getattr(cfg, "live", None) or {} + if not live_cfg.get("enabled"): + return + token = os.environ.get("BRIDGE_TOKEN") + if not token: + logger.warning("live_step %s: 影子下单启用但 BRIDGE_TOKEN 未设,跳过", account_id) + return + url = live_cfg.get("bridge_url") + if not url: + logger.warning("live_step %s: 影子下单启用但 bridge_url 未配,跳过", account_id) + return + + from .bridge_client import BridgeClient, to_bridge_code + + client = BridgeClient(url, token) + with sqlite3.connect(db_path) as conn: + conn.row_factory = sqlite3.Row + cur = conn.execute( + "SELECT id, strategy_id, symbol, direction, price, volume " + "FROM paper_trades WHERE account_id=? AND bar_date=? AND rejected=0", + (account_id, today), + ) + trades = [dict(r) for r in cur.fetchall()] + + shadowed = 0 + for t in trades: + if is_trade_shadowed(db_path, account_id, t["id"]): + continue + code = to_bridge_code(t["symbol"]) + resp = client.place_order( + code, t["direction"], t["price"], t["volume"], + reason=f"shadow:strategy:{t['strategy_id']}", + ) + ok = bool(resp and resp.get("ok")) + status = "ok" if ok else "failed" + bridge_order_id = resp.get("order_id") if resp else None + save_shadow_order(db_path, account_id, t["id"], bridge_order_id, status) + if ok: + shadowed += 1 + else: + logger.warning("live_step %s: 影子下单失败 trade=%s code=%s err=%s", + account_id, t["id"], code, + resp.get("error") if resp else "no_response") + if trades: + logger.info("live_step %s 影子下单: %d/%d ok", account_id, shadowed, len(trades)) + except Exception as e: # noqa: BLE001 影子下单绝不阻断 live_step + logger.warning("live_step %s: 影子下单异常(不阻断): %s", account_id, e) + def list_live_accounts(db_path: str) -> list[int]: """所有 mode=live & status=running 的 account_id(live_runner 遍历用)。""" diff --git a/sanguo_trader/persistence.py b/sanguo_trader/persistence.py index cf02f2c..b1fa088 100644 --- a/sanguo_trader/persistence.py +++ b/sanguo_trader/persistence.py @@ -51,6 +51,12 @@ CREATE TABLE IF NOT EXISTS paper_pending_orders ( is_market INTEGER, match_session TEXT, listing_days INTEGER, created_at TEXT ); +CREATE TABLE IF NOT EXISTS paper_shadow_orders ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + account_id INTEGER, trade_id INTEGER, bridge_order_id INTEGER, + status TEXT, posted_at TEXT, + UNIQUE(account_id, trade_id) +); """ @@ -288,3 +294,32 @@ def load_last_balance(db_path: str, account_id: int) -> dict | None: ) row = cur.fetchone() return dict(row) if row else None + + +# ---- D-3 影子下单幂等(spec §5 模式 A)---- + +def save_shadow_order(db_path: str, account_id: int, trade_id: int, + bridge_order_id: int | None, status: str) -> None: + """记录影子下单结果(幂等:UNIQUE(account_id, trade_id),重复 INSERT 被 IGNORE)。 + + 无论 ok/failed 都记录 —— 保证 scheduler 重跑不重复 POST(「不能重复下单」硬约束)。 + status 取值:'ok'(bridge 返回 ok=true)/ 'failed'(ok=false 或网络失败)。 + """ + with sqlite3.connect(db_path) as conn: + conn.execute( + """INSERT OR IGNORE INTO paper_shadow_orders + (account_id, trade_id, bridge_order_id, status, posted_at) + VALUES (?,?,?,?,?)""", + (account_id, trade_id, bridge_order_id, status, _now()), + ) + conn.commit() + + +def is_trade_shadowed(db_path: str, account_id: int, trade_id: int) -> bool: + """该成交是否已影子下单(幂等去重,避免 scheduler 重试重复 POST)。""" + with sqlite3.connect(db_path) as conn: + cur = conn.execute( + "SELECT 1 FROM paper_shadow_orders WHERE account_id=? AND trade_id=?", + (account_id, trade_id), + ) + return cur.fetchone() is not None