333 lines
15 KiB
Python
333 lines
15 KiB
Python
#!/usr/bin/env python3
|
||
"""functional_health_check.py — L1 功能健康检查
|
||
|
||
判据不是"进程活着",而是"功能是否达成":每个核心模块检查其
|
||
**输出物的新鲜度和有效性**。输出物不新鲜 = 功能未达成 = 异常。
|
||
|
||
频率:交易时段每 15 分钟(cron),输出 gateway/logs/functional_health.json。
|
||
责任矩阵 L1。修复由 L3 self_repair 读取本报告执行。
|
||
"""
|
||
import os, sys, json, sqlite3, subprocess
|
||
from datetime import datetime, timedelta
|
||
|
||
# ── 消息通道统一路由(broadcast/xmpp by delivery) ──
|
||
try:
|
||
from messenger import install_stdio_hook as _msh
|
||
_msh()
|
||
except Exception:
|
||
pass
|
||
|
||
sys.path.insert(0, '/home/hmo/MoFin')
|
||
|
||
DB = '/home/hmo/MoFin/data/mofin.db'
|
||
OUT = '/home/hmo/MoFin/gateway/logs/functional_health.json'
|
||
|
||
# ── 功能注册表:模块 → 功能判据 ──────────────────────────────
|
||
# type:
|
||
# db_freshness: DB 表 MAX(col) 距今不超过 max_age(trading=交易时段才检查, daily=每日, always=总是)
|
||
# file_freshness: 文件 mtime 距今不超过 max_age
|
||
# bot_activity: XMPP bot 活动检查(journal)
|
||
# agent_log: gateway agent.log 最近有成功 LLM 调用
|
||
# soft=True: 无输出可能属正常(如区间未触发),降级为 warn 而非 fail
|
||
REGISTRY = [
|
||
{"module": "price_monitor", "function": "实时价格写入DB(live_prices)",
|
||
"check": {"type": "db_freshness", "table": "live_prices", "col": "updated_at",
|
||
"max_age_min": 6, "when": "trading"},
|
||
"repair": {"action": "rerun_script", "script": "price_monitor.py"}},
|
||
{"module": "market_watch", "function": "市场快照采集(market_snapshots)",
|
||
"check": {"type": "db_freshness", "table": "market_snapshots", "col": "created_at",
|
||
"max_age_min": 40, "when": "trading"},
|
||
"repair": {"action": "rerun_script", "script": "market_watch.py"}},
|
||
{"module": "mtf_cache", "function": "多周期均线缓存刷新(mtf_cache)",
|
||
"check": {"type": "db_freshness", "table": "mtf_cache", "col": "updated_at",
|
||
"max_age_min": 240, "when": "trading"}, # 2026-08-24 75→240: 调度为50 7,11,14三趟/天,间隔3.5h,75min必误报
|
||
"repair": {"action": "rerun_script", "script": "refresh_mtf_cache.py"}},
|
||
{"module": "macro_context", "function": "宏观上下文刷新(macro_context_log)",
|
||
"check": {"type": "db_freshness", "table": "macro_context_log", "col": "created_at",
|
||
"max_age_min": 45, "when": "trading"},
|
||
"repair": {"action": "rerun_script", "script": "refresh_macro_context.py"}},
|
||
{"module": "health_collector", "function": "健康数据采集(mofin_health.json)",
|
||
"check": {"type": "file_freshness", "path": "/home/hmo/web-dashboard/static/mofin_health.json",
|
||
"max_age_min": 20, "when": "trading"},
|
||
"repair": {"action": "rerun_script", "script": "mofin_health.py"}},
|
||
{"module": "premarket", "function": "盘前全量重评(premarket summary)",
|
||
"check": {"type": "file_freshness", "path": "/tmp/mofin_premarket/summary.json",
|
||
"max_age_h": 26, "when": "daily"},
|
||
"repair": {"action": "rerun_script", "script": "premarket_full_review.py"}},
|
||
{"module": "reassess_daily", "function": "每日12维重评完成度(持仓当日覆盖+分析存在率)",
|
||
"check": {"type": "reassess_daily", "when": "daily"},
|
||
"repair": {"action": "rerun_script", "script": "batch_reassess.py --type holding --today"}},
|
||
{"module": "gateway_llm", "function": "LLM调用链可用(ocg-router探测)",
|
||
"check": {"type": "ocg_router", "when": "always"}, # 2026-08-24 agent.log检查过时(红线11后LLM走ocg-router直连,无gateway流量≠故障)
|
||
"repair": {"action": "llm_diagnose"}},
|
||
{"module": "sense_ocr", "function": "识图服务(SenseNova OCR配置+API可达)",
|
||
"check": {"type": "ocr_health", "when": "always"},
|
||
"repair": {"action": "llm_diagnose"}},
|
||
{"module": "xmpp_bot", "function": "XMPP消息收发(bot journal)",
|
||
"check": {"type": "bot_activity", "when": "always"},
|
||
"repair": {"action": "llm_diagnose"}},
|
||
{"module": "cron_engine", "function": "cron调度引擎本身(有job在最近10min运行)",
|
||
"check": {"type": "cron_engine", "max_age_min": 12, "when": "trading"},
|
||
"repair": {"action": "llm_diagnose"}},
|
||
]
|
||
|
||
|
||
def is_trading_now(now):
|
||
return now.weekday() < 5 and 9 <= now.hour <= 16
|
||
|
||
|
||
def check_db_freshness(conn, chk, now):
|
||
try:
|
||
row = conn.execute(f"SELECT MAX({chk['col']}) FROM {chk['table']}").fetchone()
|
||
if not row or not row[0]:
|
||
return "fail", f"表 {chk['table']} 无数据"
|
||
last = datetime.fromisoformat(str(row[0]).replace("Z", ""))
|
||
age_min = (now - last).total_seconds() / 60
|
||
limit = chk.get("max_age_min", chk.get("max_age_h", 24) * 60)
|
||
if age_min > limit:
|
||
return ("warn" if chk.get("soft") else "fail"), \
|
||
f"{chk['table']} 最新记录 {last.strftime('%m-%d %H:%M')}({age_min/60:.1f}h前,阈值{limit}min)"
|
||
return "ok", f"{age_min:.0f}min前"
|
||
except Exception as e:
|
||
return "fail", f"查询失败: {str(e)[:80]}"
|
||
|
||
|
||
def check_file_freshness(chk, now):
|
||
p = chk["path"]
|
||
if not os.path.exists(p):
|
||
return "fail", f"文件不存在: {p}"
|
||
age_min = (now.timestamp() - os.path.getmtime(p)) / 60
|
||
limit = chk.get("max_age_min", chk.get("max_age_h", 24) * 60)
|
||
if age_min > limit:
|
||
return ("warn" if chk.get("soft") else "fail"), \
|
||
f"{os.path.basename(p)} {age_min/60:.1f}h 未更新(阈值{limit}min)"
|
||
return "ok", f"{age_min:.0f}min前"
|
||
|
||
|
||
def check_agent_log(chk, now):
|
||
try:
|
||
sys.path.insert(0, '/home/hmo/MoFin')
|
||
from xmpp_logger import _scan_agent_log
|
||
r = _scan_agent_log(now.timestamp(), "zhiwei")
|
||
if r["status"] == "ok":
|
||
age = r.get("age_sec", -1)
|
||
if age > chk["max_age_min"] * 60:
|
||
return "warn", f"最近成功LLM调用在 {age//60}min 前(无新流量,可能正常)"
|
||
return "ok", f"latency={r.get('latency')}"
|
||
return "fail", f"agent.log 最近调用失败: {r.get('error','?')[:100]}"
|
||
except Exception as e:
|
||
return "fail", f"检查失败: {e}"
|
||
|
||
|
||
def check_ocg_router(chk, now):
|
||
"""LLM链路真实可用性:ocg-router /v1/models 探测(2026-08-24 替代过时 agent.log 检查)"""
|
||
try:
|
||
import urllib.request, json as _j
|
||
req = urllib.request.Request("http://127.0.0.1:19878/v1/models",
|
||
headers={"Authorization": "Bearer ocg-router-local"})
|
||
with urllib.request.urlopen(req, timeout=8) as r:
|
||
d = _j.loads(r.read().decode())
|
||
n = len(d.get("data", []))
|
||
if n > 0:
|
||
return "ok", f"ocg-router在线 {n}模型"
|
||
return "fail", "ocg-router模型列表空"
|
||
except Exception as e:
|
||
return "fail", f"ocg-router不可达: {str(e)[:80]}"
|
||
|
||
|
||
def check_bot_activity(chk, now):
|
||
try:
|
||
# 第一判据:服务当前是否 active(权威"现在在不在跑")
|
||
r2 = subprocess.run(["systemctl", "is-active", "xmpp-zhiwei"],
|
||
capture_output=True, text=True, timeout=5)
|
||
if r2.stdout.strip() != "active":
|
||
return "fail", "bot 服务非 active(已停止)"
|
||
# 第二判据:最近 journal 是否处于断线循环
|
||
r = subprocess.run(
|
||
["journalctl", "-u", "xmpp-zhiwei", "--no-pager", "--since", "30 min ago", "-o", "cat"],
|
||
capture_output=True, timeout=8, text=True)
|
||
lines = r.stdout
|
||
if "连接超时" in lines and "就绪" not in lines:
|
||
return "fail", "bot 处于断线重连循环"
|
||
if "XMPP 就绪" in lines or "已发送" in lines or "收到" in lines:
|
||
return "ok", "bot 活动正常"
|
||
return "ok", "bot 在线空闲(无新消息)"
|
||
except Exception as e:
|
||
return "fail", f"检查失败: {e}"
|
||
|
||
|
||
def check_ocr_health(chk, now):
|
||
"""识图服务健康:OCR 配置存在 + SenseNova API 可达(TCP 443)。
|
||
不发真实 OCR 请求(省钱),真实端到端验证走 K 测试。"""
|
||
import json as _json, socket
|
||
# 1. OCR 配置存在且含 key
|
||
cfg_path = '/home/hmo/.config/mofin/ocr_config.json'
|
||
if not os.path.exists(cfg_path):
|
||
return "fail", "OCR 配置文件不存在: /home/hmo/.config/mofin/ocr_config.json"
|
||
try:
|
||
cfg = _json.load(open(cfg_path))
|
||
if not cfg.get('key') or not cfg.get('base_url'):
|
||
return "fail", "OCR 配置缺 key 或 base_url"
|
||
except Exception as e:
|
||
return "fail", f"OCR 配置解析失败: {str(e)[:60]}"
|
||
# 2. SenseNova API TCP 可达
|
||
try:
|
||
host = cfg['base_url'].split('//')[1].split('/')[0].split(':')[0]
|
||
port = int(cfg['base_url'].split('//')[1].split('/')[0].split(':')[1]) if ':' in cfg['base_url'].split('//')[1].split('/')[0] else 443
|
||
s = socket.create_connection((host, port), timeout=5)
|
||
s.close()
|
||
except Exception as e:
|
||
return "fail", f"SenseNova API 不可达 ({host}): {str(e)[:60]}"
|
||
return "ok", f"配置 OK + {host}:{port} 可达"
|
||
|
||
|
||
def check_cron_engine(chk, now):
|
||
"""cron 引擎:最近 max_age_min 内是否有任何 job 运行过"""
|
||
try:
|
||
d = json.load(open('/home/hmo/.hermes/profiles/position-analyst/cron/jobs.json'))
|
||
jobs = d if isinstance(d, list) else d.get('jobs', [])
|
||
latest = None
|
||
for j in jobs:
|
||
lr = j.get('last_run_at')
|
||
if lr:
|
||
try:
|
||
t = datetime.fromisoformat(lr.replace('Z', '+00:00')).replace(tzinfo=None)
|
||
if latest is None or t > latest:
|
||
latest = t
|
||
except Exception:
|
||
pass
|
||
if not latest:
|
||
return "fail", "所有 job 均无运行记录"
|
||
age_min = (now - latest).total_seconds() / 60
|
||
if age_min > chk["max_age_min"]:
|
||
return "fail", f"最近 job 运行在 {age_min:.0f}min 前,调度引擎疑似停摆"
|
||
return "ok", f"最近 job 运行于 {age_min:.0f}min 前"
|
||
except Exception as e:
|
||
return "fail", f"检查失败: {e}"
|
||
|
||
|
||
def check_reassess_daily(conn, now):
|
||
"""每日12维重评完成度:
|
||
- 持仓:当日已重评比例(交易日13:00后应≈100%,此前按上一交易日)+ 有分析比例
|
||
- 自选:近24h重评比例(补评管道是否活着)
|
||
阈值:持仓当日覆盖<50% 或 分析存在率<70% → fail;<90% → warn
|
||
"""
|
||
try:
|
||
today = now.strftime("%Y-%m-%d")
|
||
rows = conn.execute(
|
||
"SELECT decision_type, reassessed_at, full_analysis FROM holding_strategies WHERE status='active'").fetchall()
|
||
h_total = h_today = h_fa = w_total = w_24h = 0
|
||
cutoff = now - timedelta(hours=24)
|
||
for dt, ra, fa in rows:
|
||
fa_ok = bool(fa and str(fa).strip())
|
||
ra_dt = None
|
||
if ra:
|
||
try:
|
||
ra_dt = datetime.fromisoformat(str(ra)[:19])
|
||
except Exception:
|
||
pass
|
||
if dt == '持仓策略':
|
||
h_total += 1
|
||
if fa_ok:
|
||
h_fa += 1
|
||
if ra_dt and ra_dt.strftime("%Y-%m-%d") == today:
|
||
h_today += 1
|
||
else:
|
||
w_total += 1
|
||
if ra_dt and ra_dt >= cutoff:
|
||
w_24h += 1
|
||
if h_total == 0:
|
||
return "fail", "holding_strategies 无持仓数据"
|
||
h_cov = h_today / h_total
|
||
h_fa_rate = h_fa / h_total
|
||
w_cov = w_24h / w_total if w_total else 1
|
||
detail = (f"持仓 今日重评{h_today}/{h_total}({h_cov:.0%}) 有分析{h_fa}/{h_total}({h_fa_rate:.0%}) | "
|
||
f"自选 24h重评{w_24h}/{w_total}({w_cov:.0%})")
|
||
# 12:35 前补评窗口未完成属正常,降级 warn
|
||
grace = now.hour < 13 or now.weekday() >= 5
|
||
if h_fa_rate < 0.7:
|
||
return "fail", f"分析存在率过低: {detail}"
|
||
if h_cov < 0.5 and not grace:
|
||
return "fail", f"今日持仓重评覆盖不足: {detail}"
|
||
if h_cov < 0.9 and not grace:
|
||
return "warn", f"覆盖不完整: {detail}"
|
||
if grace and (h_cov < 0.9 or w_cov < 0.2):
|
||
return "warn", f"补评窗口中: {detail}"
|
||
return "ok", detail
|
||
except Exception as e:
|
||
return "fail", f"检查异常: {e}"
|
||
|
||
|
||
def main():
|
||
now = datetime.now()
|
||
trading = is_trading_now(now)
|
||
results = []
|
||
conn = sqlite3.connect(DB, timeout=10)
|
||
|
||
for item in REGISTRY:
|
||
chk = item["check"]
|
||
when = chk.get("when", "always")
|
||
if when == "trading" and not trading:
|
||
results.append({"module": item["module"], "function": item["function"],
|
||
"status": "skip", "reason": "非交易时段"})
|
||
continue
|
||
# 开盘 grace:各管道 09:00 才启动,09:15 前隔夜数据属正常,不误报
|
||
if when == "trading" and now.weekday() < 5 and now.hour == 9 and now.minute < 15:
|
||
results.append({"module": item["module"], "function": item["function"],
|
||
"status": "skip", "reason": "开盘grace期"})
|
||
continue
|
||
if when == "daily":
|
||
# 每日类:只 fail 不 warn,且非交易日放宽到 72h
|
||
if now.weekday() >= 5:
|
||
chk = dict(chk)
|
||
chk["max_age_h"] = 72
|
||
|
||
t = chk["type"]
|
||
if t == "db_freshness":
|
||
status, reason = check_db_freshness(conn, chk, now)
|
||
elif t == "file_freshness":
|
||
status, reason = check_file_freshness(chk, now)
|
||
elif t == "agent_log":
|
||
status, reason = check_agent_log(chk, now)
|
||
elif t == "ocg_router":
|
||
status, reason = check_ocg_router(chk, now)
|
||
elif t == "bot_activity":
|
||
status, reason = check_bot_activity(chk, now)
|
||
elif t == "ocr_health":
|
||
status, reason = check_ocr_health(chk, now)
|
||
elif t == "cron_engine":
|
||
status, reason = check_cron_engine(chk, now)
|
||
elif t == "reassess_daily":
|
||
status, reason = check_reassess_daily(conn, now)
|
||
else:
|
||
status, reason = "fail", f"未知检查类型 {t}"
|
||
|
||
results.append({"module": item["module"], "function": item["function"],
|
||
"status": status, "reason": reason,
|
||
"repair": item.get("repair")})
|
||
|
||
conn.close()
|
||
|
||
fails = [r for r in results if r["status"] == "fail"]
|
||
warns = [r for r in results if r["status"] == "warn"]
|
||
report = {
|
||
"generated_at": now.strftime("%Y-%m-%d %H:%M:%S"),
|
||
"trading_hours": trading,
|
||
"summary": {"total": len(results), "ok": sum(1 for r in results if r["status"] == "ok"),
|
||
"warn": len(warns), "fail": len(fails),
|
||
"skip": sum(1 for r in results if r["status"] == "skip")},
|
||
"status": "fail" if fails else ("warn" if warns else "ok"),
|
||
"checks": results,
|
||
}
|
||
os.makedirs(os.path.dirname(OUT), exist_ok=True)
|
||
with open(OUT, "w", encoding="utf-8") as f:
|
||
json.dump(report, f, ensure_ascii=False, indent=2)
|
||
|
||
print(f"功能健康: {report['summary']} status={report['status']}")
|
||
for r in results:
|
||
if r["status"] in ("fail", "warn"):
|
||
print(f" {r['status'].upper()} {r['module']}: {r['reason']}")
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main() |