tradingSystem/app/services/ledger_service.py

1127 lines
61 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# -*- 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
from app.core import command_spec as cs
from app.core import cushion as cu
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
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"
# 该 streak 最后一次推进是哪个交易日 —— 同一天内反复对账不重复计数的判据
STREAK_YMD_KEY = "PMS_RECON_STREAK_YMD"
LIVE_INSTR = ("DISPATCHED", "JUDGE_PASSED", "RULE_PASSED")
# ================================================================ 回放
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
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 等下一跳。
"""
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
# ---- 三分流: 联调单 / 孤儿成交 / 真单。判据只查**我们自己的出口表**
# (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 = [], [], []
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
try:
od = qmt_repo.get_order(iid)
except Exception as e:
# 库抖一下**什么都不判**, 留在 processed=0 下一跳再来。原先这里把异常当"查不到"
# 并入真单分支 —— 一次连接超时就足以凭空造出一笔"外部成交"。
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)
if smoke:
seqs = [r["seq"] for r in smoke]
qmt_repo.inbox_mark(seqs, processed=qmt_repo.PROC_DIGESTED,
note="ws 联调测试单成交, 不入账 (SMOKE_)")
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)
trades = real
if not trades:
out["ok"] = not out["errors"]
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 成交已入账")
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"]:
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:
"""成交回放一跳 = 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)
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}")
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
# ================================================================ 对账
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
def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict:
"""账本 vs 下游持仓, 以下游为准修正并留痕。事实源见 positions_source()。
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
# 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、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
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 加; 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 chk.get("estimated"):
# 少数几只没成本价 —— 走 build_recon_fixes 的逐只兜底 (设计定好的行为), 不拦。
# 但必须吼一声: 估出来的成本**改不回来** (数量对得上就没有对账差异, 后续对账
# 不会再碰它), 它是一条已知的坏账, 得让人知道是哪几只。
logger.error("[对账] %s 只没有下游成本价, 将拿现价建账且**此后不会被对账修正**: %s"
"能等的话请对端补上 cost_price 再建",
len(chk["estimated"]), chk["estimated"])
if not out["coverage"]["enough"]:
logger.warning("[对账] 建账数据的情形覆盖不全: %s", out["coverage"]["hint"])
# 下游的成本价 —— 补仓位时的开仓价优先取它, 现价只兜底 (见 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}
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
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}")
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