diff --git a/README.md b/README.md index d4d4c60..bda0382 100644 --- a/README.md +++ b/README.md @@ -13,7 +13,7 @@ | `POSITION_MGMT_DESIGN.md` | 总体设计 **V0.4(定稿,开发启动)**:命令系统与管理页面/账本/仓位框架/动作引擎/两道关口/择时执行/下游通道。功能一次性开发,上线按依赖分三步切换 | | `QMT_WS_PROTOCOL.md` | **PMS ↔ QMT WebSocket 指令与回报协议 V1.0(定稿)**:传输与重连、Ed25519 签名与幂等、消息集、状态机、断线补发与对账兜底、部署前检查清单。**这是下发通道的唯一实现依据** | | `QMT_INTERFACE_REQUIREMENTS.md` | 与 QMT 侧的数据与接口需求清单 **V2.0**:A 部分(只读数据)与 C 部分(切换约定)有效;**B 部分的表通道已废止**,改由上面的 ws 协议承担 | -| `ddl_pms_v1.sql` | PMS 全部自有表建表语句(153 代理侧,**13 张**:设计 §11 的 10 张 + ws 通道 3 张) | +| `ddl_pms_v1.sql` | PMS 全部自有表建表语句(153 代理侧,**14 张**:设计 §11 的 10 张 + ws 通道 3 张 + 现金流水 1 张) | | `config/settings.py` | 配置(基础设施键名对齐 bionic;业务参数为初值,页面调参持久化到 `pms_runtime_param` 后优先) | ## 模块地图 @@ -57,8 +57,8 @@ scripts/ test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 19 例 test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 11 例 test_batch5_units.py 决策系统信号流解析与消化口径 8 例 - test_batch6_units.py ws 通道: 测试向量/签名/公钥形态/水位/DDL体检 51 例 - test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 34 例 + test_batch6_units.py ws 通道: 测试向量/签名/公钥/水位/DDL/逐笔入账 58 例 + test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 35 例 init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练; 含 DDL 体检) check_db.py 实机连通性与表结构自检 (需真实 .env) gen_keys.py ws 通道密钥: 生成 / 只取公钥(--pubkey) / PEM(--pem) / 自检(--check) @@ -199,8 +199,9 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ | # | 事项 | 状态 | |---|---|---| -| 1 | **ws 通道的账本侧改造**(清单 4~6) | 可立即开工,见下方「ws 通道实现清单」 | -| 2 | ws 通道联调(协议 §9 的 S1/S2/S3) | 需 QMT 侧配合:公钥交换 + IP 加白 | +| 1 | ~~ws 通道的账本侧改造(清单 4~6)~~ | ✅ 2026-07-29 完成 | +| 2 | **ws 通道联调**(协议 §9 的 S1/S2/S3) | 密钥已交换,可开工 | +| 2.5 | 部署便利性:`make deploy` 一句话完成 build + up + 各 profile | 待办(每次改动要敲 4 条命令太烦) | | 3 | T0 做T(二期) | 可做,设计已有,无外部依赖 | | 4 | 择时实现 A(委托决策系统盘中择时) | 阻塞:等 bionic 侧接口 | | 5 | 研判闸接通 | 阻塞:等 bionic 侧 `process_intraday_audit` 新增 PMS 请求 direction。客户端已就位,接口好了在页面填 `PMS_JUDGE_API_BASE` 即通 | @@ -213,9 +214,9 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ 1. ✅ **常驻连接进程** —— `app/ws/runner.py` + compose 的 `pms-ws`(`--profile ws`,**绝不可扩副本**,协议 §1 只准一条连接)。四个协程:存活心跳 / 收帧落库 / 5 秒 ping / 0.5 秒出口轮询。SIGTERM 优雅退出(停取新单 → 刷水位 → 最后一次 ack → 关连接 → 置 STOPPED 并清心跳),`stop_grace_period: 25s`。所有 DB 调用走 `to_thread`,不阻塞事件循环。 2. ✅ **`dispatcher._ws()`** —— 落 `pms_qmt_order` 出口表置 QUEUED 即返回,签名与发送由 ws 进程做。参数不合协议(非整百、限价缺失、代码非点式)在本地就拦下。撤单把父指令名下所有在途子单标 `cancel_state=REQUESTED`。 3. ✅ **上行消费与 seq 水位** —— `pms_qmt_inbox` 落库(seq 主键 + `trade_no` 唯一索引 = 协议 §6.1/§5.5 的双层去重),`pms_ws_state` 存连续水位,按 20 条 / 2 秒发 `ack_seq`。**落库失败绝不 ack**——落不了库就主动断线,让对端从 `last_seq+1` 重发,用协议自带的补发机制而不是自攒重试队列。 -4. 🔜 **`recon` 改为吃 `trade` 消息入账** —— 从 `pms_qmt_inbox` 里 `processed=0` 的 trade 行消费;FIFO 贪心认领退化为只处理外部/人工成交;`downstream_repo.FILLED_STATUSES` 降为旁路校验,不再作为成交判据。 -5. 🔜 **手续费口径** —— `fee` 不摊进持仓成本(对方明确逐笔费用可能有误差),只记现金流出,日终用资金快照反推校准。避免污染安全垫。 -6. 🔶 **`CANCELLED` / `EXPIRED` 自记** —— 通道层已做:`pms_qmt_order.cancel_state` 记本地是否发过撤单,与对端回的 `status` 不一致时告警。账本侧的口径随第 4 条一起落。 +4. ✅ **`recon` 改为吃 `trade` 消息入账** —— `ledger_service.consume_ws_trades()` 从 `pms_qmt_inbox` 消费 `processed=0` 的 trade 行,挂在 `replay_fills` 里(成交回放调度位从 5 分钟改为 **1 分钟**——ws 的成交是推过来的,没必要再等)。**trade 自带 `instruction_id`,是精确认领而不是 FIFO 猜**:影子期只能拿下游成交往在途指令上「凑」,同一只票挂着两条在途指令时凑错了根本看不出来,批次类型一错、卖出核销次序(T0→ADD→DCA→FILL→BASE)跟着错。ws 通了这层不确定性直接消失。`trading_order` 那条路退化为只兜外部/人工成交。 +5. ✅ **手续费口径** —— 新增 `pms_cash_flow`(第 14 张表)。逐笔 `fee` 记一条 FEE 流水(负数=流出),**绝不进 action、不摊成本**;日终 `calibrate_fees()` 用资金快照反推真实费用写一条 CALIBRATE 平掉估算误差,只调现金账不回溯改成本。影子期没有资金快照,那就只登记「待校准」不硬凑——费用既然进不了成本,晚校准没有风险。 +6. ✅ **`CANCELLED` / `EXPIRED` 自记** —— 通道层:`pms_qmt_order.cancel_state` 记本地是否发过撤单,与对端回的 `status` 不一致时告警;终态不可被后到的消息覆盖。账本侧:未成交部分由 `qty − exec_qty` 自然释放,窗口收口照旧。 部署前另有一份检查清单在协议 §10.1.1(白名单、公钥交换、Redis 持久化、密钥只走 `.env`)。 diff --git a/app/core/recon.py b/app/core/recon.py index 78ec169..ed410be 100644 --- a/app/core/recon.py +++ b/app/core/recon.py @@ -164,6 +164,72 @@ def map_fills_to_book(fills: list, open_instructions: list, known_order_ids=()) "skipped": len(fills or []) - len(fresh)} +def parent_instruction_id(child_id: str) -> str: + """子单 id → 父指令 id。子单形如 `{父指令}_D03` (executor._child_id)。""" + s = str(child_id or "") + tail = s.rsplit("_D", 1) + return tail[0] if len(tail) == 2 and tail[1].isdigit() else s + + +def map_trades_to_book(trades: list, parent_actions=None) -> dict: + """ws 逐笔成交 (协议 §5.5) → 账本动作 + 费用流水。 + + **与 map_fills_to_book 的根本差别: trade 自带 instruction_id, 是精确认领而不是 FIFO 猜。** + 影子期只能拿下游 trading_order 的成交往在途指令上"凑" —— 同一只票同时挂着两条在途指令 + 时凑错了根本看不出来, 批次类型一错, 卖出核销次序 (T0→ADD→DCA→FILL→BASE) 跟着错。 + ws 通了以后这层不确定性直接消失: 对端把 instruction_id 原样带回来了。 + + parent_actions: {父指令id: action} —— 用来定买入的批次类型 (OPEN→BASE / FILL→FILL / …)。 + 查不到就并入 BASE 并告警, 与外部成交同一口径。 + + **费用单独出一条流水, 不进 action** —— 协议 §5.5: fee 不摊进持仓成本, 只记现金流出。 + 对端明说逐笔费用是按费率估的, 估算误差一旦进了摊薄成本就会污染安全垫。 + + 返回 {"actions", "fees", "alerts", "seqs"} —— seqs 是本批消费掉的 inbox 序号。 + """ + parent_actions = parent_actions or {} + actions, fees, alerts, seqs = [], [], [], [] + for t in trades or []: + pl, seq = t.get("payload") or {}, t.get("seq") + seqs.append(seq) + code, side = pl.get("ts_code"), _side(pl.get("side")) + qty, price = int(pl.get("qty") or 0), float(pl.get("price") or 0) + iid = pl.get("instruction_id") or "" + parent = parent_instruction_id(iid) + if qty <= 0 or price <= 0 or not code: + alerts.append({"level": "WARN", "code": "BAD_TRADE", "ts_code": code, + "message": f"成交字段不完整, 已跳过: seq={seq} {pl}"[:300]}) + continue + + # amount 自洽性 (§5.5: 差异 > 0.01 元告警) —— 对不上不拦, 但要留痕 + amt = pl.get("amount") + if amt is not None and abs(float(amt) - price * qty) > 0.01: + alerts.append({"level": "WARN", "code": "AMOUNT_MISMATCH", "ts_code": code, + "message": f"成交金额不自洽: price×qty={price * qty:.2f} " + f"≠ amount={float(amt):.2f} (trade_no={pl.get('trade_no')})"}) + + action = parent_actions.get(parent) + if side == "buy" and not action: + alerts.append({"level": "WARN", "code": ALERT_EXTERNAL, "ts_code": code, + "message": f"成交带的指令 {iid} 在账本里找不到, 按外部成交并入 BASE " + f"(trade_no={pl.get('trade_no')})"}) + actions.append({"kind": "BUY" if side == "buy" else "SELL", "ts_code": code, + "qty": qty, "price": price, + "lot_type": (ACTION_TO_LOT.get(str(action or "").upper(), "BASE") + if side == "buy" else None), + "instruction_id": parent if action else None, + "order_id": pl.get("trade_no"), "alerts": []}) + + fee = pl.get("fee") + if fee is not None and float(fee) != 0: + fees.append({"ts_code": code, "amount": -abs(float(fee)), + "estimated": bool(pl.get("fee_estimated")), + "trade_no": pl.get("trade_no"), "instruction_id": parent, + "note": f"{side} {qty}股 @{price} 的手续费" + + ("(对端标注为估算, 待日终校准)" if pl.get("fee_estimated") else "")}) + return {"actions": actions, "fees": fees, "alerts": alerts, "seqs": seqs} + + def apply_sell_to_lots(lots: list, qty: int) -> dict: """卖出核销预演: 按 T0→ADD(新→旧)→DCA→FILL→BASE 分配。 账本批次不足时按可核销量分配并告警 (以下游为准, 差额留给对账修正)。""" diff --git a/app/repo/pms_repo.py b/app/repo/pms_repo.py index e0da5d4..0002ff4 100644 --- a/app/repo/pms_repo.py +++ b/app/repo/pms_repo.py @@ -467,6 +467,60 @@ def insert_ledger(*, ts_code, action, arbiter, verdict, price_at, hard_numbers=N "rsn": (reason or "")[:500], "ref": ref_id}) +# ================================================================ pms_cash_flow +def insert_cash_flow(*, ymd, kind, amount, ts_code=None, estimated=0, trade_no=None, + instruction_id=None, note=None) -> int: + """记一笔现金流水 (费用不入成本, 只入现金账 —— 协议 §5.5)。 + + (kind, trade_no) 上有唯一键做兜底去重。注意**不要拿返回值判断"是不是新插的"**: + SQLAlchemy 的 MySQL 方言默认开 CLIENT_FOUND_ROWS, 重复插入照样回 1 (同 inbox_put 的坑)。 + 调用方的去重靠 inbox 的 processed 标记, 这里的唯一键只是最后一道保险。 + """ + return execute( + "INSERT INTO pms_cash_flow (ymd, kind, ts_code, amount, estimated, trade_no, " + "instruction_id, note, created_at) VALUES (:y, :k, :code, :amt, :est, :tn, :iid, " + ":note, :ts) ON DUPLICATE KEY UPDATE id = id", + {"y": int(ymd), "k": kind, "code": ts_code, "amt": float(amount), + "est": 1 if estimated else 0, "tn": trade_no, "iid": instruction_id, + "note": (str(note)[:300] if note else None), "ts": _NOW()}) + + +def sum_cash_flow(ymd, kind=None) -> float: + sql = "SELECT COALESCE(SUM(amount), 0) AS s FROM pms_cash_flow WHERE ymd = :y" + p = {"y": int(ymd)} + if kind: + sql += " AND kind = :k" + p["k"] = kind + return float((fetch_one(sql, p) or {}).get("s") or 0) + + +def list_cash_flow(*, ymd=None, kind=None, limit: int = 200) -> list: + where, p = [], {"n": int(limit)} + if ymd: + where.append("ymd = :y") + p["y"] = int(ymd) + if kind: + where.append("kind = :k") + p["k"] = kind + sql = "SELECT * FROM pms_cash_flow" + if where: + sql += " WHERE " + " AND ".join(where) + return fetch_all(sql + " ORDER BY id DESC LIMIT :n", p) + + +def rule_rejected_today(since) -> set: + """今天已被**规则闸**拒过的 (代码, 动作) —— 自主扫描据此当日不再重复评估。 + + 评审账本是给事后判分用的 (「拒了的后来涨了多少」), 同一件事一天记一条足矣。而扫描每分钟 + 一跳, 不去重的话一条持续不通过的候选 —— 比如已超总仓上限时的补仓 —— 一天能写进 240 行 + 一模一样的记录, 把真正有信息量的行淹掉。判分锚被噪声埋了就不再是锚。 + """ + rows = fetch_all("SELECT DISTINCT ts_code, action FROM pms_action_ledger " + "WHERE verdict = 'REJECT' AND arbiter = 'rule' AND decided_at >= :d", + {"d": since}) + return {(r["ts_code"], r["action"]) for r in rows} + + def list_ledger(*, ts_code=None, limit: int = 200) -> list: sql = "SELECT * FROM pms_action_ledger" p = {"n": int(limit)} diff --git a/app/scheduler.py b/app/scheduler.py index a6e98c0..8f562f0 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -8,7 +8,7 @@ | 命令轮询 | 每 1 分钟 (全天) | 新命令解析 → 方案生成 → 任务状态机推进 | | 盘中执行 | 交易时段每 1 分钟 | 择时出手 + 自主提议扫描 (执行器下一批交付) | | 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 (下一批交付) | -| 成交回放 | 交易时段每 5 分钟 | trading_order 增量回放 + 轻对账 | +| 成交回放 | 交易时段每 1 分钟 | ws 逐笔入账 + trading_order 增量 + 轻对账 | | T 仓平回 | 14:50 (二期) | 做T强制平回 | | 日终结算 | 15:10 | 全量对账 / 除权 / 安全垫 / 命令进度日结 | | 运营日报 | 15:30 | 关注区 + 全量统计, 页面可查 | @@ -192,7 +192,9 @@ def daily_report(): celery_app.conf.beat_schedule = { "premarket": {"task": "pms.premarket", "schedule": crontab(hour=8, minute=50)}, "command_poll": {"task": "pms.command_poll", "schedule": crontab(minute="*")}, - "replay_fills": {"task": "pms.replay_fills", "schedule": crontab(minute="*/5")}, + # ws 通道的成交是推过来的, 落 inbox 后没必要再等 5 分钟才入账 —— 改成每分钟。 + # 影子期这一跳基本是空转 (inbox 为空 + 下游无新成交), 成本可以忽略。 + "replay_fills": {"task": "pms.replay_fills", "schedule": crontab(minute="*")}, "intraday_exec": {"task": "pms.intraday_exec", "schedule": crontab(minute="*")}, "signal_digest": {"task": "pms.signal_digest", "schedule": crontab(minute="*")}, "t0_close": {"task": "pms.t0_close", "schedule": crontab(hour=14, minute=50)}, diff --git a/app/services/ledger_service.py b/app/services/ledger_service.py index 6b9e535..4c19f2a 100644 --- a/app/services/ledger_service.py +++ b/app/services/ledger_service.py @@ -21,11 +21,12 @@ from datetime import datetime, timedelta from app.core import cushion as cu from app.core import recon as rc from app.core import tradedays as td -from app.repo import downstream_repo, pms_repo +from app.repo import downstream_repo, pms_repo, qmt_repo from app.services import market, param_store, portfolio logger = logging.getLogger("pms.ledger") +CF_FEE, CF_CALIBRATE = "FEE", "CALIBRATE" # pms_cash_flow.kind CURSOR_KEY = "PMS_REPLAY_CURSOR" CURSOR_ALL = "ALL" # 游标设成这个值 = 显式要求从头全量回放 (见 _seed_cursor) STREAK_KEY = "PMS_RECON_STREAK" @@ -62,9 +63,131 @@ def _seed_cursor(out: dict) -> dict: return out +def consume_ws_trades(*, limit: int = 500) -> dict: + """消费 ws 通道落在 pms_qmt_inbox 的逐笔成交 → 批次入账 + 费用流水 (协议 §5.5)。 + + 与下面的 trading_order 回放是**两条互补的路**, 都挂在 replay_fills 里: + 本函数 ws 通道推来的成交 —— 自带 instruction_id, **精确认领** + replay_fills trading_order 增量 —— 只剩外部/人工成交, FIFO 认领退化为兜底 + + 幂等靠 inbox 的 processed 标记 (0 待入账 → 1 已入账), 而不是靠返回行数 —— 那个在 + CLIENT_FOUND_ROWS 下不可信, 见 qmt_repo.inbox_put 的注释。 + """ + out = {"ok": True, "trades": 0, "actions": 0, "fees": 0, "alerts": [], "errors": []} + try: + rows = qmt_repo.inbox_pending(limit=limit) + except Exception as e: + # ws 三表没建 (影子期正常) 或库不可用 —— 不该让整个回放任务失败 + out["skipped"] = f"inbox 不可读: {type(e).__name__}" + return out + trades = [r for r in rows if r.get("msg_type") == "trade"] + if not trades: + return out + out["trades"] = len(trades) + + # 取每条成交所属父指令的 action, 用来定买入的批次类型 + parents = {rc.parent_instruction_id((r.get("payload") or {}).get("instruction_id") or "") + for r in trades} + parent_actions = {} + for pid in parents: + if not pid: + continue + try: + ins = pms_repo.get_instruction(pid) + if ins: + parent_actions[pid] = ins.get("action") + except Exception as e: + out["errors"].append(f"读指令 {pid} 失败: {type(e).__name__}: {e}") + + mapped = rc.map_trades_to_book(trades, parent_actions=parent_actions) + done = [] + for act in mapped["actions"]: + try: + _apply_action(act) + out["actions"] += 1 + except Exception as e: + logger.exception("ws 成交入账失败 %s", act) + out["errors"].append(f"{act.get('ts_code')} 入账失败: {type(e).__name__}: {e}") + ymd = td.ymd() + for f in mapped["fees"]: + try: + pms_repo.insert_cash_flow(ymd=ymd, kind=CF_FEE, amount=f["amount"], + ts_code=f["ts_code"], estimated=f["estimated"], + trade_no=f["trade_no"], + instruction_id=f["instruction_id"], note=f["note"]) + out["fees"] += 1 + except Exception as e: + out["errors"].append(f"费用入账失败 {f.get('trade_no')}: {type(e).__name__}: {e}") + + for code in {a["ts_code"] for a in mapped["actions"]}: + try: + recompute_position(code) + except Exception as e: + out["errors"].append(f"{code} 成本重算失败: {e}") + + # 只有全程无错才标已入账 —— 有错就留在 processed=0, 下一跳重试。 + # 重试是安全的: _apply_action 幂等由 trade_no 兜着 (inbox 那层已按 trade_no 去过重)。 + if not out["errors"]: + done = [s for s in mapped["seqs"] if s is not None] + if done: + try: + qmt_repo.inbox_mark(done, processed=1, note="ws 成交已入账") + except Exception as e: + out["errors"].append(f"inbox 标记失败: {type(e).__name__}: {e}") + out["alerts"] = mapped["alerts"] + for a in out["alerts"]: + logger.warning("[ws 入账告警] %s", a.get("message")) + out["ok"] = not out["errors"] + return out + + +def calibrate_fees(*, ymd=None, actual_fee=None) -> dict: + """日终用资金快照反推当日真实费用, 写一条 CALIBRATE 平掉估算误差 (协议 §5.5)。 + + 对端明说逐笔 fee 是按费率估的。协议给的校准式子是 + 当日实际费用 = 总资产变动 − 成交净额 + 资金快照来自 ws 的 funds_update / snapshot(kind=funds), 落在 pms_qmt_inbox。**影子期 + 没有这条数据**, 那就只登记一句"待校准", 不硬凑 —— 估算值本来就只影响现金账, 不影响 + 成本与安全垫, 晚校准几天没有任何风险。这也是当初把费用挡在成本之外的意义。 + + actual_fee 可显式传入 (人工按对账单校准时用), 传了就不去读快照。 + """ + ymd = int(ymd or td.ymd()) + out = {"ymd": ymd, "estimated": 0.0, "actual": None, "adjusted": 0.0, "note": ""} + try: + out["estimated"] = round(pms_repo.sum_cash_flow(ymd, CF_FEE), 2) + except Exception as e: + out["note"] = f"读当日费用流水失败: {type(e).__name__}: {e}" + return out + if actual_fee is None: + out["note"] = ("当日估算费用已入现金账; 资金快照未接通 (ws 通道未启用), " + "暂不校准 —— 费用不进成本, 晚校准无风险") + return out + actual = -abs(float(actual_fee)) + diff = round(actual - out["estimated"], 2) + out["actual"] = actual + if abs(diff) < 0.01: + out["note"] = "估算与实际一致, 无需校准" + return out + pms_repo.insert_cash_flow( + ymd=ymd, kind=CF_CALIBRATE, amount=diff, estimated=0, trade_no=None, + note=f"日终校准: 估算 {out['estimated']:.2f} → 实际 {actual:.2f}, 差额 {diff:+.2f}") + out["adjusted"] = diff + out["note"] = f"已按资金快照校准, 差额 {diff:+.2f} 元 (只调现金账, 不回溯改成本)" + logger.info("[费用校准] %s", out["note"]) + return out + + def replay_fills(*, limit: int = 500) -> dict: - """增量回放 trading_order 已成交单 → 批次入账 (每 5 分钟一跳, 幂等)。""" + """成交回放一跳 = ws 逐笔入账 + trading_order 增量回放 (幂等)。 + + 两条路互补: ws 通道的成交自带 instruction_id 精确入账; trading_order 这条在 ws 接管后 + 退化为**只兜外部/人工成交** (你在 QMT 手工下的单、别的系统下的单)。影子期只有后者。 + """ out = {"ok": True, "fills": 0, "actions": 0, "alerts": [], "errors": [], "cursor": None} + out["ws"] = consume_ws_trades(limit=limit) + if out["ws"].get("errors"): + out["errors"].extend(out["ws"]["errors"]) cursor = pms_repo.get_param(CURSOR_KEY) if not cursor: # 从未设过 (或被清空) —— 冷启动, 只对齐游标不入账 return _seed_cursor(out) @@ -94,7 +217,7 @@ def replay_fills(*, limit: int = 500) -> dict: except Exception as e: logger.exception("入账失败 %s", act) out["errors"].append(f"{act.get('ts_code')} 入账失败: {type(e).__name__}: {e}") - out["alerts"] = mapped["alerts"] + out["alerts"].extend(mapped["alerts"]) # 与 ws 那批告警合并, 日报一处看全 for code in {a["ts_code"] for a in mapped["actions"]}: try: @@ -413,6 +536,10 @@ def daily_settle() -> dict: out["steps"]["cushion"] = _settle_cushion() except Exception as e: out["errors"].append(f"安全垫结算失败: {e}") + try: + out["steps"]["fee_calibrate"] = calibrate_fees() + except Exception as e: + out["errors"].append(f"费用校准失败: {e}") try: out["steps"]["commands"] = command_service.refresh_progress() except Exception as e: diff --git a/app/services/proposal_service.py b/app/services/proposal_service.py index 62fd10e..4780331 100644 --- a/app/services/proposal_service.py +++ b/app/services/proposal_service.py @@ -51,7 +51,10 @@ def scan_and_route(*, now=None, dry_run: bool = False) -> dict: params = _scan_params(view) mkt = _market_ctx(view["held"], now) params["_mkt"] = mkt # 规则闸要用同一份 MA5, 不再重取 - skip = _inflight_keys() + # 跳过两类: ①已有在途提议或指令的 ②今天已被规则闸拒过的。 + # 后者是 2026-07-29 的教训 —— 组合已超总仓上限时, 16 只深亏票的补仓候选每分钟被拒 + # 一次, 一天往评审账本灌几千行一模一样的记录。闸门结论当天基本不会变, 记一次就够。 + skip = _inflight_keys() | _rejected_today_keys() scanned = ae.scan(positions=view["held"], params=params, market=mkt, skip=skip) except Exception as e: logger.exception("提议扫描失败") @@ -160,9 +163,41 @@ def _make_instruction(c, price, now) -> str: progress={"deadline": str(td.window_deadline(now.date(), window)), "is_command": False, "children": [], "auto": True, "reason": c["reason"]}) + bump_once_guards(c["ts_code"], c["action"], c.get("hard_numbers"), now) return iid +def bump_once_guards(ts_code: str, action: str, hard_numbers=None, now=None): + """把「只做一次」的三个计数器写上。 + + 设计 §6 写了三条一次性约束, 但它们的计数器此前**只被读、从没被写过** —— 也就是说这三条 + 纪律一直是失效的, 只是被「同一票同一动作有在途提议就不重复提」这条兜底遮住了, 而那条兜底 + 恰好在规则闸拒绝时失灵 (没生成提议 → 没东西可去重), 于是同一个候选每分钟重来一次: + + fill_count 回踩补足「每票 1 次」 + last_add_date 盈利加仓「距上次 ≥2 交易日」—— 恒 None 时永远算作"很久没加过" + dca_count 补仓「各档评估一次」—— 恒 0 时每轮都当第一次评估 + + 写入时机取「动作真的要落地」这一刻 (落指令), 而不是产出候选那一刻: 候选被闸门拦下不算 + 做过, 用掉一次名额不合理。dca_count 记的是**档位**而不是次数, 所以取本次触及的档序。 + """ + now = now or datetime.now() + fields = {} + if action == "FILL": + fields["fill_count"] = 1 + elif action == "ADD": + fields["last_add_date"] = now.date() + elif action == "DCA": + fields["dca_count"] = int((hard_numbers or {}).get("stage") or 1) + if not fields: + return + try: + pms_repo.update_position(ts_code, **fields) + except Exception as e: + # 写不上只是让这条纪律退回原样 (可能重复提), 不该把已落表的指令带崩 + logger.warning("一次性守卫计数器写入失败 %s %s: %s", ts_code, action, e) + + def _make_proposal(c, price, verdict) -> str: ttl = param_store.get_int("PMS_PROPOSAL_TTL_HOURS", 24) pid = f"PRP_{td.ymd()}_{c['ts_code'].replace('.', '')}_{c['action']}" @@ -218,6 +253,17 @@ def _tdays_between(start, today): return None +def _rejected_today_keys() -> set: + """今天已被规则闸拒过的 (代码, 动作)。读不到就返回空集 —— 去重是降噪, 不是纪律, + 读失败时宁可多记几行日志, 也不能因此漏扫一个本该评估的动作。""" + try: + return pms_repo.rule_rejected_today(datetime.now().replace( + hour=0, minute=0, second=0, microsecond=0)) + except Exception as e: + logger.warning("读当日规则闸拒绝记录失败 (按未拒过继续扫描): %s", e) + return set() + + def _inflight_keys() -> set: """已有在途提议或在途指令的 (代码, 动作) —— 同一件事不重复提。""" keys = set() diff --git a/app/web/main.py b/app/web/main.py index 9663359..82bb461 100644 --- a/app/web/main.py +++ b/app/web/main.py @@ -258,6 +258,10 @@ def api_decide(proposal_id: str, payload: dict = Body(default={})): limit_price=hn.get("price"), window_tdays=param_store.get_int("PMS_EXEC_WINDOW_TDAYS", 3), status="PROPOSED", progress={"from_proposal": proposal_id}) + # 与自主执行同一口径: 采纳即算「做过一次」, 计数器要跟着走 + # (否则页面采纳的那条动作绕开了 §6 的一次性约束) + from app.services import proposal_service + proposal_service.bump_once_guards(p["ts_code"], p["action"], hn) return {"ok": True, "decision": decision, "instruction_id": instruction_id} return ok(_decide) diff --git a/ddl_pms_v1.sql b/ddl_pms_v1.sql index ba78aa0..1ba00b3 100644 --- a/ddl_pms_v1.sql +++ b/ddl_pms_v1.sql @@ -280,3 +280,25 @@ CREATE TABLE IF NOT EXISTS pms_ws_state ( INSERT INTO pms_ws_state (id, last_seq, acked_seq, server_seq, conn_state, updated_at) VALUES (1, 0, 0, 0, 'INIT', NOW()) ON DUPLICATE KEY UPDATE id = id; + +-- 14. 现金流水 (2026-07-29 追加; 协议 QMT_WS_PROTOCOL.md §5.5 手续费口径) +-- 为什么需要它: 协议规定手续费**不摊进持仓成本**, 只记现金流出 —— 因为对端明说逐笔费用 +-- 是按费率估算的、与实际扣费有出入。估算误差一旦进了摊薄成本, 就会污染安全垫, 而安全垫是 +-- 补仓/加仓/保垫减仓共同的判断依据。所以费用必须有个**成本之外**的去处, 就是这张表。 +-- 日终再用资金快照反推当日真实费用 (总资产变动 − 成交净额), 写一条 CALIBRATE 平掉估算 +-- 误差 —— 只调现金账, 不回溯改成本。 +CREATE TABLE IF NOT EXISTS pms_cash_flow ( + id BIGINT PRIMARY KEY AUTO_INCREMENT, + ymd INT NOT NULL COMMENT 'YYYYMMDD', + kind VARCHAR(16) NOT NULL + COMMENT 'FEE(逐笔手续费, 估算) / CALIBRATE(日终资金快照反推的差额) / MANUAL(人工调整)', + ts_code VARCHAR(16) NULL COMMENT '校准行为 NULL (全组合)', + amount DECIMAL(14,2) NOT NULL COMMENT '现金变动, **流出为负**', + estimated TINYINT NOT NULL DEFAULT 0 COMMENT '1 = 对端标注 fee_estimated, 待日终校准', + trade_no VARCHAR(80) NULL COMMENT '逐笔成交号, 去重用; 校准行为 NULL', + instruction_id VARCHAR(64) NULL, + note VARCHAR(300) NULL, + created_at DATETIME NOT NULL, + UNIQUE KEY uk_kind_trade (kind, trade_no), + KEY idx_ymd_kind (ymd, kind) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='现金流水 (费用不入成本, 只入现金账)'; diff --git a/scripts/run_tests.py b/scripts/run_tests.py index 835173f..5e44b25 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -10,7 +10,7 @@ test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 (19 例) test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 (11 例) test_batch5_units.py 决策系统信号流解析与消化口径 (8 例) - test_batch6_units.py ws 通道: 协议测试向量/签名/公钥形态/水位/DDL 体检 (51 例) + test_batch6_units.py ws 通道: 测试向量/签名/公钥/水位/DDL/逐笔入账 (58 例) test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (32 例) 任一子集失败即整体失败 (退出码 1)。 """ diff --git a/scripts/test_batch6_units.py b/scripts/test_batch6_units.py index a6b8dc1..b257092 100644 --- a/scripts/test_batch6_units.py +++ b/scripts/test_batch6_units.py @@ -15,6 +15,7 @@ F. 双层去重键与成交金额自洽性。 G. 单表访问守卫: upsert (ON DUPLICATE KEY UPDATE) 不得被误判成多表。 H. dispatcher 纯助手: valid_until 归一、子单→父指令反推。 + I. ws 逐笔成交 → 账本动作: 精确认领、批次类型、费用不进成本、金额自洽性。 """ import base64 import os @@ -366,6 +367,69 @@ def run(): assert wsc.T_ORDER_UPDATE not in wsc.NEEDS_LEDGER, \ "order_update 只作状态跟踪 —— 拿它入账会和逐笔 trade 重复计数" + print("\n[I] ws 逐笔成交 → 账本动作 (清单 4/5)") + from app.core import recon as rc + + def _tr(seq, tn, iid, side="buy", qty=600, px=10.0, fee=3.21, est=True, amt=None): + return {"seq": seq, "msg_type": "trade", "payload": { + "instruction_id": iid, "trade_no": tn, "ts_code": "600000.SH", "side": side, + "qty": qty, "price": px, "amount": (px * qty if amt is None else amt), + "fee": fee, "fee_estimated": est}} + + @case("子单 id → 父指令") + def _(): + eq(rc.parent_instruction_id("INS_20260729_600000SH_OPEN_01_D07"), + "INS_20260729_600000SH_OPEN_01") + eq(rc.parent_instruction_id("INS_NO_CHILD"), "INS_NO_CHILD") + + @case("精确认领: 批次类型跟着父指令的 action 走") + def _(): + m = rc.map_trades_to_book([_tr(1, "T#1", "P1_D01"), _tr(2, "T#2", "P2_D01")], + parent_actions={"P1": "OPEN", "P2": "ADD"}) + eq([a["lot_type"] for a in m["actions"]], ["BASE", "ADD"]) + eq([a["instruction_id"] for a in m["actions"]], ["P1", "P2"]) + + @case("认领不到的指令 → 并入 BASE 且告警 (同外部成交口径)") + def _(): + m = rc.map_trades_to_book([_tr(1, "T#1", "GHOST_D01")], parent_actions={}) + eq(m["actions"][0]["lot_type"], "BASE") + eq(m["actions"][0]["instruction_id"], None) + assert any(a["code"] == "EXTERNAL_FILL" for a in m["alerts"]), m["alerts"] + + @case("手续费单独出流水, **不进 action** (协议 §5.5)") + def _(): + m = rc.map_trades_to_book([_tr(1, "T#1", "P1_D01", fee=3.21)], + parent_actions={"P1": "OPEN"}) + eq(len(m["fees"]), 1) + eq(m["fees"][0]["amount"], -3.21, "费用记为现金流出, 取负") + assert m["fees"][0]["estimated"] is True + assert "fee" not in m["actions"][0], "费用绝不能混进入账动作 —— 会污染摊薄成本" + # 入账价就是成交价本身, 不含费 + eq(m["actions"][0]["price"], 10.0) + + @case("amount 不自洽要告警但不拦 (§5.5 差异 > 0.01 元)") + def _(): + m = rc.map_trades_to_book([_tr(1, "T#1", "P1_D01", amt=9999.0)], + parent_actions={"P1": "OPEN"}) + assert any(a["code"] == "AMOUNT_MISMATCH" for a in m["alerts"]), m["alerts"] + eq(len(m["actions"]), 1, "对不上也照常入账 —— 成交是既成事实") + + @case("字段残缺的成交跳过并告警, 不带崩整批") + def _(): + m = rc.map_trades_to_book( + [_tr(1, "T#1", "P1_D01", qty=0), _tr(2, "T#2", "P1_D01")], + parent_actions={"P1": "OPEN"}) + eq(len(m["actions"]), 1) + assert any(a["code"] == "BAD_TRADE" for a in m["alerts"]) + eq(m["seqs"], [1, 2], "跳过的那条也要占住 seq, 否则 inbox 永远标不完") + + @case("卖出不带批次类型 (核销次序另算)") + def _(): + m = rc.map_trades_to_book([_tr(1, "T#1", "P1_D01", side="sell")], + parent_actions={"P1": "EXIT"}) + eq(m["actions"][0]["kind"], "SELL") + eq(m["actions"][0]["lot_type"], None) + print("\n[G] 严格单表访问守卫 (upsert 不得被误判成多表)") @case("ON DUPLICATE KEY UPDATE 的四条现存 upsert 全部放行") @@ -414,7 +478,7 @@ def run(): "DEFAULT 'NONE' COMMENT 'NONE/REQUESTED/SENT'"): eq(find_adjacent_literals(good), [], f"误报: {good[:40]}") - @case("DDL 文件本身体检通过 (13 张表 + 1 条初始行)") + @case("DDL 文件本身体检通过 (14 张表 + 1 条初始行)") def _(): sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from init_db import DDL_FILE, find_adjacent_literals, parse_statements @@ -424,7 +488,7 @@ def run(): eq(broken, [], "有认不出的残句, 说明 DDL 切分被破坏") for _, tbl, s in stmts: eq(find_adjacent_literals(s), [], f"{tbl} 有相邻字面量") - eq(len([1 for k, _, _ in stmts if k == "table"]), 13) + eq(len([1 for k, _, _ in stmts if k == "table"]), 14) eq(len([1 for k, _, _ in stmts if k == "seed"]), 1) @case("通道三表的 SQL 全部单表合规") diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py index 8a91b11..f0145bd 100644 --- a/scripts/test_wiring.py +++ b/scripts/test_wiring.py @@ -272,6 +272,11 @@ class FakeRepo: def list_ledger(self, *, ts_code=None, limit=200): return self.ledger[-limit:] + def rule_rejected_today(self, since): + # 内存桩里所有留痕都算"今天"; 只认规则闸的 REJECT (研判结论会变, 不参与当日去重) + return {(r["ts_code"], r["action"]) for r in self.ledger + if r.get("verdict") == "REJECT" and r.get("arbiter") == "rule"} + def upsert_report(self, ymd, report): self.reports[int(ymd)] = report return 1 @@ -742,6 +747,34 @@ def _(): downstream_repo.fetch_filled_orders = orig +@case("自主提议·一次性守卫真的写得进去, 且当日被拒的不再每分钟重评") +def _(): + from datetime import date + from app.repo import pms_repo + from app.services import proposal_service as ps + fake = install_fakes(positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "base_qty": 6000, "avg_cost": 10.0}]) + # 三个计数器此前只被读、从没被写过 —— 设计 §6 的三条一次性约束等于一直没生效 + ps.bump_once_guards("600000.SH", "FILL", {}) + assert fake.positions["600000.SH"]["fill_count"] == 1 + ps.bump_once_guards("600000.SH", "ADD", {}) + assert fake.positions["600000.SH"]["last_add_date"] == date.today() + ps.bump_once_guards("600000.SH", "DCA", {"stage": 2}) + assert fake.positions["600000.SH"]["dca_count"] == 2, "dca_count 记的是档序不是次数" + ps.bump_once_guards("600000.SH", "TRIM", {}) # 减持无一次性约束, 不该动任何列 + assert fake.positions["600000.SH"]["fill_count"] == 1 + + # 当日已被规则闸拒过的 (代码, 动作) 要进 skip —— 否则超上限时每分钟重评一次、 + # 每分钟往评审账本灌一行, 判分锚被自己的噪声埋掉 + pms_repo.insert_ledger(ts_code="600000.SH", action="DCA", arbiter="rule", + verdict="REJECT", price_at=10.0, reason="超总仓上限") + pms_repo.insert_ledger(ts_code="000001.SZ", action="ADD", arbiter="judge", + verdict="REJECT", price_at=10.0, reason="研判驳回") + keys = ps._rejected_today_keys() + assert ("600000.SH", "DCA") in keys, keys + assert ("000001.SZ", "ADD") not in keys, "研判驳回不该进当日去重 —— 研判结论会变" + + @case("账本服务·对账补仓位取下游成本价, 不拿现价充数") def _(): from app.core import recon as rc