tradingSystem/app/services/signal_service.py

267 lines
10 KiB
Python
Raw Permalink Normal View History

# -*- 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