tradingSystem/app/services/signal_service.py

471 lines
24 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# -*- coding: utf-8 -*-
"""
决策系统盘中信号消化 (设计 §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 action_engine as ae
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), 日切自动作废
# 买入信号留痕在评审账本里用的动作名 (2026-08-06)。
# **刻意不叫 OPEN**: 账本里 action=OPEN 的行是真的建仓评审 (放行或拒绝), 两者混在一起,
# `make t-gate` 就分不清「决策系统说这只票转多了」和「PMS 决定建这只票」。
# 动作引擎那条路读的就是这个动作名 —— 见 pms_repo.buy_signals_today。
SIGNAL_BUY_ACTION = "SIGNAL_BUY"
_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()
try:
strat_codes = pms_repo.active_strategy_codes()
except Exception:
strat_codes = set()
for db, key, parser in streams():
try:
# 试算走**只读**通道 (XREVRANGE), 绝不碰消费组 (2026-08-28 审查修):
# 消费组按 ">" 读过的消息进本消费者的待确认清单, 不确认就**永远不会再投递**
# —— 原来试算读了不确认, 恰有一条未消化的风控卖出时, 点一次试算就把它永久
# 吞掉了。正式消费则先把自己上次读了没确认的 (处理中途抛异常的) 捞回来重处理。
msgs = _peek(db, key, batch) if dry_run else _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, strat_codes)
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, strat_codes=frozenset()):
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
# **持仓票的卖出被忽略也要留痕** (2026-08-18)。digest 对"低于门槛"返回 IGNORE,
# 而 IGNORE 原来一个字不留 —— 一只你正持有、bionic 连日喊卖的票(实测 002015 一天
# 三条、300474/300433 等), 在账本里查不到, 等于风控在示警而你看不见。未持有的票
# 不留(会被全市场卖出冲垮), 只给持仓票留, 且当日按(来源,股票,动作)去重一条。
if (not dry_run and str(sig.get("action") or "").upper() == "SELL"
and pos and int(pos.get("total_qty") or 0) > 0):
k = sr.dedup_key(sig, ymd) + ":IGN"
if k not in seen:
try:
pms_repo.insert_ledger(
ts_code=code, action="SIGNAL", arbiter="rule", verdict="NOTE",
price_at=float(pos.get("price") or 0),
hard_numbers={**d["hard_numbers"], "msg_id": sig.get("msg_id")},
reason="决策系统卖出未采纳(置信度低于门槛), 持仓票请留意: " + d["reason"][:400])
seen[k] = None
out.setdefault("held_sell_noted", []).append(code)
except Exception as e:
logger.warning("[信号消化] 持仓票卖出留痕失败 %s: %s", code, e)
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
if act == sr.ACT_NOTE_BUY:
# 买入信号留痕: **持没持仓都写**。未持仓的票正是新建仓关心的那一批, 从前它们
# 连一行账本都没有 —— 「决策系统今天看多了哪几只」查不到, 事后没法复盘。
#
# 三条防冲垮的口径 (与上面 RECORD 那句「否则全市场广播会把评审账本冲垮」同一个顾虑):
# 1. 只认 BUY, HOLD 等仍走 RECORD 的老口径;
# 2. 按 (日期, 来源, 股票, 动作) 当日去重 —— 同一只票一天最多一行;
# 3. 上游那侧本来就有节流: watcher 对同股看多研判有两小时防抖锁, 昨夜已看多的票
# 直接免疫不派单, 所以 REVERSAL_BUY 是几十条的量级, 不是全市场广播。
# verdict 用 NOTE 不用 PASS —— 这是「记下来」不是「放行」, 账本里必须分得开
# (watch.py 对不认识的 verdict 有兜底符号, 不会显示异常)。
# 留痕 reason 按来源分写 (2026-09-03, 文案在 signal_rules.digest): 择时决策系统的
# 保持「决策系统盘中判该股转多」开头, 其他发送方写「盘中择时程序买入信号(来源 xxx」;
# 硬数字里另存 producer_id。新建仓的插队排序只认前者 (proposal_service._buy_signals_today)。
key = sr.dedup_key(sig, ymd)
if key in seen:
out["ignored"] += 1
return
out.setdefault("buy_notes", []).append({"ts_code": code, "reason": d["reason"]})
if dry_run:
return
px = float((pos or {}).get("price") or d["hard_numbers"].get("suggested_price") or 0)
pms_repo.insert_ledger(
ts_code=code, action=SIGNAL_BUY_ACTION, arbiter="rule", verdict="NOTE",
price_at=px, hard_numbers={**d["hard_numbers"], "msg_id": sig.get("msg_id")},
reason=d["reason"][:500])
seen[key] = None
return
key = sr.dedup_key(sig, ymd)
if key in seen:
out["ignored"] += 1
return
if _has_inflight(code):
# **确认即加速** (2026-08-18): 走到这里 act 已是 EXIT/PROPOSE —— 决策系统确认要卖。
# 若这只票已有在途清仓指令(多为你人工下的清仓命令), 不再简单跳过, 而是把它升级成
# 紧急直通(马上按现价卖、不等均价不分桶): 人工意图 + 风控独立确认 = 最强离场信号。
# 只升级已下发的卖出**指令**; 仅挂着提议(还没成指令)的不动, 那本就等你拍板。
esc = [] if dry_run else _escalate_inflight_sells(code, d["reason"])
out["ignored"] += 1
if esc:
out.setdefault("escalated", []).append({"ts_code": code, "instructions": esc})
else:
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[key] = None 再落库, 一旦 _make_exit /
# _make_proposal 抛异常 (DB 抖一下、下游拒一次), 上层 catch 住记进 out["errors"],
# 但这一天的去重键已经烧掉了 —— 同一条风控卖出信号后面再来多少次都被当成重复丢弃,
# 指令一条都不会落。失败长得像成功: 页面只多一行 error, 而该卖的票就那么留着了。
# 2026-07-31 修。
on_strategy = bool(code) and code in strat_codes
if act == sr.ACT_EXIT and not on_strategy:
iid = _make_exit(code, d, pos)
seen[key] = None
out["exits"].append({**brief, "instruction_id": iid})
else:
# 挂了策略的票: 决策系统的风控卖出只提示、不自动清仓 (强制离场会推翻你特意设的策略);
# 一律落提议进「等我拍板」由你定"维持 / 采纳即撤策略并清仓", 采纳的是全清 (as_exit)。
pid = _make_proposal(code, d, pos, sig, on_strategy=on_strategy,
as_exit=(on_strategy and act == sr.ACT_EXIT))
seen[key] = None
out["proposals"].append({**brief, "proposal_id": pid, "on_strategy": on_strategy})
if on_strategy and not dry_run:
# 同时暂停该票策略的买入这一侧 (不平仓、不动卖出、页面可恢复) —— 挡住"资金在流出、
# 网格还在逢跌买入", 又不替你做清仓这种不可逆的事。
try:
from app.services import strategy_service
affected = strategy_service.pause_buy(code, reason=d["reason"],
source=(sig.get("source") or "signal"))
if affected:
out.setdefault("strategy_buy_paused", []).extend(affected)
logger.warning("[信号消化] %s 挂着策略, 决策系统卖出只提示不自动清仓; "
"已暂停该票策略买入(页面可恢复): %s", code, affected)
except Exception as e:
logger.error("[信号消化] 暂停策略买入失败 %s: %s", code, e)
def _make_exit(code, d, pos) -> str:
"""高置信风控卖出 → 直接落卖出指令 (减持方向不设确认门槛)。"""
now = datetime.now()
# 序号用完整时分秒 (原来 %1000 同日同票同动作可能撞唯一键丢单; 2026-08-28 修)
iid = cs.make_instruction_id(td.ymd(now), code, "EXIT", int(now.strftime("%H%M%S")))
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,
# 高置信风控清仓与一键清仓同为紧急语义 (2026-08-17): 择时不做择价
# 博弈, 直通出手、限价更激进 —— 见 exec_timing.hard_gate 的 urgent 分支
"urgent": 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"])
# 清仓指令已落, 这只票还挂着的到价提议作废 (2026-09-07): 不作废它会挂到次日日结,
# 人再点采纳会对一只已经清掉的票再发一条清仓。
retired = _retire_target_proposals(code, d["reason"])
if retired:
logger.warning("[信号消化] %s 到价提议 %s 已因止损直通作废", code, retired)
return iid
def _make_proposal(code, d, pos, sig, *, on_strategy=False, as_exit=False) -> str:
ttl = param_store.get_int("PMS_PROPOSAL_TTL_HOURS", 24)
pid = f"PRP_{td.ymd()}_{code.replace('.', '')}_SIGSELL"
action = "EXIT" if as_exit else "TRIM"
reason = d["reason"]
if on_strategy:
reason = ("这只票挂着交易方案(策略): 决策系统卖出信号只提示、未自动清仓; 该票策略买入已暂停。"
"你定: 维持观察 / 采纳即撤策略并清仓。 —— " + reason)
hn = {**d["hard_numbers"], "price": float((pos or {}).get("price") or 0),
"reason": reason, "signal_source": sig.get("source"), "on_strategy": on_strategy}
try:
pms_repo.insert_proposal(proposal_id=pid, ts_code=code, action=action, qty=d["qty"],
hard_numbers=hn,
expire_at=datetime.now() + timedelta(hours=ttl),
judge_verdict=("STRATEGY_RISK" if on_strategy else None),
judge_reason=reason[:500])
except Exception:
# 确定性编号被当日已终态提议占着 → 补时间后缀重试一次 (2026-08-28 审查修)
pid = f"{pid}_{datetime.now().strftime('%H%M%S')}"
pms_repo.insert_proposal(proposal_id=pid, ts_code=code, action=action, qty=d["qty"],
hard_numbers=hn,
expire_at=datetime.now() + timedelta(hours=ttl),
judge_verdict=("STRATEGY_RISK" if on_strategy else None),
judge_reason=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
out = []
# 先捞**自己名下读了没确认的** (id="0" 只回本消费者的待确认消息): 上一跳处理中途
# 抛异常的消息会留在这里, 不捞的话它们永远不会再被投递 —— "读了不确认 = 下次还在"
# 是对消费组的误解, 2026-08-28 审查修。
try:
resp0 = c.xreadgroup(g, consumer, {key: "0"}, count=int(batch))
for _stream, messages in (resp0 or []):
out.extend(m for m in messages if m and m[1]) # 空体=已删消息, 跳过
except Exception as e:
logger.warning("[信号消化] 待确认消息回捞失败 (继续读新): %s", e)
left = int(batch) - len(out)
if left > 0:
resp = c.xreadgroup(g, consumer, {key: ">"}, count=left, block=100)
for _stream, messages in (resp or []):
out.extend(messages)
return out
def _peek(db: int, key: str, batch: int) -> list:
"""只读取样 (试算用): XREVRANGE 不进消费组、不产生任何投递痕迹。返回升序。"""
c = _client(db)
rows = c.xrevrange(key, count=int(batch)) or []
return list(reversed(rows))
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():
"""当日去重集合。用 dict 当有序集合 (保持加入顺序) —— 超上限裁剪时裁掉**最旧的**,
而不是原来 sorted() 后按字典序裁 (裁掉的与时间无关, 被裁的键再来一条就会重复消化)。"""
try:
raw = pms_repo.get_param(SEEN_KEY)
d = json.loads(raw) if raw else {}
if str(d.get("ymd")) != str(td.ymd()):
return dict()
return dict.fromkeys(d.get("keys") or [])
except Exception:
return dict()
def _save_seen(seen, ymd):
try:
pms_repo.set_param(SEEN_KEY, json.dumps({"ymd": ymd, "keys": list(seen)[-800:]}),
"system")
except Exception as e:
logger.warning("[信号消化] 去重集合写入失败: %s", e)
def _has_inflight(code: str) -> bool:
"""这只票是否已有在途的卖出, 有则高置信风控卖出不再另落指令。
2026-09-07 审查修: **到价提议不算在途**。目标价到价那条只是「等你拍板要不要止盈」,
它挂着的时候正是急转向下、风控卖出高发的时候; 原先把它算作在途, 高置信止损整条被吞
—— 不落单、不记账, 页面上只看到一张止盈提议, 看不出决策系统已判该走。
到价提议的来源写在硬数字里 (source = target_price), 按它区分, 不按动作名。
"""
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 or p["action"] not in ("TRIM", "EXIT"):
continue
if str((p.get("hard_numbers") or {}).get("source") or "") == ae.SRC_TARGET_PRICE:
continue # 等人拍板的止盈, 不挡止损
return True
except Exception as e:
logger.warning("[信号消化] 在途检查失败(按无在途继续): %s", e)
return False
def _retire_target_proposals(code: str, why: str) -> list:
"""风控高置信卖出已直通清仓后, 把这只票还挂着的到价提议作废并留痕。
不作废的话它会一直挂到次日日结, 人若再点采纳, 会对一只已经清掉的票再发一条清仓。
作废写成 DECLINED 并在账本记一行 NOTE, 复盘时看得出是被止损顶掉的, 不是人驳回的。
"""
retired = []
try:
for p in pms_repo.list_proposals(statuses=("WAIT_USER",), limit=200):
if p["ts_code"] != code:
continue
if str((p.get("hard_numbers") or {}).get("source") or "") != ae.SRC_TARGET_PRICE:
continue
if pms_repo.decide_proposal(p["proposal_id"], "DECLINED"):
pms_repo.insert_ledger(
ts_code=code, action=p["action"], arbiter="rule", verdict="NOTE",
price_at=float((p.get("hard_numbers") or {}).get("price") or 0),
hard_numbers=p.get("hard_numbers") or {}, ref_id=p["proposal_id"],
reason=f"到价提议作废: 风控高置信卖出已直通清仓 ({why})"[:500])
retired.append(p["proposal_id"])
except Exception as e:
logger.warning("[信号消化] 作废到价提议失败(清仓指令已落, 不影响): %s", e)
return retired
def _escalate_inflight_sells(code: str, reason: str) -> list:
"""把该股在途的卖出**指令**升级为紧急直通 (确认即加速, 2026-08-18)。
只动已下发的卖出指令 (progress.urgent 置 True, 下一跳 run_tick 就按现价直卖);
已是紧急的不重复动。每条升级在账本留一行, 供复盘"为什么这条清仓突然加速了"
仅挂着提议(还没成指令)的不在此列 —— 那本就等用户拍板, 不该被信号自动推成紧急。
"""
ids = []
try:
for ins in pms_repo.list_instructions(statuses=list(executor.LIVE), ts_code=code,
limit=50):
if str(ins.get("side") or "").lower() != "sell":
continue
prog = dict(ins.get("progress") or {})
if prog.get("urgent"):
continue
prog["urgent"] = True
prog["escalated_by_signal"] = str(reason)[:200]
pms_repo.update_instruction(ins["instruction_id"], progress=prog)
ids.append(ins["instruction_id"])
pms_repo.insert_ledger(
ts_code=code, action="EXIT", arbiter="rule", verdict="PASS", price_at=0,
ref_id=ins["instruction_id"],
reason="决策系统卖出确认, 在途清仓升级为紧急直通: " + str(reason)[:300])
logger.warning("[信号消化] %s 决策系统确认卖出, 在途清仓 %s 升级紧急直通",
code, ins["instruction_id"])
except Exception as e:
logger.warning("[信号消化] 升级在途清仓失败 %s: %s", code, e)
return ids
def _pos_of(view: dict, ts_code: str):
for x in view["positions"]:
if x["ts_code"] == ts_code:
return x
return None