From cd9c90dee03fbd2f7da839493c5afecd9a414bec Mon Sep 17 00:00:00 2001 From: zlt Date: Thu, 13 Aug 2026 10:45:04 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9=E9=A1=B5=E9=9D=A2=E5=B8=83?= =?UTF-8?q?=E5=B1=80=E6=A0=B7=E5=BC=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 2 +- app/services/upstream_signals.py | 186 +++++++++++++++++++++++++++++++ app/web/main.py | 8 ++ app/web/static/index.html | 98 +++++++++++++++- config/settings.py | 4 + 5 files changed, 294 insertions(+), 4 deletions(-) create mode 100644 app/services/upstream_signals.py diff --git a/.env.example b/.env.example index 35b1ad4..b81a684 100644 --- a/.env.example +++ b/.env.example @@ -16,7 +16,7 @@ DB_MYSQL_URL=mysql+pymysql://用户:密码@192.168.18.199:3306/db_gp_cj # ---- 信号/行情 Redis (208: 决策系统信号流 db2/db3 订阅 + 实时行情 db13) ---- SIGNAL_REDIS_HOST=192.168.18.208 SIGNAL_REDIS_PORT=6379 -SIGNAL_REDIS_PASSWORD=密码 +SIGNAL_REDIS_PASSWORD=wlkj2018 # ---- PMS 自己的 Celery 总线 (150 Redis, db8; 决策系统用 db7, 逻辑隔离) ---- PMS_REDIS_URL=redis://:密码@192.168.16.150:6379/8 diff --git a/app/services/upstream_signals.py b/app/services/upstream_signals.py new file mode 100644 index 0000000..cafaf68 --- /dev/null +++ b/app/services/upstream_signals.py @@ -0,0 +1,186 @@ +# -*- 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 diff --git a/app/web/main.py b/app/web/main.py index aad8f36..7bef926 100644 --- a/app/web/main.py +++ b/app/web/main.py @@ -498,6 +498,14 @@ def api_open_scan(): return ok(proposal_service.disposition_snapshot) +@app.get("/api/upstream-signals") +def api_upstream_signals(): + """上游信号只读快照 (三系统盘中信号: 决策系统买卖 / 择时层 BUY / 双向告警 / 资金异动 / 实况分)。 + 只展示、不进任何下单或决策逻辑; 内部一律 XREVRANGE/ZREVRANGE, 不建消费组、不 ACK、不写。""" + from app.services import upstream_signals + return ok(upstream_signals.snapshot) + + @app.get("/api/upstream/plan") def api_upstream_plan(limit: int = Query(30), date: str = Query(None), bucket: str = Query("main")): diff --git a/app/web/static/index.html b/app/web/static/index.html index 82b8ebe..be40c1d 100644 --- a/app/web/static/index.html +++ b/app/web/static/index.html @@ -106,6 +106,13 @@ pre.json{background:var(--surface-2);border:1px solid var(--hair);border-radius: .exp-led{padding:5px 0;border-bottom:1px solid var(--hair);font-size:12.5px;line-height:1.5;} .exp-led:last-child{border-bottom:none;} @media(max-width:900px){.pos-exp{grid-template-columns:1fr;gap:14px;}} +.up-grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(400px,1fr));gap:16px;} +.up-cell{border:1px solid var(--hair);border-radius:11px;padding:12px 14px;background:var(--surface-2);} +.up-h{font-size:13px;font-weight:640;color:var(--ink-2);margin-bottom:8px;} +.up-h .src-err{color:var(--stop);font-weight:400;font-size:11.5px;} +.mr-two{display:grid;grid-template-columns:1fr 1fr;gap:16px;font-size:12.5px;} +.mr-row{padding:3px 0;border-bottom:1px solid var(--hair);} +.mr-row:last-child{border-bottom:none;} @@ -564,6 +571,83 @@ pre.json{background:var(--surface-2);border:1px solid var(--hair);border-radius: + +
+

上游信号 · 只读参考,不进下单逻辑 {{ showUpstream?'▲ 收起':'▼ 展开' }}

+ +
+
@@ -1259,6 +1343,8 @@ createApp({ const showSettings = ref(false); // 设置默认折叠 const lastRefresh = ref(''); const autoRefresh = ref(true); + const upstream = ref({ sources: {} }); + const showUpstream = ref(true); const openScan = ref({ by_code: {}, notes: [] }); const insTab = ref('live'); const stratDlg = reactive({ visible:false, busy:false, reasons:[], ts_code:'', price:0, @@ -1684,6 +1770,11 @@ createApp({ const d = await call('get', '/api/open-scan'); openScan.value = (d && d.by_code) ? d : { by_code: {}, notes: [] }; } + const usrc = k => (upstream.value.sources || {})[k]; + async function loadUpstream() { + const d = await call('get', '/api/upstream-signals'); + upstream.value = (d && d.sources) ? d : { sources: {} }; + } async function loadPlans(cid) { plansOf.value = cid; const d = await call('get', '/api/plans?command_id=' + cid); plans.value = d.data || []; @@ -1718,7 +1809,7 @@ createApp({ loading.value = true; err.value = ''; await Promise.all([loadOverview(), loadParams(), loadCatalog(), loadCommands(), loadPositions(), loadInstructions(), loadLedger(), loadProposals(), - loadDispatchMode(), loadWs(), loadStrategies(), loadOpLog(), loadOpenScan()]); + loadDispatchMode(), loadWs(), loadStrategies(), loadOpLog(), loadOpenScan(), loadUpstream()]); await loadNames(); lastRefresh.value = (health.value.now || '').slice(11, 19); loading.value = false; @@ -1918,7 +2009,7 @@ createApp({ let refreshTimer = null, refreshTick = 0; async function refreshLive() { await Promise.all([loadOverview(), loadPositions(), loadInstructions(), loadLedger(), - loadProposals(), loadStrategies(), loadOpLog(), loadOpenScan()]); + loadProposals(), loadStrategies(), loadOpLog(), loadOpenScan(), loadUpstream()]); await loadNames(); lastRefresh.value = (health.value.now || '').slice(11, 19); } @@ -1953,7 +2044,8 @@ createApp({ openStrategy, validateStrategy, attachStrategy, setStrategyStatus, stratStateText, resumeBuy, showLedger, onPosExpand, openScan, insTab, todayYmd, insLive, insDone, insEnd, dispOf, loadOpenScan, denyToday, - showOpLog, showSettings, watchCandidates, marketOpen, marketState, lastRefresh, autoRefresh }; + showOpLog, showSettings, watchCandidates, marketOpen, marketState, lastRefresh, autoRefresh, + upstream, showUpstream, usrc, loadUpstream }; } }).use(ElementPlus).mount('#app'); diff --git a/config/settings.py b/config/settings.py index 267b60f..7c794f8 100644 --- a/config/settings.py +++ b/config/settings.py @@ -40,6 +40,10 @@ class Settings(BaseSettings): SIGNAL_REDIS_SOCKET_TIMEOUT: int = 5 PMS_REDIS_URL: str = "redis://:pass@192.168.16.150:6379/8" + # 上游信号(mtf 实况分)只读源 —— 214 db0。密钥只从 .env 注入(见下「密钥」纪律), 代码默认留空。 + # .env 里加一行: UP_MR_REDIS_URL=redis://:<密码>@192.168.16.214:36379/0 + # (208 上的决策系统/择时层/资金异动几条流复用现有 SIGNAL_REDIS_*, 无需另配。) + UP_MR_REDIS_URL: str = "" # PMS 自己的 Celery 总线, 独立 db=8 (决策系统用 db7, 物理同机逻辑隔离) PMS_WEB_HOST: str = "0.0.0.0" # 管理页面 (FastAPI + 单页)