tradingSystem/app/services/proposal_service.py

934 lines
54 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 -*-
"""
自主提议: 扫描 → 规则闸 → 研判闸 → 按自主档位分流 (设计 §6 / §7)
==================================================================
分流规则 (设计原文):
full 闸门与研判通过即执行 → 直接落指令
propose_only 增持类全部待用户确认 (一期默认) → 落 pms_proposal 队列
off 不扫描
两条无条件覆盖档位的规矩:
* **规则算出来的减持不设确认门槛** —— TRIM 保垫减仓属纯规则自动执行, 任何档位都直接落
指令。但这一条覆盖的是**档位**, 不覆盖「强制入人工队列」这个标记 (2026-09-03 修): 带了
这个标记的减持一律交人, 方向是卖也不例外。哪种减持带标记由候选的来源说了算, 见
`action_engine.source_confirm_why` 与 `_route_one` 里那段说明。
* **15% 及更深的补仓永远需用户确认** —— 即便档位是 full, 也强制入队。
研判闸不可用时 (决策系统未接通/超时), 按设计自动降级为 propose_only + ERROR 告警,
**绝不把「研判拿不到」当成「研判通过」**。
2026-08-06 加了第五类动作「新建仓 OPEN」: 除了管已有持仓, 也从上游候选池里挑新票建底仓,
走的是同一条 规则闸 → 研判闸 → 档位分流。三处与已有四类不同的地方, 都在本文件里:
* 档位走自己的 `PMS_OPEN_AUTONOMY` (默认 full), 不跟随 `PMS_AUTONOMY` —— 「新建仓要不要
人点头」和「加仓要不要人点头」是两个决定。但 `PMS_AUTONOMY=off` 仍然是总闸, 它一关,
连新建仓一起停: off 的意思就是「自主动作全部停下」, 不该有子开关能绕过它。
* 研判驳回对新建仓做当日去重 (加仓类**有意**不去重, 那条纪律不动), 理由见
`pms_repo.judge_rejected_today`。
* 单轮扫描给研判一个时间预算, 用尽的候选留到下一跳 —— 见 `_judge_budget_left` 的说明。
"""
from __future__ import annotations
import logging
import time
from datetime import datetime, timedelta
from app.core import action_engine as ae
from app.core import copy
from app.core import command_spec as cs
from app.core import rule_gate
from app.core import signal_rules as sr
from app.core import tradedays as td
from app.repo import pms_repo
from app.services import (command_service, executor, industry, judge, market, param_store,
plan_feed, portfolio)
logger = logging.getLogger("pms.proposal")
AUTONOMY_FULL, AUTONOMY_PROPOSE, AUTONOMY_OFF = "full", "propose_only", "off"
def _macro_gate() -> dict:
"""宏观偏热闸状态 (macro_service 维护, 这里只读; MACRO_TIMING_PLAN.md §5.2)。
读不到一律按不生效 —— 闸的安全方向是「不额外拦」: 买入本身另有规则闸把关,
不能让宏观层故障把整个自主引擎摁死。"""
try:
from app.services import macro_service
return macro_service.gate_state() or {}
except Exception as e: # noqa: BLE001
logger.warning("[提议] 读宏观闸状态失败 (按不生效): %s", e)
return {}
def scan_and_route(*, now=None, dry_run: bool = False) -> dict:
"""自主提议扫描一轮。dry_run=True 只出候选与判定, 不落任何表。"""
now = now or datetime.now()
out = {"ok": True, "autonomy": None, "open_autonomy": None, "candidates": 0,
"open_candidates": 0, "executed": [], "queued": [],
"rejected": [], "skipped": [], "errors": [], "degraded": False, "dry_run": dry_run}
autonomy = param_store.get("PMS_AUTONOMY", AUTONOMY_PROPOSE)
out["autonomy"] = autonomy
if autonomy == AUTONOMY_OFF:
out["skipped"].append({"why": "自主档位 off, 不扫描 (新建仓一并停 —— off 是总闸)"})
return out
if param_store.get_bool("PMS_GLOBAL_EXEC_HALT", False):
out["skipped"].append({"why": "全局暂停执行 (休假模式)"})
return out
open_autonomy = param_store.get("PMS_OPEN_AUTONOMY", AUTONOMY_FULL)
out["open_autonomy"] = open_autonomy
try:
view = portfolio.positions_view()
params = _scan_params(view)
# 个股参数命令的当前值 (黑名单、目标价、止损价…)。取数提到扫描之前, 因为动作引擎
# 现在要用它评「目标价到价」那条动作; 后面分流与规则闸用的是同一份, 不重复查。
stock_params = command_service.effective_stock_params()
mkt = _market_ctx(view["held"], now)
params["_mkt"] = mkt # 规则闸要用同一份 MA5, 不再重取
# 跳过三类: ①已有在途提议或指令的 ②今天已被规则闸拒过的 ③今天已被研判闸驳回的新建仓。
# ② 是 2026-07-29 的教训 —— 组合已超总仓上限时, 16 只深亏票的补仓候选每分钟被拒
# 一次, 一天往评审账本灌几千行一模一样的记录。闸门结论当天基本不会变, 记一次就够。
# ③ 只挡新建仓: 加仓类对研判驳回**有意**不去重 (契约里那句「研判结论会变」), 那条
# 纪律不动; 而候选池是几十只的量级, 不挡的话 ② 那个洞会原样从研判闸重来一遍。
#
# 三份合并成 `{(代码, 动作): 原因}`, **在途的排在最后, 覆盖被拒的** (2026-08-06):
# 这四种处境完全不同 —— 在途的明天照样被挡, 被拒的日切就重新评估。从前它们在
# skipped 里长得一模一样, 看的人判断不了这只票明天还会不会再被评估。
# 一只票既有在途指令又今天被拒过时, 显示"有在途"更贴近它此刻的实际状态。
skip = {**_judge_rejected_open_keys(), **_rejected_today_keys(),
**_declined_today_keys(), **_inflight_keys()}
# 挂了 ACTIVE 交易方案(策略)的票交策略层接管, 动作引擎不再对它自动提议
# (读库失败按空集 —— 宁可这轮不排除, 也不能因读不到把全体持仓都排除)
try:
strategy_codes = pms_repo.active_strategy_codes()
except Exception:
strategy_codes = set()
scanned = ae.scan(positions=view["held"], params=params, market=mkt, skip=skip,
strategy_codes=strategy_codes, stock_params=stock_params)
except Exception as e:
logger.exception("提议扫描失败")
return {**out, "ok": False, "errors": [f"扫描失败: {type(e).__name__}: {e}"]}
# 宏观偏热闸: 只拦已有持仓四类动作里的**买入侧** (FILL/ADD/DCA), TRIM 保垫减仓照常。
# 落点选在扫描产出之后、分流之前: 不碰 action_engine, 不产生指令自然过不了后面任何一层。
mg = _macro_gate()
if mg.get("active"):
kept = []
for c in scanned["candidates"]:
if c.get("side") == "buy":
out["skipped"].append({"ts_code": c.get("ts_code"), "action": c.get("action"),
"why": f"宏观偏热闸生效 ({mg.get('why') or '股汇对冲指数偏热'}), "
f"自主增持暂停"})
else:
kept.append(c)
scanned["candidates"] = kept
out["candidates"] = len(scanned["candidates"])
out["skipped"].extend(scanned["skipped"])
brake_active = td.ymd() < param_store.get_int("PMS_BRAKE_UNTIL", 0)
# ---- 新建仓: 候选池那一路 (取不到候选池不影响上面四类, 反之亦然) ----
open_cands = []
if open_autonomy == AUTONOMY_OFF:
out["skipped"].append({"action": ae.A_OPEN, "why": "新建仓档位 off, 不扫描候选池"})
else:
try:
open_cands = _scan_open(view, params, stock_params, skip, mkt, out)
except Exception as e:
logger.exception("新建仓扫描失败")
out["errors"].append(f"新建仓扫描失败: {type(e).__name__}: {e}")
out["open_candidates"] = len(open_cands)
# 研判的时间预算从这一刻起算。**先跑已有持仓的四类, 再跑新建仓** —— 预算真用尽时,
# 被推到下一跳的一定是新建仓, 已有持仓的动作行为与 2026-08-06 之前一致。
deadline = time.monotonic() + max(
0, param_store.get_int("PMS_JUDGE_TICK_BUDGET_SEC", 150))
for c in list(scanned["candidates"]) + open_cands:
try:
_route_one(c, view, params, stock_params, brake_active, now, dry_run, out,
deadline=deadline)
except Exception as e:
logger.exception("提议分流失败 %s", c.get("ts_code"))
out["errors"].append(f"{c.get('ts_code')} {c.get('action')}: "
f"{type(e).__name__}: {e}")
out["ok"] = not out["errors"]
return out
# ================================================================ 新建仓的取数与筛选
def _scan_open(view, params, stock_params, skip, mkt, out) -> list:
"""候选池 → 新建仓候选。产出的候选数天生不超过剩余名额 (名额与金额在纯逻辑里边走边扣)。
与命令驱动那条路 (command_service._candidates) 有意不同的两点:
1. **只认上游计划接口这一个来源。** buy_plan 那张旧表在目标架构下没有明确的写入方,
白名单是给人下命令用的 —— 无人值守地建仓, 事实源必须单一。
2. **价格只认实时价。** 那边用 market.plan_price, 拿不到实时价会回落昨收; 它自己的
注释就写着「拿昨收当现价去做不追高这类判断会出错, 只给规划期定量用」。
这条路是要真下单的, 所以走 market.get_prices, 取不到就整只跳过并留痕。
"""
# 宏观偏热闸: 新建仓全是买入, 整轮直接不扫 (MACRO_TIMING_PLAN.md §5.2)。
# 放在函数入口: disposition_snapshot 复用本函数, 页面"系统在盯的候选"会如实显示闸生效。
mg = _macro_gate()
if mg.get("active"):
out["skipped"].append({"action": ae.A_OPEN,
"why": f"宏观偏热闸生效 ({mg.get('why') or '股汇对冲指数偏热'}), "
f"本轮不自动新建仓"})
return []
src = (param_store.get("PMS_CANDIDATE_SOURCE", plan_feed.SRC_PLAN_API)
or plan_feed.SRC_PLAN_API).strip()
if src == plan_feed.SRC_BUY_PLAN:
out["skipped"].append({"action": ae.A_OPEN,
"why": f"候选池来源是 {src}, 自主新建仓只认上游计划接口"})
return []
black = {c for c, d in (stock_params or {}).items() if d.get("black")}
held = [x["ts_code"] for x in view["held"]]
try:
# 按判决分流开着时, 上游判为「仅展示」的票在候选阶段就剔掉 (理由见
# plan_feed.select_candidates 的说明: 不剔的话它们会白占 top_n 的名额)。
sel = plan_feed.candidates(held=held, black=black,
route_by_verdict=bool(params.get("open_route_by_verdict")))
except plan_feed.PlanFeedError as e:
# 拿不到 ≠ 今天没票可买。显式留痕, 本轮不产新建仓候选, 已有持仓的四类照常。
logger.error("[新建仓] 候选池取数失败, 本轮不建仓: %s", e)
out["skipped"].append({"action": ae.A_OPEN, "why": f"候选池取不到, 本轮不建仓: {e}"})
return []
# 仅展示的票逐只写跳过原因: disposition_snapshot 复用本函数, 页面「系统在盯的候选」
# 由此显示「选股系统判为仅展示」而不是一个说不清的空白。
for row in (sel.get("display_only") or []):
# 2026-09-04 起这里是带判决依据的行; 早先只是一串代码, 兼容着读。
row = row if isinstance(row, dict) else {"ts_code": row}
out["skipped"].append({"ts_code": row.get("ts_code"), "action": ae.A_OPEN,
"why": ae.display_only_why(row)})
items = list(sel.get("items") or [])
if not items:
out["skipped"].append({"action": ae.A_OPEN, "why": _no_candidate_why(sel)})
return []
codes = [x["ts_code"] for x in items]
prices = market.get_prices(codes)
# 行业名一次批量取, 本轮复用。滚动扣减要算行业集中度, 所以必须先有它 ——
# caps_ctx 对**从没持仓过**的票带不出行业名, 不显式传的话 sizer.check_caps 遇到
# sector 为空会整段跳过, 行业集中度那道硬拦截就静默失效了。
# (industry.get_many 走 gp_stock_category 时是逐只查库、没有缓存; 给它加按日缓存能
# 省掉这几十次往返, 但那会改到既有函数的时序行为, 单独提、单独拍板, 这次不夹带。)
sectors = industry.get_many(codes) if codes else {}
# 今天被决策系统盘中判过转多的票 (signal_service 落的留痕)。**只用来排序, 不改资格**:
# 不在候选池里的票不会因为有信号就被建仓, 候选层那一整套过滤一道都不绕。
# 它的增量是实打实的 —— 候选按分数降序取, 名额只剩两个时第 25 名永远轮不上;
# 而「此刻转多」是盘中才有的新信息, 昨夜算出来的分数与买入区间都表达不了它。
sig_buy = _buy_signals_today()
cands = [{**x, "price": prices.get(x["ts_code"]), "sector": sectors.get(x["ts_code"]),
"sig_buy": sig_buy.get(x["ts_code"])} for x in items]
hit = [c["ts_code"] for c in cands if c.get("sig_buy")]
if hit:
logger.info("[新建仓] 候选池里今天被决策系统判过转多的: %s", hit)
t = view["totals"]
p = view["params"]
# 名额与金额都要扣掉**在途的建仓承诺** (2026-08-28 审查修): names_count 只数已成交
# 持仓, 在途 OPEN 指令与等拍板的 OPEN 提议两头都不占 —— 候选没进买入区间时指令滞留
# 在途, 下一跳名额照旧又放两只, 几分钟就把持仓数上限击穿。同票去重挡不住换一只票。
inflight = _inflight_open_commitments(view)
slots = int(p["max_names"] or 0) - int(t["names_count"] or 0) - len(inflight["names"])
room = float(p["portfolio_cap"] or 0) * float(p["scale"] or 0) \
- float(t["portfolio_mv"] or 0) - inflight["amount"]
if inflight["names"]:
out["skipped"].append({"action": ae.A_OPEN,
"why": f"在途建仓已占 {len(inflight['names'])} 个名额、"
f"{inflight['amount']:,.0f}"
f"({sorted(inflight['names'])[:6]}), 本轮按余量扫描"})
# 真实可用资金封顶。**这一条比规则闸严一档, 是刻意的**: 规则闸那条「拿不到 ws 资金快照
# 只告警不拦」是为**已经排好的命令**设计的 —— 通道故障不该升级成业务停摆。而无人值守地
# 从零建仓完全可以等一等, 07-30 那次 (scale 200 万 / 账户实际 98 万, 方案一路放行到
# 下游才被拒, 而拒了不自动重发) 不该换个入口重演。
if str(t.get("cash_source") or "") == portfolio.CASH_WS and t.get("cash_avail") is not None:
room = min(room, float(t["cash_avail"]))
elif param_store.get_bool("PMS_OPEN_REQUIRE_WS_CASH", True):
out["skipped"].append({"action": ae.A_OPEN,
"why": f"读不到券商账户的可用资金,这一轮不开新仓;"
f"已持有的票照常处理。"
f"原因:{t.get('cash_why') or '原因未知'}"})
return []
res = ae.scan_open(candidates=cands, params=params, caps=portfolio.caps_ctx(view),
room_amt=room, slots=slots, skip=skip)
out["skipped"].extend(res["skipped"])
# 当日行情快照只对**真的产出了候选**的票取 —— 候选数已被名额与金额扣到很小,
# 不必为整池几十只票各拉一遍分钟线。规则闸的「不追高」靠的就是这一份。
for c in res["candidates"]:
code = c["ts_code"]
d = {"ma5": market.get_ma5(code)}
try:
d["day"] = market.day_snapshot(code) or {}
except Exception as e:
logger.warning("[新建仓] 取当日快照失败 %s: %s", code, e)
d["day"] = {}
mkt[code] = d
return res["candidates"]
# 候选一只都不剩时, 这一栏要说清是被哪几道筛掉的。键名是代码里的, 翻成人话再显示。
# 顺序就是筛选实际发生的顺序, 照着念下来正好是一条票走过的路。
_DROP_CN = [
("held", "已经持有"),
("black", "在黑名单里"),
("display", "选股系统判定今天不买"),
("st", "是 ST 类"),
("tier", "档位不够"),
("score", "分数不够"),
("sources", "传导源数不够"),
("upside", "券商预期空间不够"),
("theme", "同一主题已经选够"),
("dup", "和已选的重复"),
("capped", "超出今天的名额"),
]
def _no_candidate_why(sel) -> str:
"""今天一只候选都没有时的那句说明。
2026-09-04 改过一次。改之前是把落选计数那个字典直接拼进字符串, 页面上显示成
「落选明细 {'held': 0, 'black': 0, ... 'display': 83}」—— 键全是英文, 值为零的
也全列出来, 读的人要在十一个数字里自己找出哪个不是零。
现在只念不为零的那几项, 按筛选实际发生的顺序, 用中文写成一句话。
"""
n = sel.get("considered")
dropped = sel.get("dropped") or {}
parts = [f"{cn} {dropped[k]}" for k, cn in _DROP_CN
if isinstance(dropped.get(k), int) and dropped[k] > 0]
head = f"今天考察了 {n} 只票,一只都没留下" if n else "今天没有票可考察"
if not parts:
return head + ""
return head + "" + "".join(parts) + ""
def _inflight_open_commitments(view) -> dict:
"""在途的新建仓承诺: {names: {代码}, amount: 预估占用金额}。
口径: 在途 (LIVE) 的 OPEN 指令 + 等拍板 (WAIT_USER) 的 OPEN 提议, 只算**尚未持有**
的票 (已成交的部分 names_count 里已经有了)。金额按 数量×现价 估 (指令下发前没有限价),
现价取不到时按提议里存的价兜底。读库失败按空 —— 与其它在途读取同一口径, 宁可这轮
多放也不能因读不到把新建仓全停 (规则闸和真实资金封顶仍在后面兜)。"""
names, amount = set(), 0.0
held = {x["ts_code"] for x in view.get("held") or []}
rows = []
try:
for i in pms_repo.list_instructions(statuses=list(executor.LIVE), limit=300):
if i.get("action") == ae.A_OPEN and i["ts_code"] not in held:
rows.append((i["ts_code"], int(i.get("qty") or 0), None))
except Exception as e:
logger.warning("读在途 OPEN 指令失败 (按无在途): %s", e)
try:
for pr in pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200):
if pr.get("action") == ae.A_OPEN and pr["ts_code"] not in held:
hn = pr.get("hard_numbers") or {}
rows.append((pr["ts_code"], int(pr.get("qty") or 0),
float(hn.get("price") or 0) or None))
except Exception as e:
logger.warning("读在途 OPEN 提议失败 (按无在途): %s", e)
if not rows:
return {"names": names, "amount": 0.0}
try:
prices = market.get_prices(sorted({c for c, _, _ in rows}))
except Exception:
prices = {}
for code, qty, px_hint in rows:
if code in names:
continue
names.add(code)
px = float(prices.get(code) or 0) or float(px_hint or 0)
amount += qty * px
return {"names": names, "amount": round(amount, 2)}
def _plain_check(c) -> str:
"""规则闸的 failed 串形如 "T1_UNAVAILABLE: 卖出…" —— 去掉给运维看的大写代码前缀, 只留中文说明。
实现搬到 app/core/copy.py 了: 同一套口径后端两处 + 前端一份共用, 不再各写各的。
这里留一层薄壳, 是因为本文件里已经有多处按名字调它。
"""
return copy.strip_code(c)
def _today_open_reject_reasons(now):
"""今天各票**新建仓**被驳回的实因, 供候选处置显示真正的「为什么」:
判 = 决策系统给的原话 (评审账本 arbiter=judge 那行的 reason);
合规 = 规则闸具体未过的检查 (failed_checks, 去掉大写代码前缀)。
只读评审账本、按票取当天最新一条 (list_ledger 按 id 倒序); 读不到就返回空, 调用方退回概述。"""
judge, rule = {}, {}
today = now.strftime("%Y-%m-%d")
try:
for r in pms_repo.list_ledger(limit=500):
if r.get("verdict") != "REJECT" or r.get("action") != ae.A_OPEN:
continue
if str(r.get("decided_at"))[:10] != today:
continue
code = r.get("ts_code")
if not code:
continue
if r.get("arbiter") == "judge" and code not in judge:
judge[code] = (r.get("reason") or "").strip()
elif r.get("arbiter") == "rule" and code not in rule:
fc = r.get("failed_checks") or []
rule[code] = "".join(_plain_check(x) for x in fc) if fc else (r.get("reason") or "").strip()
except Exception:
logger.exception("读当日新建仓驳回实因失败 (退回概述)")
return judge, rule
def disposition_snapshot(now=None) -> dict:
"""候选池处置快照 (只读, 无副作用): 候选池里每只票今天为什么下单 / 没下单。
页面「系统在盯的候选」原来只铺上游榜单, 看不出每只票的去向。这里补上去向与原因:
名额满 / 可投金额不足 / 超上限 / 不足一手 / 在途 / 已持有 / 今日已拒 —— 这些原因每分钟的
scan_and_route 都算了, 但算完即弃 (评审账本只留规则闸/研判闸的驳回, 这一层查不到)。
实现上**直接复用** `_scan_open`: 把它写 out / mkt 的两个收集器换成丢弃用的临时容器,
一个字都不改扫描逻辑, 口径与实盘扫描逐字一致 (不另写一份, 免得日后分叉)。_scan_open 内部
只读库与行情、不落任何表, 所以这样调用是安全的。
「已下单 / 已生成提议」交给前端用它已在手的 instructions / proposals 按股票对齐, 这里不重复查。
"""
now = now or datetime.now()
res = {"ok": True, "ready": True, "at": now.isoformat(timespec="seconds"),
"by_code": {}, "notes": [],
"autonomy": param_store.get("PMS_AUTONOMY", AUTONOMY_PROPOSE),
"open_autonomy": param_store.get("PMS_OPEN_AUTONOMY", AUTONOMY_FULL)}
if res["autonomy"] == AUTONOMY_OFF:
res["notes"].append("自主档位 off —— 本轮不扫描候选池 (off 是总闸)")
return res
if param_store.get_bool("PMS_GLOBAL_EXEC_HALT", False):
res["notes"].append("全局暂停执行 (休假模式) —— 本轮不扫描候选池")
return res
if res["open_autonomy"] == AUTONOMY_OFF:
res["notes"].append("新建仓档位 off —— 不扫描候选池")
return res
try:
view = portfolio.positions_view()
params = _scan_params(view)
stock_params = command_service.effective_stock_params()
# 给交易员的白话原因: 用同一套去重键 (口径与实盘扫描一致), 但值不含内部闸门术语、也不指引
# 任何终端命令。能取到实因就带上实因 (决策系统原话 / 合规具体未过项), 取不到才退回概述。
# 落库顺序同实盘 (在途最后, 覆盖被拒)。
_jr, _rr = _today_open_reject_reasons(now)
skip = {}
for _k in _judge_rejected_open_keys():
_w = _jr.get(_k[0])
skip[_k] = ("决策系统判暂不建仓:" + _w) if _w else "决策系统今天判过暂不建仓,明天开盘会重新评估"
for _k in _rejected_today_keys():
_w = _rr.get(_k[0])
skip[_k] = ("未通过合规检查:" + _w) if _w else "今天没通过合规检查,明天开盘会重新评估"
for _k in _inflight_keys(): # 在途; 前端多半已按指令/提议标成「已下单/待你确认」
skip[_k] = "系统已经在处理这只票(见『等我拍板』或『今日在办』)"
sink = {"skipped": []} # 丢弃用: _scan_open 只往里 append/extend, 不读它
would = _scan_open(view, params, stock_params, skip, {}, sink)
except Exception as e:
logger.exception("候选处置快照失败")
return {"ok": False, "ready": False, "error": f"{type(e).__name__}: {e}",
"by_code": {}, "notes": []}
for c in would:
res["by_code"][c["ts_code"]] = {
"disp": "would", "why": "名额与资金都够, 本轮将建底仓 (下一跳落单)"}
for s in sink["skipped"]:
code = s.get("ts_code")
if not code or code == "*": # 池级原因 (名额满/资金为零/候选池空…) 挂到 notes
if s.get("why"):
res["notes"].append(s["why"])
else: # 单票原因; 不覆盖已判 would 的
res["by_code"].setdefault(
code, {"disp": "deny", "why": s.get("why") or "未产出候选"})
return res
def _route_one(c, view, params, stock_params, brake_active, now, dry_run, out,
*, deadline=None):
code, action, side = c["ts_code"], c["action"], c["side"]
is_open = (action == ae.A_OPEN)
pos = _pos_of(view, code)
# 新建仓的现价取候选自带的那份。**不能取 pos 的**: 从没持仓过的票在账本里没有行,
# `_pos_of` 回的是空壳、没有 price, 取它会是 0, 规则闸一上来就判 PRICE_MISSING。
price = float(c.get("price") or 0) if is_open else float(pos.get("price") or 0)
# ---- 一级: 规则闸 (自主动作受刹车约束, is_command=False) ----
gate = rule_gate.check(
side=side, action=action, qty=c["qty"], price=price,
ctx={"ts_code": code, "position": pos,
# 当日行情走 _market_ctx 取到的真实快照。原来这里是
# {"vwap": price, "day_chg_from_open": None}
# 两个都是编的: vwap 拿现价顶等于"现价恰好等于均价", day_chg 给 None 让
# 「不追高」整道闸静默跳过。取不到就留空, 让规则闸记 *_MISSING 警告。
"day": {**((params.get("_mkt") or {}).get(code, {}).get("day") or {}),
"price": price,
"ma5": (params.get("_mkt") or {}).get(code, {}).get("ma5")},
"params": {"no_chase_ma5": params.get("no_chase_ma5"),
"buy_halt_dayup": params.get("buy_halt_dayup"),
"sector_source_ready": view["sector_ready"]},
# 新建仓必须**显式**把 is_new_name 与行业名传进去。caps_ctx 只在
# view["positions"] 里找得到该票时才带得出行业, 而新票根本不在里面 ——
# 不传的话 sizer.check_caps 遇到 sector 为空会整段跳过, 行业集中度那道硬拦截
# 就静默失效了 (不报错、不少数据, 就是不生效)。
"caps": (portfolio.caps_ctx(view, ts_code=code, is_new_name=True,
sector=c.get("sector")) if is_open
else (portfolio.caps_ctx(view, ts_code=code) if side == "buy" else None)),
"flags": {"buy_halt": params.get("buy_halt"), "exec_halt": params.get("exec_halt"),
"brake_active": brake_active,
"blacklisted": bool(stock_params.get(code, {}).get("black")),
# 用户设的止损价, 与黑名单同一份事实源 (命令表)。规则闸拿它只告警
# 不拦截, 见 rule_gate 里那段说明。
"stop_price": stock_params.get(code, {}).get("stop_price"),
"is_command": False}})
if not gate["passed"]:
out["rejected"].append({"ts_code": code, "action": action, "by": "rule",
"failed": gate["failed"]})
if not dry_run:
pms_repo.insert_ledger(ts_code=code, action=action, arbiter="rule", verdict="REJECT",
price_at=price, hard_numbers=c["hard_numbers"],
failed_checks=gate["failed"], reason="自主提议未通过合规检查")
return
# ---- 二级: 研判闸 (补足/加仓/补仓/新建仓; 减持不送研判) ----
verdict = {"verdict": judge.PASS, "reason": "", "degraded": False}
if c.get("judge_required"):
if not _judge_budget_left(deadline):
# 本轮研判时间用尽 —— 整条跳过, 下一跳重来。
# **不当成「研判不可用」入人工队列**: 那会在自动档位下凭空造出一个人工确认队列,
# 与「不用每天靠人」这件事正好相反。规则闸白跑一次的成本是毫秒级。
out["skipped"].append({"ts_code": code, "action": action,
"why": "本轮研判时间预算用尽, 下一跳继续 "
f"(PMS_JUDGE_TICK_BUDGET_SEC="
f"{param_store.get_int('PMS_JUDGE_TICK_BUDGET_SEC', 150)}s)"})
return
verdict = judge.request(c, context={"position": _judge_ctx(pos),
"recent_ledger": _recent_ledger(code)})
# 新建仓的低把握驳回改交人 (2026-09-07, 第二道保险)。择时决策系统的提示词已改成
# 「证据不足以判断 → 不可用」, 但模型未必每次守得住; 它给的把握度是现成的读数,
# 低于阈值的驳回按「不可用」处理 —— 进人的待确认队列, 不记驳回、不杀提议。
# 只对新建仓: 加仓类的驳回口径不动。阈值一次定死, 不按复盘读数回调。
if is_open and verdict["verdict"] == judge.REJECT:
conf = verdict.get("confidence")
floor_conf = param_store.get_int("PMS_JUDGE_REJECT_CONF_MIN", 60)
if conf is not None and float(conf) < floor_conf:
verdict = {**verdict, "verdict": judge.UNAVAILABLE, "degraded": True,
"reason": (f"决策系统驳回但把握度只有 {int(conf)} (低于 {floor_conf}), "
f"按证据不足交人: {verdict.get('reason') or ''}")[:500]}
if verdict["verdict"] == judge.REJECT:
out["rejected"].append({"ts_code": code, "action": action, "by": "judge",
"failed": [verdict.get("reason") or "决策系统驳回"]})
if not dry_run:
pms_repo.insert_ledger(ts_code=code, action=action, arbiter="judge",
verdict="REJECT", price_at=price,
hard_numbers=c["hard_numbers"],
reason=verdict.get("reason") or "决策系统驳回")
return
if verdict.get("degraded"):
out["degraded"] = True
# ---- 三级: 按档位分流 ----
# 新建仓走自己的档位 (PMS_OPEN_AUTONOMY), 不跟随全局 —— 「新建仓要不要人点头」和
# 「加仓要不要人点头」是两个不同的决定。全局 off 已经在 scan_and_route 入口拦掉了,
# 走到这里说明总闸是开的。三条无条件覆盖档位的规矩对新建仓一样有效: 规则算出来的减持
# 自动执行、深档补仓强制确认、研判不可用一律入队。
autonomy = (out.get("open_autonomy") or AUTONOMY_FULL) if is_open else out["autonomy"]
# 强制入人工队列是**一票否决**, 排在方向与档位前面 (2026-09-03 修)。原来这一行是
# auto_exec = (side == "sell") or (autonomy == AUTONOMY_FULL and not force_queue)
# 卖出方向在或运算的左边, 把 force_queue 整个短路了 —— 方向是卖, 「强制入人工队列」
# 这个标记就不起作用。今天还没出事, 是因为走到这里的卖出候选只有保垫减仓一种, 它既不
# 强制确认也不送研判, force_queue 恒为假; 但设计里明确要求「研究证据走弱触发的减持必须
# 交人裁决、绝不自动卖」, 那类减持一旦接上来, 结果会是自动卖出。
#
# 为什么选一票否决而不是按来源开白名单: 白名单要求「谁可以自动卖」有一份完整清单, 而
# 这条路上真正需要自动卖的两条 —— 决策系统高置信风控卖出的自动止损、用户命令驱动的清仓
# —— 根本不经过提议分流 (前者是 signal_service 直接落卖出指令, 后者是命令服务 → 方案
# 生成器 → 执行器), 白名单在这里会是一份空转的清单。而 force_queue 这个名字本来就承诺了
# 「强制入人工队列」, 让它对所有方向都算数, 是把这个名字兑现, 不是新加一条规矩。
# 来源标记仍然要有, 但它的职责是让新来源能声明自己必须交人 (见 ae.source_confirm_why),
# 不是去给已有的自动止损发通行证。
#
# 保留下来的行为: 规则算出来的保垫减仓照旧自动执行 (它的 force_queue 是假), 风控高置信
# 卖出与命令清仓两条路一个字没碰。
src_why = ae.source_confirm_why(c.get("source"))
force_queue = (bool(c.get("needs_user_confirm")) or bool(src_why)
or bool(verdict.get("degraded")))
auto_exec = (not force_queue) and (side == "sell" or autonomy == AUTONOMY_FULL)
# 自动执行开关 (2026-09-03, PMS_OPEN_AUTO_EXEC_ON_VERDICT): 新建仓档位是 propose_only 时,
# 「判决候选 + 决策系统研判真回了通过 + 规则闸通过 (走到这里就是通过了) + 上游风险列表
# 为空」四条齐, 这一条按 full 处理。强制入队的 (关注 / 深档补仓 / 研判不可用) 永远不走。
auto_why = None
if is_open and not auto_exec and not force_queue and autonomy == AUTONOMY_PROPOSE:
auto_why = _verdict_auto_exec_why(c, verdict)
if auto_why:
auto_exec = True
if dry_run:
(out["executed"] if auto_exec else out["queued"]).append(
{**_brief(c), "route": "auto" if auto_exec else "queue",
"judge": verdict["verdict"], "dry_run": True,
**({"why": auto_why} if auto_why else {})})
return
if auto_exec:
iid = _make_instruction(c, price, now)
reason = verdict.get("reason") or c["reason"]
if auto_why:
reason = f"判决候选自动执行: {reason}"
pms_repo.insert_ledger(ts_code=code, action=action,
arbiter="judge" if c.get("judge_required") else "rule",
verdict="PASS", price_at=price, hard_numbers=c["hard_numbers"],
ref_id=iid, reason=reason[:500])
out["executed"].append({**_brief(c), "instruction_id": iid,
"why": ("减持自动执行 (规则触发, 来源没有要求交人裁决)"
if side == "sell"
else (auto_why or ("新建仓档位 full" if is_open
else "档位 full")))})
else:
pid = _make_proposal(c, price, verdict)
# 交人的原因按从具体到笼统取: 候选自带的 (关注判决等, 见 action_engine.verdict_confirm_why)
# → 来源强制的 (研究证据走弱那类减持) → 深档补仓那条老规矩 → 研判不可用 → 档位。
why = (c.get("confirm_why") or src_why
or ("深档补仓强制确认" if c.get("needs_user_confirm") else None)
or ("研判不可用, 降级人工确认" if verdict.get("degraded") else None)
or ("新建仓档位 propose_only" if is_open else "档位 propose_only"))
out["queued"].append({**_brief(c), "proposal_id": pid, "why": why})
def _verdict_auto_exec_why(c, verdict):
"""自动执行开关的四条件判定: 全部成立回一句原因 (写进 executed 与账本), 否则 None。
四条: ① 开关 PMS_OPEN_AUTO_EXEC_ON_VERDICT 为真; ② 上游判决是「候选」; ③ 决策系统的
研判**真的回了通过** —— 只认带应答体 (raw) 的 PASS, 「动作不在研判范围即放行」那种
没有问过决策系统的 PASS 不算, 研判不可用更不算; ④ 上游风险列表为空 (None 与 [] 都算
空: 判为候选本身就意味着上游没标硬风险, 缺键是旧字段布局)。规则闸通过是调用方保证的
(未过早就 return 了)。任一条不满足就回 None, 调用方照旧入队 —— 这条路只放宽不收紧。
"""
if not param_store.get_bool("PMS_OPEN_AUTO_EXEC_ON_VERDICT", False):
return None
hn = c.get("hard_numbers") or {}
if hn.get("verdict") != ae.VERDICT_CANDIDATE:
return None
if (not c.get("judge_required") or verdict.get("verdict") != judge.PASS
or verdict.get("degraded") or not isinstance(verdict.get("raw"), dict)):
return None
if hn.get("risk"):
return None
return ("判决候选自动执行 (判决候选 + 研判通过 + 规则闸通过 + 风险列表为空; "
"PMS_OPEN_AUTO_EXEC_ON_VERDICT=True)")
def _judge_budget_left(deadline) -> bool:
"""这一轮还够不够再送一次研判。
**这不是节流, 是让一次心跳做得完。** judge.request 是同步阻塞的, 单次上限
PMS_JUDGE_TIMEOUT (默认 90 秒), 而 scheduler 给所有调度任务设的软超时是 240 秒 ——
三只票送研判就顶破了, 任务被 celery 打死在中途, 而且是在已经落了一部分表之后。
以前送研判的只有已有持仓那几只、多数轮次还被去重挡掉, 所以一直没撞上; 新建仓上线后
冷启动那天会有十来条候选, 第一跳就会捅穿。
预算不够时调用方整条跳过并留痕, 下一分钟的心跳接着做 —— 一条候选都不丢, 也没有任何
按天计的上限。决策系统那侧的裁决按「日期+股票+动作」缓存半小时, 所以真正慢的只有
缓存过期后的第一跳, 之后同一批候选都是秒回。
"""
if deadline is None:
return True
need = max(1, param_store.get_int("PMS_JUDGE_TIMEOUT", 90))
return (time.monotonic() + need) <= deadline
# ================================================================ 落表
def _make_instruction(c, price, now) -> str:
ymd = td.ymd(now)
# 序号用完整时分秒 (2026-08-28 修): 原来 %1000 只留「分钟个位+秒」, 同日同票同动作
# 第二条指令约 1/600 概率撞唯一键丢单一跳。撞了再退一步逐秒加一重试。
seq = int(now.strftime("%H%M%S"))
window = param_store.get_int("PMS_EXEC_WINDOW_TDAYS", 3)
iid = None
last_err = None
for salt in range(3):
try_iid = cs.make_instruction_id(ymd, c["ts_code"], c["action"], seq + salt)
try:
pms_repo.insert_instruction(
instruction_id=try_iid, origin_type="proposal", origin_id=None,
ts_code=c["ts_code"],
action=c["action"], side=c["side"], qty=c["qty"], limit_price=None,
window_tdays=window, status=executor.ST_PROPOSED,
progress={"deadline": str(td.window_deadline(now.date(), window)),
"is_command": False, "children": [], "auto": True,
"reason": c["reason"]})
iid = try_iid
break
except Exception as e:
last_err = e
if iid is None:
raise last_err
g = bump_once_guards(c["ts_code"], c["action"], c.get("hard_numbers"), now) or {}
if not g.get("ok"):
# 指令已经落表了, 不回滚; 但要在评审账本上留一条痕, 否则这条纪律失效没有任何记录
try:
pms_repo.insert_ledger(
ts_code=c["ts_code"], action=c["action"], arbiter="rule", verdict="WARN",
price_at=float(price or 0), ref_id=iid,
hard_numbers={"once_guard_fields": {k: str(v) for k, v in
(g.get("fields") or {}).items()}},
reason=f"一次性守卫计数器未写入 ({g.get('error')}) —— "
f"该票 {c['action']} 的「只做一次」本轮失效, 留意重复出手")
except Exception:
logger.exception("一次性守卫失败留痕也没写上 %s", c["ts_code"])
return iid
def bump_once_guards(ts_code: str, action: str, hard_numbers=None, now=None):
"""把「只做一次」的三个计数器写上。
设计 §6 写了三条一次性约束, 但它们的计数器此前**只被读、从没被写过** —— 也就是说这三条
纪律一直是失效的, 只是被「同一票同一动作有在途提议就不重复提」这条兜底遮住了, 而那条兜底
恰好在规则闸拒绝时失灵 (没生成提议 → 没东西可去重), 于是同一个候选每分钟重来一次:
fill_count 回踩补足「每票 1 次」
last_add_date 盈利加仓「距上次 ≥2 交易日」—— 恒 None 时永远算作"很久没加过"
dca_count 补仓「各档评估一次」—— 恒 0 时每轮都当第一次评估
写入时机取「动作真的要落地」这一刻 (落指令), 而不是产出候选那一刻: 候选被闸门拦下不算
做过, 用掉一次名额不合理。dca_count 记的是**档位**而不是次数, 所以取本次触及的档序。
"""
now = now or datetime.now()
fields = {}
if action == "FILL":
fields["fill_count"] = 1
elif action == "ADD":
fields["last_add_date"] = now.date()
elif action == "DCA":
fields["dca_count"] = int((hard_numbers or {}).get("stage") or 1)
if not fields:
return {"ok": True, "fields": {}}
try:
n = pms_repo.update_position(ts_code, **fields)
except Exception as e:
# 写不上不该把已落表的指令带崩, 但**必须回成败**: 计数器没写上, 这条一次性纪律
# 就退回失效状态, 同一票同一动作在首条指令离开 LIVE 之后会再来一次。
# (2026-07-31 静默失败专项: 原来这里 return None, 调用方拿不到任何信号)
logger.error("一次性守卫计数器写入失败 %s %s: %s —— 该票这条一次性纪律本轮失效",
ts_code, action, e)
return {"ok": False, "fields": fields, "error": f"{type(e).__name__}: {e}"}
if n == 0:
logger.error("一次性守卫计数器没写到任何行 %s %s (持仓行不存在?) —— "
"该票这条一次性纪律本轮失效", ts_code, action)
return {"ok": False, "fields": fields, "error": "持仓行不存在, 影响 0 行"}
return {"ok": True, "fields": fields}
def _make_proposal(c, price, verdict) -> str:
ttl = param_store.get_int("PMS_PROPOSAL_TTL_HOURS", 24)
pid = f"PRP_{td.ymd()}_{c['ts_code'].replace('.', '')}_{c['action']}"
hn = {**(c.get("hard_numbers") or {}), "price": price, "reason": c["reason"],
"needs_user_confirm": c.get("needs_user_confirm", False),
# 候选来源 (2026-09-03): 人在「等我拍板」里要看得出这条减持是规则算的还是研究
# 证据走弱推来的 —— 两者该不该点头是两回事。没有来源的按动作引擎自身处理。
"source": c.get("source") or ae.SRC_ENGINE,
# 研判应答的结论与置信度 (2026-09-03): 人裁决时要看得见决策系统怎么说、有多确定。
# judge_reason 另有一列, 这两项进硬数字是为了随账本走 (采纳/驳回时原样落账)。
"judge_verdict": verdict.get("verdict"), "judge_conf": verdict.get("confidence")}
try:
pms_repo.insert_proposal(
proposal_id=pid, ts_code=c["ts_code"], action=c["action"], qty=c["qty"],
hard_numbers=hn, expire_at=datetime.now() + timedelta(hours=ttl),
judge_verdict=verdict.get("verdict"),
judge_reason=(verdict.get("reason") or "")[:500])
except Exception:
# 确定性编号被当日已终态 (过期/驳回) 的旧提议占着 —— 补时间后缀重试一次,
# 不让唯一键冲突把整轮扫描刷成 error (2026-08-28 审查修; 用户驳回的当日
# 不重提另有 _declined_today_keys 挡在扫描入口, 这里兜的是过期重生成)
pid = f"{pid}_{datetime.now().strftime('%H%M%S')}"
pms_repo.insert_proposal(
proposal_id=pid, ts_code=c["ts_code"], action=c["action"], qty=c["qty"],
hard_numbers=hn, expire_at=datetime.now() + timedelta(hours=ttl),
judge_verdict=verdict.get("verdict"),
judge_reason=(verdict.get("reason") or "")[:500])
return pid
# ================================================================ 上下文
def _scan_params(view: dict) -> dict:
p = dict(view["params"])
p.update({
"cushion_solid": param_store.get_float("PMS_CUSHION_SOLID", 0.03),
"trim_peak": param_store.get_float("PMS_TRIM_PEAK", 0.06),
"trim_giveback": param_store.get_float("PMS_TRIM_GIVEBACK", 0.5),
"dca_triggers": param_store.get_tuple_floats("PMS_DCA_TRIGGERS", (-0.08, -0.15)),
"dca_deep_confirm": param_store.get_float("PMS_DCA_DEEP_CONFIRM", -0.15),
"dca_max_ratio": param_store.get_float("PMS_DCA_MAX_RATIO", 0.5),
"no_chase_ma5": param_store.get_float("PMS_NO_CHASE_MA5", 0.06),
"buy_halt_dayup": param_store.get_float("PMS_BUY_HALT_DAYUP", 0.05),
"build_window_tdays": param_store.get_int("PMS_BUILD_WINDOW_TDAYS", 10),
"fill_max_loss": param_store.get_float("PMS_FILL_MAX_LOSS", -0.03),
"open_signal_priority": param_store.get_bool("PMS_OPEN_SIGNAL_PRIORITY", True),
"open_route_by_verdict": param_store.get_bool("PMS_PLAN_ROUTE_BY_VERDICT", True),
})
return p
def _market_ctx(held: list, now) -> dict:
"""每票的 MA5 / 5日高点 / 建仓天数 / 距上次加仓天数 (交易日口径)。"""
out = {}
today = now.date() if hasattr(now, "date") else now
for p in held or []:
code = p["ts_code"]
d = {"ma5": market.get_ma5(code), "high5": market.get_high5(code)}
# 当日行情快照 —— **规则闸的「不追高(当日涨幅)」全靠它**。
# 原来这里没取, `_route_one` 给规则闸硬编码 `day_chg_from_open: None`, 而规则闸对
# None 是整道跳过 → 动作引擎产出的每一笔自主买单 (FILL/ADD/DCA) 从来没过过这道闸,
# 账本里也查不到痕迹。executor 那条路一直是拿 day_snapshot 的, 只有这里漏了。
# 取不到就留空, 由规则闸记 DAYUP_MISSING 警告 —— 不再无声无息。
try:
d["day"] = market.day_snapshot(code) or {}
except Exception as e: # 行情不可用不该让整轮扫描崩掉
logger.warning("[提议] 取当日快照失败 %s: %s", code, e)
d["day"] = {}
opened, last_add = p.get("opened_date"), p.get("last_add_date")
d["tdays_since_open"] = _tdays_between(opened, today)
d["tdays_since_last_add"] = _tdays_between(last_add, today)
out[code] = d
return out
def _tdays_between(start, today):
"""start(含) 到 today 的交易日数; start 为空返回 None (调用方按「无约束」处理)。"""
if not start:
return None
try:
return max(0, td.trade_days_left(today, start) - 1)
except Exception:
return None
def _rejected_today_keys() -> dict:
"""今天已被规则闸拒过的 {(代码, 动作): 原因}。读不到就返回空 —— 去重是降噪, 不是纪律,
读失败时宁可多记几行日志, 也不能因此漏扫一个本该评估的动作。"""
try:
keys = pms_repo.rule_rejected_today(datetime.now().replace(
hour=0, minute=0, second=0, microsecond=0))
except Exception as e:
logger.warning("读当日规则闸拒绝记录失败 (按未拒过继续扫描): %s", e)
return {}
return {k: "今日已被规则闸拒过 (日切后重新评估; 拒因看 make t-gate)" for k in keys}
def _buy_signals_today() -> dict:
"""今天**择时决策系统**判过盘中转多的留痕。读不到就返回空 —— 它只影响候选的先后,
不影响资格, 所以读失败时降级成「按分数排」就够了, 不该因此让整轮新建仓停摆。
只认发送方以 bionic 开头的那些 (2026-09-03): db2 那条流盘中择时程序也在写, 它的入场
触发不是决策系统的结论, 不该拿来插队。区分靠留痕 reason 的固定开头 (signal_rules 里
两边共用同一个常量), 账本查询本身一个字不改。同一只票同一天两家都发过时, 查询按
最新一条取 reason —— 最新那条若是择时程序的, 决策系统更早的那条会被盖掉, 这只票就
不插队; 方向是少插一次队, 不是多插。
"""
try:
rows = pms_repo.buy_signals_today(datetime.now().replace(
hour=0, minute=0, second=0, microsecond=0))
except Exception as e:
logger.warning("[新建仓] 读当日转多留痕失败 (本轮按纯分数排序): %s", e)
return {}
return {code: v for code, v in (rows or {}).items()
if sr.is_bionic_buy_note((v or {}).get("reason"))}
def _judge_rejected_open_keys() -> dict:
"""今天已被研判闸驳回的**新建仓** {(代码, 动作): 原因}。读不到就返回空 (同上: 去重是降噪)。
只取 OPEN 那些键。加仓类对研判驳回**有意**不做当日去重 —— 契约里那句「研判结论会变」,
节流责任放在决策系统那侧的半小时缓存上, 这条纪律一个字不动。
"""
try:
keys = pms_repo.judge_rejected_today(datetime.now().replace(
hour=0, minute=0, second=0, microsecond=0))
except Exception as e:
logger.warning("读当日研判驳回记录失败 (按未驳回过继续扫描): %s", e)
return {}
return {(code, act): "今日已被研判闸驳回 (日切后重新评估; 驳回理由看 make t-gate)"
for code, act in keys if act == ae.A_OPEN}
def _declined_today_keys() -> dict:
"""今天已被用户**驳回**的提议 {(代码, 动作): 原因} —— 当日不再重提 (2026-08-28)。
两个理由: ① 你上午刚说过"", 触发条件没变的话下午每分钟再问一遍是烦人不是尽责;
② 提议编号是「日期+代码+动作」的确定性编号, 驳回的行还占着编号, 同日重建必撞
唯一键, 整轮 scan_and_route 会被这个报错刷屏。日切自动解除, 明天重新评估。"""
keys = {}
today = str(td.ymd())
try:
for p in pms_repo.list_proposals(statuses=("DECLINED",), limit=200):
at = str(p.get("decided_at") or p.get("created_at") or "")
if at[:10].replace("-", "") == today or at[:10] == f"{today[:4]}-{today[4:6]}-{today[6:]}":
keys[(p["ts_code"], p["action"])] = "今天已被你驳回过, 当日不再重提 (明天重新评估)"
except Exception as e:
logger.warning("读当日已驳回提议失败 (按无): %s", e)
return keys
def _inflight_keys() -> dict:
"""已有在途提议或在途指令的 {(代码, 动作): 原因} —— 同一件事不重复提。
提议与指令分开写原因, 并且**指令写在后面覆盖提议**: 一条提议被采纳之后会变成指令,
这时候「有在途指令」比「有提议在等确认」更贴近它此刻的状态。原因里带上单号与状态,
是为了让人从 skipped 那一行就能接着往下查, 不必再去翻两张表。
"""
keys = {}
# 减持侧按**代码**去重, 不按 (代码, 动作) (2026-09-07 审查修): 一只票有一条清仓在等人拍板时,
# 保垫减仓若按另一个动作名单独放行, 会抢在人前面自动卖掉一部分。所以一条减持在途,
# 这只票两种减持一起记进跳过集合。动作引擎那边还有一道同样的闸, 两处互为保险。
def _mark(code, action, why):
keys[(code, action)] = why
if action in ae.SELL_SIDE_ACTIONS:
for other in ae.SELL_SIDE_ACTIONS:
keys.setdefault((code, other), why + " (同票另一种减持一并让路)")
try:
for p in pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200):
_mark(p["ts_code"], p["action"], f"已有在途提议 {p.get('proposal_id')} 在等人确认")
except Exception as e:
logger.warning("读提议队列失败: %s", e)
try:
for i in pms_repo.list_instructions(statuses=list(executor.LIVE), limit=300):
_mark(i["ts_code"], i.get("action"),
f"已有在途指令 {i.get('instruction_id')} ({i.get('status')})")
except Exception as e:
logger.warning("读在途指令失败: %s", e)
return keys
def _pos_of(view: dict, ts_code: str) -> dict:
for x in view["positions"]:
if x["ts_code"] == ts_code:
return x
return {"ts_code": ts_code, "total_qty": 0, "avail_qty": 0, "frozen_reason": "NONE"}
def _judge_ctx(pos: dict) -> dict:
"""送研判的账本快照 (设计 §7: PMS 备齐硬数字与账本流水作为 context)。"""
keep = ("ts_code", "price", "avg_cost", "total_qty", "base_qty", "add_qty", "dca_qty",
"cushion_pct", "cushion_peak", "pct_of_scale", "target_pct", "support_ref",
"pressure_ref", "stop_ref", "ref_source", "sector", "neg_cushion_days")
return {k: pos.get(k) for k in keep}
def _recent_ledger(ts_code: str, limit: int = 10) -> list:
try:
rows = pms_repo.list_ledger(ts_code=ts_code, limit=limit)
except Exception:
return []
# price_at 是数据库的 DECIMAL 列, 直接带出去会是 Decimal —— 那是 2026-08-06 研判闸
# 第一次真发请求就全军覆没的原因 (json 序列化不了)。judge.jsonable 已经在出口统一拦了,
# 这里再就地转成 float, 是为了让日志、页面、留痕里的数字也是干净的, 不必依赖出口那一道。
return [{"at": str(r.get("decided_at")), "action": r.get("action"),
"arbiter": r.get("arbiter"), "verdict": r.get("verdict"),
"price": float(r.get("price_at") or 0), "reason": r.get("reason")} for r in rows]
def _brief(c: dict) -> dict:
return {"ts_code": c["ts_code"], "action": c["action"], "side": c["side"],
"qty": c["qty"], "reason": c["reason"]}