diff --git a/README.md b/README.md index 2ed2143..870bd13 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,7 @@ app/ rule_gate.py 规则闸终检: 上限/一手/可卖/冻结/刹车/行业/不追高 (减持只放行不阻拦) exec_timing.py 择时实现B: 分日配额 / 分笔 / 买卖出手判定 / 14:45 兜底 / 窗口收口 action_engine.py 动作引擎: FILL 回踩补足 / ADD 盈利加仓 / DCA 补仓 / TRIM 保垫减仓 + signal_rules.py 决策系统两条信号流的解析与消化口径 (含置信度尺度归一) tradedays.py 交易日历: 调度守卫与执行窗口计算 db/session.py 三库连接 + **严格单表访问守卫** (JOIN/逗号连表/跨表子查询一律拒绝) repo/ 单表数据访问: pms_repo (自有 10 表) / downstream_repo (下游只读三表) @@ -39,6 +40,7 @@ app/ dispatcher.py 下发通道三适配器: shadow(默认) / plan_x / channel_y proposal_service.py 自主提议: 扫描→规则闸→研判闸→按自主档位分流 (执行/入队) judge.py 研判闸客户端 (决策系统未接通时自动降级为人工确认) + signal_service.py 盘中信号订阅 (db2 广播 + db3 风控卖出) → 卖出指令或提议 ledger_service.py 成交回放 / 对账 / 除权 / 盘前 / 日终结算 / 运营日报 market.py 行情 (Redis db13) 与参考位 (决策系统主口径 + 兜底自算) industry.py 行业划分可插拔适配器 (custom_table / gp_stock_category / 停用) @@ -50,7 +52,8 @@ scripts/ test_batch2_units.py 命令 / 方案 / 回放对账 纯逻辑 35 例 test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 18 例 test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 11 例 - test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 30 例 + test_batch5_units.py 决策系统信号流解析与消化口径 8 例 + test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 32 例 init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练) check_db.py 实机连通性与表结构自检 (需真实 .env) ``` @@ -101,7 +104,7 @@ git pull && docker compose build && docker compose up -d | 命令轮询 | 每 1 分钟(全天) | 新命令解析 → 方案生成 → 状态机推进 | ✅ | | 成交回放 | 交易时段每 5 分钟 | `trading_order` 增量回放 + 盘中轻对账 | ✅ | | 盘中执行 | 交易时段每 1 分钟 | 方案转指令 → 自主提议扫描 → 择时出手(规则闸终检 → 下发 → 记子单) | ✅ | -| 信号消化 | 交易时段每 1 分钟 | 订阅决策系统盘中信号 | 🔜 下一批 | +| 信号消化 | 交易时段每 1 分钟 | 订阅 db2 盘中广播 + db3 风控卖出 → 卖出指令或提议 | ✅ | | T 仓平回 | 14:50 | 做T强制平回 | 🔜 二期(现只自证 T 仓为 0) | | 日终结算 | 15:10 | 除权检测 / 全量对账 / 安全垫 / 命令进度日结 | ✅ | | 运营日报 | 15:30 | 关注区 + 全量统计(页面「日报」按钮可查) | ✅ | @@ -112,9 +115,11 @@ git pull && docker compose build && docker compose up -d | 模式 | 行为 | 什么时候用 | |---|---|---| -| `shadow`(默认) | 指令照常过规则闸、照常置 DISPATCHED,但**不写下游**。你在 QMT 侧人工执行,成交由回放按 FIFO 认领回账本 | 通道协商完成前的一期口径(设计 §9:命令类降仓/清仓由用户人工执行、PMS 记账跟踪) | -| `plan_x` | 买入写 `trading_buy_plan`(`is_active=6` 待挂单、署名 `approved_by='pms'`);**卖出无对应通道,自动退回影子** | QMT 侧确认沿用旧通道过渡时 | -| `channel_y` | 写统一指令表 `pms_order_request`(DDL 见需求清单 B1) | B1 协商落地、表建好之后 | +| `shadow`(默认) | 指令照常过规则闸、照常置 DISPATCHED,但**不写下游**。你在 QMT 侧人工执行,成交由回放按 FIFO 认领回账本 | 直连服务就绪前的一期口径(设计 §9:命令类降仓/清仓由用户人工执行、PMS 记账跟踪) | +| `plan_x` | 买入写 `trading_buy_plan` | ⚠️ **已作废**:架构已定 trading_service 全量退出业务,此模式不再使用(代码暂留,勿在实盘开启) | +| `channel_y` | 写统一指令表 `pms_order_request` | 新 QMT 直连服务就绪后启用,由它消费本表 | + +> **目标架构(2026-07-28 已定)**:`trading_service` 全量退出业务,只保留看板与统计展示;新写一个 QMT 直连服务承担挂单与订单/持仓/资金回写;PMS 只管决策与账本。三者之间的数据接口待协定后另行成文。 影子模式下的完整闭环:页面下命令 → 方案落表 → 方案转指令 → 择时按日配额给出「今天该出多少、什么价」→ 你照着在 QMT 下单 → 5 分钟一次的回放把成交认领回批次账本 → 命令进度自动推进。整条链路除了「人手下单」这一步,其余与实盘接管后完全一致。 @@ -130,9 +135,9 @@ git pull && docker compose build && docker compose up -d ## 已实现 / 待开发 -**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口);三模式下发通道;**动作引擎四类自主动作 + 研判闸客户端 + 提议分流**;管理页面四块 + 运维/日报抽屉;调度器八个调度位;单测 108 例。 +**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口);三模式下发通道;**动作引擎四类自主动作 + 研判闸客户端 + 提议分流**;**决策系统信号消化**(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器八个调度位;单测 118 例。 -**待开发(下一批)**:决策系统盘中信号订阅(风控 SELL / 止盈 / 反转 → 卖出方案或提议)、择时实现 A(委托决策系统盘中择时)、T0 做T(二期)。研判闸客户端已就位,等 bionic 侧 `process_intraday_audit` 新增 PMS 请求 direction 后,在页面填 `PMS_JUDGE_API_BASE` 即接通。 +**待开发**:T0 做T(二期)、择时实现 A(委托决策系统盘中择时,等 bionic 侧接口)、新 QMT 直连服务的对接(等接口协定)。研判闸客户端已就位,等 bionic 侧 `process_intraday_audit` 新增 PMS 请求 direction 后,在页面填 `PMS_JUDGE_API_BASE` 即接通。 **待外部协商**:`QMT_INTERFACE_REQUIREMENTS.md` 的 A/B/C/D 各项——尤其 A1(`trading_position` 完整 DDL 与可用数量列)、A2(`trading_order` 状态枚举与**来源标识**)、B1(统一指令通道)。在来源标识到位前,回放按「同股同向 + 下发早于成交 + FIFO」贪心认领,认领不上即判外部成交并告警;持仓数量列用候选名探测,探测结果可经页面「运维 → 导出下游表结构」查看,也是回填 D1 的现成材料。 diff --git a/app/core/signal_rules.py b/app/core/signal_rules.py new file mode 100644 index 0000000..2865f3f --- /dev/null +++ b/app/core/signal_rules.py @@ -0,0 +1,146 @@ +# -*- 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 / HOLD 信号 —— **不产生买入动作**。买什么、买多少是 PMS 自己的命令与动作引擎说了算 + (设计: 持仓系统管「做什么、多少」)。这类信号只作为择时参考落痕, 不越权。 +""" +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" + + +def _num(v, d=0.0): + try: + return float(v) + except (TypeError, ValueError): + return d + + +def _norm_conf(v) -> float: + """置信度归一到 0~1。两条流一条给 0~1 一条给 0~100, 大于 1 的一律按百分制处理。""" + c = _num(v) + if c > 1: + c = c / 100.0 + return max(0.0, min(1.0, c)) + + +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 ""} + + +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(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")} + + if not code: + return _r(ACT_IGNORE, 0, "信号缺少股票代码", hard) + if act != "SELL": + # 买入/持有类信号不产生动作 —— 买什么买多少由 PMS 的命令与动作引擎决定 + 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 + 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')}" diff --git a/app/scheduler.py b/app/scheduler.py index d64e3e8..a6e98c0 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -146,8 +146,9 @@ def intraday_exec(): @celery_app.task(name="pms.signal_digest") @guard(trade_day=True, session=True) def signal_digest(): - """信号消化: 订阅决策系统盘中信号 (风控 SELL/止盈/反转) → 卖出方案或提议。下一批交付。""" - return {"consumed": 0, "note": "决策系统信号订阅为下一批交付 (设计 §10 信号消化)"} + """信号消化: 订阅决策系统盘中信号 (db2 盘中广播 + db3 风控 LLM 卖出) → 卖出指令或提议。""" + from app.services import signal_service + return signal_service.consume() @celery_app.task(name="pms.t0_close") diff --git a/app/services/param_store.py b/app/services/param_store.py index b596a54..dd9f65b 100644 --- a/app/services/param_store.py +++ b/app/services/param_store.py @@ -80,6 +80,11 @@ DESC = { "PMS_T0_CLOSE_TIME": "T仓强制平回时点", "PMS_T0_STOCK_DAY_LOSS": "单票当日T亏熔断", "PMS_T0_GLOBAL_DAY_LOSS": "全局当日T亏熔断", "PMS_REPLAY_INTERVAL_MIN": "成交回放间隔 (分钟)", "PMS_RECON_ALARM_DAYS": "连续不一致升级天数", + "PMS_SIGNAL_ENABLED": "是否消化决策系统盘中信号", + "PMS_SIGNAL_GROUP": "信号消费组名 (独立于 trading_service, 互不抢消息)", + "PMS_SIGNAL_SELL_CONF_MIN": "卖出信号消化门槛 (低于此不动)", + "PMS_SIGNAL_AUTO_EXIT_CONF": "卖出信号直接清仓门槛 (之间则落提议)", + "PMS_SIGNAL_TRIM_RATIO": "中等置信度卖出信号的减仓比例", } # loaded 标记必不可少: 不能用「data 是否为空」判断缓存是否有效 —— diff --git a/app/services/signal_service.py b/app/services/signal_service.py new file mode 100644 index 0000000..2def840 --- /dev/null +++ b/app/services/signal_service.py @@ -0,0 +1,266 @@ +# -*- coding: utf-8 -*- +""" +决策系统盘中信号消化 (设计 §10「信号消化」) +============================================= +订阅两条流, 转成 PMS 自己的卖出动作或提议: + + db2 `intraday_signals:{YYYY-MM-DD}` 盘中 BUY/SELL/HOLD 广播 (每日一条流) + db3 `bionic:signals:llm_sell_actions` 风控 LLM 卖出动作 (固定 key) + +**用独立消费组** (`pms_signal_consumer`), 与 trading_service 的 `qmt_main_activator` / +`qmt_sell_activator` 互不抢消息 —— Redis Stream 的消费组之间各自看到全量消息, +所以 PMS 可以和现有消费者并行订阅, 迁移期两边都能跑。 + +不做常驻进程: 由调度器 `signal_digest` 每分钟拉一批, 与其余任务同一套守卫和降级口径。 +解析与消化规则在 core/signal_rules.py (纯逻辑), 本模块只管连 Redis、落表、留痕。 +""" +from __future__ import annotations + +import json +import logging +from datetime import datetime, timedelta + +from config.settings import settings +from app.core import command_spec as cs +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 executor, param_store, portfolio + +logger = logging.getLogger("pms.signal") + +SEEN_KEY = "PMS_SIGNAL_SEEN" # 当日去重集合 (JSON), 日切自动作废 +_clients = {} + + +def _client(db: int): + """Redis 客户端。强制 RESP2 —— 与行情库同一个坑 (服务端 <6.0 不认 HELLO)。""" + if db in _clients: + return _clients[db] + import redis + kw = dict(host=settings.SIGNAL_REDIS_HOST, port=settings.SIGNAL_REDIS_PORT, + password=settings.SIGNAL_REDIS_PASSWORD or None, db=db, + decode_responses=True, socket_timeout=settings.SIGNAL_REDIS_SOCKET_TIMEOUT) + try: + c = redis.Redis(protocol=2, **kw) + except TypeError: + c = redis.Redis(**kw) + _clients[db] = c + return c + + +def group_name() -> str: + return param_store.get("PMS_SIGNAL_GROUP", "pms_signal_consumer") or "pms_signal_consumer" + + +def streams() -> list: + """[(db, key, parser)] —— 盘中流按日期拼 key, 风控流是固定 key。""" + ymd = datetime.now().strftime("%Y-%m-%d") + tpl = param_store.get("PMS_SIGNAL_STREAM_INTRADAY", "intraday_signals:{ymd}") + sell_key = param_store.get("PMS_SIGNAL_STREAM_SELL", "bionic:signals:llm_sell_actions") + return [(settings.SIGNAL_REDIS_DB_INTRADAY, tpl.format(ymd=ymd), sr.parse_intraday), + (settings.SIGNAL_REDIS_DB_ACTIONS, sell_key, sr.parse_risk_sell)] + + +def status() -> dict: + """页面用: 两条流的连通性与积压情况。""" + out = {"enabled": param_store.get_bool("PMS_SIGNAL_ENABLED", True), + "group": group_name(), "streams": []} + for db, key, _ in streams(): + item = {"db": db, "key": key} + try: + c = _client(db) + item["length"] = c.xlen(key) + groups = c.xinfo_groups(key) + mine = [g for g in groups if g.get("name") == group_name()] + item["pending"] = mine[0].get("pending") if mine else None + item["group_ready"] = bool(mine) + item["other_groups"] = [g.get("name") for g in groups + if g.get("name") != group_name()] + except Exception as e: + item["error"] = f"{type(e).__name__}: {e}" + out["streams"].append(item) + return out + + +# ================================================================ 消费 +def consume(*, batch: int = 50, dry_run: bool = False) -> dict: + """拉一批信号并消化。每分钟一跳, 幂等 (消费组 ACK + 当日去重)。""" + out = {"ok": True, "read": 0, "exits": [], "proposals": [], "recorded": 0, + "ignored": 0, "errors": [], "dry_run": dry_run} + if not param_store.get_bool("PMS_SIGNAL_ENABLED", True): + out["skipped"] = "信号消化已关闭 (PMS_SIGNAL_ENABLED=False)" + return out + + try: + view = portfolio.positions_view() + except Exception as e: + return {**out, "ok": False, "errors": [f"读账本失败: {type(e).__name__}: {e}"]} + + prm = {"sell_conf_min": param_store.get_float("PMS_SIGNAL_SELL_CONF_MIN", 0.75), + "auto_exit_conf": param_store.get_float("PMS_SIGNAL_AUTO_EXIT_CONF", 0.85), + "trim_ratio": param_store.get_float("PMS_SIGNAL_TRIM_RATIO", 1 / 3)} + seen = _load_seen() + ymd = td.ymd() + + for db, key, parser in streams(): + try: + msgs = _read(db, key, batch) + except Exception as e: + out["errors"].append(f"{key} 读取失败: {type(e).__name__}: {e}") + continue + out["read"] += len(msgs) + for msg_id, fields in msgs: + try: + sig = parser(fields, msg_id=msg_id) + _handle(sig, view, prm, seen, ymd, dry_run, out) + if not dry_run: + _ack(db, key, msg_id) + except Exception as e: + logger.exception("信号处理失败 %s", msg_id) + out["errors"].append(f"{msg_id}: {type(e).__name__}: {e}") + + if not dry_run: + _save_seen(seen, ymd) + out["ok"] = not out["errors"] + return out + + +def _handle(sig, view, prm, seen, ymd, dry_run, out): + code = sig.get("ts_code") + pos = _pos_of(view, code) if code else None + d = sr.digest(sig, pos, prm) + act = d["action"] + + if act == sr.ACT_IGNORE: + out["ignored"] += 1 + return + if act == sr.ACT_RECORD: + # 只给持有的票留痕, 否则全市场广播会把评审账本冲垮 + if pos and int(pos.get("total_qty") or 0) > 0 and not dry_run: + pms_repo.insert_ledger(ts_code=code, action="SIGNAL", arbiter="rule", + verdict="PASS", price_at=float(pos.get("price") or 0), + hard_numbers={**d["hard_numbers"], "msg_id": sig.get("msg_id")}, + reason=d["reason"][:500]) + out["recorded"] += 1 + return + + key = sr.dedup_key(sig, ymd) + if key in seen: + out["ignored"] += 1 + return + if _has_inflight(code): + out["ignored"] += 1 + out.setdefault("skipped_inflight", []).append(code) + return + + brief = {"ts_code": code, "qty": d["qty"], "confidence": d["hard_numbers"]["confidence"], + "reason": d["reason"]} + if dry_run: + (out["exits"] if act == sr.ACT_EXIT else out["proposals"]).append( + {**brief, "dry_run": True}) + return + + seen.add(key) + if act == sr.ACT_EXIT: + iid = _make_exit(code, d, pos) + out["exits"].append({**brief, "instruction_id": iid}) + else: + pid = _make_proposal(code, d, pos, sig) + out["proposals"].append({**brief, "proposal_id": pid}) + + +def _make_exit(code, d, pos) -> str: + """高置信风控卖出 → 直接落卖出指令 (减持方向不设确认门槛)。""" + now = datetime.now() + iid = cs.make_instruction_id(td.ymd(now), code, "EXIT", int(now.strftime("%H%M%S")) % 1000) + window = param_store.get_int("PMS_EXEC_WINDOW_TDAYS", 3) + pms_repo.insert_instruction( + instruction_id=iid, origin_type="system", origin_id=d["hard_numbers"].get("source"), + ts_code=code, action="EXIT", side="sell", qty=d["qty"], limit_price=None, + window_tdays=window, status=executor.ST_PROPOSED, + progress={"deadline": str(td.window_deadline(now.date(), window)), + "is_command": False, "children": [], "from_signal": True, + "reason": d["reason"]}) + pms_repo.insert_ledger(ts_code=code, action="EXIT", arbiter="rule", verdict="PASS", + price_at=float((pos or {}).get("price") or 0), + hard_numbers=d["hard_numbers"], ref_id=iid, + reason=d["reason"][:500]) + logger.warning("[信号消化] %s 转清仓指令 %s —— %s", code, iid, d["reason"]) + return iid + + +def _make_proposal(code, d, pos, sig) -> str: + ttl = param_store.get_int("PMS_PROPOSAL_TTL_HOURS", 24) + pid = f"PRP_{td.ymd()}_{code.replace('.', '')}_SIGSELL" + hn = {**d["hard_numbers"], "price": float((pos or {}).get("price") or 0), + "reason": d["reason"], "signal_source": sig.get("source")} + pms_repo.insert_proposal(proposal_id=pid, ts_code=code, action="TRIM", qty=d["qty"], + hard_numbers=hn, + expire_at=datetime.now() + timedelta(hours=ttl), + judge_verdict=None, judge_reason=d["reason"][:500]) + return pid + + +# ================================================================ Redis 细节 +def _read(db: int, key: str, batch: int) -> list: + c = _client(db) + g, consumer = group_name(), param_store.get("PMS_SIGNAL_CONSUMER", "pms_1") + try: + c.xgroup_create(key, g, id="$", mkstream=True) # 只消化新消息, 不回溯历史 + logger.info("[信号消化] 建消费组 %s @ %s", g, key) + except Exception as e: + if "BUSYGROUP" not in str(e): + raise + resp = c.xreadgroup(g, consumer, {key: ">"}, count=int(batch), block=100) + out = [] + for _stream, messages in (resp or []): + out.extend(messages) + return out + + +def _ack(db: int, key: str, msg_id: str): + try: + _client(db).xack(key, group_name(), msg_id) + except Exception as e: + logger.warning("[信号消化] ACK 失败 %s: %s", msg_id, e) + + +# ================================================================ 去重与助手 +def _load_seen() -> set: + try: + raw = pms_repo.get_param(SEEN_KEY) + d = json.loads(raw) if raw else {} + if str(d.get("ymd")) != str(td.ymd()): + return set() + return set(d.get("keys") or []) + except Exception: + return set() + + +def _save_seen(seen: set, ymd): + try: + pms_repo.set_param(SEEN_KEY, json.dumps({"ymd": ymd, "keys": sorted(seen)[-500:]}), + "system") + except Exception as e: + logger.warning("[信号消化] 去重集合写入失败: %s", e) + + +def _has_inflight(code: str) -> bool: + try: + for i in pms_repo.list_instructions(statuses=list(executor.LIVE), ts_code=code, limit=20): + if str(i.get("side")).lower() == "sell": + return True + for p in pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200): + if p["ts_code"] == code and p["action"] in ("TRIM", "EXIT"): + return True + except Exception as e: + logger.warning("[信号消化] 在途检查失败(按无在途继续): %s", e) + return False + + +def _pos_of(view: dict, ts_code: str): + for x in view["positions"]: + if x["ts_code"] == ts_code: + return x + return None diff --git a/app/web/main.py b/app/web/main.py index dd1a94c..8e68871 100644 --- a/app/web/main.py +++ b/app/web/main.py @@ -316,6 +316,19 @@ def api_scan_proposals(dry_run: bool = Query(False)): return ok(proposal_service.scan_and_route, dry_run=dry_run) +@app.post("/api/ops/digest-signals") +def api_digest_signals(dry_run: bool = Query(False)): + """消化一批决策系统盘中信号。dry_run=true 只解析判定, 不落表也不 ACK。""" + from app.services import signal_service + return ok(signal_service.consume, dry_run=dry_run) + + +@app.get("/api/signal-status") +def api_signal_status(): + from app.services import signal_service + return ok(signal_service.status) + + @app.get("/api/dispatch-mode") def api_dispatch_mode(): from app.services import dispatcher, judge diff --git a/config/settings.py b/config/settings.py index 7e773c7..46ec998 100644 --- a/config/settings.py +++ b/config/settings.py @@ -111,6 +111,16 @@ class Settings(BaseSettings): PMS_T0_STOCK_DAY_LOSS: float = 0.003 # 单票当日T亏熔断 (占规模) PMS_T0_GLOBAL_DAY_LOSS: float = 0.01 # 全局当日T亏熔断 + # --- 决策系统信号消化 (设计 §10) --- + PMS_SIGNAL_ENABLED: bool = True + PMS_SIGNAL_GROUP: str = "pms_signal_consumer" # 独立消费组, 不与 trading_service 抢消息 + PMS_SIGNAL_CONSUMER: str = "pms_1" + PMS_SIGNAL_STREAM_INTRADAY: str = "intraday_signals:{ymd}" # db2, 每日一条流 + PMS_SIGNAL_STREAM_SELL: str = "bionic:signals:llm_sell_actions" # db3, 固定 key + PMS_SIGNAL_SELL_CONF_MIN: float = 0.75 # 低于此置信度的卖出信号不消化 + PMS_SIGNAL_AUTO_EXIT_CONF: float = 0.85 # 高于此置信度直接转清仓指令, 之间则落提议 + PMS_SIGNAL_TRIM_RATIO: float = 0.3333 # 中等置信度时的减仓比例 + # --- 对账与回放 --- PMS_REPLAY_INTERVAL_MIN: int = 5 PMS_RECON_ALARM_DAYS: int = 3 # 连续不一致 N 日升级 ERROR diff --git a/scripts/run_tests.py b/scripts/run_tests.py index d56c911..4776b0a 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -9,7 +9,8 @@ test_batch2_units.py 命令状态机 / 方案生成器 / 回放对账纯逻辑 (35 例) test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 (18 例) test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 (11 例) - test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (30 例) + test_batch5_units.py 决策系统信号流解析与消化口径 (8 例) + test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (32 例) 任一子集失败即整体失败 (退出码 1)。 """ import os @@ -19,7 +20,7 @@ import sys HERE = os.path.dirname(os.path.abspath(__file__)) ROOT = os.path.dirname(HERE) SUITES = ["test_core_units.py", "test_batch2_units.py", "test_batch3_units.py", - "test_batch4_units.py", "test_wiring.py"] + "test_batch4_units.py", "test_batch5_units.py", "test_wiring.py"] def main(): diff --git a/scripts/test_batch5_units.py b/scripts/test_batch5_units.py new file mode 100644 index 0000000..7ddcd4e --- /dev/null +++ b/scripts/test_batch5_units.py @@ -0,0 +1,141 @@ +# -*- coding: utf-8 -*- +""" +第五批模块单测 (实机运行, 零外部依赖) +====================================== +运行: 在 tradingSystem 仓库根目录执行 python scripts/test_batch5_units.py +覆盖: signal_rules 两条信号流的解析 (含两条流置信度尺度不同这个坑) 与消化口径。 +""" +import json +import os +import sys +import traceback + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +from app.core import signal_rules as sr # noqa: E402 + +RESULTS = [] + + +def case(name): + def deco(fn): + RESULTS.append((name, fn)) + return fn + return deco + + +PRM = {"sell_conf_min": 0.75, "auto_exit_conf": 0.85, "trim_ratio": 1 / 3} +POS = {"ts_code": "600000.SH", "total_qty": 6000, "avail_qty": 6000, "price": 10.0} + + +# ================================================================ 解析 +@case("解析·盘中流 (db2): 扁平字段 + component_scores 内嵌 JSON, 置信度本就是 0~1") +def _(): + s = sr.parse_intraday({"ts_code": "600000.SH", "action": "buy", "confidence": "0.83", + "component_scores": json.dumps({"minute_qrs": 2.4}), + "suggested_price": "10.25"}, msg_id="1-1") + assert s["source"] == sr.SRC_INTRADAY and s["ts_code"] == "600000.SH" + assert s["action"] == "BUY" and abs(s["confidence"] - 0.83) < 1e-9 + assert abs(s["minute_qrs"] - 2.4) < 1e-9 and s["suggested_price"] == 10.25 + # component_scores 给成 dict 或坏 JSON 都不能炸 + assert sr.parse_intraday({"component_scores": {"minute_qrs": 3}})["minute_qrs"] == 3.0 + assert sr.parse_intraday({"component_scores": "{坏"})["minute_qrs"] == 0.0 + + +@case("解析·风控流 (db3): 外层套 data JSON, 置信度是 0~100 要归一") +def _(): + inner = {"ts_code": "600000.SH", "action": "SELL", "confidence": 88, + "dominant_signal": "破位", "llm_reason": "跌破关键支撑且量能背离"} + s = sr.parse_risk_sell({"data": json.dumps(inner, ensure_ascii=False)}, msg_id="2-1") + assert s["source"] == sr.SRC_RISK_SELL and s["action"] == "SELL" + assert abs(s["confidence"] - 0.88) < 1e-9, s # 88 → 0.88, 两条流尺度不同 + assert s["dominant_signal"] == "破位" and "支撑" in s["reason"] + # 已经是 0~1 的也不会被再除一次 + assert abs(sr.parse_risk_sell({"data": json.dumps({"ts_code": "x", "action": "SELL", + "confidence": 0.9})})["confidence"] + - 0.9) < 1e-9 + # 坏 JSON → 明确标记, 不抛异常 + bad = sr.parse_risk_sell({"data": "{不是JSON"}) + assert bad["ts_code"] == "" and "解析失败" in bad["parse_error"] + + +# ================================================================ 消化 +@case("消化·高置信风控卖出 → 直接清仓 (减持不设确认门槛)") +def _(): + s = sr.parse_risk_sell({"data": json.dumps({"ts_code": "600000.SH", "action": "SELL", + "confidence": 90, "llm_reason": "逻辑走坏"})}) + d = sr.digest(s, POS, PRM) + assert d["action"] == sr.ACT_EXIT and d["qty"] == 6000, d + assert "清仓" in d["reason"] and "逻辑走坏" in d["reason"], d + assert d["hard_numbers"]["confidence"] == 0.9 + + +@case("消化·中等置信 → 落提议待确认 (按比例减)") +def _(): + s = sr.parse_risk_sell({"data": json.dumps({"ts_code": "600000.SH", "action": "SELL", + "confidence": 80})}) + d = sr.digest(s, POS, PRM) + assert d["action"] == sr.ACT_PROPOSE and d["qty"] == 2000, d # 6000 的 1/3, 整百 + assert "待确认" in d["reason"] + + +@case("消化·低置信 / 未持有 / 缺代码 一律不动") +def _(): + low = sr.parse_risk_sell({"data": json.dumps({"ts_code": "600000.SH", "action": "SELL", + "confidence": 60})}) + assert sr.digest(low, POS, PRM)["action"] == sr.ACT_IGNORE + hi = sr.parse_risk_sell({"data": json.dumps({"ts_code": "600000.SH", "action": "SELL", + "confidence": 95})}) + assert sr.digest(hi, {"total_qty": 0}, PRM)["action"] == sr.ACT_IGNORE + assert sr.digest(hi, None, PRM)["action"] == sr.ACT_IGNORE + assert sr.digest({"ts_code": "", "action": "SELL", "confidence": 0.95}, + POS, PRM)["action"] == sr.ACT_IGNORE + + +@case("消化·BUY/HOLD 信号只留痕不买 (买什么买多少由 PMS 自己决定)") +def _(): + b = sr.parse_intraday({"ts_code": "600000.SH", "action": "BUY", "confidence": "0.95"}) + d = sr.digest(b, POS, PRM) + assert d["action"] == sr.ACT_RECORD and d["qty"] == 0, d + assert "不据此买入" in d["reason"], d + h = sr.parse_intraday({"ts_code": "600000.SH", "action": "HOLD", "confidence": "0.99"}) + assert sr.digest(h, POS, PRM)["action"] == sr.ACT_RECORD + + +@case("消化·盘中流里的 SELL 也照样消化 (两条流同一套口径)") +def _(): + s = sr.parse_intraday({"ts_code": "600000.SH", "action": "SELL", "confidence": "0.92"}) + d = sr.digest(s, POS, PRM) + assert d["action"] == sr.ACT_EXIT and d["qty"] == 6000, d + + +@case("去重键·同票同源同动作当日只算一次") +def _(): + s = {"source": sr.SRC_RISK_SELL, "ts_code": "600000.SH", "action": "SELL"} + k1 = sr.dedup_key(s, 20260728) + assert k1 == sr.dedup_key(dict(s), 20260728) + assert k1 != sr.dedup_key(s, 20260729) + assert k1 != sr.dedup_key({**s, "source": sr.SRC_INTRADAY}, 20260728) + + +# ---------------------------------------------------------------- runner +def main(): + passed, failed = 0, 0 + for name, fn in RESULTS: + try: + fn() + print(f" PASS {name}") + passed += 1 + except Exception: + print(f" FAIL {name}") + traceback.print_exc() + failed += 1 + print("-" * 60) + if failed: + print(f"FAILED: {failed} / {passed + failed}") + sys.exit(1) + print(f"ALL PASS ({passed} cases)") + + +if __name__ == "__main__": + main() diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py index d2b516a..90e9e85 100644 --- a/scripts/test_wiring.py +++ b/scripts/test_wiring.py @@ -1004,6 +1004,119 @@ def _(): assert r2["verdict"] == judge.PASS and r2["degraded"] is False, r2 # TRIM 不在研判范围 +class FakeRedis: + """Redis Stream 的最小替身 (消费组 + xreadgroup + ack)。""" + + def __init__(self, msgs=None): + self.msgs = dict(msgs or {}) + self.acked, self.groups = [], [] + + def xgroup_create(self, key, group, id="$", mkstream=False): + self.groups.append((key, group)) + + def xreadgroup(self, group, consumer, streams, count=10, block=0): + key = list(streams)[0] + m = self.msgs.pop(key, []) + return [(key, m)] if m else [] + + def xack(self, key, group, msg_id): + self.acked.append(msg_id) + + def xlen(self, key): + return len(self.msgs.get(key, [])) + + def xinfo_groups(self, key): + return [{"name": g, "pending": 0} for k, g in self.groups if k == key] + + +def _install_signal_fakes(fake, sell_msgs=None, intraday_msgs=None): + """把两条流的假客户端装上, 返回 {db: FakeRedis}。""" + import json as _json + from datetime import datetime as _dt + from app.services import signal_service as ss + from config.settings import settings as _st + ymd = _dt.now().strftime("%Y-%m-%d") + r2 = FakeRedis({f"intraday_signals:{ymd}": list(intraday_msgs or [])}) + r3 = FakeRedis({"bionic:signals:llm_sell_actions": list(sell_msgs or [])}) + by_db = {_st.SIGNAL_REDIS_DB_INTRADAY: r2, _st.SIGNAL_REDIS_DB_ACTIONS: r3} + ss._client = lambda db: by_db[db] + return by_db + + +def _sell_msg(mid, code, conf, reason="逻辑走坏"): + import json as _json + return (mid, {"data": _json.dumps({"ts_code": code, "action": "SELL", + "confidence": conf, "llm_reason": reason}, + ensure_ascii=False)}) + + +@case("信号消化·高置信风控卖出转清仓指令; 中等置信落提议; 当日去重") +def _(): + from app.services import signal_service as ss + fake = install_fakes(prices={"600000.SH": 10.0, "000001.SZ": 8.0}, + params={"PMS_TOTAL_SCALE": "2000000"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 10.0}, + {"ts_code": "000001.SZ", "total_qty": 3000, + "avail_qty": 3000, "avg_cost": 8.0}]) + by_db = _install_signal_fakes(fake, sell_msgs=[ + _sell_msg("3-1", "600000.SH", 92), # 高置信 → 清仓 + _sell_msg("3-2", "000001.SZ", 80), # 中置信 → 提议 + _sell_msg("3-3", "600519.SH", 95), # 没持仓 → 忽略 + ]) + r = ss.consume() + assert r["ok"], r + assert [x["ts_code"] for x in r["exits"]] == ["600000.SH"], r + assert r["exits"][0]["qty"] == 6000 + assert [x["ts_code"] for x in r["proposals"]] == ["000001.SZ"], r + assert r["proposals"][0]["qty"] == 1000 # 3000 的 1/3 + assert r["ignored"] >= 1, r + assert len(by_db[3].acked) == 3, by_db[3].acked # 三条都 ACK + + ins = [i for i in fake.instructions.values() if i["action"] == "EXIT"] + assert ins and ins[0]["side"] == "sell" and ins[0]["progress"]["from_signal"] is True + assert any(x["arbiter"] == "rule" and x["action"] == "EXIT" for x in fake.ledger) + prop = [p for p in fake.proposals.values() if p["action"] == "TRIM"] + assert prop and prop[0]["hard_numbers"]["signal_source"] == "risk_sell", prop + + # 同一条信号再来一次: 当日去重 + 在途检查, 不重复下指令 + n_ins = len(fake.instructions) + _install_signal_fakes(fake, sell_msgs=[_sell_msg("3-4", "600000.SH", 92)]) + r2 = ss.consume() + assert not r2["exits"] and len(fake.instructions) == n_ins, r2 + + +@case("信号消化·BUY 只留痕不买; 关闭开关即不消化; 试算不落表不ACK") +def _(): + from app.services import signal_service as ss + fake = install_fakes(prices={"600000.SH": 10.0}, params={"PMS_TOTAL_SCALE": "2000000"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 10.0}]) + _install_signal_fakes(fake, intraday_msgs=[ + ("1-1", {"ts_code": "600000.SH", "action": "BUY", "confidence": "0.95"})]) + r = ss.consume() + assert r["recorded"] == 1 and not r["exits"], r + assert not any(i["side"] == "buy" for i in fake.instructions.values()) + assert any(x["action"] == "SIGNAL" for x in fake.ledger), fake.ledger + + fake2 = install_fakes(prices={"600000.SH": 10.0}, + params={"PMS_TOTAL_SCALE": "2000000", "PMS_SIGNAL_ENABLED": "false"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 10.0}]) + by = _install_signal_fakes(fake2, sell_msgs=[_sell_msg("3-9", "600000.SH", 95)]) + r2 = ss.consume() + assert "skipped" in r2 and not fake2.instructions, r2 + + fake3 = install_fakes(prices={"600000.SH": 10.0}, params={"PMS_TOTAL_SCALE": "2000000"}, + positions=[{"ts_code": "600000.SH", "total_qty": 6000, + "avail_qty": 6000, "avg_cost": 10.0}]) + by3 = _install_signal_fakes(fake3, sell_msgs=[_sell_msg("3-10", "600000.SH", 95)]) + r3 = ss.consume(dry_run=True) + assert r3["exits"] and r3["exits"][0]["dry_run"] is True, r3 + assert not fake3.instructions and not fake3.ledger + assert by3[3].acked == [], "试算不该 ACK" + + @case("装配·执行相关路由与调度接线到位") def _(): from app.web.main import app @@ -1011,12 +1124,13 @@ def _(): paths = {r.path for r in app.routes} for p in ("/api/ops/materialize", "/api/ops/exec-tick", "/api/ops/sweep-windows", "/api/instructions/{instruction_id}/cancel", "/api/dispatch-mode", - "/api/ops/scan-proposals"): + "/api/ops/scan-proposals", "/api/ops/digest-signals", "/api/signal-status"): assert p in paths, p import inspect src = inspect.getsource(sch.intraday_exec) assert "executor" in src and "run_tick" in src, "调度器未接执行器" assert "proposal_service" in src, "调度器未接自主提议扫描" + assert "signal_service" in inspect.getsource(sch.signal_digest), "调度器未接信号消化" # ---------------------------------------------------------------- runner