diff --git a/deploy/bot/xmpp_agent_core.py b/deploy/bot/xmpp_agent_core.py index 94266f6f..401a802e 100644 --- a/deploy/bot/xmpp_agent_core.py +++ b/deploy/bot/xmpp_agent_core.py @@ -385,6 +385,14 @@ def _inbound_loop(bot): body = _process_image_message(img_url) log.info(f"🖼️ OCR完成: {body[:100]}") reply = call_hermes(body) + if not reply: + # 网关不可用/超时:绝不静默吞掉用户消息——重试一次,仍失败则兜底告知 + log.warning("call_hermes 无响应(网关可能重启中),10s后重试一次") + time.sleep(10) + reply = call_hermes(body) + if not reply: + reply = "⚠️ 网关暂时不可用(可能在重启/网络波动),你的消息我已收到但未能处理。请几分钟后重发一次,或稍候我会自行恢复。" + log.error("call_hermes 重试仍失败,已回兜底消息") ack_mgr.ack(body[:40]) if reply: with _outbound_lock: diff --git a/deploy/profile-scripts/alert_helper.py b/deploy/profile-scripts/alert_helper.py index 08605dbd..cdf850a3 100644 --- a/deploy/profile-scripts/alert_helper.py +++ b/deploy/profile-scripts/alert_helper.py @@ -30,6 +30,7 @@ LOG = "/home/hmo/MoFin/gateway/logs/alert_helper.log" INFO_MIN_INTERVAL = 1800 # 同类 info 30 分钟最多 1 条 INFO_MAX_LINES = 8 # info 篇幅上限 ACTION_MAX_LINES = 30 # action 篇幅上限(宽松但不失控) +ACTION_DEDUP_SEC = 300 # action 相同内容 5 分钟内不重复发(防同秒双发/竞态) DEDUP_SEC = 24 * 3600 # 相同内容 24h 不重复 @@ -103,6 +104,10 @@ def notify(category, body, level=INFO): entry["suppressed"] = 0 prefix = f"📟【MoFin系统·{category}】(非知微本人)" else: + # ACTION: 不限速,但相同内容 5 分钟内去重(防竞态双发) + if body_hash == entry.get("last_hash") and now - entry.get("last_ts", 0) < ACTION_DEDUP_SEC: + _log(f"[{category}] ACTION 内容重复(5min内),静默") + return False lines = body.splitlines() if len(lines) > ACTION_MAX_LINES: body = "\n".join(lines[:ACTION_MAX_LINES]) + f"\n…(共{len(lines)}行)" diff --git a/deploy/profile-scripts/intraday_health_check.py b/deploy/profile-scripts/intraday_health_check.py index ea93960e..2927de72 100644 --- a/deploy/profile-scripts/intraday_health_check.py +++ b/deploy/profile-scripts/intraday_health_check.py @@ -65,22 +65,7 @@ def db_today_count(table, date_col): return -1 -def check_xiaoguo(): - """小果管道:进程/scanner有数据/API可达(降级不报错)""" - # 进程 — 不一定有常驻进程(no_agent cron模式) - # 数据 — 今日有扫描记录 - scans_today = db_today_count("xiaoguo_scan_tracker", "last_scanned_at") - if scans_today <= 0: - # 可能是小果离线了,不报严重,记录即可 - return - # API — 用socket快速检测可达性(3s超时) - try: - s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - s.settimeout(3) - s.connect(("node122", 18003)) - s.close() - except: - pass +# check_xiaoguo 已删除(小果 2026-07-20 全线退役,残留检查只会制造假警报) def check_price_monitor(): @@ -152,37 +137,27 @@ def check_price_monitor(): def check_bots(): zhiwei = subprocess.run(["systemctl", "is-active", "xmpp-zhiwei.service"], capture_output=True, text=True, timeout=5).stdout.strip() == "active" - xiaoguo = subprocess.run(["systemctl", "is-active", "xmpp-xiaoguo.service"], - capture_output=True, text=True, timeout=5).stdout.strip() == "active" log(zhiwei, "知微XMPP Bot离线") - log(xiaoguo, "小果XMPP Bot离线") + # 小果已全线退役(2026-07-20),不再检查 xmpp-xiaoguo / :8645 def check_gateways(): - log(check_port(8643), "知微Gateway :8643 未监听") - log(check_port(8645), "小果Gateway :8645 未监听") + log(check_port(8643), "知微Gateway :8643 未监听") def check_signal_pipeline(): - """信号从xiaoguo_scanner→signal_news→consumer是否通畅""" - unproc = 0 + """信号管道通畅性(小果已退役,只查全量积压)""" total_unproc = 0 try: conn = get_conn() - # xiaoguo 信号堆积(4h以内时效) - r = conn.execute("SELECT COUNT(*) FROM signal_news WHERE source LIKE 'xiaoguo%' AND (processed=0 OR processed IS NULL) AND created_at > datetime('now', '-4 hours')").fetchone() - unproc = r[0] # 全量未处理信号(跨来源) r2 = conn.execute("SELECT COUNT(*) FROM signal_news WHERE (processed=0 OR processed IS NULL)").fetchone() total_unproc = r2[0] conn.close() except: pass - log(unproc < 30, f"xiaoguo信号堆积: {unproc}条未处理(需<30)") - # 其他来源信号积压预警 - other = total_unproc - unproc - if other > 50: - log(False, f"其它来源信号积压: {other}条未处理(divergence_watch/trend等无consumer)") + if total_unproc > 50: + log(False, f"信号积压: {total_unproc}条未处理") # 宏观风险状态检查 try: @@ -248,7 +223,7 @@ def main(): check_bots() check_gateways() - check_xiaoguo() + # check_xiaoguo 已移除(小果已退役) if 9 <= now.hour < 16: # 开盘前10分钟(9:00-9:10)跳过价格新鲜度检查 # price_monitor 从 09:00 才开始启动,09:01 检查时数据还未更新(前一天收盘数据) diff --git a/deploy/profile-scripts/price_monitor.py b/deploy/profile-scripts/price_monitor.py index f28e7896..652e1d9e 100644 --- a/deploy/profile-scripts/price_monitor.py +++ b/deploy/profile-scripts/price_monitor.py @@ -37,7 +37,8 @@ XMPP_USER = "hmo@yoin.fun" XMPP_BRIDGE = "http://127.0.0.1:5805/" def push_to_xmpp(text): - """通过知微 HTTP bridge 推送到Dad私信""" + """原始直推(已废弃直用)——保留给极少数必须原样的场景。 + 新代码请用 _push_action/_push_digest。""" if not text.strip(): return try: @@ -51,6 +52,25 @@ def push_to_xmpp(text): except Exception as e: print(f"[XMPP推送失败] {e}", file=sys.stderr) + +# ── 分级推送(2026-07-21 信噪比纪律,红线#12)── +# ACTION: 破止损/重评确认的操作信号 — 直通不限速 +# INFO: 未确认的进区提示 — 聚合成摘要,30min 限 1 条 +def _push_action(category, text): + try: + from alert_helper import notify, ACTION + notify(category, text, ACTION) + except Exception as e: + print(f"[ACTION推送失败] {e}", file=sys.stderr) + + +def _push_digest(category, text): + try: + from alert_helper import notify, INFO + notify(category, text, INFO) + except Exception as e: + print(f"[INFO推送失败] {e}", file=sys.stderr) + # ── 批量拉取价格 ────────────────────────────────────────────────────────── def fetch_all_prices(codes): @@ -460,6 +480,9 @@ def run_once(round_label=""): # 批量拉取这些股票的价格 prices = fetch_all_prices(list(check_codes)) + # 本轮进区事件收集(聚合成一条摘要推送,替代逐条轰炸) + _zone_entries = [] + for d in active: code = d["code"] trig = d.get("trigger", {}) @@ -496,7 +519,7 @@ def run_once(round_label=""): if _budget_low: outputs.append(f" 📨 止损触发(超时跳过重评)→已推送Dad") if _can_push(code, "stop_loss"): - push_to_xmpp(f"⚠️ {name}({code}) {price} → 跌破止损{hi}!") + _push_action("止损告警", f"⚠️ {name}({code}) {price} → 跌破止损{hi}!") else: try: cost = d.get("cost", 0) or 0 @@ -512,7 +535,7 @@ def run_once(round_label=""): rr = result.get("rr_ratio", 0) if _can_push(code, "stop_loss"): msg = f"🔔 {name}({code}) 价{price}→触发操作区间{max(buy_lo,0):.2f}~{buy_hi:.2f},已触发重评|RR={rr}" - push_to_xmpp(msg) + _push_action("操作信号", msg) outputs.append(f" 📨 止损重评→已推送Dad: {action}") except Exception as e: outputs.append(f" ⚠️ 止损重评失败: {e}") @@ -529,11 +552,11 @@ def run_once(round_label=""): extra = f"({act})" outputs.append(f"⚡ {name}({code}) {price} → 进入{label}{lo}~{hi}{extra}") record_event(code, name, "entry_zone", price, f"{lo}~{hi}", label) - # 进入区间 → 立即重评并推送给Dad(时间不够则跳过重评直接推原始告警) + # 进入区间 → 立即重评并推送给Dad(时间不够则记入摘要,不逐条轰炸) if _budget_low: if _can_push(code, key): - push_to_xmpp(f"⚡ {name}({code}) {price} → 进入{label}{lo}~{hi}") - outputs.append(f" 📨 区间触发(超时跳过重评)→已推送Dad") + _zone_entries.append(f"{name}({code}) {price}→{label}{lo}~{hi}") + outputs.append(f" 📨 区间触发(超时)→记入摘要") else: try: cost = d.get("cost", 0) or 0 @@ -552,7 +575,7 @@ def run_once(round_label=""): rr = result.get("rr_ratio", 0) if _can_push(code, key): msg = f"🔔 {name}({code}) 价{price}→触发{zone_desc},已触发重评|RR={rr}" - push_to_xmpp(msg) + _push_action("操作信号", msg) outputs.append(f" 📨 区间触发重评→已推送Dad: {action}") else: reason = f"重评结果:{timing_signal},不构成操作建议" @@ -568,6 +591,12 @@ def run_once(round_label=""): state[code][key] = False state_updated = True + # === 第二步收尾:进区事件聚合成一条摘要推送(INFO级,30min限1条+截断)=== + if _zone_entries: + _digest = f"📋 {len(_zone_entries)}只进入操作区:\n" + "\n".join(f"• {e}" for e in _zone_entries) + _push_digest("盘中触发", _digest) + outputs.append(f"📨 进区摘要({len(_zone_entries)}只)→已按INFO策略推送") + # === 第三步:买入区偏离检测 + 自动重评 === reassesed_codes = [] # 先做急跌检测(仅持仓,自选股不推送暴跌告警) @@ -595,7 +624,7 @@ def run_once(round_label=""): stop_loss = d.get("stop_loss", 0) sl_note = f" 止损{stop_loss}" if stop_loss else "" msg = f"🔻 {name}({code}) {price} 暴跌{cp:.1f}%!{sl_note}" - push_to_xmpp(msg) + _push_action("急跌告警", msg) outputs.append(msg) state.setdefault(code, {})["__sharp_decline_triggered"] = True state_updated = True