Merge branch 'session-work'

# Conflicts:
#	deploy/profile-scripts/intraday_health_check.py
This commit is contained in:
知微
2026-07-21 10:43:41 +08:00
3 changed files with 50 additions and 8 deletions
+8
View File
@@ -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:
+5
View File
@@ -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)}行)"
+37 -8
View File
@@ -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