feat(live): D-3 sanguo实盘分支(影子下单)+D期设计文档

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期完整设计)
This commit is contained in:
2026-07-11 00:03:46 +08:00
parent eff9ed2ae9
commit ff84b3d4b0
6 changed files with 380 additions and 0 deletions
+7
View File
@@ -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影子)
@@ -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 期模拟盘已端到端跑通(PaperEngineraw/qfq 双源 + 总账/分户双层记账 + 软限额 + 占用成本 + 分红送股)。D 期目标:**接入实盘**,实现「模拟→实盘」同引擎切换。
核心约束(决定架构):
- **miniQMT 仅 Windows 桌面**,必须登录常驻,提供 `xtquant` Python API
- **sanguo 跑 NAS Linux Docker 容器**`sanguo_vnpy_v2`
- **Windows 与 NAS 在不同网络**(异地),需跨网打通
- 实走为**日线级**(每日 20:30 单根 bar 推进),对延迟不敏感
→ 结论:唯一可行路径是 **miniQMTxtquant**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 监听确认)|
| frpsVPS| 无 allowPorts 限制 | ✅ |
| CaddyVPS| `bridge.mysanguo.top { reverse_proxy 127.0.0.1:18765 }` | ✅ validate + reload |
| DNSNameSilo| `bridge A 43.133.235.218` TTL 3600 | ✅ 生效 |
| HTTPS 证书 | Caddy 自动 ACMETLS1.3| ✅ |
| 全链路验证 | `curl https://bridge.mysanguo.top``502 server: Caddy` | ✅ 502=隧道通到 Windows8765 待写)|
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=<ts>` | `[{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: <random>`bridge 校验,不符 401。token 走环境变量(sanguo + Windows bridge 两端同值),**不进 git**
2. **最小接口**:只放 §4.2 的接口,不暴露 xtquant 全能力
3. **限速**FastAPI middleware 限流(防爆破)
4. **可选 IP 白名单**sanguo 经 VPS 反代,bridge 看到的源 IP 是 VPS43.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]] wikiFRP/Caddy 基建)
- vps-access skill
- 前序:`docs/superpowers/specs/2026-07-07-phase3c-paper-trading-design.md`
- 网络层落地:NAS `/volume1/stock/frp_windows/`frpc + README
+2
View File
@@ -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"),
)
+100
View File
@@ -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 codesh600000 / 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):
沪市 SSE60/68/51/56/58 开头 → sh
深市 SZSE00/30/15 开头 → sz
默认 → shguess_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")
+63
View File
@@ -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 完成 @%spending=%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 立即 returnlive_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_idlive_runner 遍历用)。"""
+35
View File
@@ -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