diff --git a/app/services/param_store.py b/app/services/param_store.py index ffe1387..d54cc5b 100644 --- a/app/services/param_store.py +++ b/app/services/param_store.py @@ -55,6 +55,9 @@ RUNTIME_EXTRA = { "PMS_UPSTREAM_ALERT_SCAN": (300, int, "告警扫描窗: 每轮从告警流取最近 N 条再分类限量 (窗要大于各类上限之和, 防 price_notice 等高频源把风控告警挤出最新N)"), "PMS_UPSTREAM_ALERT_PER_CAT": (20, int, "告警每类上限: 同一类型最多显示最新 N 条 (砍单类刷屏, 不挤别的类)"), "PMS_UPSTREAM_ALERT_MAX_AGE_MIN": (60, int, "告警时间窗(分钟): 超过此龄判为陈旧并在页面置灰标注(不删除, 风控不丢信息), 0=不限龄。与每类上限正交, 先按龄标记再每类限量"), + # 2026-09-09: 点股票看它今天全部信号 (台账 055)。快照口径只留每类最新 N 条, 到下午 + # 一只活跃票早盘的信号早被挤掉; 按代码回扫是另一条路径, 上限单独可调。 + "PMS_UPSTREAM_BYCODE_SCAN": (5000, int, "单票信号回扫上限: 点开一只票时每条上游流最多倒序回扫 N 条 (分页 500/次)。只在点击时触发, 不进轮询; 调大更全但更慢"), # 2026-08-18 宏观择时层: 闸状态由 macro_service 每次扫描写入, 扫描层与策略层只读。 # 用运行参数承载 (不进策略 state、不另建表) —— 与 PMS_STRATEGY_BUYPAUSE 同一手法。 "PMS_MACRO_GATE_STATE": ("", str, "宏观闸状态 (macro_service 维护的 JSON: active/ymd/原因), 页面与扫描层只读, 勿手改"), @@ -442,6 +445,7 @@ _RANGES = { # 任务会先被 celery 打死 (而且是在已经落了一部分表之后)。 "PMS_UPSTREAM_ALERT_SCAN": (1, 20000), "PMS_UPSTREAM_ALERT_PER_CAT": (1, 1000), "PMS_UPSTREAM_ALERT_MAX_AGE_MIN": (0, 1440), # 0=不限龄, 上限 24 小时 + "PMS_UPSTREAM_BYCODE_SCAN": (100, 50000), # 下限 100 保证至少扫得到一页 # 宏观择时层 "PMS_MACRO_HOT_TH": (0, 100), "PMS_MACRO_COLD_TH": (-100, 0), "PMS_MACRO_EXIT_BAND": (0, 100), "PMS_MACRO_CONFIRM_DAYS": (1, 10), diff --git a/app/services/upstream_signals.py b/app/services/upstream_signals.py index c5237a6..870a34f 100644 --- a/app/services/upstream_signals.py +++ b/app/services/upstream_signals.py @@ -282,3 +282,308 @@ def snapshot(force: bool = False) -> dict: _cache["at"], _cache["data"] = now, out return out + +# ================================================================ 按代码回扫 (2026-09-09) +# 信号栏一行 = 上游流里的一条原始消息, 同一只票一天出现很多次 (实测活跃票一天十来条)。 +# 这里按代码把五个源汇到一起, 给页面的「这只票今天的信号」抽屉用。 +# +# 与 snapshot() 的分工: 那个是**全场最近 N 条**的快照 (告警每类只留最新 20 条), 到下午 +# 一只活跃票早盘的信号早被挤掉了; 这个是**单票回扫到当天开盘**, 所以才叫"全部"。 +# 两份缓存不合并 —— 口径不同, 合并会让单票结果被全场快照污染。 +# +# 只读纪律与 snapshot() 逐字相同: 一律 XREVRANGE / ZREVRANGE, 不建消费组、不 ACK、不写。 +# **这个接口只能点击触发, 绝不能加进页面轮询** —— 加进去就是每 30 秒对上游做一次全流扫描, +# 把只读观察变成压力源。 +_BYCODE_CACHE = {} +_BYCODE_MAX = 32 +_SCAN_PAGE = 500 + + +def _bycode_scan() -> int: + return max(100, param_store.get_int("PMS_UPSTREAM_BYCODE_SCAN", 5000)) + + +def _code6(v) -> str: + """任意写法里抽六位数字。抽不到回空串。""" + import re + m = re.search(r"(\d{6})", str(v or "")) + return m.group(1) if m else "" + + +def _norm_code(v) -> str: + """SH600000 / 600000.SH / sh600000 / 600000 -> 600000.SH; 认不出后缀就只回六位。""" + s = str(v or "").strip().upper() + six = _code6(s) + if not six: + return "" + for ex in ("SH", "SZ", "BJ"): + if ex in s: + return "%s.%s" % (six, ex) + return six + + +def _same_code(a, b) -> bool: + """两边都带交易所后缀时要求后缀一致; 只有一边带就退回六位匹配。 + 六位在 A 股跨市场唯一, 够用; 这道二次校验是给将来可能出现的别的品种留的。""" + sa, sb = str(a or "").upper(), str(b or "").upper() + if _code6(sa) != _code6(sb) or not _code6(sa): + return False + exa = next((x for x in ("SH", "SZ", "BJ") if x in sa), "") + exb = next((x for x in ("SH", "SZ", "BJ") if x in sb), "") + return (not exa) or (not exb) or exa == exb + + +def _id_ms(sid) -> int: + try: + return int(str(sid).split("-")[0]) + except Exception: + return 0 + + +def _prev_id(sid) -> str: + """比 sid 严格小的边界 id。不用 Redis 6.2 的 `(id` 排他区间 —— 本模块已经因为上游 + 可能低于 6.0 而强制走 RESP2, 排他区间同样不能假设。手工减一在所有版本上都成立。""" + try: + ms, seq = str(sid).split("-") + ms, seq = int(ms), int(seq) + except Exception: + return "-" + if seq > 0: + return "%d-%d" % (ms, seq - 1) + if ms > 0: + return "%d-18446744073709551615" % (ms - 1) + return "-" + + +def _scan_stream(c, key, *, cap, page=_SCAN_PAGE, stop_before_ms=None): + """分页倒序回扫一条流, 返回 (条目列表, 是否还没见底)。 + + 一次 XREVRANGE COUNT 5000 会把几兆塞进一个应答且没法中途停, 所以分页。 + stop_before_ms 给**不按日期分键**的流用 (风控卖出、资金异动): 一旦某条早于今天零点 + 就立刻收手, 且不算截断 —— 今天的已经扫全了。 + """ + entries, cur, truncated = [], "+", False + while len(entries) < cap: + want = min(page, cap - len(entries)) + batch = c.xrevrange(key, max=cur, min="-", count=want) + if not batch: + break + stop = False + for sid, f in batch: + if stop_before_ms is not None and _id_ms(sid) < stop_before_ms: + stop = True + break + entries.append((sid, f)) + if stop: + return entries, False + if len(batch) < want: + break + cur = _prev_id(batch[-1][0]) + if cur == "-": + break + else: + truncated = True + return entries, truncated + + +def _today_start_ms() -> int: + t = time.localtime() + return int(time.mktime((t.tm_year, t.tm_mon, t.tm_mday, 0, 0, 0, 0, 0, -1)) * 1000) + + +def _row(src, src_label, ms, **kw) -> dict: + """时间线一行的统一形状。缺的字段给 None 不缺键 —— 模板最怕键时有时无。""" + base = {"src": src, "src_label": src_label, "ts": int(ms or 0), "time": _hm(ms), + "ymd": _ymd_of(ms), "cat": None, "cat_label": None, "direction": None, + "level": None, "value": None, "price": None, "target": None, "confidence": None, + "entry_score": None, "z_dd": None, "window_net": None, "window_ret": None, + "action": None, "dominant_signal": None, "reason": None, "stale": False, + "age_min": None, "dup": 1, "ts_from": "trigger"} + base.update(kw) + return base + + +def _dedup(rows) -> list: + """同源 + 同分钟 + 同值的合并成一条并累加 dup; **跨源永不合并** —— + 决策系统与择时层在同一分钟都看多, 是两条独立证据, 合并就是删信息。""" + seen, out = {}, [] + for r in rows: + ident = { + "alert": (r.get("cat"), r.get("level"), r.get("value")), + "decision": (r.get("action"), r.get("price")), + "timing": (r.get("price"), r.get("entry_score")), + "sell": (r.get("dominant_signal"), r.get("confidence")), + "metrics": (r.get("z_dd"), r.get("window_net")), + }.get(r.get("src"), (r.get("value"),)) + key = (r.get("src"), r.get("time"), ident) + if key in seen: + seen[key]["dup"] += 1 + continue + seen[key] = r + out.append(r) + return out + + +def _bycode_intraday(c, ymd, code, cap): + ent, trunc = _scan_stream(c, "intraday_signals:%s" % ymd, cap=cap) + rows = [] + for _sid, f in ent: + if not _same_code(f.get("ts_code"), code): + continue + producer = str(f.get("producer_id") or "") + is_itd = ("intraday_timing" in producer) or (f.get("entry_score") is not None) + rows.append(_row("timing" if is_itd else "decision", + "盘中择时层入场" if is_itd else "决策系统广播", + f.get("trigger_time"), + cat=f.get("signal_type"), + cat_label="择时层入场" if is_itd else ("建议买入" if str(f.get("action") or "").upper() == "BUY" else "建议卖出"), + direction="看涨" if str(f.get("action") or "").upper() == "BUY" else "看跌", + action=f.get("action"), price=_num(f.get("suggested_price")), + target=_num(f.get("target_price")), confidence=_num(f.get("confidence")), + entry_score=_num(f.get("entry_score")))) + return rows, len(ent), trunc + + +def _bycode_alerts(c, ymd, code, cap): + ent, trunc = _scan_stream(c, "intraday_alerts:%s" % ymd, cap=cap) + max_age, now, rows = _alert_max_age_sec(), time.time(), [] + for _sid, f in ent: + if not _same_code(f.get("ts_code"), code): + continue + src = str(f.get("source") or "") + meta = f.get("metadata") + if not isinstance(meta, dict): + try: + meta = json.loads(meta) + except Exception: + meta = {} + dir_fixed, cat, cat_label = _ALERT_META.get(src, (None, "other", "其他")) + if dir_fixed is None: + mdir = str((meta or {}).get("direction") or "").lower() + d = "看跌" if mdir in ("down", "short", "bear") else ("看涨" if mdir in ("up", "long", "bull") else "—") + else: + d = dir_fixed + try: + age_sec = now - float(f.get("trigger_time")) / 1000.0 + age_min = int(age_sec // 60) + except (TypeError, ValueError): + age_sec, age_min = None, None + rows.append(_row("alert", "盘中告警", f.get("trigger_time"), cat=cat, cat_label=cat_label, + direction=d, level=f.get("level"), value=_num(f.get("value")), + age_min=age_min, + stale=bool(max_age) and age_sec is not None and age_sec > max_age, + reason=str(f.get("source") or ""))) + return rows, len(ent), trunc + + +def _bycode_sell(c, code, cap): + ent, trunc = _scan_stream(c, "bionic:signals:llm_sell_actions", cap=cap, + stop_before_ms=_today_start_ms()) + rows = [] + for _sid, f in ent: + raw = f.get("data") + d = f + if raw: + try: + d = json.loads(raw) if isinstance(raw, str) else dict(raw) + except Exception: + d = {} + if not _same_code(d.get("ts_code") or f.get("ts_code"), code): + continue + ts = d.get("timestamp") or d.get("trigger_time") or f.get("ts") + rows.append(_row("sell", "风控卖出动作", ts, cat="risk_sell", cat_label="风控卖出", + direction="看跌", action=d.get("action"), + confidence=_num(d.get("confidence")), + dominant_signal=d.get("dominant_signal"), + reason=d.get("llm_reason") or d.get("reason") or f.get("reason"))) + return rows, len(ent), trunc + + +def _bycode_metrics(c, code, cap): + """资金异动进时间线 (2026-09-09): 时间取流 id 的毫秒部分, 那是**写入时刻**不是上游触发时刻, + 所以每行标 ts_from=stream_id, 页面写「写入 13:41」而不是「触发 13:41」。""" + ent, trunc = _scan_stream(c, "mtf:intraday:stream:metrics", cap=cap, + stop_before_ms=_today_start_ms()) + rows = [] + for sid, f in ent: + cd = f.get("ts_code") or f.get("code") + payload = f.get("data") + if isinstance(payload, str): + try: + payload = json.loads(payload) + except Exception: + payload = {} + if not isinstance(payload, dict): + payload = {} + if not _same_code(cd or payload.get("ts_code"), code): + continue + g = lambda k, dv=None: payload.get(k, f.get(k, dv)) + rows.append(_row("metrics", "盘中资金异动", _id_ms(sid), cat="fund_flow", cat_label="资金异动", + direction=g("direction"), level=g("level"), + value=_num(g("value") if g("value") is not None else g("net")), + z_dd=_num(g("z_dd")), window_net=_num(g("window_net")), + window_ret=_num(g("window_ret")), ts_from="stream_id")) + return rows, len(ent), trunc + + +def by_code(ts_code: str, force: bool = False) -> dict: + """一只票今天的全部上游信号。点击触发, 不进任何轮询。""" + code = _norm_code(ts_code) + if not code: + return {"ok": False, "error": "认不出的股票代码: %s" % ts_code, "timeline": [], + "snapshot": {"metrics": [], "mr": {"watch": None, "avoid": None}}} + six = _code6(code) + ttl = max(2, param_store.get_int("PMS_UPSTREAM_CACHE_SEC", 8)) + now = time.time() + hit = _BYCODE_CACHE.get(six) + if not force and hit and (now - hit[0]) < ttl: + return hit[1] + + ymd = datetime.now().strftime("%Y-%m-%d") + cap = _bycode_scan() + out = {"ok": True, "ts_code": code, "code6": six, "ymd": ymd, + "at": datetime.now().isoformat(timespec="seconds"), + "timeline": [], "snapshot": {"metrics": [], "mr": {"watch": None, "avoid": None}}, + "scanned": {}, "truncated": {}, "sources": {}} + rows = [] + + def _try(name, fn): + try: + got, scanned, trunc = fn() + rows.extend(got) + out["scanned"][name] = scanned + out["truncated"][name] = bool(trunc) + out["sources"][name] = {"ok": True} + except Exception as e: + logger.warning("[单票信号] %s %s 读取失败: %s", code, name, e) + out["sources"][name] = {"ok": False, "error": "%s: %s" % (type(e).__name__, e)} + + dbi, dba = settings.SIGNAL_REDIS_DB_INTRADAY, settings.SIGNAL_REDIS_DB_ACTIONS + _try("intraday", lambda: _bycode_intraday(_c208(dbi), ymd, code, cap)) + _try("alerts", lambda: _bycode_alerts(_c208(dbi), ymd, code, cap)) + _try("sell_actions", lambda: _bycode_sell(_c208(dba), code, cap)) + _try("metrics", lambda: _bycode_metrics(_c208(dbi), code, cap)) + + # 实况分是一个按分数排的有序集合, 连写入时刻都没有 —— **绝不给它编时间**, + # 排进时间线的假时间戳比放在快照段里说"没有时间"危险得多。 + try: + mr = _read_mr(_c214()) + out["snapshot"]["mr"] = { + "watch": next((x for x in mr.get("watch") or [] if _same_code(x.get("ts_code"), code)), None), + "avoid": next((x for x in mr.get("avoid") or [] if _same_code(x.get("ts_code"), code)), None)} + out["sources"]["mr"] = {"ok": True} + except Exception as e: + logger.warning("[单票信号] %s mr 读取失败: %s", code, e) + out["sources"]["mr"] = {"ok": False, "error": "%s: %s" % (type(e).__name__, e)} + + # 资金异动同时进时间线与快照段: 时间线给"什么时候进的钱", 快照段给"现在是什么状态"。 + out["snapshot"]["metrics"] = [r for r in rows if r["src"] == "metrics"][:5] + rows = _dedup(rows) + rows.sort(key=lambda r: (r.get("ts") or 0), reverse=True) + out["timeline"] = rows + + _BYCODE_CACHE[six] = (now, out) + if len(_BYCODE_CACHE) > _BYCODE_MAX: + for k in sorted(_BYCODE_CACHE, key=lambda x: _BYCODE_CACHE[x][0])[:-_BYCODE_MAX]: + _BYCODE_CACHE.pop(k, None) + return out diff --git a/app/web/main.py b/app/web/main.py index be95f68..8d70c29 100644 --- a/app/web/main.py +++ b/app/web/main.py @@ -692,6 +692,20 @@ def api_upstream_signals(): return ok(upstream_signals.snapshot) +@app.get("/api/upstream/signals/{ts_code}") +def api_upstream_signals_by_code(ts_code: str): + """一只票**今天的全部**上游信号 (2026-09-09, 台账 055)。 + + 与 /api/upstream-signals 的分工: 那个是全场最近 N 条的快照, 告警每类只留最新 20 条, + 到下午一只活跃票早盘的信号早被挤掉; 这个按代码倒序回扫到当天开盘, 所以才叫"全部"。 + + **只能点击触发, 绝不能进页面轮询** —— 进轮询就是每 30 秒对上游做一次全流扫描, + 把只读观察变成压力源。内部纪律与快照相同: 一律 XREVRANGE/ZREVRANGE, 不建消费组、不 ACK、不写。 + """ + from app.services import upstream_signals + return ok(upstream_signals.by_code, ts_code) + + # ================================================================ 宏观择时 (MACRO_TIMING_PLAN.md) @app.get("/api/macro/status") def api_macro_status(): diff --git a/app/web/static/index.html b/app/web/static/index.html index 8799c3c..ccf8281 100644 --- a/app/web/static/index.html +++ b/app/web/static/index.html @@ -70,6 +70,10 @@ body{margin:0;background:var(--plane);color:var(--ink); .panel h3,.tblk h3{margin:0 0 10px;font-size:15px;color:var(--ink);font-weight:660;} .sub2{color:var(--muted);font-size:12.5px;margin:-4px 0 12px;} .banner{margin:10px 0;} +/* 2026-09-09: 顶部那条横幅在文档流里, 浮层(2010)、遮罩(2005)、弹层都压得住它, 往下滚一屏 + 就彻底看不见 —— 于是"报错了"在页面上表现成"点了没反应"。再浮一份固定在顶上。 */ +.banner-float{position:fixed;top:8px;left:50%;transform:translateX(-50%);z-index:99995; + width:min(760px,92vw);box-shadow:0 8px 28px rgba(17,17,16,.28);} .attn{color:var(--warn);padding:4px 0;} .tcard{background:var(--surface-2);border:1px solid var(--hair);border-radius:11px;padding:12px 14px;margin-bottom:10px;} .feed-row{padding:5px 0;border-bottom:1px solid var(--hair);font-size:13px;} @@ -148,7 +152,10 @@ pre.json{background:var(--surface-2);border:1px solid var(--hair);border-radius: .sig-ev{display:grid;grid-template-columns:repeat(auto-fit,minmax(300px,1fr));gap:14px;margin-top:14px;} .sig-ev .up-cell{margin:0;} .sig-none{color:var(--ink-2);padding:10px 0;font-size:12.5px;} -.up-h .src-err{color:var(--stop);font-weight:400;font-size:11.5px;} +/* 2026-09-09: 原来只写了 .up-h, 而信号栏的分组头是 .srch —— 角标即使渲染出来也没有颜色, + 在风控面板上看起来就跟普通文字一样。两处都要。 */ +.up-h .src-err,.srch .src-err{color:var(--stop);font-weight:400;font-size:11.5px;} +.sig-fail{color:var(--stop);font-size:11.5px;padding:3px 0;} .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;} @@ -365,6 +372,26 @@ body.dock-r:not(.r-fold) .side-r .strip{display:none;} .oprow2 .b{flex:none;width:92px;} .oprow2 .b .el-button{width:100%;} /* 登录遮罩与登录人胶囊 (2026-08 加权限) */ +/* ── 单票信号抽屉 (2026-09-09): 点一只票看它今天的全部信号 ── + 层级 3000/2995 是显式写死的, 不能删也不能改小: 见 .banner-float 那条注释。 */ +.cs-scrim{position:fixed;inset:0;background:rgba(17,17,16,.34);z-index:2995;} +.cs-drawer{position:fixed;top:0;right:0;bottom:0;width:min(620px,94vw);z-index:3000; + background:var(--surface);border-left:1px solid var(--hair);display:flex;flex-direction:column; + box-shadow:-12px 0 40px rgba(17,17,16,.22);} +.cs-h{display:flex;align-items:center;gap:8px;padding:11px 13px;border-bottom:1px solid var(--hair);flex:none;flex-wrap:wrap;} +.cs-h .ttl{font-size:14px;font-weight:660;} +.cs-b{flex:1;overflow:auto;padding:10px 13px 16px;} +.cs-sect{margin:0 0 12px;} +.cs-sect>.h{font-size:12px;font-weight:660;color:var(--muted);margin:0 0 5px; + border-bottom:1px dashed var(--hair);padding-bottom:4px;} +.cs-row{display:flex;align-items:baseline;gap:7px;padding:4px 0;font-size:12.5px; + border-bottom:1px solid rgba(0,0,0,.04);flex-wrap:wrap;} +.cs-row .t{font-size:11.5px;color:var(--muted);width:40px;flex:none;} +.cs-row .s{font-size:11px;padding:0 5px;border-radius:5px;background:var(--chip);color:var(--muted);flex:none;} +.cs-row.old{opacity:.5;} +.cs-dup{font-size:11px;color:var(--muted);} +.cs-warn{background:#fff8e6;border:1px solid #f0d9a0;border-radius:8px;padding:7px 9px; + font-size:12px;margin:0 0 10px;line-height:1.5;} .login-mask{position:fixed;inset:0;z-index:9999;display:flex;align-items:center;justify-content:center; background:rgba(20,20,18,.55);backdrop-filter:blur(4px);} .login-card{width:320px;background:var(--surface);border:1px solid var(--hair);border-radius:14px; @@ -457,6 +484,10 @@ body.dock-r:not(.r-fold) .side-r .strip{display:none;} + + +