# -*- 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 recon as rc from app.core import tradedays as td 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" 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 的注释。 """ 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, 分不出联调还是真单; # 出口表是本端写的, 骗不了自己。为什么必须挡, 见 qmt_repo.SMOKE_PREFIX 的注释。 # 标 processed=2 (已消化) 而不是 1 (已入账): 两者语义不同, 事后翻 inbox 一眼能看出 # 这笔是被有意跳过的, 而不是入账入丢了。 smoke, real = [], [] for r in trades: iid = (r.get("payload") or {}).get("instruction_id") or "" try: od = qmt_repo.get_order(iid) if iid else None except Exception as e: # 查不到就按真单走 —— 宁可多入账也不能漏账 out["errors"].append(f"出口表反查 {iid} 失败: {type(e).__name__}: {e}") od = None (smoke if od and qmt_repo.is_smoke(od.get("parent_instruction_id")) else real).append(r) if smoke: seqs = [r["seq"] for r in smoke] qmt_repo.inbox_mark(seqs, processed=2, note="ws 联调测试单成交, 不入账 (SMOKE_)") out["smoke_skipped"] = len(smoke) logger.warning("[联调] 跳过 %s 笔测试单成交, 不入账: seq=%s", len(smoke), seqs) trades = real if not trades: return out out["trades"] = len(trades) # 取每条成交所属父指令的 action, 用来定买入的批次类型 parents = {rc.parent_instruction_id((r.get("payload") or {}).get("instruction_id") or "") for r in trades} parent_actions = {} for pid in parents: if not pid: continue try: ins = pms_repo.get_instruction(pid) if ins: parent_actions[pid] = ins.get("action") except Exception as e: out["errors"].append(f"读指令 {pid} 失败: {type(e).__name__}: {e}") mapped = rc.map_trades_to_book(trades, parent_actions=parent_actions) done = [] for act in mapped["actions"]: try: _apply_action(act) out["actions"] += 1 except Exception as e: logger.exception("ws 成交入账失败 %s", act) out["errors"].append(f"{act.get('ts_code')} 入账失败: {type(e).__name__}: {e}") ymd = td.ymd() for f in mapped["fees"]: try: pms_repo.insert_cash_flow(ymd=ymd, kind=CF_FEE, amount=f["amount"], ts_code=f["ts_code"], estimated=f["estimated"], trade_no=f["trade_no"], instruction_id=f["instruction_id"], note=f["note"]) out["fees"] += 1 except Exception as e: out["errors"].append(f"费用入账失败 {f.get('trade_no')}: {type(e).__name__}: {e}") for code in {a["ts_code"] for a in mapped["actions"]}: try: recompute_position(code) except Exception as e: out["errors"].append(f"{code} 成本重算失败: {e}") # 只有全程无错才标已入账 —— 有错就留在 processed=0, 下一跳重试。 # 重试是安全的: _apply_action 幂等由 trade_no 兜着 (inbox 那层已按 trade_no 去过重)。 if not out["errors"]: done = [s for s in mapped["seqs"] if s is not None] if done: try: qmt_repo.inbox_mark(done, processed=1, note="ws 成交已入账") except Exception as e: out["errors"].append(f"inbox 标记失败: {type(e).__name__}: {e}") out["alerts"] = mapped["alerts"] for a in out["alerts"]: logger.warning("[ws 入账告警] %s", a.get("message")) out["ok"] = not out["errors"] return out def calibrate_fees(*, ymd=None, actual_fee=None) -> dict: """日终用资金快照反推当日真实费用, 写一条 CALIBRATE 平掉估算误差 (协议 §5.5)。 对端明说逐笔 fee 是按费率估的。协议给的校准式子是 当日实际费用 = 总资产变动 − 成交净额 资金快照来自 ws 的 funds_update / snapshot(kind=funds), 落在 pms_qmt_inbox。**影子期 没有这条数据**, 那就只登记一句"待校准", 不硬凑 —— 估算值本来就只影响现金账, 不影响 成本与安全垫, 晚校准几天没有任何风险。这也是当初把费用挡在成本之外的意义。 actual_fee 可显式传入 (人工按对账单校准时用), 传了就不去读快照。 """ ymd = int(ymd or td.ymd()) out = {"ymd": ymd, "estimated": 0.0, "actual": None, "adjusted": 0.0, "note": ""} try: out["estimated"] = round(pms_repo.sum_cash_flow(ymd, CF_FEE), 2) except Exception as e: out["note"] = f"读当日费用流水失败: {type(e).__name__}: {e}" return out if actual_fee is None: out["note"] = ("当日估算费用已入现金账; 资金快照未接通 (ws 通道未启用), " "暂不校准 —— 费用不进成本, 晚校准无风险") return out actual = -abs(float(actual_fee)) diff = round(actual - out["estimated"], 2) out["actual"] = actual if abs(diff) < 0.01: out["note"] = "估算与实际一致, 无需校准" return out pms_repo.insert_cash_flow( ymd=ymd, kind=CF_CALIBRATE, amount=diff, estimated=0, trade_no=None, note=f"日终校准: 估算 {out['estimated']:.2f} → 实际 {actual:.2f}, 差额 {diff:+.2f}") out["adjusted"] = diff out["note"] = f"已按资金快照校准, 差额 {diff:+.2f} 元 (只调现金账, 不回溯改成本)" logger.info("[费用校准] %s", out["note"]) return out def replay_fills(*, limit: int = 500) -> dict: """成交回放一跳 = 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 def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict: """账本 vs 下游持仓, 以下游为准修正并留痕。 force=True 绕过爆炸半径限制 (见 _recon_blast_guard) —— 只在人工确认下游读数确实 正确之后使用。 """ out = {"ok": True, "diffs": [], "fixes": [], "errors": [], "columns": {}, "severity": rc.SEV_OK} try: ds = downstream_repo.fetch_positions() except Exception as e: out.update({"ok": False, "errors": [f"读 trading_position 失败: {type(e).__name__}: {e}"]}) return out out["columns"] = ds["columns"] if ds["rows"] and ds["columns"].get("qty") is None: out.update({"ok": False, "errors": [ "下游持仓表未识别出数量列 —— 请按 QMT_INTERFACE_REQUIREMENTS A1/D1 取得 DDL 后, " "把列名补进 downstream_repo.QTY_CANDIDATES"]}) 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 streak = param_store.get_int(STREAK_KEY, 0) streak = streak + 1 if diffs else 0 # 必须走 ParamStore 写入: 直接写库不会失效缓存, 会导致连续天数一直读到旧值 param_store.set_param(STREAK_KEY, streak, "system") 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) # 下游的成本价 —— 补仓位时的开仓价优先取它, 现价只兜底 (见 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 _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 未配置"}) 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