224 lines
12 KiB
Python
224 lines
12 KiB
Python
# -*- 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')}"
|