"""sanguo-data-monitor — VPS 数据面监控班次(spec §20.15.3, 23:50 evening / 08:05 morning). 只读体检不修数据; 告警写 data/alerts.db(单写者), 前端/晨检消费 health JSON。 日志源=akshare_static_download 自身 FileHandler(data/static/logs/, utf-8 确定性) ——migration_logs 是 PS 重定向副本身受编码坑, 不作解析源。 退出码: 0=绿/黄, 3=有红(wrapper 恒转 0, 红态由 alerts/health 承载)。 """ import argparse import datetime as dt import glob import json import os import subprocess import sys import traceback _HERE = os.path.dirname(os.path.abspath(__file__)) _REPO = os.path.dirname(os.path.dirname(_HERE)) for _p in (_HERE, _REPO): if _p not in sys.path: sys.path.insert(0, _p) from sanguo_data import alerts_store as store # noqa: E402 import monitor_checks as chk # noqa: E402 import monitor_registry as reg # noqa: E402 _TODAY_OVERRIDE = None # 测试注入口 _HEALTH_DIR = "monitor" _LOG_GLOB = os.path.join("static", "logs", "akshare_static_*.log") # 过程层日志窗=26h(两班通用): 覆盖晨班接昨夜+周六 08:05 接 03:00 周班; # 不用长窗——旧失败会每日重炒(09-26 实证: 上周五 RemoteDisconnected 连带红) _LOG_WINDOW_H = 26 _OK_LAST_RESULTS = {"0", "3", "267009", "267011"} # 0=成功 3=vintage 类「有发现」设计退出码(非故障) 267009=运行中 267011=未跑过 def _now(): if _TODAY_OVERRIDE is not None: return dt.datetime.combine(_TODAY_OVERRIDE, dt.time(23, 50)) return dt.datetime.now() def _event(key, severity, title, detail, evidence, host="vps", source="data", today=None): return {"alert_id": f"{key}-{today}", "check_key": key, "source": source, "host": host, "severity": severity, "title": title, "detail": detail, "evidence": evidence} def _ro_conn(path): import sqlite3 return sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=10) # ---------- 结果层: 逐注册表项(log_only 归过程层, 不入此处) ---------- def check_entry(e, ctx): """注册表项分发; 单检抛错=黄降级不阻断全班(D-2: 对齐 strategy 侧 strategy_monitor.check_entry 先例——旧码无隔离, 任一检 KeyError/TypeError =整班崩+wrapper 恒 0=SCHTASKS 年龄判定失明, 最坏 ~22h 仅靠 peer 反盯)。""" try: return _check_entry_impl(e, ctx) except Exception as exc: # noqa: BLE001 return {"key": f"data-{e['name']}-stale", "status": "yellow", "expected": e.get("title", e["name"]), "actual": f"error: {exc!r}", "evidence": {"error": repr(exc)}} def _check_entry_impl(e, ctx): # root 语义: "static"→/static/; "data"→/ # (09-26 修正: 首版误拼 data/data/... 致三域恒黄; fixture 曾同镜像故测试假绿) if e.get("root") == "data": root = ctx["data_root"] else: root = os.path.join(ctx["data_root"], "static") name, kind = e["name"], e["kind"] ev = {"kind": kind} def done(check_kind, status, expected, actual, evidence=None): return {"key": f"data-{name}-{check_kind}", "status": status, "expected": expected, "actual": actual, "evidence": evidence or ev} if kind == "per_date": exp = chk.expected_daily(ctx["days"], ctx["now"], e.get("ready", "20:30"), e.get("lag", 0)) got = chk.latest_stem(os.path.join(root, e["subdir"]), e["pattern"]) latest = got[0].split("_")[0] if got else None if exp is None: # 无日历(db 不可达): 只做存在性弱检查 return done("stale", "green" if got else "yellow", "n/a", latest) if latest is None or latest < exp: # 落后/缺失才告; 领先期望=绿 sev = "yellow" if exp not in ctx["confirmed_days"] \ else e.get("severity", "red") return done("stale", sev, exp, latest, {"producer": e.get("producer"), "unconfirmed_day": exp not in ctx["confirmed_days"]}) rows = chk.parquet_rows(got[1]) if got else None # 滚动分位带(P2-6, 10-01 拍板): row_band=(floor, fallback_hi), 带内 # 判据=P1~P99 滚动带(窗口=最近 30 个 prior 文件, 剔 0 行真空; 观测 # <10 回退静态整带)——根治静态上限被真实增长涨破的手动校准循环 lo, hi, n_obs, band_src = chk.rolling_rows_band( os.path.join(root, e["subdir"]), e["pattern"], got[0] if got else None, floor=e["row_band"][0], fallback_hi=e["row_band"][1]) # 未确认日 0 行空文件(节假日闸门产物)不告 row-band (block_trade # band 下限>0 对 0925 空文件误黄实案); 交易日 0 行仍黄保持感知. # 0928 修正: 豁免按文件日期(latest)判而非 exp——lag 型(margin/ # margin_sse T-1) exp=ok[-1-lag]=确认日, 旧条件对闸门空写的非 # 确认日空文件永不豁免(0927 margin_sse 20260925 空文件假黄实案) if not (rows == 0 and latest not in ctx["confirmed_days"]) \ and (rows is None or not (lo <= rows <= hi)): return done("rows", "yellow", f"{lo:.0f}~{hi:.0f}", rows, {"band_source": band_src, "n_obs": n_obs, "lo_eff": lo, "hi_eff": hi}) return done("stale", "green", exp, latest) if kind == "weekly_friday": exp = chk.expected_weekly(ctx["days"], ctx["now"], e.get("ready", "06:00"), confirmed=ctx["confirmed_days"]) got = chk.latest_stem(os.path.join(root, e["subdir"]), e["pattern"]) latest = got[0].split("_")[0] if got else None if exp is None: return done("stale", "green" if got else "yellow", "n/a", latest) if latest is None or latest < exp: # 周窗锚落后才告(周六班可能超前=绿) sev = "yellow" if exp not in ctx["confirmed_days"] \ else e.get("severity", "red") return done("stale", sev, exp, latest) lo, hi, n_obs, band_src = chk.rolling_rows_band( os.path.join(root, e["subdir"]), e["pattern"], got[0] if got else None, floor=e["row_band"][0], fallback_hi=e["row_band"][1]) rows = chk.parquet_rows(got[1]) if got else None # 同 per_date 0928 修正: 豁免按文件日期判(假周五闸门空写 0 行文件) if not (rows == 0 and latest not in ctx["confirmed_days"]) \ and (rows is None or not (lo <= rows <= hi)): return done("rows", "yellow", f"{lo:.0f}~{hi:.0f}", rows, {"band_source": band_src, "n_obs": n_obs, "lo_eff": lo, "hi_eff": hi}) return done("stale", "green", exp, latest) if kind == "month_bucket": got = chk.latest_stem(os.path.join(root, e["subdir"]), e["pattern"]) if not got: return done("stale", e.get("severity", "yellow"), "any", None) ym = got[0].split("_")[0] try: d0 = dt.date(int(ym[:4]), int(ym[4:6]), 1) except ValueError: return done("stale", "yellow", "YYYYMM", ym) age = (ctx["now"].date() - d0).days if age > e["max_age_days"]: return done("stale", "yellow", f"<={e['max_age_days']}d", f"{age}d") lo, hi = e["row_band"] rows = chk.parquet_rows(got[1]) if rows is None or not (lo <= rows <= hi): # D-6(10-05): 损坏件 rows=None 不再静默绿——对齐 per_date 判定 return done("rows", "yellow", f"{lo}~{hi}", rows) return done("stale", "green", ym, ym) if kind == "year_file": year = str(ctx["now"].year) p = os.path.join(root, e["subdir"], f"{year}.parquet") if not os.path.exists(p): return done("stale", e.get("severity", "yellow"), year, None) mtime = dt.date.fromtimestamp(os.path.getmtime(p)).isoformat() exp = chk.expected_daily(ctx["days"], ctx["now"], e.get("ready", "22:30")) want = exp.replace("-", "") if exp else None got_d = mtime.replace("-", "") if want and got_d < want: return done("stale", e.get("severity", "yellow"), want, got_d) return done("stale", "green", want, got_d) if kind == "dir_mtime": d = os.path.join(root, e["subdir"]) if not os.path.isdir(d): return done("stale", e.get("severity", "yellow"), "exists", None) age_h = (ctx["now"] - dt.datetime.fromtimestamp( chk.dir_latest_mtime(d))).total_seconds() / 3600 if age_h > e["max_age_days"] * 24: return done("stale", "yellow", f"<={e['max_age_days']}d", f"{age_h:.0f}h") return done("stale", "green", f"<={e['max_age_days']}d", f"{age_h:.0f}h") if kind == "db_sample": exp = chk.expected_daily(ctx["days"], ctx["now"], e.get("ready", "22:30"), e.get("lag", 0)) worst, readable = None, False for sym, exc in e["samples"]: got_d = None try: conn = _ro_conn(os.path.join(ctx["data_root"], "quant_trading.db")) try: row = conn.execute( "SELECT MAX(datetime) FROM dbbardata WHERE symbol=? AND " "exchange=? AND interval='d'", (sym, exc)).fetchone() finally: conn.close() # D-10(10-05): close 进 finally, 异常路径不泄漏 got_d = (row[0] or "").replace("-", "") except Exception: # noqa: BLE001 只读探测, 库坏=黄红由调用方判 got_d = None if got_d: readable = True if exp and (got_d is None or got_d < exp) and worst is None: worst = {"symbol": f"{sym}.{exc}", "expected": exp, "actual": got_d} if exp and worst: sev = "yellow" if exp not in ctx["confirmed_days"] \ else e.get("severity", "red") return done("stale", sev, exp, worst["actual"], worst) if exp is None: # D-5(10-05): 日历不可达不再无条件绿(fail-open)——弱检查=样本可读 # (对齐 per_date 存在性弱检), 库也不可达=黄 return done("stale", "green" if readable else "yellow", "n/a", "sample-readable" if readable else "db-unreadable") return done("stale", "green", exp, "ok") if kind == "vintage_json": p = os.path.join(root, e["subdir"]) try: with open(p, encoding="utf-8") as f: data = json.load(f) except (OSError, ValueError): return done("verdict", "yellow", "readable", None) checked = data.get("checked_at", "") try: stale = not checked or (ctx["now"] - dt.datetime.fromisoformat( checked)).total_seconds() > 26 * 3600 except (ValueError, TypeError): stale = True # D-9(10-05): 非字符串 fromisoformat 抛 TypeError 不炸 if stale: return done("verdict", "yellow", "<=26h", checked) # 结论无顶层字段(static_vintage_check 契约): 挖洞藏在 tables.*.holes_* holes = {} for t, entry in (data.get("tables") or {}).items(): n_perm = len(entry.get("holes_permanent", [])) n_back = len(entry.get("holes_backfillable", [])) if n_perm or n_back: holes[t] = {"permanent": n_perm, "backfillable": n_back} if holes: return done("verdict", "yellow", "no-holes", f"{len(holes)} 表有洞") return done("verdict", "green", "no-holes", "clean") if kind == "static_exists": d = os.path.join(root, e["subdir"]) exists = os.path.isdir(d) and any( f.endswith(".parquet") for f in os.listdir(d)) return done("exists", "green" if exists else e.get("severity", "yellow"), "exists", exists) if kind == "disk": # D-8(10-05): 阈值消费注册表(旧硬编码 8/20 与注册表脱钩, crit_gb 全仓 # 无消费=改表不改行为的校准陷阱) r = chk.check_disk(ctx["data_root"], warn_gb=e.get("warn_gb", 20), crit_gb=e.get("crit_gb", 8)) return {"key": "infra-disk-data", "status": r["status"], "expected": f">={e['warn_gb']}G", "actual": f"{r['free_gb']}G", "evidence": r} return done("unknown", "yellow", "n/a", None) # ---------- 过程层: 日志统计行+签名分诊 ---------- def collect_log_layer(data_root, now, window_h=_LOG_WINDOW_H): logs = sorted(glob.glob(os.path.join(data_root, _LOG_GLOB))) keep = [p for p in logs if (now - dt.datetime.fromtimestamp(os.path.getmtime(p)) ).total_seconds() / 3600 <= window_h] keep.sort(key=os.path.getmtime) # D-1: mtime 序 → 后到运行整行胜出 stats, sigs_by_type = {}, {} for p in keep: try: with open(p, encoding="utf-8", errors="replace") as f: text = f.read() except OSError: continue for t, s in chk.parse_stat_lines(text).items(): # D-1(10-05): 旧 max(rows) 归并吞掉后到运行的 failed + 假 # resolve(早班 rows=100/failed=0 压住后到 rows=50/failed=3 → # 漏报+既有告警被假自愈双失效)——改后到运行(mtime 新)整行胜出 stats[t] = s # 10-05: 签名按 [type] 归因(跨类型连坐修), 跨文件按 sig 去重合并 for t, ss in chk.signatures_by_type(text).items(): bucket = sigs_by_type.setdefault(t, []) have = {x["sig"] for x in bucket} for s in ss: if s["sig"] not in have: bucket.append(s) have.add(s["sig"]) return stats, sigs_by_type def log_layer_results(stats, sigs_by_type, today): """统计行 failed>0 → 按「该类型自己的」ERROR 行签名分级(10-05 归因修: 他类型反爬签名不再抬本类型 failed); 断路器独立键(全局); 其余 ok → passed.""" results, passed = [], [] breaker = any(s["sig"] == "breaker" for ss in sigs_by_type.values() for s in ss) for t, s in sorted(stats.items()): if s.get("failed", 0) > 0: tsigs = sigs_by_type.get(t, []) worst = max((x["severity"] for x in tsigs), default=None) ev = {"stats": s, "signatures": tsigs} sev = "red" if (worst == "red" or breaker) else "yellow" results.append(_event(f"data-{t}-failed", sev, f"{t} 类型级 failed={s['failed']}", f"rows={s.get('rows')}", ev, today=today)) else: passed.append(f"data-{t}-failed") if breaker: results.append(_event("data-akshare-breaker", "red", "断路器触发(熔断连累: 后续类型整批跳过)", "exit code 2 会话", {"signatures": sigs_by_type}, today=today)) else: passed.append("data-akshare-breaker") return results, passed def log_only_results(entries, stats, sigs_by_type, today): """valuation 等 log_only 项: 统计行存在且 failed=0(周班项平日无行属正常). 10-05: red 判定用该类型自己的签名(归因修).""" results, passed = [], [] for e in entries: key = f"data-{e['name']}-stat" s = stats.get(e["name"]) red_sig = any(x["severity"] == "red" for x in sigs_by_type.get(e["name"], [])) if s is None: if e["log_cadence"] == "weekly": continue results.append(_event(key, "yellow", f"{e['name']} 无当日统计行", "任务未跑或日志缺失", {}, today=today)) elif s.get("failed", 0) > 0: results.append(_event(key, "red" if red_sig else "yellow", f"{e['name']} 统计 failed={s['failed']}", f"rows={s.get('rows')}", {"stats": s}, today=today)) else: passed.append(key) return results, passed # ---------- 调度层: schtask 在岗(spec §20.15.1 表后注; retry 一次性死案即此类) ---------- def schtask_results(snapshot, now): results = [] today = now.date().isoformat().replace("-", "") # weekday=72h: 周五晚班→周一晚班次跨度 69h+容差(0928 早班 ak-events-retry # 57h/qmt-probe0915 71h 假黄实案); 真死 >72h 仍黄, next_run=N/A→red 兜底 limits = {"daily": 26 if now.weekday() < 5 else 80, "weekday": 72, "weekly": 8 * 24, "monthly": 40 * 24, "always": None} for name, cadence in reg.SCHTASKS: key = f"infra-task-{name}" info = snapshot.get(name) if info is None: results.append(_event(key, "red", f"schtask {name} 缺失", "schtasks /query 无此任务", {}, source="infra", today=today)) continue next_run = info.get("next_run", "") last_run = info.get("last_run", "") last_result = info.get("last_result", "") # 1010 二审 P2-6: 主动禁用(Disabled)不告——qmt-relogin 0929/30 连两日 # 红实案=运维检修暂停(防与手动登录打架), N/A→red 兜底只应抓真死 st = info.get("status", "") if "禁用" in st or "disabled" in st.lower(): continue bad = [] if "N/A" in next_run and cadence != "always": bad.append("next_run=N/A(一次性死)") if last_result and last_result not in _OK_LAST_RESULTS: bad.append(f"last_result={last_result}") limit_h = limits[cadence] if cadence == "weekday" and now.weekday() >= 5: limit_h = None age_h = _parse_schtask_age(last_run, now) if limit_h is not None and age_h is not None and age_h > limit_h: bad.append(f"last_run {age_h:.0f}h 前(限 {limit_h}h)") if bad: results.append(_event(key, "red" if "N/A" in next_run else "yellow", f"schtask {name} 异常", "; ".join(bad), {"next_run": next_run, "last_run": last_run, "last_result": last_result}, source="infra", today=today)) return results def _parse_schtask_age(last_run, now): if not last_run or "N/A" in last_run: return None for fmt in ("%Y/%m/%d %H:%M:%S", "%Y-%m-%d %H:%M"): try: t = dt.datetime.strptime(last_run.strip(), fmt) return (now - t).total_seconds() / 3600 except ValueError: continue return None # ---------- 衍生层: NAS inbox 摄取(单写者纪律: NAS 不写远端 sqlite) ---------- def ingest_inbox(conn, data_root, now): """摄取 NAS 事件+绿键回传(09-26 补: NAS 侧绿检查须经 green_keys 回传, 否则 VPS 侧 NAS 告警永不转绿). 返回最新 generated_at(静默检测用).""" inbox = os.path.join(data_root, _HEALTH_DIR, "inbox") latest_ts = None if not os.path.isdir(inbox): return latest_ts processed = os.path.join(inbox, "processed") os.makedirs(processed, exist_ok=True) green_keys = [] for p in sorted(glob.glob(os.path.join(inbox, "*.json"))): try: with open(p, encoding="utf-8") as f: payload = json.load(f) for ev in payload.get("events", []): store.upsert_alert(conn, ev) green_keys.extend(payload.get("green_keys", [])) ts = payload.get("generated_at") if ts: latest_ts = max(latest_ts or "", ts) os.replace(p, os.path.join(processed, os.path.basename(p))) except (OSError, ValueError) as exc: print(f"[inbox] skip bad file {p}: {exc}") if green_keys: store.resolve_passed(conn, green_keys) return latest_ts # ---------- 主流程 ---------- def write_health(data_root, health): d = os.path.join(data_root, _HEALTH_DIR) os.makedirs(d, exist_ok=True) for fn in (f"health_{health['date']}_{health['shift']}.json", "health_latest.json"): tmp = os.path.join(d, fn + ".tmp") with open(tmp, "w", encoding="utf-8") as f: json.dump(health, f, ensure_ascii=False, indent=1) os.replace(tmp, os.path.join(d, fn)) def _calendar_days(confirmed_days, now_date, lookback=15): """结果层日历 = 确认集(dbbardata 000300 实证) ∪ 近 lookback 日工作日假定集 (09-25 首跑实锤的循环依赖修: dbbardata 晚间源滞后/漏日会把全体期望塌一天, 误红鲜活域)。10-03 假黄根治: 假定集剔 EXCHANGE_HOLIDAYS 休市日(已知假日= 零期望); 表缺年=不扣(黄照旧绝不静默转绿)。确认集恒保留——数据本体是真相, 表错不吞真相。语义: 缺失遇「未确认日」一律降黄(未登记休市/合并窗/源滞后 三态容忍); 确认日缺失才红(跨域互证: 某域有该日数据=真是交易日)。""" union = set(confirmed_days) d0 = now_date - dt.timedelta(days=lookback) while d0 <= now_date: hol = reg.EXCHANGE_HOLIDAYS.get(d0.year) if d0.weekday() < 5 and (hol is None or d0.strftime("%Y%m%d") not in hol): union.add(d0.isoformat()) d0 += dt.timedelta(days=1) return sorted(union) def main(argv=None): ap = argparse.ArgumentParser(description=__doc__) ap.add_argument("--shift", choices=["evening", "morning"], default="evening") ap.add_argument("--data-root", default="data") ap.add_argument("--db", default=store.default_alerts_db_path()) ap.add_argument("--skip-tasks", action="store_true") ap.add_argument("--skip-inbox", action="store_true") args = ap.parse_args(argv) db_path = args.db if args.db == "data/alerts.db" and args.data_root != "data": db_path = os.path.join(args.data_root, "alerts.db") # D-2(10-05): 整班顶层兜底 —— 裸崩溃时 wrapper 恒 0 + last_run 新鲜 = # SCHTASKS 年龄判定失明(审计推演最坏 ~22h 仅 peer 反盯); 兜底写 red # health + infra-monitor-crash 告警后 exit 3, 静默窗收敛到本班内 try: return _run_shift(args, db_path) except Exception: # noqa: BLE001 return _crash_exit(args, db_path) def _run_shift(args, db_path): now = _now() today = now.date().isoformat().replace("-", "") days = chk.trading_days( os.path.join(args.data_root, "quant_trading.db"), now.date()) confirmed = {d.replace("-", "") for d in days} days = _calendar_days(days, now.date()) # 语义见函数 docstring(10-03 加休市表剔除) ctx = {"data_root": args.data_root, "days": days, "now": now, "confirmed_days": confirmed} conn = store.connect(db_path) checks, events, passed = [], [], [] for e in reg.REGISTRY: if not e.get("enabled", True) or e["kind"] == "log_only": continue # log_only 归过程层 log_only_results r = check_entry(e, ctx) checks.append(r) if r["status"] == "green": passed.append(r["key"]) # kind 归并修复: 绿路径统一返 stale key, rows key 只在黄时存在 # → rows 类黄一旦发生绿班永不 resolve (sw_industry/block_trade # 悬死实案); 绿班恒收 rows key(幂等, 无告警时 resolve 零命中) if e["kind"] in ("per_date", "weekly_friday", "month_bucket"): passed.append(f"data-{e['name']}-rows") else: kindword = r["key"].rsplit("-", 1)[-1] # stale/rows/verdict/exists events.append(_event( r["key"], r["status"], f"{e['name']} {kindword} 异常", f"actual={r['actual']} expected={r['expected']}", {"evidence": r.get("evidence", {}), "kind": e["kind"], "producer": e.get("producer", "vps")}, source="infra" if r["key"].startswith("infra-") else "data", today=today)) stats, sigs_by_type = collect_log_layer(args.data_root, now) log_events, log_passed = log_layer_results(stats, sigs_by_type, today) lo_events, lo_passed = log_only_results( [e for e in reg.REGISTRY if e["kind"] == "log_only" and e.get("enabled", True)], stats, sigs_by_type, today) events += log_events + lo_events passed += log_passed + lo_passed if not args.skip_tasks and os.name == "nt": # schtasks 控制台输出=OEM 代码页(zh-CN=GBK); wrapper 以 -X utf8 启动时 # text=True 会按 UTF-8 误解码致中文字段名全乱(09-25 首跑 14 任务全假红 # 实锤) → 字节捕获+手动 GBK 解码(GBK 兼容 ASCII, 英文系统不受影响)。 raw = subprocess.run(["schtasks", "/query", "/fo", "LIST", "/v"], capture_output=True, timeout=180).stdout task_events = schtask_results( chk.parse_schtasks_list(raw.decode("gbk", errors="replace")), now) events += task_events flagged = {ev["check_key"] for ev in task_events} passed += [f"infra-task-{n}" for n, _ in reg.SCHTASKS if f"infra-task-{n}" not in flagged] nas_silent_key = "infra-nas-monitor-silent" if not args.skip_inbox: nas_ts = ingest_inbox(conn, args.data_root, now) if nas_ts: try: age_h = (now - dt.datetime.fromisoformat(nas_ts) ).total_seconds() / 3600 except (ValueError, TypeError): # D-9(10-05): TypeError 同防 age_h = None if age_h is not None and age_h > 30: events.append(_event(nas_silent_key, "yellow", "NAS 监控静默>30h", f"last={nas_ts}", {"last_inbox_ts": nas_ts}, source="infra", today=today)) else: passed.append(nas_silent_key) for ev in events: store.upsert_alert(conn, ev) store.resolve_passed(conn, passed) store.purge_expired(conn) conn.close() counts = {s: sum(1 for c in checks if c["status"] == s) for s in ("green", "yellow", "red")} red_events = sum(1 for ev in events if ev["severity"] == "red") yellow_events = sum(1 for ev in events if ev["severity"] == "yellow") # 拆分防晨检误导航(🟢-7): checks 绿而 alerts 层 26h 窗旧红回声时, # overall=red 会引导去查已自愈项; 两字段各自可读, overall 保兼容取最坏 checks_overall = "red" if counts["red"] else ( "yellow" if counts["yellow"] else "green") alerts_overall = "red" if red_events else ( "yellow" if yellow_events else "green") overall = "red" if (counts["red"] or red_events) else ( "yellow" if (counts["yellow"] or yellow_events) else "green") health = {"date": now.date().isoformat(), "shift": args.shift, "host": "vps", "generated_at": now.isoformat(timespec="seconds"), "overall": overall, "checks_overall": checks_overall, "alerts_overall": alerts_overall, "counts": counts, "alerts_emitted": len(events), "checks": checks} write_health(args.data_root, health) print(f"[data-monitor] {args.shift} overall={overall} checks={len(checks)}" f" alerts={len(events)}") return 3 if overall == "red" else 0 def _crash_exit(args, db_path): """D-2(10-05): 整班崩溃兜底出口 —— red health + infra-monitor-crash 告警 + exit 3(告警/health 任一写失败也不掩原始崩溃事实)。""" now = _now() today = now.date().isoformat().replace("-", "") detail = traceback.format_exc() sys.stderr.write(f"[data-monitor] CRASH 整班崩溃(D-2 兜底):\n{detail}\n") try: conn = store.connect(db_path) store.upsert_alert(conn, _event( "infra-monitor-crash", "red", "data-monitor 整班崩溃(顶层兜底)", detail[-800:], {"data_root": args.data_root, "shift": args.shift}, source="infra", today=today)) conn.close() except Exception: # noqa: BLE001 崩溃中的崩溃 pass try: write_health(args.data_root, { "date": now.date().isoformat(), "shift": args.shift, "host": "vps", "generated_at": now.isoformat(timespec="seconds"), "overall": "red", "checks_overall": "red", "alerts_overall": "red", "counts": {"green": 0, "yellow": 0, "red": 0}, "alerts_emitted": 1, "checks": [], "crash": detail[-400:]}) except Exception: # noqa: BLE001 pass return 3 if __name__ == "__main__": sys.exit(main())