# -*- coding: utf-8 -*- """ 决策系统盘中信号的解析与消化规则 (纯逻辑, 无外部依赖, 可单测) ================================================================ 设计 POSITION_MGMT_DESIGN.md §10「信号消化」与 §1: 「风控 SELL、盘中 ENTRY/EXIT 广播照常产出, 但**下游停止直接执行**, 改由 PMS 订阅消化后统一决定卖出指令。」 两条流的真实格式 (2026-07-28 从 trading_service 的两个消费者实测确认): 买入/盘中信号 Redis db2, key = `intraday_signals:{YYYY-MM-DD}`, 每日一条流 扁平字段: ts_code / action(BUY|SELL|HOLD) / confidence(**0~1**) / component_scores(JSON 字符串, 内含 minute_qrs) / suggested_price 风控卖出信号 Redis db3, key = `bionic:signals:llm_sell_actions`, 固定 key 外层含 data(JSON 字符串), 内层: ts_code / action(SELL) / confidence(**0~100**) / dominant_signal / llm_reason 两条流的置信度**尺度不同**(0~1 与 0~100), 这是最容易踩的坑, 统一在 parse 里归一到 0~1。 消化口径 (PMS 侧): * SELL 信号 —— 只对**持有的票**有意义。置信度够高即转卖出动作 (减持方向不设确认门槛, 与保垫减仓同一口径); 置信度中等则落提议队列等用户裁决。 * BUY 信号 —— **本模块仍然不产生买入动作**, 但 2026-08-06 起**无论持没持仓都要留痕**。 * HOLD 等其余 —— 只给持仓票留痕 (口径未变)。 2026-08-06 补的那条留痕, 由来值得写下来 ------------------------------------------ 决策系统盘中判出 REVERSAL_BUY 时会往 db2 这条流广播 (action=BUY / signal_type=ENTRY), PMS 一直订阅得到, 但走到 digest 就被归进 ACT_RECORD, 而 signal_service 的 RECORD 分支 **只给持仓票写账本**。于是「决策系统今天看多了某只没持仓的票」这件事, PMS 收到了、计了个数, 然后一个字都不留 —— 账本查不到、页面看不见, 事后复盘问「那天系统看见了吗」答不上来。 这与 2026-08-05 那次「同一条结论一边当判决一边当摆设」是同一个形状, 只是换了个入口。 当时那句注释的理由 (「买什么买多少由 PMS 的命令与动作引擎决定」) 在动作引擎**没有新建仓 动作**的时候是成立的 —— PMS 确实没有能力消化一个买入信号。动作引擎补上 OPEN 之后这条理由 就失效了, 所以先把留痕补上: 本模块只负责「记下来」, 要不要据此建仓由动作引擎那条路决定 (proposal_service 读账本里这些痕迹, 给候选排序时优先, 见那边的说明)。 **分工没有变**: 信号消化不下买单, 买什么买多少仍然归动作引擎。 """ from __future__ import annotations import json SRC_INTRADAY, SRC_RISK_SELL = "intraday", "risk_sell" ACT_EXIT, ACT_PROPOSE, ACT_RECORD, ACT_IGNORE = "EXIT", "PROPOSE", "RECORD", "IGNORE" # 买入信号留痕: 不产生任何买入动作, 但**持没持仓都要写账本** (与 ACT_RECORD 的差别就在这) ACT_NOTE_BUY = "NOTE_BUY" # 信号来源 (2026-09-03 起区分)。db2 那条盘中流**两家都写**: 择时决策系统的大脑广播 # (producer_id 形如 bionic_brain_intraday_v2.0) 与盘中择时程序的入场信号 (producer_id 形如 # intraday_timing_v0.1.0), 格式相近而语义不同 —— 前者是知情二审的结论, 后者是量化程序的 # 触发。此前 PMS 一律记成「决策系统盘中判该股转多」, 于是新建仓的插队排序把择时程序的 # 触发也当成了决策系统的结论。现在按 producer_id 前缀分开: 只有以 bionic 开头的才算 # 择时决策系统; 缺失一律记作 unknown, 不猜。 PRODUCER_UNKNOWN = "unknown" BIONIC_PRODUCER_PREFIX = "bionic" # 留痕 reason 的固定开头。动作引擎那条路 (proposal_service._buy_signals_today) 靠这个开头 # 认出「哪几条是择时决策系统判的」—— 账本查询只回 reason 不回硬数字, 用开头判最省一次改表。 BUY_NOTE_PREFIX_BIONIC = "决策系统盘中判该股转多" BUY_NOTE_PREFIX_OTHER = "盘中择时程序买入信号" def is_bionic_producer(producer_id) -> bool: return str(producer_id or "").strip().lower().startswith(BIONIC_PRODUCER_PREFIX) def is_bionic_buy_note(reason) -> bool: """账本里一条 SIGNAL_BUY 留痕是不是择时决策系统判的 (看 reason 开头)。""" return str(reason or "").startswith(BUY_NOTE_PREFIX_BIONIC) def _num(v, d=0.0): try: return float(v) except (TypeError, ValueError): return d def _norm_conf(v) -> float: """盘中流 (0~1 制) 的置信度归一。大于 1 的按百分制兜底处理。""" c = _num(v) if c > 1: c = c / 100.0 return max(0.0, min(1.0, c)) def _norm_conf_pct(v) -> float: """风控流 (db3) 的置信度归一: **文档口径固定 0~100 制, 一律除以 100**。 原来两条流共用 _norm_conf 的"大于 1 才除"启发式, 0~100 制里 (0,1] 的取值 (即 0~1% 这种噪声级置信度) 不会被缩放, 1 会被当成 100% 而直接触发自动清仓 —— 方向恰好反了: 越低的置信被解释得越高 (2026-08-28 审查修)。 既然文档定死了百分制, 就不做"像不像小数"的猜测; 上游若真发 0~1 制小数, 会被折成不足 1% 而落进"低于门槛只留痕"一档, 错的方向是漏报不是误清仓。 """ c = _num(v) return max(0.0, min(100.0, c)) / 100.0 def parse_intraday(fields: dict, *, msg_id: str = None) -> dict: """买入/盘中信号流 (db2) 的一条消息 → 归一结构。""" f = fields or {} scores = {} raw_scores = f.get("component_scores") if raw_scores: try: scores = json.loads(raw_scores) if isinstance(raw_scores, str) else dict(raw_scores) except (json.JSONDecodeError, TypeError, ValueError): scores = {} return {"source": SRC_INTRADAY, "msg_id": msg_id, "ts_code": (f.get("ts_code") or "").strip(), "action": (f.get("action") or "").strip().upper(), "confidence": _norm_conf(f.get("confidence")), "minute_qrs": _num(scores.get("minute_qrs")), "suggested_price": _num(f.get("suggested_price")) or None, "reason": f.get("reason") or f.get("llm_reason") or "", "dominant_signal": f.get("dominant_signal") or "", # 发送方标识 (2026-09-03 起保留): 缺就是 unknown, 不猜是谁发的 "producer_id": str(f.get("producer_id") or "").strip() or PRODUCER_UNKNOWN} def parse_risk_sell(fields: dict, *, msg_id: str = None) -> dict: """风控卖出信号流 (db3) 的一条消息 → 归一结构。外层套一层 data JSON 字符串。""" f = fields or {} inner = f raw = f.get("data") if raw: try: inner = json.loads(raw) if isinstance(raw, str) else dict(raw) except (json.JSONDecodeError, TypeError, ValueError): return {"source": SRC_RISK_SELL, "msg_id": msg_id, "ts_code": "", "action": "", "confidence": 0.0, "parse_error": "内层 data JSON 解析失败", "reason": "", "dominant_signal": ""} return {"source": SRC_RISK_SELL, "msg_id": msg_id, "ts_code": (inner.get("ts_code") or "").strip(), "action": (inner.get("action") or "").strip().upper(), "confidence": _norm_conf_pct(inner.get("confidence")), "dominant_signal": inner.get("dominant_signal") or "", "reason": (inner.get("llm_reason") or inner.get("reason") or "")[:500], "suggested_price": _num(inner.get("suggested_price")) or None} def digest(signal: dict, position: dict, params: dict) -> dict: """一条信号 → 一个消化结论。 position: PMS 账本里这只票的快照 (无持仓传 None 或 total_qty=0) params: {sell_conf_min, auto_exit_conf, trim_ratio} 返回 {"action": EXIT|PROPOSE|RECORD|IGNORE, "qty", "reason", "hard_numbers"} """ code = (signal or {}).get("ts_code") or "" act = (signal or {}).get("action") or "" conf = _num((signal or {}).get("confidence")) hard = {"source": signal.get("source"), "confidence": round(conf, 4), "dominant_signal": signal.get("dominant_signal"), "minute_qrs": signal.get("minute_qrs"), "suggested_price": signal.get("suggested_price")} if not code: return _r(ACT_IGNORE, 0, "信号缺少股票代码", hard) if act == "BUY": # 买入信号**仍然不产生买入动作** —— 买什么买多少归动作引擎。 # 但无论持没持仓都要留痕: 未持仓的票正是新建仓关心的那一批, 从前它们连账本都没有。 # 留痕文案按来源分写 (2026-09-03): 择时决策系统的保持原句「决策系统盘中判该股转多」; # 其他发送方 (盘中择时程序、未知) 写「盘中择时程序买入信号(来源 xxx)」—— # 动作引擎的插队排序只认前者, 所以这个开头就是两类的分界线, 别改。 held_now = int((position or {}).get("total_qty") or 0) hard["held"] = held_now producer = str(signal.get("producer_id") or "").strip() or PRODUCER_UNKNOWN hard["producer_id"] = producer head = (BUY_NOTE_PREFIX_BIONIC if is_bionic_producer(producer) else f"{BUY_NOTE_PREFIX_OTHER}(来源 {producer})") return _r(ACT_NOTE_BUY, 0, f"{head} (置信度 {conf:.0%}" + (f", 建议价 {signal.get('suggested_price')}" if signal.get("suggested_price") else "") + f"){'; 该股当前有持仓' if held_now > 0 else '; 该股当前无持仓'} —— " f"只留痕, 买不买由动作引擎按闸门决定: " + (signal.get("reason") or signal.get("dominant_signal") or "未给理由"), hard) if act != "SELL": # HOLD 等其余类型: 口径未变, 只给持仓票留痕 return _r(ACT_RECORD, 0, f"{act or '未知'} 信号仅作择时参考留痕, PMS 不据此买入", hard) held = int((position or {}).get("total_qty") or 0) if held <= 0: return _r(ACT_IGNORE, 0, "未持有该票, 卖出信号无对象", hard) conf_min = _num(params.get("sell_conf_min"), 0.75) auto_conf = _num(params.get("auto_exit_conf"), 0.85) if conf < conf_min: return _r(ACT_IGNORE, 0, f"置信度 {conf:.0%} < 消化门槛 {conf_min:.0%}, 不动", hard) avail = int((position or {}).get("avail_qty") or 0) hard.update({"total_qty": held, "avail_qty": avail}) why = signal.get("reason") or signal.get("dominant_signal") or "决策系统风控卖出" if conf >= auto_conf: # 高置信风控卖出 = 清仓。减持方向不设确认门槛 (与保垫减仓同一口径) return _r(ACT_EXIT, held, f"风控 SELL 置信度 {conf:.0%} ≥ {auto_conf:.0%}, 清仓 {held} 股 —— {why}", hard) ratio = _num(params.get("trim_ratio"), 1 / 3) # 四舍五入到一手, 不用向下取整: 配置里写 0.3333 还是 1/3 不该让 3000 股的三分之一 # 一会儿算成 1000 一会儿算成 900。不足一手时退化为全卖 (一手是最小可操作单位)。 qty = int(round(held * ratio / 100)) * 100 if qty <= 0: qty = held # 科创板部分卖出最少 200 股: 量够就抬到 200 (偏保守多卖一点), 不够就退化为全卖 if code.startswith(("688", "689")) and 0 < qty < 200 and qty != held: qty = 200 if held >= 200 else held return _r(ACT_PROPOSE, qty, f"风控 SELL 置信度 {conf:.0%} 介于 {conf_min:.0%}~{auto_conf:.0%}, " f"提议减 {qty} 股待确认 —— {why}", hard) def _r(action, qty, reason, hard): return {"action": action, "qty": int(qty), "reason": reason, "hard_numbers": hard} def dedup_key(signal: dict, ymd) -> str: """当日去重键: 同一只票、同一来源、同一动作, 一天只消化一次。""" return f"{ymd}:{signal.get('source')}:{signal.get('ts_code')}:{signal.get('action')}"