# -*- 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 consensus 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 tech_rules from app.core import tradedays as td from app.repo import pms_repo from app.services import (candidate_pub, command_service, consensus_service, executor, industry, judge, market, param_store, reask_service, plan_feed, portfolio, strategy_service) 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, # 系统自决本轮自动采纳的新建仓只数 (每日上限累计用; observe 读数从 self_decide_seen 取) "_sd_accepts": 0} 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, 不再重取 # 系统自决每日新建仓上限的基数: 今天已自动采纳几只 (评审账本按前缀数)。只读一次, # 本轮内 _route_one 每采纳一只再往 out["_sd_accepts"] 加, 两者相加即当日累计。off 时不查库。 if ae.self_decide_on(params): try: _t0 = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0) params["_sd_base_count"] = pms_repo.system_pass_count_today(_t0) except Exception as e: # noqa: BLE001 logger.warning("[系统自决] 读当日已采纳数失败 (本轮按 0 起算): %s", e) params["_sd_base_count"] = 0 # 跳过三类: ①已有在途提议或指令的 ②今天已被规则闸拒过的 ③今天已被研判闸驳回的新建仓。 # ② 是 2026-07-29 的教训 —— 组合已超总仓上限时, 16 只深亏票的补仓候选每分钟被拒 # 一次, 一天往评审账本灌几千行一模一样的记录。闸门结论当天基本不会变, 记一次就够。 # ③ 只挡新建仓: 加仓类对研判驳回**有意**不去重 (契约里那句「研判结论会变」), 那条 # 纪律不动; 而候选池是几十只的量级, 不挡的话 ② 那个洞会原样从研判闸重来一遍。 # # 三份合并成 `{(代码, 动作): 原因}`, **在途的排在最后, 覆盖被拒的** (2026-08-06): # 这四种处境完全不同 —— 在途的明天照样被挡, 被拒的日切就重新评估。从前它们在 # skipped 里长得一模一样, 看的人判断不了这只票明天还会不会再被评估。 # 一只票既有在途指令又今天被拒过时, 显示"有在途"更贴近它此刻的实际状态。 # 研判驳回那一份先过一遍解锁判定 (2026-09-09, 台账 055): 满足解锁条件的从跳过集合里 # 剔掉, 让它今天能被重新评估一次。**只剔研判闸这一份** —— 规则闸拒、用户驳回、在途 # 三份一个字不动 (合并顺序天然保证这点)。没有任何票解锁时行为与今天逐字相同。 _jk = _judge_rejected_open_keys() _unlocked = {} try: _unlocked = reask_service.evaluate_all(_jk, now) except Exception as e: logger.warning("[新建仓] 解锁判定出错 (本轮当日闸照旧全挡): %s", e) for _k in _unlocked: _jk.pop(_k, None) out["reask_unlocked"] = [{"ts_code": c, "action": a, "why": v.get("why"), "hits": v.get("hits")} for (c, a), v in _unlocked.items()] # 系统自决当日已放弃的新建仓也并入跳过集合 (改动六): 排在用户驳回之后、在途之前 —— # 「今日系统已判放弃,明日重评」。off 时 _system_declined_today_keys 返回空, 逐字回旧。 skip = {**_jk, **_rejected_today_keys(), **_declined_today_keys(), **_system_declined_today_keys(), **_inflight_keys()} # 挂了 ACTIVE 交易方案(策略)的票交策略层接管, 动作引擎不再对它自动提议 # (读库失败按空集 —— 宁可这轮不排除, 也不能因读不到把全体持仓都排除) try: strategy_codes = pms_repo.active_strategy_codes() except Exception: strategy_codes = set() # 逻辑状态挂到持仓行上 (2026-09-07 第三件): 早上拉计划时查回的当日态。读不到就没有这个键, # 动作引擎按「没有读数」处理, 绝不折成存疑; 挂不上只记日志, 不拖垮扫描。 try: from app.services import logic_state_service logic_state_service.attach(view["held"]) except Exception as e: # noqa: BLE001 logger.warning("[逻辑状态] 挂到持仓行失败, 本轮按没有读数: %s", e) # 三源合议增持门 (2026-09-11 工作包二): 给持仓行装配合议四块供动作引擎的增持门读。持仓票 # 的买方评析从逻辑状态映射里带出 (计划里读不到持仓票)。开关 PMS_TECH_GATE_INCREASE 关掉 # 时整段不做, 持仓行不挂 consensus_blocks, 增持门不设 —— 行为与接入前逐字相同。 # 装配失败只记日志、不拦扫描 (合议是加不改)。 if params.get("tech_gate_increase"): try: from app.services import logic_state_service as _lss _lmap = _lss.state_map() _crs = {code: (st or {}).get("company_review") for code, st in (_lmap or {}).items()} # 持仓行的择时票也带上当日盘中转多留痕 (2026-09-14 加固, 评审发现四): 来源与候选 # 相同 —— 决策系统判过盘中转多的留痕 (_buy_signals_today)。没这留痕的行照旧只按 # 昨夜定性。持仓行昨夜中性而当日转多时, 择时票升为看多, 增持门的方向判断才跟得上。 _sig = _buy_signals_today() _flip = {code: (v or {}).get("at") for code, v in (_sig or {}).items() if (v or {}).get("at")} consensus_service.decorate_positions(view["held"], company_reviews=_crs, flip_map=_flip) except Exception as e: # noqa: BLE001 logger.warning("[合议] 持仓增持门装配失败, 本轮按无合议 (不设门): %s", e) # 技术面转空自动离场 (2026-09-11 工作包三): 给持仓行挂技术面状态供 eval_tech_exit 读, # 并备好「本次翻空已处理」去重集 (清掉 SAR 已不在空头的, 那是新一轮翻向)。挂不上只记日志、 # 本轮不评转空离场 (无读数弃权)。离场档位 off 时 tech_exit_on 已为假, 评估器自然不产出。 if params.get("tech_exit_on"): try: from app.services import tech_service as _ts _tech_states = _ts.state_map() _ts.attach(view["held"], _tech_states) _done = _load_tech_exit_done() params["tech_exit_done"] = _keep_tech_exit_done(_done, _tech_states) except Exception as e: # noqa: BLE001 logger.warning("[技术面离场] 挂技术面状态失败, 本轮不评转空离场: %s", e) params["tech_exit_done"] = 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, now=now) except Exception as e: logger.exception("新建仓扫描失败") out["errors"].append(f"新建仓扫描失败: {type(e).__name__}: {e}") out["open_candidates"] = len(open_cands) # 候选池发布给择时决策系统 (2026-09-09, 台账 004): 它的盘中强势扫描拿这份当覆盖面, # 读不到就本轮不扫、绝不回退全市场。这里只写代码不写别的; 写失败只记日志, 不影响本轮。 # 落点选在候选产出之后、分流之前: 发布的是"今天该盯哪几只", 与它们最后买没买无关。 if not dry_run and open_cands: try: candidate_pub.publish([x.get("ts_code") for x in open_cands]) except Exception as e: logger.warning("[新建仓] 候选池发布失败 (不影响本轮): %s", e) # 研判的时间预算从这一刻起算。**先跑已有持仓的四类, 再跑新建仓** —— 预算真用尽时, # 被推到下一跳的一定是新建仓, 已有持仓的动作行为与 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: _rk = _unlocked.get((c.get("ts_code"), c.get("action"))) if _rk: c["reask"] = _rk _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}") # 技术面转空离场去重集写回 (2026-09-11 工作包三): 本轮真处理过 (执行或入队, 未被闸拒) 的 # 转空离场记进去重集, 下次同一翻空不再重复处理。试算不写。清掉 SAR 已翻多的那步在扫描前做过了。 if not dry_run: try: _rej = {(r.get("ts_code"), r.get("action")) for r in out["rejected"]} _acted = {c["ts_code"] for c in scanned["candidates"] if c.get("source") == ae.SRC_TECH_EXIT and (c["ts_code"], c["action"]) not in _rej} if _acted: _save_tech_exit_done(set(params.get("tech_exit_done") or set()) | _acted) except Exception as e: # noqa: BLE001 logger.warning("[技术面离场] 去重集写回失败 (下轮可能重评一次): %s", e) # 三源合议观察读数落表 (2026-09-14 观察读数包, 台账 013): 把本轮合议相关的判定归类落表, # 供台账 006 到 011 的复核。只加不改行为 —— 失败只记警告、不抛 (SOFT_FAIL); dry_run 内部会跳过。 # 返回值接住 (test_batch10 [A1] 查丢弃返回值)。 if not dry_run: try: from app.services import consensus_stats rec = consensus_stats.record_round(out, scanned, params, now) if not rec.get("ok"): logger.warning("[合议观察] 本轮落表未全成: %s", rec.get("error")) except Exception as e: # noqa: BLE001 logger.warning("[合议观察] 落表出错 (不影响扫描): %s", e) out["ok"] = not out["errors"] return out _TECH_EXIT_DONE_KEY = "PMS_TECH_EXIT_DONE" def _keep_tech_exit_done(done, tech_states) -> set: """本次翻空去重集里, 本轮该保留 (不清除) 的代码集。 2026-09-14 加固 (评审发现三): 原来是「SAR 方向是空才保留」, 而「不是空」把「映射里没有 这只票」也算了进去 —— 停牌或决策系统那晚漏算的持仓票缺读数, 会被当成翻多而清出去重集; 若同一轮另有别的票处理了转空离场, 清过的集合会被写回, 第二天读数回来、翻空仍在两天内就 再减一次 (多减三分之一)。改成**只在 SAR 方向明确为多时才清除**: 翻空与缺读数一律保留。""" return {c for c in (done or set()) if (tech_states.get(c) or {}).get("sar_side") != "多"} def _load_tech_exit_done() -> set: """本次翻空已处理过的代码集 (跨轮去重, 不动表结构; 与 tech_state_map 同一手法用运行参数承载)。 读不到或坏了返回空集。清除 (SAR 翻多=新一轮) 在 scan_and_route 扫描前按当轮技术面映射做。""" import json raw = param_store.get(_TECH_EXIT_DONE_KEY, "") or "" if not raw: return set() try: m = json.loads(raw) except (TypeError, ValueError): return set() return set(m.get("codes") or []) if isinstance(m, dict) else set() def _save_tech_exit_done(codes) -> None: import json r = param_store.set_param(_TECH_EXIT_DONE_KEY, json.dumps( {"codes": sorted(codes), "at": datetime.now().strftime("%Y-%m-%d %H:%M:%S")}, ensure_ascii=False), "system") if not r.get("ok"): logger.error("[技术面离场] 去重集写入失败: %s —— 同一翻空下轮可能再处理一次", r.get("error")) # ================================================================ 新建仓的取数与筛选 def _scan_open(view, params, stock_params, skip, mkt, out, now=None) -> 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 [] # 三源合议装配 (2026-09-11 工作包二): 开着时给每只候选凑齐基本面 / 技术面 / 择时三票, 合成方向 # 与路由, 挂到候选上供 scan_open 分流与定档。一次批量取昨夜定性与技术面映射, 逐只装配。 # 开关关掉整段不做, cands 原样进 scan_open —— 与接入前逐字相同。装配失败只记日志、不拦扫描 # (合议是加不改: 装不上就当没有合议, scan_open 里 con 为 None 自然退回旧路)。 if params.get("consensus_route"): _attach_consensus(cands, out, now=now) 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"] def _intraday_cfg() -> dict | None: """盘中确认包的开关与阈值 (2026-09-14, 台账 012)。总闸关掉返回 None —— 整段不走, 逐字回旧。 cparams 是合议参数, 收口突破改立场后要用它重新合议。""" if not param_store.get_bool("PMS_TECH_INTRADAY_CONFIRM", True): return None return {"breakout_on": param_store.get_bool("PMS_TECH_INTRADAY_BREAKOUT", True), "vol_min": param_store.get_float("PMS_TECH_VOL_CONFIRM", 1.5), "window": param_store.get("PMS_TECH_INTRADAY_WINDOW", "0945-1430") or "0945-1430", "cparams": consensus_service._params()} def _intraday_market(code, cache) -> tuple: """一只票的盘中读数 (现价/均价/时段折算量比)。本轮按代码缓存, 同一只不取两次。 量比走与解锁重问相同的算法 (reask_service._vol_ratio)。任一取数失败传 None (按取不到)。""" if code in cache: return cache[code] price = vwap = vratio = None try: day = market.day_snapshot(code) or {} price, vwap = day.get("price"), day.get("vwap") except Exception as e: # noqa: BLE001 logger.warning("[盘中确认] 取当日快照失败 %s: %s", code, e) try: vratio = reask_service._vol_ratio(code) except Exception as e: # noqa: BLE001 logger.warning("[盘中确认] 取量比失败 %s: %s", code, e) cache[code] = (price, vwap, vratio) return cache[code] def _apply_intraday(code, blocks, cfg, now, cache) -> str | None: """对一只候选跑盘中确认, 原地改 blocks。返回 intraday kind (或 None, 表示相位无关未跑)。 收口突破改技术面立场并重新合议; 开口未确认/时段外把路由改观察带 wait_confirm; 都在硬数字 tech 块记 intraday。""" tech = blocks.get("tech") or {} if tech.get("phase") not in ("收口等待", "开口向上"): return None price, vwap, vratio = _intraday_market(code, cache) r = tech_rules.intraday_confirm(tech, price=price, vwap=vwap, upper_prev=tech.get("boll_upper"), vol_ratio=vratio, now=now, window=cfg["window"], vol_min=cfg["vol_min"], breakout_on=cfg["breakout_on"]) kind = r.get("kind") if kind == "breakout": new_tech = r["state"] blocks["tech"] = new_tech cp = cfg["cparams"] blocks["consensus"] = consensus.decide( (blocks.get("fund") or {}).get("stance") or "无读数", new_tech.get("stance") or "无读数", (blocks.get("timing") or {}).get("stance") or "无读数", tech_phase=new_tech.get("phase"), fund_required=(cp.get("fund_required") and cp.get("consensus_route")), weak_confirm=cp.get("weak_confirm", True), timing_dissent=cp.get("timing_dissent", "confirm")) new_tech["intraday"] = {"kind": kind, "reads": r["reads"]} elif kind in ("opened_wait", "outside_window"): con = blocks.get("consensus") or {} con["route"] = "观察" con["route_reason"] = r["why"] con["disp"] = ae.DISP_WAIT_CONFIRM blocks["consensus"] = con tech["intraday"] = {"kind": kind, "reads": r["reads"]} elif kind == "opened_ok": tech["intraday"] = {"kind": kind, "reads": r["reads"]} return kind def _attach_consensus(cands, out=None, now=None) -> None: """给每只候选装配三源合议, 挂三样到候选上 (原地改): consensus —— 合议块 (direction / votes / strength / reason / route / route_reason), scan_open 分流读它; consensus_blocks —— 四块意见 (fund / tech / timing / consensus), scan_open 定档时取 fund/tech 喂 advise_v2; consensus_hard —— 提议硬数字六键 (方案附录丁), 进评审账本与提议卡, 不进送研判名单。 批量取一次昨夜定性与技术面映射; 盘中转多留痕从候选自带的 sig_buy 取时刻。任何一步失败都不抛 —— 装不上就当这只票没有合议 (scan_open 里 con 为 None 自然退回旧路), 合议是加不改。 盘中确认 (2026-09-14 盘中确认包): 装配完四块后, 若总闸开着且传了 now, 对收口等待/开口向上的 候选跑一遍 intraday_confirm —— 收口突破改立场重新合议, 开口未确认/时段外把路由改观察带 wait_confirm。 总闸关掉 (_intraday_cfg 为 None) 或没传 now 时整段不走, 与接入前逐字相同。 每装配一只候选往 out["consensus_seen"] 追加一条紧凑记录 (观察读数/候选处置快照用): ts_code / fund / tech / timing / direction / route / phase, 盘中突破的另记 intraday=breakout。 out 为 None 或缺该键时只装配不留痕。""" if not cands: return seen = out.setdefault("consensus_seen", []) if isinstance(out, dict) else None try: codes = [c.get("ts_code") for c in cands if c.get("ts_code")] nightly = consensus_service.nightly_map(codes) tech_states = consensus_service.state_map() except Exception as e: # noqa: BLE001 logger.warning("[合议] 批量取数失败, 本轮候选不装配合议: %s", e) return icfg = _intraday_cfg() if now is not None else None icache: dict = {} for c in cands: code = cs.normalize_code(c.get("ts_code") or "") if not code: continue try: flip_at = (c.get("sig_buy") or {}).get("at") or None blocks = consensus_service.assemble( c, nightly=nightly.get(code), tech_states=tech_states, flip_at=flip_at) ikind = _apply_intraday(code, blocks, icfg, now, icache) if icfg else None c["consensus"] = blocks["consensus"] c["consensus_blocks"] = blocks c["consensus_hard"] = consensus_service.hard_keys(blocks) if seen is not None: con, fund, tech, tm = (blocks["consensus"], blocks["fund"], blocks["tech"], blocks["timing"]) rec = {"ts_code": code, "fund": fund.get("stance"), "tech": tech.get("stance"), "timing": tm.get("stance"), "direction": con.get("direction"), "route": con.get("route"), "phase": tech.get("phase"), # 候选栏芯片用: 基本面质地档 (好/中/差), 从候选自带的买方评析取 "overall": (c.get("company_review") or {}).get("overall")} if ikind: rec["intraday"] = ikind seen.append(rec) except Exception as e: # noqa: BLE001 logger.warning("[合议] 装配失败 %s (本票按无合议): %s", code, e) # 候选一只都不剩时, 这一栏要说清是被哪几道筛掉的。键名是代码里的, 翻成人话再显示。 # 顺序就是筛选实际发生的顺序, 照着念下来正好是一条票走过的路。 _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 = {} _dj = _judge_rejected_open_keys() # 解锁重问 (2026-09-09): 判定是只读的, 这里跟扫描那边用同一个函数、同一份口径, # 所以页面上写的"已解锁"与扫描下一跳真会重问的是同一批票, 不会两处说法打架。 _un = {} try: _un = reask_service.evaluate_all(_dj) or {} except Exception as e: logger.warning("[处置] 解锁判定出错 (按未解锁显示): %s", e) for _k in _dj: _w = _jr.get(_k[0]) base = ("决策系统判暂不建仓:" + _w) if _w else "决策系统今天判过暂不建仓,明天开盘会重新评估" if _k in _un: skip[_k] = (base + "。不过盘中输入变了,这只票今天已解锁一次重问:" + str(_un[_k].get("why") or "") + "。重问放行的话会进你的待确认队列。") else: skip[_k] = base for _k in _rejected_today_keys(): _w = _rr.get(_k[0]) skip[_k] = ("未通过合规检查:" + _w) if _w else "今天没通过合规检查,明天开盘会重新评估" _sd_keys = _system_declined_today_keys() # 系统自决今日已放弃 (明日重评) for _k in _sd_keys: skip[_k] = "今日系统已判放弃,明日重评" for _k in _inflight_keys(): # 在途; 前端多半已按指令/提议标成「已下单/待你确认」 skip[_k] = "系统正在处理这只票(见『待确认提议』或『当日指令』)" sink = {"skipped": []} # 丢弃用: _scan_open 只往里 append/extend, 不读它 would = _scan_open(view, params, stock_params, skip, {}, sink, now=now) 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 的 # 合议判观察的候选带 disp=wait_tech/wait_confirm, 其余不产出的按 deny。 res["by_code"].setdefault( code, {"disp": s.get("disp") or "deny", "why": s.get("why") or "未产出候选"}) # 系统自决今日放弃的候选打处置词 sys_decline (09-14 评审补): 页面显示「系统已放弃」而不是「未建仓」。 for (_code, _act) in _sd_keys: if _act == ae.A_OPEN and _code in res["by_code"]: res["by_code"][_code]["disp"] = "sys_decline" # 候选栏三小块 (2026-09-14 页面收尾包): 每只候选的合议摘要, 从 consensus_seen 来。 # 基本面立场与质地档 / 技术面立场与相位 / 合议方向与路由。挂到对应代码的 by_code 行上。 for rec in (sink.get("consensus_seen") or []): code = rec.get("ts_code") if code and code in res["by_code"]: res["by_code"][code]["consensus"] = { "fund": rec.get("fund"), "overall": rec.get("overall"), "tech": rec.get("tech"), "phase": rec.get("phase"), "timing": rec.get("timing"), # 择时那一票 (2026-09-15: 候选栏三枚维度标签用) "direction": rec.get("direction"), "route": rec.get("route")} 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} # 系统自决 decline/full 档: 逻辑存疑的新建仓不送研判 (09-14 评审补) —— 自决必然放弃它, 送研判 # 只是白占一次调用与本轮时间预算。record 档不跳: 那一档要照旧进人工队列, 提议上得有真研判结论。 _sd_mode = str((params or {}).get("self_decide") or "off").strip() if (is_open and c.get("decide_kind") == ae.DK_LOGIC_DOUBT and _sd_mode in ("decline", "full")): verdict = {"verdict": "SKIPPED", "reason": "逻辑存疑,不送研判(系统自决直接放弃)", "degraded": False, "raw": None, "confidence": None} elif 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 # 解锁重问 (2026-09-09): 名额在**真的要发请求**这一刻才算用掉 —— 走到这里说明规则闸 # 已经过了, 候选也真的产出来了。解锁之后被名额或金额挡住、被规则闸拦下的那些, 什么都 # 没问过, 不该消耗这只票当天唯一的一次机会。试算一律不落账、不带强制标记。 _rk = c.get("reask") if not dry_run else None _force, _rtext = False, None if _rk: if reask_service.commit(code, _rk, price): _force = True _rtext = reask_service.reask_text(_rk) else: c.pop("reask", None) # 名额没落上就不算重问, 这一跳按普通候选走 verdict = judge.request(c, context={"position": _judge_ctx(pos), "recent_ledger": _recent_ledger(code)}, force=_force, reask_text=_rtext) # 新建仓的低把握驳回改交人 (2026-09-07, 第二道保险)。择时决策系统的提示词已改成 # 「证据不足以判断 → 不可用」, 但模型未必每次守得住; 它给的把握度是现成的读数, # 低于阈值的驳回按「不可用」处理 —— 进人的待确认队列, 不记驳回、不杀提议。 # 只对新建仓: 加仓类的驳回口径不动。阈值一次定死, 不按复盘读数回调。 # **系统自决开着时不做这个转换** (2026-09-14, §3.5): 自决下 REJECT 一律放弃并记账 judge, # 不再"低把握改不可用交人"——低把握的 REJECT 也走下面的放弃路, 解锁重问对它照常生效。 if is_open and verdict["verdict"] == judge.REJECT and not ae.self_decide_on(params): 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: # 驳回也把把握度记进硬数字 (2026-09-08): 低把握驳回转交人那道保险 (PMS_JUDGE_REJECT_CONF_MIN) # 是否在起作用, 只能从驳回行的把握度分布看出来; 此前只有提议行记它, 驳回行没有, 复核无据。 # 基线快照 (2026-09-09, 台账 055): 驳回那一刻的价、两个口径的当日涨幅、 # 时段折算量比、昨夜定性与支撑压力。解锁判定要拿它跟此刻的读数比 —— # 没有这个块, "相对驳回那一刻变强了没有"就无从判起。取数失败写 None, # 不阻断驳回落账: 记账优先于留痕。 _hn = {**(c["hard_numbers"] or {}), "judge_conf": verdict.get("confidence")} try: _hn["reask_base"] = reask_service.base_snapshot( code, price, judge_conf=verdict.get("confidence")) except Exception as e: logger.warning("[新建仓] 基线快照取不到 [%s] (驳回照常落账): %s", code, e) _rsn = verdict.get("reason") or "决策系统驳回" if c.get("reask"): # 重问后仍驳回: **静默记账**, 页面不单列 (2026-09-09 拍板)。它就是账本里 # 又一行"决策系统驳回建新仓", 与别的驳回长得一样; 下一跳 G2 会挡住, # 当日不再重问, 这条路彻底闭合。 _hn["reask_seq"] = 1 _rsn = "重问后仍驳回:" + _rsn pms_repo.insert_ledger(ts_code=code, action=action, arbiter="judge", verdict="REJECT", price_at=price, hard_numbers=_hn, reason=_rsn[:500]) 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")) # 重问放行的一律交人 (2026-09-09 拍板)。三条理由: # 一, 重问放行的票多半现价已经高于择时买入区间上沿, 自动落指令会一路等到窗口作废, # 期间还占着名额与预估金额, 把真正买得进的票挤掉 —— 这是实打实的损失。 # 二, "输入变了"和"证据更强了"是两回事: 解锁由代码按硬条件判, 定性由模型判, 而 # "这个变化够不够动用一个持仓名额"是边界情形, 按既有分工归人。 # 三, 这次要解决的是**连提议都不出**, 不是"一定要买到"。把牌摊在人面前, 人拍板; # 拍了板买不上是择时区间和追高原则在起作用, 那两条一个字没改。 reask_why = ("盘中重问放行,需人工确认:" + str((c.get("reask") or {}).get("why") or "") + "。能不能买到由择时区间决定,可能当日买不上。") if c.get("reask") else None # 上游标了硬风险的票必须交人 (2026-09-10 补的洞)。 # # 这一条原先只写在下面那条更窄的旁路里 (_verdict_auto_exec_why 的第四条「上游风险 # 列表为空」), 而一票否决这条链完全不看 risk。后果是反的: 档位是 propose_only 时 # 旁路会拦住带风险的票, 一旦档位改成 full, 那条旁路根本不走, 带硬风险的候选反而 # 畅通无阻直接落指令 —— 越放开越不设防。 # # 硬风险是选股系统在候选卡上标的三类 (akg-factor-bridge/card.py): 昨夜给出 # 卖出/回避/已剔除信号、用的传导快照日与计划日不符、吸筹为高位派发。这三类都是 # 「有人看见了不对劲」, 不是评分低, 所以归人裁决而不是归档位。 _hn_risk = (c.get("hard_numbers") or {}).get("risk") risk_why = None if _hn_risk: _rl = _hn_risk if isinstance(_hn_risk, (list, tuple)) else [_hn_risk] risk_why = "选股系统标了风险,需人工确认:" + ";".join(str(x) for x in _rl if x) force_queue = (bool(c.get("needs_user_confirm")) or bool(src_why) or bool(verdict.get("degraded")) or bool(reask_why) or bool(risk_why)) auto_exec = (not force_queue) and (side == "sell" or autonomy == AUTONOMY_FULL) # ---- 系统自决 (2026-09-14 台账 015): 新建仓/深亏补仓/目标价由合议+研判自动决定 ---- # 位置在 force_queue 之后、四条件旁路之前: 自决非 off 时它先决定, 旁路不再插手。 # 重问放行也归自决 (试探仓): 先给它打上 decide_kind, 再交 _self_decide 判。 if c.get("reask") and ae.self_decide_on(params) and not c.get("decide_kind"): c["decide_kind"] = ae.DK_REASK sd = _self_decide(c, verdict, params, params.get("_sd_base_count", 0) + out.get("_sd_accepts", 0)) if sd.get("would") and sd["would"] != "none": _con = _sd_consensus_brief(c) out.setdefault("self_decide_seen", []).append( {"ts_code": code, "action": action, "decide_kind": sd.get("rule"), "would": sd["would"], "decision": sd.get("decision"), "tier": sd.get("tier"), "why": sd.get("why"), "direction": _con.get("direction"), "route": _con.get("route"), "conf": verdict.get("confidence")}) sd_decision = sd.get("decision", "none") if sd_decision == "accept": auto_exec, force_queue = True, False if is_open: out["_sd_accepts"] = out.get("_sd_accepts", 0) + 1 elif sd_decision == "decline": # 自动放弃: 落账 (评审方 system) + 追加 rejected by system + 当日不重评 + 返回 if dry_run: out["rejected"].append({**_brief(c), "by": "system", "route": "decline", "dry_run": True, "why": sd.get("why")}) return _sd_decline_ledger(code, action, c, verdict, price, sd) out["rejected"].append({"ts_code": code, "action": action, "by": "system", "failed": [sd.get("why") or "系统自决放弃"]}) return elif sd_decision == "confirm": auto_exec, force_queue = False, True c["confirm_why"] = sd.get("why") or c.get("confirm_why") # 自动执行开关 (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: sd_accept = sd_decision == "accept" strat_cancelled, buys_cancelled = [], [] if sd_accept and side == "sell": # 目标价自动清仓两道前置 (2026-09-14): 按当前可卖量夹紧; 该票挂着方案先撤, # 免得清仓单与网格买入腿对倒 (附录乙那处冲突)。可卖量为零则不落, 明日重评。 posrow = pms_repo.get_position(code) or {} clamped = ae.clamp_sell_qty(int(c.get("qty") or 0), posrow) if clamped <= 0: out["skipped"].append({"ts_code": code, "action": action, "why": "系统到价自动清仓: 当前可卖量为零, 本轮不清 (明日重评)"}) return c["qty"] = clamped # 停买入侧三件 (09-14 评审补): 撤方案、驳回买入提议、撤在途买单 —— 不只在有活动方案时, # 一张在途的补仓/网格买单同样会与清仓单对倒。撤不成只记日志, 清仓照落。 try: r = strategy_service.cancel_by_code( code, reason="系统到价自动清仓: 先撤该票交易方案与在途买单, 免得清仓单与买入腿对倒") strat_cancelled = r.get("cancelled") or [] buys_cancelled = r.get("instructions") or [] if r.get("errors"): logger.warning("[自决清仓] 停买入侧有未成项 %s: %s", code, r["errors"]) except Exception as e: # noqa: BLE001 logger.warning("[自决清仓] 停买入侧失败 %s (清仓照落): %s", code, e) elif sd_accept and is_open: # 试探仓类采纳: 把数量重算成建议档 (弱基本面那类还带上紧止盈标记 advice.tight_trail)。 _apply_advice_tier(c, c.setdefault("hard_numbers", {}), verdict, force_trial=(sd.get("tier") == "trial")) if sd_accept: # 指令与账本理由一致。到价清仓把原句 (现价/目标价/股数) 接在后面, 复盘看得见到的是什么价 (09-14 评审补)。 _why = sd.get("why") or c["reason"] if sd.get("rule") == ae.DK_TARGET_PRICE and c.get("reason"): _why = f"{_why};{c['reason']}" c["reason"] = _why iid = _make_instruction(c, price, now) if sd_accept: reason = c["reason"] _dc = params.get("_sd_base_count", 0) + out.get("_sd_accepts", 0) hard = {**(c.get("hard_numbers") or {}), "self_decide": _self_decide_hard(c, verdict, sd, params, strat_cancelled, _dc, buys_cancelled=buys_cancelled)} else: reason = verdict.get("reason") or c["reason"] if auto_why: reason = f"判决候选自动执行: {reason}" hard = c["hard_numbers"] 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=hard, ref_id=iid, reason=reason[:500]) out["executed"].append({**_brief(c), "instruction_id": iid, "why": (sd.get("why") if sd_accept else ("减持自动执行 (规则触发, 来源没有要求交人裁决)" 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) → 来源强制的 (研究证据走弱那类减持) → # 深档补仓那条老规矩 → 研判不可用 → 档位。 # 重问放行仍排在最前 (2026-09-09 拍板): 它是这张提议最需要人知道的那件事 # (曾被驳回、什么变了、可能买不上)。硬风险 (2026-09-10 加) 紧随其后 —— 两者 # 实际上几乎不会同时出现: 上游标了硬风险的票判决多半不是「候选」, 走不到重问那一步, # 所以这个次序之争是空的, 不必为它推翻上一次拍板。 why = (reask_why or risk_why or 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 _apply_advice_tier(c, hn, verdict, *, force_trial=False): """按建议档位补研判一档、重算数量与仓位 (原 _make_proposal 内联那段抽出, 2026-09-14)。 自动路在落指令前也调它, 让自动采纳的试探仓真下的是试探仓的股数与分批, 弱基本面那类还带上 紧止盈标记 (advice.tight_trail)。force_trial=True 时把建议直接压到试探仓 —— 系统自决按 decide_kind 判要试探仓时, 研判把握度未必低, apply_judge 不会降档, 所以显式压一次。 原地改 hn 与 c["qty"], 返回 (hn, qty, is_trial)。建议算不出不拦 (与从前一样吞掉异常)。 """ is_trial = False try: from app.services import advice_service as _adv adv2 = _adv.apply_judge(hn.get("advice"), verdict, param_store.get_int("PMS_ADVICE_TRIAL_CONF_MIN", 60)) if force_trial and isinstance(adv2, dict) \ and adv2.get("tier") not in (_adv.TIER_NONE, _adv.TIER_TRIAL): adv2 = {**adv2, "tier": _adv.TIER_TRIAL, "factor": 1.0 / 6.0, "splits": (1.0,)} if isinstance(adv2, dict): hn["advice"], hn["advice_text"] = adv2, adv2.get("text") is_trial = adv2.get("tier") == _adv.TIER_TRIAL if adv2.get("needs_user_confirm"): hn["needs_user_confirm"] = True fix, resize_why = _adv.resize_to_advice(adv2, hn, c["ts_code"]) if fix: old_qty = c.get("qty") c["qty"] = fix["qty"] hn["target_pct"] = fix["target_pct"] hn["target_amount"] = fix["target_amount"] hn["base_amount"] = fix["base_amount"] hn["batch_scheme"] = fix["batch_scheme"] hn["advice_resized"] = {"from_qty": old_qty, "to_qty": fix["qty"], "tier": adv2.get("tier"), "note": fix.get("note") or ""} logger.info("提议按建议档位重算数量 %s: %s → %s 股 (%s)%s", c["ts_code"], old_qty, fix["qty"], adv2.get("tier"), "; " + fix["note"] if fix.get("note") else "") elif resize_why and "档位没变" not in resize_why: # 算不出来要留痕, 让人在卡上看得见「为什么数量还是原来那个」, # 而不是默默按旧数量下单。 hn["advice_resize_note"] = resize_why except Exception: # noqa: BLE001 —— 建议方案出错不拦 pass return hn, c.get("qty"), is_trial def _self_decide(c, verdict, params, today_auto_count) -> dict: """系统自决 (2026-09-14 台账 015): 一条候选归不归自决管、自决怎么判。 返回 {decision, tier, why, rule, would}。decision 取 accept / decline / confirm / none: none 表示不归自决管 (走旧路)。would 是未经档位降级 (record/decline) 的原判, 供观察读数记 「本会怎么判」。判定次序按规格书 §3.6 改动三, 决定表见 §3.5。 """ mode = str((params or {}).get("self_decide") or "off").strip() action, dk = c.get("action"), c.get("decide_kind") is_open = action == ae.A_OPEN is_dca_deep = action == ae.A_DCA and dk == ae.DK_DCA_DEEP is_target = action == ae.A_EXIT and dk == ae.DK_TARGET_PRICE # 一, 不归自决管: 关掉, 或不是新建仓/深亏补仓/目标价这三类 if mode == "off" or not (is_open or is_dca_deep or is_target): return {"decision": "none"} conf = verdict.get("confidence") trial_min = param_store.get_int("PMS_ADVICE_TRIAL_CONF_MIN", 60) dk_cn = ae.DK_CN.get(dk, "") def _mode(res): """按档位把原判降级: record 只记不做 (全变 none); decline 只放弃、不采纳 (accept 变 none)。 would 存原判供读数。""" res = {"tier": None, "rule": dk or "consensus_strong", **res} res["would"] = res["decision"] if mode == "record": res["decision"] = "none" elif mode == "decline" and res["decision"] == "accept": res["decision"] = "none" return res # 目标价到价: 不进研判范围, 直接自动清仓 (卖出方向, 落指令前再按可卖量夹紧、先撤方案) if is_target: return _mode({"decision": "accept", "rule": ae.DK_TARGET_PRICE, "why": "目标价到价自动清仓"}) # 二, 前置四种交人: 研判不可用 / 上游硬风险 / 合议装不上 / 判决词认不出 if verdict.get("degraded"): return _mode({"decision": "confirm", "why": "研判不可用,降级人工确认"}) if (c.get("hard_numbers") or {}).get("risk"): return _mode({"decision": "confirm", "why": "上游标了硬风险,交人裁决"}) if is_open and not (c.get("consensus") or (c.get("hard_numbers") or {}).get("consensus")): return _mode({"decision": "confirm", "why": "三源合议装不上,交人裁决"}) if dk == ae.DK_VERDICT_UNKNOWN: return _mode({"decision": "confirm", "why": "判决词认不出,交人裁决"}) # 三, 逻辑存疑: 自动放弃, 评审方 system if dk == ae.DK_LOGIC_DOUBT: return _mode({"decision": "decline", "why": "研究证据走弱(逻辑存疑),自动放弃"}) # 四, 真通过 (规格 §3.5): 只认带应答体的 PASS。动作不在研判范围时研判直接回 PASS 且没有应答体, # 那不是决策系统的结论, 不能据此自动入场 —— 交人 (09-14 评审补; 把 OPEN 从研判范围里拿掉时 # 全部新建仓都会落到这里, 方向是保守)。 if not isinstance(verdict.get("raw"), dict): return _mode({"decision": "confirm", "why": "研判没有真正应答(动作不在研判范围或无应答体),交人裁决"}) # 深亏补仓: 研判必答「杀逻辑还是杀情绪」, 契约规定杀逻辑或说不清一律 REJECT (REJECT 已在前面 # 早退当放弃)。所以走到这里的 PASS 即杀情绪: 把握度不低于 60 自动补, 否则放弃。 if is_dca_deep: if verdict.get("verdict") == judge.PASS and conf is not None and float(conf) >= trial_min: return _mode({"decision": "accept", "rule": ae.DK_DCA_DEEP, "why": f"深亏补仓:研判通过(把握度 {int(conf)}),按研判必答题自动补仓"}) return _mode({"decision": "decline", "rule": ae.DK_DCA_DEEP, "why": "深亏补仓:把握度不足,自动放弃"}) # 新建仓 (走到这: verdict 是真 PASS —— REJECT 与 degraded 都已在前面处理) conf_ok = conf is not None and float(conf) >= trial_min is_trial = (dk in ae.TRIAL_DECIDE_KINDS) or (not conf_ok) if is_trial: if not (params or {}).get("self_decide_trial"): return _mode({"decision": "confirm", "why": f"{dk_cn or '把握度偏低'}:试探仓自动采纳未开,交人裁决"}) why = ("系统采纳试探仓:" + (dk_cn or "把握度偏低") + (f";研判通过(把握度 {int(conf)})" if conf is not None else "")) res = {"decision": "accept", "tier": "trial", "why": why} else: res = {"decision": "accept", "tier": "standard", "why": "系统采纳:合议放行" + (f",研判通过(把握度 {int(conf)})" if conf is not None else "")} # 五, 每日自动新建仓上限: 超过转交人 (0 = 不限) dmax = int((params or {}).get("self_decide_daily_max") or 0) if res["decision"] == "accept" and is_open and dmax and today_auto_count >= dmax: return _mode({"decision": "confirm", "why": f"今日自动新建仓已达上限 {dmax} 只,本只转交人裁决"}) return _mode(res) def _sd_consensus_brief(c) -> dict: """从候选取合议方向/强弱/路由三项 (self_decide 块用)。硬数字里的合议六键或候选自带的都行。""" con = (c.get("hard_numbers") or {}).get("consensus") or c.get("consensus") or {} return {"direction": con.get("direction"), "strength": con.get("strength"), "route": con.get("route")} def _self_decide_hard(c, verdict, sd, params, strat_cancelled, daily_count, *, buys_cancelled=None) -> dict: """账本硬数字的 self_decide 块 (数据契约 §3.8): 决定种类、档、仓位档、合议三项、研判、当日计数、 撤掉的方案与在途买单。""" tier_cn = {"trial": "试探仓", "standard": "标准仓", "half": "减半仓"}.get(sd.get("tier")) return {"rule": sd.get("rule"), "mode": str((params or {}).get("self_decide") or "off"), "tier": tier_cn or sd.get("tier"), "consensus": _sd_consensus_brief(c), "judge": {"verdict": verdict.get("verdict"), "confidence": verdict.get("confidence")}, "daily_count": daily_count, "daily_max": int((params or {}).get("self_decide_daily_max") or 0), "strategy_cancelled": list(strat_cancelled or []), "buys_cancelled": list(buys_cancelled or [])} def _sd_decline_ledger(code, action, c, verdict, price, sd) -> None: """系统自决放弃落账 (改动六): 评审方 system, 结论 REJECT, 硬数字带原硬数字、合议六键、 研判把握度与 self_decide 块。研判 REJECT 引起的放弃走的是研判闸那条 (评审方 judge), 不到这里。""" hard = {**(c.get("hard_numbers") or {}), "judge_conf": verdict.get("confidence"), "self_decide": {"rule": sd.get("rule"), "decision": "decline", "consensus": _sd_consensus_brief(c), "judge": {"verdict": verdict.get("verdict"), "confidence": verdict.get("confidence")}}} pms_repo.insert_ledger(ts_code=code, action=action, arbiter="system", verdict="REJECT", price_at=price, hard_numbers=hard, reason=(sd.get("why") or "系统自决放弃")[:500]) 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"), # 两个期限的头 (2026-09-08 量价研判链): 决策系统并列给的 5 日与 20 日判断, 只显示不触发。 "judge_pv_heads": verdict.get("pv_heads")} # 建议方案补研判一档 (2026-09-09 接入方案): 把握度低或不可用 → 试探仓, 只在人工确认后下。 # # 2026-09-10 补上数量重算。原先这里只改档位文字、不改数量 —— 候选产出时 advise() # 的档位是进过数量的, 但研判回来之后这一降档没跟上, 于是卡上写「试探仓, 目标 1.0%」 # 而采纳后真下的是标准仓 6% 的第一批。实测 8 张待确认提议里 7 张是这个样子, 差三倍。 # 采纳按钮那边直接抄提议里的股数原样落单 (只有卖出侧会按可卖量重算), 所以不在这里 # 改, 后面就没有第二次机会了。 _apply_advice_tier(c, hn, verdict) 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), # 逻辑状态四态的三个旋钮 (2026-09-07 第三件): 分流开关 (持仓停增持 + 新建仓强制确认)、 # 减持提议开关 (默认关)、减持比例 (默认与保垫减仓相同的三分之一)。 "logic_state_route": param_store.get_bool("PMS_LOGIC_STATE_ROUTE", True), "open_route_by_logic": param_store.get_bool("PMS_LOGIC_STATE_ROUTE", True), # 建议方案直接作为新建仓的仓位档与分批 (2026-09-09 接入方案第六之二节; 模拟侧默认开) "plan_by_advice": param_store.get_bool("PMS_PLAN_BY_ADVICE", True), "logic_doubt_trim": param_store.get_bool("PMS_LOGIC_DOUBT_TRIM", False), "logic_doubt_trim_ratio": param_store.get_float("PMS_LOGIC_DOUBT_TRIM_RATIO", 1.0 / 3), # 三源合议 (2026-09-11 工作包二): 合议参与新建仓分流 (scan_open 读它), 持仓增持门 (scan 读它)。 # 两个都关掉时, 合议既不装配也不分流、增持门不设 —— 行为与接入前逐字相同。 "consensus_route": param_store.get_bool("PMS_CONSENSUS_ROUTE", True), "tech_gate_increase": param_store.get_bool("PMS_TECH_GATE_INCREASE", True), # 系统自决 (2026-09-14 台账 015): 档位 off/record/decline/full 由 action_engine 与 # _self_decide 读; off 时 scan_open/eval_dca/eval_target 的自决分支一行不走, 逐字回旧。 "self_decide": str(param_store.get("PMS_SELF_DECIDE", "off") or "off").strip(), "self_decide_trial": param_store.get_bool("PMS_SELF_DECIDE_TRIAL", False), "self_decide_daily_max": param_store.get_int("PMS_SELF_DECIDE_DAILY_MAX", 3), # 技术面转空自动离场 (2026-09-11 工作包三, 台账 010)。档位 off/propose_only/full: # off 不评 (tech_exit_on 为假)、propose_only 交人、full 自动执行 (与保垫减仓同档)。 # tech_exit_done 去重集在 scan_and_route 里按当轮技术面映射填 (这里给空占位)。 "tech_exit_on": (str(param_store.get("PMS_TECH_EXIT_AUTONOMY", "full") or "full").strip() in ("full", "propose_only")) and param_store.get_bool("PMS_TECH_ENABLED", True), "tech_exit_propose_only": str(param_store.get("PMS_TECH_EXIT_AUTONOMY", "full") or "full").strip() == "propose_only", "tech_exit_trim_ratio": param_store.get_float("PMS_TECH_EXIT_TRIM_RATIO", 1.0 / 3), "tech_exit_fresh_days": param_store.get_int("PMS_TECH_FLIP_FRESH_DAYS", 2), "tech_exit_done": set(), }) 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 _system_declined_today_keys() -> dict: """今天已被**系统自决放弃**过的新建仓等 {(代码, 动作): 原因} —— 当日不再重评, 明日重评。 评审方 system 的 REJECT (逻辑存疑、合议引起的自动放弃) 就是这类。研判 REJECT 引起的放弃 评审方是 judge, 走 _judge_rejected_open_keys 那条 (还要参与解锁重问), 不在此列。 off 时账本里没有 system 放弃行, 自然返回空 —— 逐字回旧。读不到就返回空 (去重是降噪)。""" try: keys = pms_repo.system_decided_today(datetime.now().replace( hour=0, minute=0, second=0, microsecond=0)) except Exception as e: # noqa: BLE001 logger.warning("读当日系统自决放弃记录失败 (按未放弃过继续扫描): %s", e) return {} return {k: "今日系统已判放弃,明日重评" for k in keys} 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"]}