# -*- coding: utf-8 -*- """上游信号只读快照 (2026-08-12: 只展示、不进任何下单/决策逻辑)。 把三系统盘中信号取最近若干条给页面看: · 决策系统(bionic) 买卖广播 + 风控卖出动作 —— 208 db2 intraday_signals / db3 llm_sell_actions · 盘中择时层(intraday_timing) BUY 入场 + 双向风控告警 —— 208 db2 intraday_signals / intraday_alerts · mtf 资金异动 + 实况分关注/回避榜 —— 208 db2 stream:metrics / 214 db0 mr:board **严格只读**: 一律 XREVRANGE / ZREVRANGE 取最近, 不建消费组、不 ACK、不写 —— 碰不到上游的消费与 投递, 也不和 PMS 现有信号消化(signal_service 的消费组)抢消息。连不通的源单独降级 (ok=False + error), 不影响其余与整页。取数带进程内缓存 (默认 8s), 不高频打上游。 `intraday_signals` 一条流两家都写、格式不同: 按 producer_id / entry_score 字段特征分成 「决策系统」与「择时层」两类展示。 """ from __future__ import annotations import json import logging import time from datetime import datetime from config.settings import settings from app.services import param_store logger = logging.getLogger("pms.upstream") _clients = {} _cache = {"at": 0.0, "data": None} def _to() -> int: return max(1, param_store.get_int("PMS_UPSTREAM_TIMEOUT_SEC", 3)) def _limit() -> int: return max(1, param_store.get_int("PMS_UPSTREAM_LIMIT", 40)) def _c208(db: int): """208(与 SIGNAL_REDIS 同实例)的只读客户端, 按库缓存。""" key = ("208", db) if key in _clients: return _clients[key] 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=_to(), socket_connect_timeout=_to()) try: c = redis.Redis(protocol=2, **kw) # RESP2: 服务端 <6.0 不认 HELLO (与行情库同一个坑) except TypeError: c = redis.Redis(**kw) _clients[key] = c return c def _c214(): """214(mtf 实况分)的只读客户端, 走 .env 里的 UP_MR_REDIS_URL。""" if "214" in _clients: return _clients["214"] import redis url = (getattr(settings, "UP_MR_REDIS_URL", "") or "").strip() if not url: raise RuntimeError("UP_MR_REDIS_URL 未配置 (请在 .env 加, 见 config/settings.py 注释)") try: c = redis.Redis.from_url(url, protocol=2, decode_responses=True, socket_timeout=_to(), socket_connect_timeout=_to()) except TypeError: c = redis.Redis.from_url(url, decode_responses=True, socket_timeout=_to(), socket_connect_timeout=_to()) _clients["214"] = c return c def _hm(ms): """epoch 毫秒 → HH:MM (本地时区)。拿不到就空串。""" try: return time.strftime("%H:%M", time.localtime(int(ms) // 1000)) except Exception: return "" def _num(v): try: return round(float(v), 4) except (TypeError, ValueError): return v # ================================================================ 各源解析 (只读) def _read_intraday(c, ymd): """一条 intraday_signals 流两家都写: 有 producer_id=intraday_timing / entry_score 的归「择时层 BUY」, 其余归「决策系统买卖」(bionic 广播 ENTRY/EXIT · BUY/SELL)。""" key = "intraday_signals:%s" % ymd decision, timing = [], [] for _id, f in c.xrevrange(key, count=_limit()): producer = str(f.get("producer_id") or "") is_itd = ("intraday_timing" in producer) or (f.get("entry_score") is not None) row = {"ts_code": f.get("ts_code"), "action": f.get("action"), "signal_type": f.get("signal_type"), "price": _num(f.get("suggested_price")), "target": _num(f.get("target_price")), "confidence": _num(f.get("confidence")), "entry_score": _num(f.get("entry_score")), "time": _hm(f.get("trigger_time"))} (timing if is_itd else decision).append(row) return {"key": key, "decision": decision, "timing": timing} def _read_alerts(c, ymd): key = "intraday_alerts:%s" % ymd out = [] for _id, f in c.xrevrange(key, count=_limit()): src = str(f.get("source") or "") meta = f.get("metadata") if not isinstance(meta, dict): try: meta = json.loads(meta) except Exception: meta = {} mdir = str((meta or {}).get("direction") or "").lower() low = src.lower() # 源名都含 "intensity", 不能用 "_in" 这种松匹配 (会误配 out_intensity)。用精确 token, 先判看跌。 if mdir in ("down", "short", "bear") or any(x in low for x in ("_down", "out_intensity", "capitulation")): d = "看跌" elif mdir in ("up", "long", "bull") or any(x in low for x in ("_up", "in_intensity", "breakout", "buy_emitted")): d = "看涨" else: d = "—" out.append({"ts_code": f.get("ts_code"), "source": src, "direction": d, "level": f.get("level"), "value": _num(f.get("value")), "time": _hm(f.get("trigger_time"))}) return {"key": key, "items": out} def _read_metrics(c): """mtf 资金异动: payload schema 未在契约里写全, 挑常见字段, 其余整体带过去供页面兜底显示。""" key = "mtf:intraday:stream:metrics" out = [] for _id, f in c.xrevrange(key, count=_limit()): code = f.get("ts_code") or f.get("code") out.append({"ts_code": code, "z_dd": _num(f.get("z_dd")), "direction": f.get("direction"), "level": f.get("level"), "value": _num(f.get("value") if f.get("value") is not None else f.get("net")), "raw": {k: v for k, v in f.items() if k not in ("ts_code", "code")}}) return {"key": key, "items": out} def _read_sell_actions(c): key = "bionic:signals:llm_sell_actions" out = [] for _id, f in c.xrevrange(key, count=_limit()): out.append({"ts_code": f.get("ts_code"), "action": f.get("action"), "confidence": _num(f.get("confidence")), "reason": f.get("reason"), "time": _hm(f.get("trigger_time") or f.get("ts"))}) return {"key": key, "items": out} def _read_mr(c): """mtf 实况分 关注/回避榜 (zset, score=实况分 0~1): 高分=关注、低分=回避, 各取前 10。""" key = "mtf:mr:intraday:board" top = c.zrevrange(key, 0, 9, withscores=True) bottom = c.zrange(key, 0, 9, withscores=True) fmt = lambda pairs: [{"ts_code": k, "score": round(float(v), 3)} for k, v in pairs] return {"key": key, "watch": fmt(top), "avoid": fmt(bottom)} # ================================================================ 快照 def snapshot(force: bool = False) -> dict: ttl = max(2, param_store.get_int("PMS_UPSTREAM_CACHE_SEC", 8)) now = time.time() if not force and _cache["data"] is not None and (now - _cache["at"]) < ttl: return _cache["data"] ymd = datetime.now().strftime("%Y-%m-%d") out = {"ok": True, "at": datetime.now().isoformat(timespec="seconds"), "sources": {}} def _try(name, fn): try: out["sources"][name] = {"ok": True, **fn()} except Exception as e: logger.warning("[上游信号] %s 读取失败: %s", name, e) out["sources"][name] = {"ok": False, "error": f"{type(e).__name__}: {e}"} dbi, dba = settings.SIGNAL_REDIS_DB_INTRADAY, settings.SIGNAL_REDIS_DB_ACTIONS _try("intraday", lambda: _read_intraday(_c208(dbi), ymd)) _try("alerts", lambda: _read_alerts(_c208(dbi), ymd)) _try("metrics", lambda: _read_metrics(_c208(dbi))) _try("sell_actions", lambda: _read_sell_actions(_c208(dba))) _try("mr", lambda: _read_mr(_c214())) _cache["at"], _cache["data"] = now, out return out