tradingSystem/app/services/ledger_service.py

1120 lines
61 KiB
Python
Raw Normal View History

# -*- coding: utf-8 -*-
"""
账本服务: 成交回放 · 对账 · 除权 · 日终结算 (设计 §4生命线)
================================================================
纯逻辑在 app/core/recon.py, 本模块只负责取数落库与状态推进
摊薄成本口径 ( core/cushion.PositionCost 等价, 但数据来自批次表):
cum_buy = Σ (剩余数量 + 已核销数量) × 开仓价
cum_sell = Σ 已核销数量 × 核销均价
avg_cost = max(0, (cum_buy cum_sell) / 当前持股数)
做T利润通过 T0 批次的买卖流水自然摊入上式, **不再另行扣减** realized_t_profit
(该列仅作展示统计, 重复扣减会把成本做低两次)
铁律: 对账以下游为准; 修正一律走 RECON 批次留痕; 连续 N 日不一致升级 ERROR
"""
from __future__ import annotations
import logging
from datetime import datetime, timedelta
2026-07-30 08:42:06 +08:00
from app.core import command_spec as cs
from app.core import cushion as cu
2026-07-31 13:46:15 +08:00
from app.core import rebuild_check as rbc
from app.core import recon as rc
from app.core import tradedays as td
from app.core import ws_codec as wsc
2026-07-29 10:49:56 +08:00
from app.repo import downstream_repo, pms_repo, qmt_repo
from app.services import market, param_store, portfolio
logger = logging.getLogger("pms.ledger")
2026-07-29 10:49:56 +08:00
CF_FEE, CF_CALIBRATE = "FEE", "CALIBRATE" # pms_cash_flow.kind
CURSOR_KEY = "PMS_REPLAY_CURSOR"
2026-07-29 10:01:00 +08:00
CURSOR_ALL = "ALL" # 游标设成这个值 = 显式要求从头全量回放 (见 _seed_cursor)
STREAK_KEY = "PMS_RECON_STREAK"
# 该 streak 最后一次推进是哪个交易日 —— 同一天内反复对账不重复计数的判据
STREAK_YMD_KEY = "PMS_RECON_STREAK_YMD"
LIVE_INSTR = ("DISPATCHED", "JUDGE_PASSED", "RULE_PASSED")
# ================================================================ 回放
2026-07-29 10:01:00 +08:00
def _seed_cursor(out: dict) -> dict:
"""首次回放: 把游标对齐到当前最新成交, **不追认历史**。
`trading_order` 里躺着旧系统多年的成交记录PMS 刚上线时一条在途指令都没有, 若从头
回放, 每一条历史成交都会被判成外部成交并入 BASE 批次 账本上凭空长出一堆早就
清掉的持仓, 摊薄成本与安全垫全错, 日志还刷几百条 EXTERNAL_FILL 告警而安全垫是补仓
盈利加仓保垫减仓共同的判断依据, 它一错整条纪律链跟着错**账本必须从干净的起点开始**
要补历史有两条路 (页面参数设置 PMS_REPLAY_CURSOR):
* 设成某个 order_id 从它之后开始回放
* 设成 ALL 从头全量回放
"""
try:
anchor = downstream_repo.latest_filled_order_id()
except Exception as e:
out.update({"ok": False, "errors": [f"读下游最新成交失败: {type(e).__name__}: {e}"]})
return out
if not anchor:
out["note"] = "下游暂无已成交单, 游标待下次再对齐"
return out
pms_repo.set_param(CURSOR_KEY, anchor, updated_by="system")
out.update({"cursor": anchor, "seeded": True})
out["note"] = (f"首次回放: 游标已对齐到当前最新成交 {anchor}, **不追认历史** —— "
f"账本从现在起跟踪。需要补历史请在页面把 PMS_REPLAY_CURSOR 改成某个 "
f"order_id (从它之后开始) 或 {CURSOR_ALL} (从头全量)")
logger.warning("[replay_fills] %s", out["note"])
return out
2026-07-29 10:49:56 +08:00
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 的注释
每笔成交先按出口表反查分三路 (见下方注释): 真单入账(1) / 联调单跳过(2) /
**孤儿成交挂起(3, 不入账)**; 反查本身报错的一笔都不判, 留在 0 等下一跳
2026-07-29 10:49:56 +08:00
"""
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"]
2026-07-29 14:50:04 +08:00
if not trades:
return out
# ---- 三分流: 联调单 / 孤儿成交 / 真单。判据只查**我们自己的出口表**
# (pms_qmt_order.parent_instruction_id 以 SMOKE_ 开头), 不看对端字段 —— 对端只回子
# instruction_id, 分不出联调还是真单; 出口表是本端写的, 骗不了自己。
# 联调单标 processed=2 (已消化) 而不是 1 (已入账): 两者语义不同, 事后翻 inbox 一眼能看出
# 这笔是被有意跳过的, 而不是入账入丢了。见 qmt_repo.SMOKE_PREFIX 与 PROC_* 的注释。
#
# **孤儿成交 (出口表反查不到 instruction_id) 一律挂起, 绝不入账。**
# 这条是 2026-07-30 用真实事故换来的。ws 这条路上的 trade 必带 instruction_id (§5.5),
# 而那个 id 是 PMS 自己生成、自己写进出口表的, 所以查不到只有一种解释: **数据不一致** ——
# 上下游停机后 inbox 里留着的旧行、出口表被清过、或换过环境。它**不可能**是"有人在 QMT
# 手工下了单": 手工成交走 trading_order 那条路, 根本不进 inbox。
# 原先这里沿用 replay_fills 的口径「查不到就按真单走, 宁可多入账也不能漏账」, 结果 11 笔
# 停机遗留的联调成交被判成外部成交并入 BASE, 账本上凭空长出 600000.SH 1100 股 @9.273 ——
# 而它们对应的出口表行早就不在了, SMOKE_ 那道闸压根没机会生效。
# 两种错的代价根本不对称: 漏账看得见 (inbox 留着行、有 ERROR、通道状态报数字), 错账看不见
# (凭空长出来的持仓与真持仓同形, 摊薄成本和安全垫却已经全错, 而补仓/加仓/保垫减仓全挂在
# 安全垫上)。所以宁可挂起等人工, 不猜。
smoke, orphan, real = [], [], []
2026-07-29 14:50:04 +08:00
for r in trades:
iid = (r.get("payload") or {}).get("instruction_id") or ""
if not iid: # §5.5 要求 trade 必带 instruction_id, 没有就是坏数据
orphan.append(r)
continue
2026-07-29 14:50:04 +08:00
try:
od = qmt_repo.get_order(iid)
except Exception as e:
# 库抖一下**什么都不判**, 留在 processed=0 下一跳再来。原先这里把异常当"查不到"
# 并入真单分支 —— 一次连接超时就足以凭空造出一笔"外部成交"。
2026-07-29 14:50:04 +08:00
out["errors"].append(f"出口表反查 {iid} 失败: {type(e).__name__}: {e}")
continue
if od is None:
orphan.append(r)
elif qmt_repo.is_smoke(od.get("parent_instruction_id")):
smoke.append(r)
else:
real.append(r)
2026-07-29 14:50:04 +08:00
if smoke:
seqs = [r["seq"] for r in smoke]
qmt_repo.inbox_mark(seqs, processed=qmt_repo.PROC_DIGESTED,
note="ws 联调测试单成交, 不入账 (SMOKE_)")
2026-07-29 14:50:04 +08:00
out["smoke_skipped"] = len(smoke)
logger.warning("[联调] 跳过 %s 笔测试单成交, 不入账: seq=%s", len(smoke), seqs)
if orphan:
seqs = [r["seq"] for r in orphan]
codes = sorted({(r.get("payload") or {}).get("ts_code") or "?" for r in orphan})
qmt_repo.inbox_mark(seqs, processed=qmt_repo.PROC_ORPHAN,
note="出口表查不到对应委托, 挂起待人工确认, 未入账")
out["orphan_held"] = len(orphan)
msg = (f"{len(orphan)} 笔 ws 成交在出口表 pms_qmt_order 里查不到对应委托, 已挂起"
f"**未入账** (processed={qmt_repo.PROC_ORPHAN}): seq={seqs}, 标的={codes}"
f"ws 的 trade 必带 PMS 自己发出的 instruction_id, 查不到=数据不一致 "
f"(多为上下游停机后遗留的旧 inbox 行)。**确认这些成交真该入账**, 再把这些行的 "
f"processed 改回 0 让它们重新过一遍; 若是遗留垃圾, 保持挂起即可。")
out["alerts"].append({"type": "ORPHAN_WS_TRADE", "message": msg,
"seqs": seqs, "ts_codes": codes})
logger.error("[ws 入账] %s", msg)
2026-07-29 14:50:04 +08:00
trades = real
2026-07-29 10:49:56 +08:00
if not trades:
out["ok"] = not out["errors"]
2026-07-29 10:49:56 +08:00
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=qmt_repo.PROC_BOOKED,
note="ws 成交已入账")
2026-07-29 10:49:56 +08:00
except Exception as e:
out["errors"].append(f"inbox 标记失败: {type(e).__name__}: {e}")
# extend 而不是赋值 —— 上面孤儿成交的告警已经在 out["alerts"] 里了, 直接赋值会把它冲掉,
# 于是"挂起了一笔账"这件事只剩日志里一行, 页面和调度返回值里全看不见。
out["alerts"].extend(mapped["alerts"])
for a in mapped["alerts"]:
2026-07-29 10:49:56 +08:00
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:
2026-07-29 10:49:56 +08:00
"""成交回放一跳 = ws 逐笔入账 + trading_order 增量回放 (幂等)。
两条路互补: ws 通道的成交自带 instruction_id 精确入账; trading_order 这条在 ws 接管后
退化为**只兜外部/人工成交** (你在 QMT 手工下的单别的系统下的单)影子期只有后者
"""
out = {"ok": True, "fills": 0, "actions": 0, "alerts": [], "errors": [], "cursor": None}
2026-07-29 10:49:56 +08:00
out["ws"] = consume_ws_trades(limit=limit)
if out["ws"].get("errors"):
out["errors"].extend(out["ws"]["errors"])
2026-07-29 10:01:00 +08:00
cursor = pms_repo.get_param(CURSOR_KEY)
if not cursor: # 从未设过 (或被清空) —— 冷启动, 只对齐游标不入账
return _seed_cursor(out)
if str(cursor).strip().upper() == CURSOR_ALL:
cursor = None # 显式要求从头全量回放
try:
fills = downstream_repo.fetch_filled_orders(since_id=cursor, limit=limit)
except Exception as e:
out.update({"ok": False, "errors": [f"读 trading_order 失败: {type(e).__name__}: {e}"]})
return out
out["fills"] = len(fills)
if not fills:
out["cursor"] = cursor
return out
try:
instrs = _open_instructions()
except Exception as e:
instrs = []
out["errors"].append(f"读在途指令失败(按外部成交处理): {e}")
mapped = rc.map_fills_to_book(fills, instrs)
for act in mapped["actions"]:
try:
_apply_action(act)
out["actions"] += 1
except Exception as e:
logger.exception("入账失败 %s", act)
out["errors"].append(f"{act.get('ts_code')} 入账失败: {type(e).__name__}: {e}")
2026-07-29 10:49:56 +08:00
out["alerts"].extend(mapped["alerts"]) # 与 ws 那批告警合并, 日报一处看全
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}")
try:
nc = rc.next_cursor(fills, cursor)
pms_repo.set_param(CURSOR_KEY, nc, "system")
out["cursor"] = nc
except Exception as e:
out["errors"].append(f"游标推进失败: {e}")
out["ok"] = not out["errors"]
for a in out["alerts"]:
logger.warning("[回放告警] %s", a.get("message"))
return out
def _open_instructions() -> list:
"""在途指令 (回放认领的候选池)。
dispatched_at 下发时点 建单时点, **不用 updated_at**:
updated_at 每次部分成交回写都会往前跳, 拿它当下发时点会让同一条指令的后续成交
被时间守卫挡在门外, 误判成外部成交建单时点是天然的下界, 宁松勿紧
"""
rows = pms_repo.list_instructions(statuses=list(LIVE_INSTR), limit=500)
out = []
for r in rows:
prog = r.get("progress") or {}
out.append({"instruction_id": r["instruction_id"], "ts_code": r["ts_code"],
"side": r.get("side"), "qty": int(r.get("qty") or 0),
"exec_qty": int(r.get("exec_qty") or 0), "action": r.get("action"),
"dispatched_at": str(prog.get("dispatched_at")
or r.get("created_at") or "")})
return out
def _apply_action(act: dict):
code, qty, px = act["ts_code"], int(act["qty"]), float(act["price"] or 0)
if qty <= 0:
return
pms_repo.ensure_position(code)
if act["kind"] == "BUY":
pms_repo.insert_lot(ts_code=code, lot_type=act.get("lot_type") or "BASE", qty=qty,
open_price=px, open_date=datetime.now().date(),
instruction_id=act.get("instruction_id"),
note="外部成交并入 BASE" if not act.get("instruction_id") else None)
pms_repo.bump_position_qty(code, total_delta=qty, avail_delta=0) # T+1: 当日买入不可卖
else:
lots = pms_repo.list_lots(code, status="OPEN")
res = rc.apply_sell_to_lots(
[{"lot_id": l["id"], "lot_type": l["lot_type"], "qty": int(l["qty"]),
"open_date": _date_key(l["open_date"])} for l in lots], qty)
by_id = {l["id"]: l for l in lots}
t_profit = 0.0
for a in res["alloc"]:
lot = by_id[a["lot_id"]]
pnl = (px - float(lot["open_price"] or 0)) * a["qty"]
pms_repo.close_lot_qty(a["lot_id"], qty=a["qty"], close_price=px, realized_pnl=pnl)
if lot["lot_type"] == "T0":
t_profit += pnl
pms_repo.bump_position_qty(code, total_delta=-qty, avail_delta=-qty)
if t_profit:
pos = pms_repo.get_position(code) or {}
pms_repo.update_position(
code, realized_t_profit=float(pos.get("realized_t_profit") or 0) + t_profit)
for al in res["alerts"]:
logger.warning("[卖出核销] %s %s", code, al.get("message"))
if act.get("instruction_id"):
try:
pms_repo.add_instruction_exec(act["instruction_id"], qty)
_settle_instruction(act["instruction_id"])
except Exception as e:
logger.warning("指令进度更新失败 %s: %s", act["instruction_id"], e)
def _settle_instruction(instruction_id: str):
"""成交量达到指令数量 → 置 CONFIRMED (部分成交保持在途, 由窗口/过期规则收口)。"""
ins = pms_repo.get_instruction(instruction_id)
if ins and int(ins.get("exec_qty") or 0) >= int(ins.get("qty") or 0) > 0:
pms_repo.update_instruction(instruction_id, status="CONFIRMED")
def _date_key(d):
try:
return int(str(d).replace("-", "")[:8])
except (TypeError, ValueError):
return 0
def recompute_position(ts_code: str) -> dict:
"""由批次表重算持仓数量/摊薄成本/安全垫/垫子峰值。"""
lots = pms_repo.list_lots(ts_code, status=None, limit=2000)
qty = sum(int(l["qty"] or 0) for l in lots)
cum_buy = sum((int(l["qty"] or 0) + int(l["closed_qty"] or 0)) * float(l["open_price"] or 0)
for l in lots)
cum_sell = sum(int(l["closed_qty"] or 0) * float(l["close_avg_price"] or 0) for l in lots)
by_type = {}
for l in lots:
if int(l["qty"] or 0) > 0:
by_type[l["lot_type"]] = by_type.get(l["lot_type"], 0) + int(l["qty"])
avg_cost = max(0.0, (cum_buy - cum_sell) / qty) if qty > 0 else None
px = market.get_price(ts_code) or avg_cost or 0
cp = (px / avg_cost - 1.0) if (avg_cost and avg_cost > 0 and px) else None
solid = param_store.get_float("PMS_CUSHION_SOLID", 0.03)
pos = pms_repo.get_position(ts_code) or {}
peak = max(float(pos.get("cushion_peak") or 0), cp or 0)
scale = param_store.get_float("PMS_TOTAL_SCALE", 0)
fields = {
"total_qty": qty, "base_qty": by_type.get("BASE", 0) + by_type.get("RECON", 0),
"fill_qty": by_type.get("FILL", 0), "add_qty": by_type.get("ADD", 0),
"dca_qty": by_type.get("DCA", 0), "t0_qty": by_type.get("T0", 0),
"avg_cost": round(avg_cost, 3) if avg_cost else None,
"cushion_pct": round(cp, 4) if cp is not None else None,
"cushion_state": cu.cushion_state(cp, solid), "cushion_peak": round(peak, 4),
"pct_of_scale": round(qty * px / scale, 4) if scale > 0 else None,
"status": "CLOSED" if qty <= 0 else (pos.get("status") or "HOLDING"),
}
if qty > 0 and (pos.get("status") in (None, "PLANNED", "CLOSED")):
fields["status"] = "HOLDING"
if qty > 0 and not pos.get("opened_date"):
fields["opened_date"] = datetime.now().date()
pms_repo.update_position(ts_code, **fields)
return fields
# ================================================================ 对账
2026-07-29 16:55:45 +08:00
def _recon_blast_guard(book: list, ds: dict, diffs: list, *, force: bool = False):
"""对账修正的爆炸半径限制。返回 None 表示放行, 否则返回拦截原因。
只拦大面积重写, 不拦日常漂移 后者正是对账存在的意义两条判据:
1. **下游读回空而本端有持仓** 直接拦真的一夜清仓与读空在数据上完全同形,
但前者极罕见后者(重启/连错库/权限/网络)很常见, 且代价不对称: 拦错了只是晚一天
修账, 放错了是 22 个持仓的批次结构不可逆地没了
2. **受影响的股票数既超绝对值又超占比** 两个条件同时满足才拦, 避免持仓很少时
(比如只剩 3 ) 一点正常漂移就被占比判据挡住
首次建账 (本端无持仓下游有) 不受限制: 那时 held=0, 两条判据都不成立
"""
if force:
return None
held = [b for b in book if int(b.get("total_qty") or 0) > 0]
if not held:
return None # 空账本 → 建账/认领, 放行
if not (ds.get("rows") or []):
return {"why": "下游持仓读回空集, 而本端有持仓 —— 视为读数异常而非真实清仓",
"held": len(held), "diffs": len(diffs)}
max_names = param_store.get_int("PMS_RECON_MAX_FIX_NAMES", 5)
max_ratio = param_store.get_float("PMS_RECON_MAX_FIX_RATIO", 0.34)
ratio = len(diffs) / len(held)
if len(diffs) > max_names and ratio > max_ratio:
return {"why": f"差异面过大: {len(diffs)} 只 / 持仓 {len(held)}"
f"({ratio:.0%} > {max_ratio:.0%} 且 > {max_names} 只)",
"held": len(held), "diffs": len(diffs)}
return None
SRC_WS, SRC_TABLE, SRC_NONE = "ws", "table", "none"
def positions_source() -> dict:
"""对账的持仓事实源:**ws 快照为主、`trading_position` 表为兜底**2026-07-30 拍板)。
为什么要有这一层
----------------
协议 §6.2 定的对账事实源是 ws `query_positions` `snapshot{kind:"positions"}`;
而目标架构里 `trading_service` 全量退出业务只留看板, **那张表在新架构下没有明确的
写入方**07-30 实测: QMT 已切到模拟仓且功能正常, `trading_position` 却是空的
(`fetch_positions` columns 三个 None 就是"表里一行都没有"的铁证)
只认表的话, 账本永远建不起来
三条仲裁规则**两个源不一致时不许静默挑一个** 那会变成"两个同名不同物":
1. ws 快照新鲜 ws若表也非空且与 ws 对不上, **照样用 ws, 但记一条告警**
列出差异只数 那说明表的写入方与 QMT 已经不同步, 是要修的事, 不是噪音
2. ws 快照缺失/过期/字段不认 退回表, 并说明退回的原因 (三种原因处理起来完全不同:
没接通要找对端, 过期要看 pms-ws 活没活, 字段不认要补 ws_codec 的候选名)
3. 两个源都拿不到 `source=none`**这不等于"清仓"**, 上层必须据此拒绝改账
新鲜度按本端 `received_at` , 不用 payload `as_of` 理由见 qmt_repo.latest_snapshot
"""
mode = param_store.get("PMS_RECON_SOURCE", "ws_first") or "ws_first"
max_age = param_store.get_int("PMS_RECON_WS_SNAPSHOT_MAX_AGE_SEC", 900)
out = {"source": SRC_NONE, "rows": [], "columns": {"qty": None, "avail": None, "cost": None},
"raw_count": 0, "as_of": 0, "age_sec": None, "alerts": [], "mode": mode}
ws_snap, ws_why = None, ""
if mode in ("ws_first", "ws_only"):
try:
snap = qmt_repo.latest_snapshot("positions")
except Exception as e:
snap, ws_why = None, f"读 ws 快照失败: {type(e).__name__}: {e}"
if not snap:
ws_why = ws_why or ("ws 从未回过 positions 快照 —— 确认 pms-ws 在跑、"
"PMS_QMT_WS_ENABLED 已开, 且对端实现了 query_positions")
elif snap.get("age_sec") is not None and snap["age_sec"] > max_age:
ws_why = (f"ws 快照已过期 ({snap['age_sec']:.0f}s > {max_age}s) —— "
f"pms-ws 可能没在跑, 或对端不再回应 query_positions")
else:
parsed = wsc.parse_positions_snapshot(snap["payload"])
if parsed["raw_count"] and parsed["columns"].get("qty") is None:
ws_why = (f"ws 快照有 {parsed['raw_count']} 个条目却认不出数量列 —— "
f"把对端实际字段名补进 ws_codec._SNAP_QTY")
else:
# 归一到 PMS 内部点式口径。协议 §8 规定就是点式, 但对端给前缀式 (SH600000)
# 时若不转, diff 会拿 "SH600000" 去比账本里的 "600000.SH" —— 结果是**每一只
# 都对不上**: 账本那只判"下游没有了"要核销, ws 那只判"新持仓"要补。
# 一次代码格式不一致就能造出一轮双向全量重写, 比读空还狠。
for r in parsed["rows"]:
r["ts_code"] = downstream_repo.to_dot(r["ts_code"])
ws_snap = {**parsed, "age_sec": snap.get("age_sec"), "seq": snap.get("seq")}
tbl, tbl_err = None, ""
if mode in ("ws_first", "table_only"):
try:
tbl = downstream_repo.fetch_positions()
except Exception as e:
tbl_err = f"读 trading_position 失败: {type(e).__name__}: {e}"
if ws_snap:
out.update({"source": SRC_WS, "rows": ws_snap["rows"], "columns": ws_snap["columns"],
"raw_count": ws_snap["raw_count"], "as_of": ws_snap["as_of"],
"age_sec": ws_snap["age_sec"], "seq": ws_snap.get("seq")})
if tbl and (tbl.get("rows") or []):
ws_map = {r["ts_code"]: int(r.get("qty") or 0) for r in ws_snap["rows"]}
tb_map = {r["ts_code"]: int(r.get("qty") or 0) for r in tbl["rows"]}
gap = [c for c in set(ws_map) | set(tb_map) if ws_map.get(c, 0) != tb_map.get(c, 0)]
if gap:
out["alerts"].append(
{"level": "WARN", "code": "SOURCE_DISAGREE",
"message": f"ws 快照与 trading_position 对不上 {len(gap)}"
f"(ws {len(ws_map)} 只 / 表 {len(tb_map)} 只): "
f"{sorted(gap)[:8]}。**已按 ws 为准**; 表的写入方与 QMT "
f"不同步, 需查明是谁在写那张表"})
return out
if tbl is not None and not tbl_err:
# **「应答了空」≠「没应答」。** 表查询成功返回 0 行, 是下游给出的一个**有效应答**
# (「我没有持仓」) —— 它可能是真清仓、也可能是它自己读空, 该交给爆炸半径闸按老规矩
# 处理, 人工确认后 force 能放行。而"两个源都取不到答案"是压根不知道账户状态, 那才
# 该在对账入口就拒掉。
# 一开始把两者都归成 source=none, 结果 force 被入口那道拒绝挡死, 07-29 那个
# 「下游读空绝不清账」用例的 force 分支当场挂了 —— 空集是数据, 不是缺数据。
if tbl["rows"] and tbl["columns"].get("qty") is None:
out["alerts"].append({"level": "ERROR", "code": "TABLE_NO_QTY_COL",
"message": "trading_position 未识别出数量列 —— 按 "
"QMT_INTERFACE_REQUIREMENTS A1/D1 取 DDL 后把列名"
"补进 downstream_repo.QTY_CANDIDATES"})
return out # 认不出数量列 = 读不到, 不是没持仓
out.update({"source": SRC_TABLE, "rows": tbl["rows"], "columns": tbl["columns"],
"raw_count": tbl.get("raw_count") or len(tbl["rows"])})
if mode == "ws_first" and ws_why:
out["alerts"].append({"level": "WARN", "code": "WS_SNAPSHOT_UNAVAILABLE",
"message": f"退回 trading_position 表作为事实源: {ws_why}"})
return out
out["alerts"].append({"level": "WARN", "code": "NO_POSITION_SOURCE",
"message": f"两个事实源都拿不到持仓 —— ws: {ws_why or '未启用'}; "
f"表: {tbl_err or '空集'}。**这不等于清仓**"})
return out
2026-07-29 16:55:45 +08:00
def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict:
"""账本 vs 下游持仓, 以下游为准修正并留痕。事实源见 positions_source()。
2026-07-29 16:55:45 +08:00
force=True 绕过爆炸半径限制 ( _recon_blast_guard) 只在人工确认下游读数确实
正确之后使用
"""
out = {"ok": True, "diffs": [], "fixes": [], "errors": [], "columns": {}, "severity": rc.SEV_OK}
src = positions_source()
ds = {"rows": src["rows"], "columns": src["columns"], "raw_count": src["raw_count"]}
out.update({"source": src["source"], "source_mode": src["mode"],
"source_age_sec": src.get("age_sec"), "source_alerts": src["alerts"],
"columns": src["columns"]})
for a in src["alerts"]:
(logger.error if a["level"] == "ERROR" else logger.warning)(
"[对账·事实源] %s", a["message"])
if src["source"] == SRC_NONE:
# 走到这里意味着**两个源都没有给出应答**(ws 无新鲜快照 + 表查询异常/未启用), 不是
# "应答了空集"——后者是有效数据, 归 SRC_TABLE 交给爆炸半径闸, 见 positions_source。
#
# 按账本有没有持仓分两级, 不是一律 ERROR: 账本也空时 (刚清账、等对端装持仓) 本来就
# 无账可对, 而盘中轻对账每分钟一跳, 刷 ERROR 只会把真告警埋掉 —— 与补发期告警限流
# 同一个道理。账本有持仓却拿不到任何事实源, 那才是真要停下来的事。
#
# **这一级连 force 都不放行**, 与爆炸半径闸不同。force 的语义是「人工已确认下游读数
# 正确」, 而这里根本没有读数可供确认 —— 没有任何数字, 人也无从确认。真要清账走
# scripts/reset_ledger.py 那条明路, 别拿一个空壳子当"下游事实"去核销批次。
held = [p for p in pms_repo.list_positions() if int(p.get("total_qty") or 0) > 0]
if held:
out.update({"severity": rc.SEV_ERROR, "fixes": [],
"blocked": {"why": "两个事实源都没有应答 (不是应答了空集), 而本端有"
"持仓 —— 拒绝对账。读不到 ≠ 清仓; force 在此不放行, "
"要清账用 scripts/reset_ledger.py",
"held": len(held), "force_ignored": bool(force)}})
logger.error("[对账] 拒绝对账: 无事实源应答而本端有 %s 只持仓 (force=%s 不放行)",
len(held), force)
else:
out["note"] = "两个事实源都没有应答, 账本也空 —— 无可对之账 (等对端装持仓)"
return out
2026-07-30 08:42:06 +08:00
# 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、
# 基金、脏数据), 而对账是**会照着它改账本**的, 放进来就变成真持仓。
# 2026-07-29: 下游只剩一只 `999999.SH` 2000 股 —— 我们自己测拒绝路径用的假代码 ——
# 若照单全收, 明天日终就会给账本凭空补出 2000 股不存在的票。
ds_rows, junk = [], []
for r in ds["rows"]:
(ds_rows if cs.is_stock_code(r.get("ts_code")) else junk).append(r)
if junk:
out["junk_codes"] = [{"ts_code": r.get("ts_code"), "qty": r.get("qty")} for r in junk]
logger.error("[对账] 下游有 %s 条非 A 股个股代码, 已忽略不入账: %s"
"若其中确有真实标的, 说明 command_spec.STOCK_PREFIXES 缺了代码段",
len(junk), out["junk_codes"])
ds = {**ds, "rows": ds_rows}
book = [{"ts_code": r["ts_code"], "total_qty": int(r.get("total_qty") or 0)}
for r in pms_repo.list_positions()]
diffs = rc.diff_positions(book, [{"ts_code": r["ts_code"], "qty": r["qty"]}
for r in ds["rows"]])
out["diffs"] = diffs
# 连续不一致计数 —— **按交易日推进, 且只有权威那一趟才写**。
# 盘中轻对账 (apply_fix=False) 每分钟跑一次, 让它推进的话「连续 3 日」就成了「连续
# 3 分钟」, 日报还会写出「连续 175 日」这种数。详见 rc.advance_streak 的注释。
streak = param_store.get_int(STREAK_KEY, 0)
if apply_fix:
st = rc.advance_streak(streak, param_store.get_int(STREAK_YMD_KEY, 0),
int(datetime.now().strftime("%Y%m%d")), bool(diffs))
if st["changed"] or st["ymd"] != param_store.get_int(STREAK_YMD_KEY, 0):
# 必须走 ParamStore 写入: 直接写库不会失效缓存, 连续天数会一直读到旧值
param_store.set_param(STREAK_KEY, st["streak"], "system")
param_store.set_param(STREAK_YMD_KEY, st["ymd"], "system")
streak = st["streak"]
out["streak"] = streak
out["severity"] = rc.recon_severity(streak if diffs else 0,
param_store.get_int("PMS_RECON_ALARM_DAYS", 3))
if not diffs or not apply_fix:
return out
2026-07-29 16:55:45 +08:00
blast = _recon_blast_guard(book, ds, diffs, force=force)
if blast:
# 差异面太大, 不自动改账。**这一条是 2026-07-29 的血的教训**: 下游一次读空
# (对端模拟环境重启), 对账照着「以下游为准」把 22 个持仓 126 万市值全核销了,
# 全程没有任何拦截 —— 上面那道数量列校验写的是 `if ds["rows"] and ...`,
# 空结果集连它都不走。
# 「以下游为准」是对的, 但它的前提是「读到的确实是下游的真实状态」。读空、读漏、
# 连错库都会长成同一个样子, 而代价是不可逆的批次核销。所以给它加个爆炸半径:
# 小幅漂移照常自动修 (那正是对账的价值), 大面积重写一律停下来等人。
out.update({"severity": rc.SEV_ERROR, "fixes": [], "blocked": blast})
logger.error("[对账] 拒绝自动修正: %s。差异 %s 项, 持仓 %s 只。"
"确认下游读数无误后, 用 apply_fix=true&force=true 手工放行",
blast["why"], len(diffs), blast["held"])
return out
codes = [d["ts_code"] for d in diffs]
prices = market.get_prices(codes)
2026-07-31 13:46:15 +08:00
# 首次建账的成本价闸 (2026-07-31 加; QMT_SIDE_S3_CLOSEOUT.md §5)
# ---------------------------------------------------------------
# 只在**账本为空**时生效 —— 那一刻是不可逆的: 这一批 RECON 批次的开仓价会定死每只票
# 的摊薄成本, 而摊薄成本是安全垫的分母, 安全垫又是补仓/加仓/保垫减仓的共同判据。
# 日常漂移不走这道闸: 那时已有批次各自带着自己的成本, 补几百股用什么价影响有限,
# 而每天拦一次对账才是真的坏事。
#
# 为什么非要有这道闸: `rc.build_recon_fixes` 在下游给不出成本价时**静默退回现价**,
# 只在 note 里留一句"现价兜底"。于是账本建起来了、页面一切正常、每个数都长得像真的,
# 只有安全垫齐刷刷是 0 —— 而没人会盯着一个"看起来就该是 0"的字段。
# 漏账看得见 (有行、有 ERROR、有数字), 错账看不见。所以拿不准就不建。
if not [b for b in book if int(b.get("total_qty") or 0) > 0] and \
param_store.get_bool("PMS_RECON_REQUIRE_COST", True):
chk = rbc.check_costs(ds.get("rows") or [], prices)
out["cost_check"] = {k: chk[k] for k in ("n", "counts", "eq_ratio", "blocking",
"reasons", "hint")}
out["coverage"] = rbc.coverage(chk["rows"])
if chk["blocking"] and not force:
out.update({"severity": rc.SEV_ERROR, "fixes": [],
"blocked": {"why": "首次建账被成本价闸拦下: " + "; ".join(chk["reasons"]),
"hint": chk["hint"],
"how": "请对端把 trading_position 的 cost_price 填成真实"
"成本价再重试; 确认这份数据就是对的可用 force=true "
"放行 (那意味着你接受安全垫从 0 起算)"}})
logger.error("[对账] 首次建账被成本价闸拦下: %s%s", chk["reasons"], chk["hint"])
return out
if chk["blocking"]:
logger.warning("[对账] 成本价闸本应拦下 (%s), 但 force=true 放行 —— "
"安全垫将从 0 起算, 补仓/加仓/保垫减仓这一轮判不准", chk["reasons"])
if not out["coverage"]["enough"]:
logger.warning("[对账] 建账数据的情形覆盖不全: %s", out["coverage"]["hint"])
2026-07-29 10:12:12 +08:00
# 下游的成本价 —— 补仓位时的开仓价优先取它, 现价只兜底 (见 rc.build_recon_fixes 注释:
# 拿现价当成本会让安全垫齐刷刷归零, 整条纪律链跟着失灵)
costs = {r["ts_code"]: r.get("cost") for r in (ds.get("rows") or [])}
lots_map = {c: [{"lot_id": l["id"], "lot_type": l["lot_type"], "qty": int(l["qty"]),
"open_date": _date_key(l["open_date"])}
for l in pms_repo.list_lots(c, status="OPEN")] for c in codes}
2026-07-29 10:12:12 +08:00
fixes = rc.build_recon_fixes(diffs, price_map=prices, lots_map=lots_map, cost_map=costs)
for f in fixes:
try:
_apply_fix(f)
recompute_position(f["ts_code"])
except Exception as e:
logger.exception("对账修正失败 %s", f)
out["errors"].append(f"{f['ts_code']} 修正失败: {type(e).__name__}: {e}")
out["fixes"] = fixes
out["ok"] = not out["errors"]
if out["severity"] == rc.SEV_ERROR:
logger.error("[对账] 连续 %s 日不一致, 升级 ERROR 待人工: %s 项差异", streak, len(diffs))
return out
2026-07-31 13:46:15 +08:00
def rebuild_preflight() -> dict:
"""账本重建的**只读**预检 (README 待办 #4)。一个字都不写, 随时可跑。
回答三个问题, 顺序就是它们该被回答的顺序:
1. 事实源给不给得出持仓? (给不出就没什么可谈的)
2. 这份持仓的成本价能不能用? (`rebuild_check` 的成本价闸, 这是重建的成败所在)
3. 建完之后验不验得到纪律? (§5.2 的四种情形; 不阻断, 但缺了就是白建一轮)
再加一段开关现状清账期间关掉的东西, 重建完必须记得打开 尤其 `PMS_SIGNAL_ENABLED`:
账本空时信号消化判 IGNORE 也照样 ACK, 卖出信号会被消费组静默吃掉且跨日拿不回来
"""
out = {"ok": True, "ready": False, "steps": [], "switches": {}, "hint": ""}
book = [p for p in pms_repo.list_positions() if int(p.get("total_qty") or 0) > 0]
out["book_held"] = len(book)
out["first_build"] = not book
src = positions_source()
out["source"] = {"source": src["source"], "mode": src["mode"], "age_sec": src.get("age_sec"),
"rows": len(src["rows"]), "columns": src["columns"],
"alerts": src["alerts"]}
if src["source"] == SRC_NONE:
out["steps"].append({"step": "事实源", "ok": False,
"why": "ws 快照与 trading_position 都没有应答 —— 不是空集, 是没读到。"
"先确认 pms-ws 在跑且对端实现了 query_positions"})
out["hint"] = "拿不到事实源, 谈不上重建 (读不到 ≠ 清仓)"
return out
rows = [r for r in src["rows"] if cs.is_stock_code(r.get("ts_code"))]
out["steps"].append({"step": "事实源", "ok": bool(rows),
"why": f"{src['source']} 给出 {len(rows)} 只持仓"
+ ("" if rows else " —— 对端还没装持仓, 等它")})
if not rows:
out["hint"] = "下游一只持仓都没有 —— 等对端装好再来"
return out
prices = market.get_prices([r["ts_code"] for r in rows])
chk = rbc.check_costs(rows, prices)
out["cost_check"] = chk
out["steps"].append({"step": "成本价体检", "ok": not chk["blocking"], "why": chk["hint"]})
cov = rbc.coverage(chk["rows"])
out["coverage"] = cov
out["steps"].append({"step": "情形覆盖", "ok": cov["enough"], "why": cov["hint"],
"blocking": False})
out["switches"] = {
"PMS_DISPATCH_MODE": param_store.get("PMS_DISPATCH_MODE", "shadow"),
"PMS_AUTONOMY": param_store.get("PMS_AUTONOMY", "propose_only"),
"PMS_SIGNAL_ENABLED": param_store.get_bool("PMS_SIGNAL_ENABLED", True),
"PMS_EXEC_HALT": param_store.get_bool("PMS_EXEC_HALT", False),
"PMS_SECTOR_SOURCE": param_store.get("PMS_SECTOR_SOURCE", ""),
}
out["after"] = [
"重建完立刻把 PMS_SIGNAL_ENABLED 打开 —— 账本空时信号消化判 IGNORE 也照样 ACK, "
"卖出信号被消费组静默吃掉且跨日拿不回来",
"确认账本无误后再重开 pms-beat (15:10 日终结算会自动认领 trading_position, "
"对端装到一半时跨过 15:10 会拿半成品建账)",
"跑一次 ledger_service.rebuild_accept() 看安全垫分布与行业集中度",
]
out["ready"] = not chk["blocking"]
out["hint"] = (chk["hint"] if chk["blocking"] else
f"可以重建: {len(rows)} 只持仓、成本价可用。" +
("" if cov["enough"] else "注意 " + cov["hint"]))
return out
def rebuild_accept() -> dict:
"""重建之后的**只读**判收。回答"这本账建对了没有"
最硬的一条判据是**安全垫分布**: 如果重建后每一只票的安全垫都是 0, 那就是踩了
拿现价当成本那个坑 页面上每个数都合理, 只有这一处露馅所以专门数它
顺带回答 README 待办 #3 留的那个问题: `PMS_SECTOR_MAX_RATIO=40%` 是当初按"二级或更粗"
的粒度定的, 换到 gp_hybk 三级 (884*) 之后偏不偏松 有了真实持仓分布才算得出来
"""
out = {"ok": True, "checks": [], "hint": ""}
view = portfolio.positions_view()
held = view.get("positions") or []
t = view.get("totals") or {}
out["held"] = len(held)
if not held:
out.update({"ok": False, "hint": "账本还是空的 —— 重建没跑, 或跑了但被闸拦下了"})
return out
cush = [x.get("cushion_pct") for x in held if x.get("cushion_pct") is not None]
zero = [x for x in held if abs(float(x.get("cushion_pct") or 0)) < 1e-9]
out["cushion"] = {"n": len(cush), "zero": len(zero),
"min": round(min(cush), 4) if cush else None,
"max": round(max(cush), 4) if cush else None,
"solid": t.get("solid_names"), "neg": t.get("neg_names")}
all_zero = bool(held and len(zero) == len(held))
out["checks"].append({
"check": "安全垫分布", "ok": not all_zero,
"why": ("**每一只的安全垫都是 0** —— 这正是拿现价当成本的样子。摊薄成本等于当天价, "
"补仓/加仓/保垫减仓这一整条纪律链会全程判不出来。请核对下游 cost_price"
if all_zero else
f"{len(cush)} 只有安全垫, 区间 {out['cushion']['min']:+.1%} ~ "
f"{out['cushion']['max']:+.1%}, 其中 {len(zero)} 只为 0")})
lots = {}
for x in held:
for l in pms_repo.list_lots(x["ts_code"], status="OPEN"):
lots[l["lot_type"]] = lots.get(l["lot_type"], 0) + 1
out["lots"] = lots
out["checks"].append({"check": "批次账", "ok": bool(lots),
"why": f"批次分布 {lots or '(空)'}" +
("" if lots else " —— 持仓有行但没有批次, 账本结构不完整")})
ready = bool(t.get("sector_ready", view.get("sector_ready")))
names, mv = t.get("sector_names") or {}, t.get("sector_mv") or {}
# **分母必须与规则闸一致**: sizer.check_caps 的行业判据是
# (sector_mv + add) / port_after > sector_max_ratio
# 即「占组合持仓市值」, 不是占总规模 PMS_SCALE。两边用不同分母的话, 这里报"没超"而
# 规则闸拦人 (或反过来), 而两个数字都自称是"行业集中度"——同名不同物最难查。
# (settings.py 里那句注释写的是"占总仓", 容易被读成占 PMS_SCALE, 已在此写明。)
port_mv = float(t.get("portfolio_mv") or 0.0) or sum(
float(x.get("market_value") or 0.0) for x in held)
max_ratio = param_store.get_float("PMS_SECTOR_MAX_RATIO", 0.40)
max_names = param_store.get_int("PMS_SECTOR_MAX_NAMES", 4)
top = sorted(((s, (mv[s] / port_mv if port_mv else 0.0), names.get(s, 0)) for s in mv),
key=lambda x: -x[1])[:5]
over_ratio = [x for x in top if x[1] > max_ratio + 1e-9]
over_names = [s for s, n in names.items() if n > max_names]
out["sector"] = {"ready": ready, "max_ratio": max_ratio, "max_names": max_names,
"denominator": "组合持仓市值 (与 sizer.check_caps 一致)",
"portfolio_mv": round(port_mv, 2),
"over_ratio": [s for s, _, _ in over_ratio], "over_names": over_names,
"top": [{"sector": s, "ratio": round(r, 4), "names": n} for s, r, n in top]}
worst = top[0][1] if top else 0.0
if not ready:
why = "行业源没就绪 —— 集中度约束此刻等同未配置且会静默失效"
elif over_ratio or over_names:
# 接管进来的持仓超限**不算重建失败** —— 账本忠实反映了下游的真实状态, 那才是它的职责。
# 但必须说出来: 规则闸从此会挡住这些行业的加仓, 不说的话下次加不进去会以为是 bug。
why = (f"接管进来的持仓**已经超限**: "
+ (f"占比超上限 {max_ratio:.0%} 的有 {[s for s, _, _ in over_ratio]}"
f" (最大 {worst:.1%})" if over_ratio else "")
+ (f" 只数超上限 {max_names} 的有 {over_names}" if over_names else "")
+ "。这不是重建出错 —— 账本忠实反映了下游真实持仓; 但规则闸从此会挡住这几个"
"行业的加仓, 心里要有数 (减持方向不受影响)")
else:
why = f"最大行业占组合 {worst:.1%} (上限 {max_ratio:.0%}), 只数上限 {max_names}, 未超限"
if worst < max_ratio * 0.5:
why += ("。阈值是当初按二级或更粗的粒度定的, 现在是三级 884* —— "
"拿这份真实分布回看它偏不偏松 (参考项目三级用 20%)")
out["checks"].append({"check": "行业集中度", "ok": ready, "why": why})
out["ok"] = all(c["ok"] for c in out["checks"])
bad = [c["check"] for c in out["checks"] if not c["ok"]]
out["hint"] = "账本重建判收通过" if out["ok"] else "这几项没过: " + " / ".join(bad)
return out
def _apply_fix(f: dict):
code = f["ts_code"]
pms_repo.ensure_position(code)
if f["op"] == "ADD_RECON_LOT":
px = float(f.get("price") or 0)
pms_repo.insert_lot(ts_code=code, lot_type="RECON", qty=int(f["qty"]),
open_price=px, open_date=datetime.now().date(),
note=f["note"] + ("" if px > 0 else " [缺现价, 成本待人工核]"))
else:
lots = {l["id"]: l for l in pms_repo.list_lots(code, status="OPEN")}
px = float(f.get("price") or 0)
for a in f.get("alloc") or []:
lot = lots.get(a["lot_id"])
pnl = (px - float(lot["open_price"] or 0)) * a["qty"] if (lot and px) else 0.0
pms_repo.close_lot_qty(a["lot_id"], qty=a["qty"], close_price=px or
float(lot["open_price"] or 0), realized_pnl=pnl)
pms_repo.insert_ledger(ts_code=code, action="RECON", arbiter="rule", verdict="PASS",
price_at=float(f.get("price") or 0),
hard_numbers={"op": f["op"], "qty": f["qty"]},
reason=f["note"])
# ================================================================ 除权
def detect_and_apply_ex_right() -> dict:
"""用昨日结算快照与今日持仓/价格比对, 识别送转股并按比例调整批次。"""
out = {"checked": 0, "ex_rights": [], "mismatches": [], "errors": []}
prev = _prev_snapshot()
if not prev:
out["errors"].append("无昨日结算快照, 本次跳过除权检测 (次日起生效)")
return out
for pos in pms_repo.list_positions(only_open=True):
code = pos["ts_code"]
old = prev.get(code)
if not old:
continue
out["checked"] += 1
px = market.get_price(code)
r = rc.detect_ex_right(int(old.get("qty") or 0), int(pos.get("total_qty") or 0),
float(old.get("price") or 0), float(px or 0))
if not r:
continue
if r["kind"] == "EX_RIGHT":
try:
lots = pms_repo.list_lots(code, status="OPEN")
for l in rc.apply_ex_right(lots, r["ratio"]):
pms_repo.update_lot(l["id"], qty=l["qty"], open_price=l["open_price"],
note=l["note"])
recompute_position(code)
pms_repo.insert_ledger(ts_code=code, action="RECON", arbiter="rule",
verdict="PASS", price_at=px or 0, hard_numbers=r,
reason=f"除权调整 ×{r['ratio']}")
out["ex_rights"].append({"ts_code": code, **r})
except Exception as e:
out["errors"].append(f"{code} 除权调整失败: {e}")
else:
out["mismatches"].append({"ts_code": code, **r})
logger.error("[除权] %s 比例不吻合, 待人工: %s", code, r.get("reason"))
return out
def _prev_snapshot() -> dict:
"""取最近一份日终快照 {code: {qty, price}} (存在 pms_daily_report 里, 不新增表)。"""
r = pms_repo.latest_report()
if not r:
return {}
snap = (r.get("report") or {}).get("snapshot") or {}
return snap if isinstance(snap, dict) else {}
# ================================================================ 盘前 / 日终
def premarket() -> dict:
"""盘前准备 (08:50): T+1 可卖重置、参考位取数、刹车结算。"""
out = {"ok": True, "avail_reset": 0, "refs": 0, "brake": None, "errors": []}
try:
out["avail_reset"] = pms_repo.reset_avail_all()
except Exception as e:
out["errors"].append(f"可卖量重置失败: {e}")
for pos in pms_repo.list_positions(only_open=True):
code = pos["ts_code"]
try:
refs = market.get_refs(code, base_cost=pos.get("avg_cost"))
pms_repo.update_position(code, support_ref=refs.get("support"),
pressure_ref=refs.get("pressure"),
stop_ref=refs.get("stop"),
ref_source=refs.get("source"))
out["refs"] += 1
except Exception as e:
out["errors"].append(f"{code} 参考位取数失败: {e}")
try:
out["brake"] = _settle_brake()
except Exception as e:
out["errors"].append(f"刹车结算失败: {e}")
out["ok"] = not out["errors"]
return out
def _settle_brake() -> dict:
"""组合刹车: 自高水位回撤 ≥ 阈值 → 自主增持停 N 个交易日 (命令类不受限)。"""
v = portfolio.positions_view()
mv = v["totals"]["portfolio_mv"]
hw = param_store.get_float("PMS_HIGH_WATER", 0.0)
dd_limit = param_store.get_float("PMS_BRAKE_DRAWDOWN", 0.05)
days = param_store.get_int("PMS_BRAKE_DAYS", 3)
until = param_store.get_int("PMS_BRAKE_UNTIL", 0)
today = td.ymd()
if mv > hw:
param_store.set_param("PMS_HIGH_WATER", mv, "system")
hw = mv
drawdown = (1 - mv / hw) if hw > 0 else 0.0
if hw > 0 and drawdown >= dd_limit and today >= until:
until = td.ymd(td.next_trade_day(datetime.now().date(), days))
param_store.set_param("PMS_BRAKE_UNTIL", until, "system")
logger.warning("[刹车] 自高水位回撤 %.1f%%%.0f%%, 自主增持暂停至 %s",
drawdown * 100, dd_limit * 100, until)
return {"high_water": hw, "portfolio_mv": mv, "drawdown": round(drawdown, 4),
"brake_until": until, "active": today < until}
def daily_settle() -> dict:
"""日终结算 (15:10): 除权检测 → 全量对账 → 垫子峰值/连负天数 → 命令进度日结 → 快照留存。"""
from app.services import command_service
out = {"ok": True, "steps": {}, "errors": []}
try:
out["steps"]["ex_right"] = detect_and_apply_ex_right()
except Exception as e:
out["errors"].append(f"除权检测失败: {e}")
try:
out["steps"]["recon"] = reconcile()
except Exception as e:
out["errors"].append(f"对账失败: {e}")
try:
out["steps"]["cushion"] = _settle_cushion()
except Exception as e:
out["errors"].append(f"安全垫结算失败: {e}")
2026-07-29 10:49:56 +08:00
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:
out["errors"].append(f"命令进度结算失败: {e}")
try:
pms_repo.expire_proposals()
except Exception as e:
out["errors"].append(f"提议过期处理失败: {e}")
out["ok"] = not out["errors"]
return out
def _settle_cushion() -> dict:
"""更新垫子峰值与「安全垫连续为负天数」(清弱票判定所需)。"""
v = portfolio.positions_view()
streak = portfolio.neg_streak_map()
updated = 0
for x in v["held"]:
code, cp = x["ts_code"], x["cushion_pct"]
streak[code] = (int(streak.get(code, 0)) + 1) if (cp is not None and cp < 0) else 0
peak = max(float(x["cushion_peak"] or 0), cp or 0)
pms_repo.update_position(code, cushion_pct=cp, cushion_peak=round(peak, 4),
cushion_state=x["cushion_state"])
updated += 1
held = {x["ts_code"] for x in v["held"]}
portfolio.save_neg_streak({k: v2 for k, v2 in streak.items() if k in held})
return {"updated": updated,
"neg_streak": {k: v2 for k, v2 in streak.items() if v2 > 0 and k in held}}
# ================================================================ 日报
def build_daily_report(ymd: int = None) -> dict:
"""运营日报 (15:30): 关注区 + 全量统计 + 当日快照 (快照供次日除权检测)。"""
from app.services import command_service
ymd = int(ymd or td.ymd())
v = portfolio.positions_view()
t = v["totals"]
try:
recon_state = {"streak": param_store.get_int(STREAK_KEY, 0)}
except Exception:
recon_state = {}
cmds = pms_repo.list_commands(statuses=["EXECUTING", "PARTIAL", "PENDING", "PLANNING"],
limit=100)
proposals = pms_repo.list_proposals(statuses=("WAIT_USER",), limit=100)
live_ins = pms_repo.list_instructions(statuses=list(LIVE_INSTR), limit=200)
attention = []
for c in cmds:
prog = c.get("progress") or {}
attention.append({"type": "命令进度", "command_id": c["command_id"],
"cmd_type": c["cmd_type"], "status": c["status"],
"done": prog.get("done_amount"), "target": prog.get("target_amount"),
"deadline": prog.get("deadline")})
if proposals:
attention.append({"type": "待确认提议", "count": len(proposals)})
if recon_state.get("streak"):
attention.append({"type": "对账差异", "streak": recon_state["streak"],
"severity": rc.recon_severity(
recon_state["streak"],
param_store.get_int("PMS_RECON_ALARM_DAYS", 3))})
brake_until = param_store.get_int("PMS_BRAKE_UNTIL", 0)
if brake_until > ymd:
attention.append({"type": "组合刹车", "until": brake_until})
if v["price_missing"]:
attention.append({"type": "行情缺失", "codes": v["price_missing"]})
if not v["sector_ready"]:
attention.append({"type": "行业约束停用", "hint": "PMS_SECTOR_SOURCE 未配置"})
# 总规模与账户实际总资产的偏差。**这条必须每天摆出来**: 三道仓位闸 (总仓/单股/只数) 全部
# 以 PMS_TOTAL_SCALE 为基准算, 而它是用户命令参数, 不是账户余额。两者一旦差得远, 规划器
# 会排出账户根本执行不了的方案, 却一路过闸 —— 到 QMT 那儿才被 INSUFFICIENT_CASH 拒。
# 07-30 就是这个情形: scale 200 万 vs 模拟账户 98 万, 差一倍而页面上看不出任何异常。
ta, scale_v = t.get("total_asset"), float(t.get("scale") or 0)
if ta and scale_v > 0:
gap = float(ta) / scale_v - 1
if abs(gap) > param_store.get_float("PMS_SCALE_GAP_ALARM", 0.10):
attention.append({"type": "规模与账户不符", "scale": scale_v,
"total_asset": round(float(ta), 2), "gap": round(gap, 4),
"hint": "仓位上限按 scale 算, 资金校验按账户算 —— 差得远时方案会"
"过闸但下不出去。改 PMS_TOTAL_SCALE 或核对账户"})
if t.get("cash_source") != "ws":
attention.append({"type": "资金快照未接通", "hint": t.get("cash_why") or "",
"note": "买入前的资金校验已降级 (只按 scale市值 估算), 见规则闸的 "
"CASH_ESTIMATED 告警"})
if td.calendar_degraded():
attention.append({"type": "交易日历降级", "hint": "未安装 chinesecalendar, 节假日不可辨"})
report = {
"ymd": ymd, "generated_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
"totals": t, "attention": attention,
"commands": [{"command_id": c["command_id"], "cmd_type": c["cmd_type"],
"status": c["status"], "progress": c.get("progress")} for c in cmds],
"proposals": len(proposals), "live_instructions": len(live_ins),
"positions": [{"ts_code": x["ts_code"], "qty": x["total_qty"], "price": x["price"],
"avg_cost": x["avg_cost"], "cushion_pct": x["cushion_pct"],
"cushion_state": x["cushion_state"], "mv": x["market_value"],
"pct_of_scale": x["pct_of_scale"]} for x in v["held"]],
# 次日除权检测用的快照 (数量 + 收盘价)
"snapshot": {x["ts_code"]: {"qty": x["total_qty"], "price": x["price"]}
for x in v["held"]},
"recon": recon_state,
}
try:
pms_repo.upsert_report(ymd, report)
except Exception as e:
logger.error("日报落表失败: %s", e)
report["save_error"] = str(e)
return report
def expire_stale_instructions() -> int:
"""下发后长时间未被接受的指令置过期 (不自动重发 —— 设计 §13)。"""
mins = param_store.get_int("PMS_DISPATCH_EXPIRE_MIN", 30)
cut = datetime.now() - timedelta(minutes=mins)
n = 0
for r in pms_repo.list_instructions(statuses=["DISPATCHED"], limit=500):
try:
if str(r.get("updated_at") or "") and str(r["updated_at"]) < str(cut):
pms_repo.update_instruction(r["instruction_id"], status="EXPIRED")
n += 1
except Exception:
continue
return n