1181 lines
65 KiB
Python
1181 lines
65 KiB
Python
# -*- 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 写入: 直接写库不会失效缓存, 连续天数会一直读到旧值。
|
||
# **返回值必须接**: set_param 写不进去时只回 {"ok": False, "error": ...} 而不抛
|
||
# 异常 —— 2026-07-31 就是因为把它丢了, STREAK_YMD 键不在白名单里被静默拒写,
|
||
# prev_ymd 永远读回 0, 按日推进的修复形同虚设而单测全绿 (单测只测纯函数)。
|
||
for k, v in ((STREAK_KEY, st["streak"]), (STREAK_YMD_KEY, st["ymd"])):
|
||
w = param_store.set_param(k, v, "system")
|
||
if not (w or {}).get("ok"):
|
||
logger.error("[对账] 连续天数写入失败 %s=%s: %s —— "
|
||
"按日推进将失效, 计数会退化成按次累加", k, v,
|
||
(w or {}).get("error"))
|
||
out.setdefault("warnings", []).append(
|
||
f"连续天数未能落库 ({k}): {(w or {}).get('error')}")
|
||
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()
|
||
# 「该踩没踩上」必须冒到盘前准备的 errors 里 —— 它不是一条 warning, 是风控没生效
|
||
for w in (out["brake"].get("warnings") or []):
|
||
(out["errors"] if "刹车未生效" in w else out.setdefault("warnings", [])).append(w)
|
||
except Exception as e:
|
||
out["errors"].append(f"刹车结算失败: {e}")
|
||
out["ok"] = not out["errors"]
|
||
return out
|
||
|
||
|
||
def _settle_brake() -> dict:
|
||
"""组合刹车: 自高水位回撤 ≥ 阈值 → 自主增持停 N 个交易日 (命令类不受限)。
|
||
|
||
两处 set_param 的返回值都必须接 (2026-07-31 静默失败专项):
|
||
* PMS_BRAKE_UNTIL 写不进去 = 刹车根本没踩。proposal_service 每一跳重新去
|
||
param_store 读这个键判 brake_active, 读不到就是 False, 自主增持照跑。而本函数
|
||
原来照样 return {"brake_until": until, "active": True} 并打日志"暂停至 X" ——
|
||
日报、页面、日志三处都说停了, 实际一路买穿整个回撤。
|
||
* PMS_HIGH_WATER 写不进去 = 高水位永远停在旧值, 回撤按陈年高点算, 阈值再也够不着。
|
||
风控是慢性失效, 没有任何一处会报错。
|
||
"""
|
||
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()
|
||
warnings = []
|
||
if mv > hw:
|
||
w = param_store.set_param("PMS_HIGH_WATER", mv, "system") or {}
|
||
if w.get("ok"):
|
||
hw = mv
|
||
else:
|
||
logger.error("[刹车] 高水位写入失败 (仍按旧值 %s 算回撤): %s", hw, w.get("error"))
|
||
warnings.append(f"高水位未能落库, 回撤按旧高点 {hw:,.0f} 计算: {w.get('error')}")
|
||
drawdown = (1 - mv / hw) if hw > 0 else 0.0
|
||
engaged = False
|
||
if hw > 0 and drawdown >= dd_limit and today >= until:
|
||
nxt = td.ymd(td.next_trade_day(datetime.now().date(), days))
|
||
w = param_store.set_param("PMS_BRAKE_UNTIL", nxt, "system") or {}
|
||
if w.get("ok"):
|
||
until, engaged = nxt, True
|
||
logger.warning("[刹车] 自高水位回撤 %.1f%% ≥ %.0f%%, 自主增持暂停至 %s",
|
||
drawdown * 100, dd_limit * 100, until)
|
||
else:
|
||
logger.error("[刹车] **该踩刹车但没踩上** —— 回撤 %.1f%% ≥ %.0f%%, "
|
||
"PMS_BRAKE_UNTIL 写入失败: %s。自主增持仍在放行",
|
||
drawdown * 100, dd_limit * 100, w.get("error"))
|
||
warnings.append(f"**刹车未生效**: 回撤 {drawdown:.1%} 已达阈值但 "
|
||
f"PMS_BRAKE_UNTIL 写不进去 ({w.get('error')}), 自主增持仍在放行")
|
||
out = {"high_water": hw, "portfolio_mv": mv, "drawdown": round(drawdown, 4),
|
||
"brake_until": until, "active": today < until, "engaged": engaged}
|
||
if warnings:
|
||
out["warnings"] = warnings
|
||
return out
|
||
|
||
|
||
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:
|
||
r = out["steps"]["recon"] = reconcile()
|
||
# reconcile 的"我拒绝对账"是**返回值**不是异常 (爆炸半径拦截 / 两源无应答 /
|
||
# 成本价体检不过 → ok=False 或 severity=ERROR)。原来这里只 catch 异常, 于是
|
||
# 一整天没对上账的日子, daily_settle 照样 ok=True 回给调度器和运维页。
|
||
if not r.get("ok"):
|
||
out["errors"].append("对账未完成: " + "; ".join(
|
||
str(x) for x in (r.get("errors") or [r.get("blocked") or "原因见 recon 步骤"])))
|
||
elif r.get("severity") == "ERROR":
|
||
out["errors"].append(
|
||
f"对账连续不一致已达 ERROR (连续 {r.get('streak')} 日), 需人工介入")
|
||
for w in (r.get("warnings") or []):
|
||
out.setdefault("warnings", []).append(f"对账: {w}")
|
||
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"]}
|
||
w = portfolio.save_neg_streak({k: v2 for k, v2 in streak.items() if k in held}) or {}
|
||
out = {"updated": updated,
|
||
"neg_streak": {k: v2 for k, v2 in streak.items() if v2 > 0 and k in held}}
|
||
if not w.get("ok"):
|
||
# 写不进去要让日终结算整体报失败, 不能只在本函数里留个字段等人去翻
|
||
raise RuntimeError(f"安全垫连负天数未能落库 ({w.get('error')}) —— "
|
||
f"降仓命令的「清弱票」判据会停在旧值")
|
||
return out
|
||
|
||
|
||
# ================================================================ 日报
|
||
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
|