#!/usr/bin/env python3 """ Palladium Z1 Monitor — LD逻辑板使用监控系统 解析 `test_server -short` 输出,展示 Palladium Z1 仿真器的 逻辑板 (LD) 状态、逻辑板利用率、作业信息和 T-Pod 可用性。 支持两种数据来源: 1. 手动粘贴 — 在 Web 界面粘贴 test_server -short 输出 2. API 推送 — 在 Palladium 主机运行 collect.py 远程推送 """ from __future__ import annotations # PEP 563: 兼容 Python 3.7+ (避免 PEP 585 泛型语法) import os import re import io import csv import json import subprocess import threading import time from datetime import datetime, timezone, timedelta from flask import Flask, render_template, request, jsonify, redirect, url_for, Response from flask_sqlalchemy import SQLAlchemy # ═══════════════════════════════════════════════════════════════════════════ # Config # ═══════════════════════════════════════════════════════════════════════════ BASE_DIR = os.path.dirname(os.path.abspath(__file__)) INSTANCE_DIR = os.path.join(BASE_DIR, "instance") os.makedirs(INSTANCE_DIR, exist_ok=True) app = Flask(__name__) app.config["SECRET_KEY"] = os.environ.get("PALLADIUM_SECRET", "palladium-z1-monitor-2026") app.config["SQLALCHEMY_DATABASE_URI"] = f"sqlite:///{os.path.join(INSTANCE_DIR, 'palladium.db')}" app.config["SQLALCHEMY_TRACK_MODIFICATIONS"] = False db = SQLAlchemy(app) # ═══════════════════════════════════════════════════════════════════════════ # Domain State Constants # ═══════════════════════════════════════════════════════════════════════════ DOMAIN_AVAILABLE = "available" DOMAIN_DOWNLOADED = "downloaded" DOMAIN_DISABLED = "disabled" DOMAIN_NOT_EXIST = "not_exist" DOMAIN_COLORS = { DOMAIN_AVAILABLE: "#28a745", DOMAIN_DOWNLOADED: "#007bff", DOMAIN_DISABLED: "#dc3545", DOMAIN_NOT_EXIST: "#4a4a5a", } DOMAIN_LABELS = { DOMAIN_AVAILABLE: "可用", DOMAIN_DOWNLOADED: "已加载", DOMAIN_DISABLED: "已禁用", DOMAIN_NOT_EXIST: "不存在", } # ═══════════════════════════════════════════════════════════════════════════ # Models # ═══════════════════════════════════════════════════════════════════════════ class ScanSnapshot(db.Model): __tablename__ = "snapshots" id = db.Column(db.Integer, primary_key=True) timestamp = db.Column(db.DateTime, default=lambda: datetime.now(timezone.utc), index=True) emulator = db.Column(db.String(100), default="") hardware = db.Column(db.String(100), default="") configmgr = db.Column(db.String(200), default="") system_status = db.Column(db.String(20), default="") source = db.Column(db.String(50), default="manual") raw_output = db.Column(db.Text) # Summary stats (denormalised for quick queries) total_boards = db.Column(db.Integer, default=0) online_boards = db.Column(db.Integer, default=0) offline_boards = db.Column(db.Integer, default=0) total_domains = db.Column(db.Integer, default=0) available_domains = db.Column(db.Integer, default=0) downloaded_domains = db.Column(db.Integer, default=0) disabled_domains = db.Column(db.Integer, default=0) not_exist_domains = db.Column(db.Integer, default=0) active_jobs = db.Column(db.Integer, default=0) used_boards = db.Column(db.Integer, default=0) utilization = db.Column(db.Float, default=0.0) boards = db.relationship("BoardStatus", backref="snapshot", cascade="all, delete-orphan", lazy=True) jobs = db.relationship("JobInfo", backref="snapshot", cascade="all, delete-orphan", lazy=True) pods = db.relationship("PodInfo", backref="snapshot", cascade="all, delete-orphan", lazy=True) class BoardStatus(db.Model): __tablename__ = "board_statuses" id = db.Column(db.Integer, primary_key=True) snapshot_id = db.Column(db.Integer, db.ForeignKey("snapshots.id"), nullable=False, index=True) rack = db.Column(db.Integer) cluster = db.Column(db.Integer) ld_index = db.Column(db.Integer) status = db.Column(db.String(20)) ccd = db.Column(db.String(20)) d0 = db.Column(db.String(20), default="-") d1 = db.Column(db.String(20), default="-") d2 = db.Column(db.String(20), default="-") d3 = db.Column(db.String(20), default="-") d4 = db.Column(db.String(20), default="-") d5 = db.Column(db.String(20), default="-") d6 = db.Column(db.String(20), default="-") d7 = db.Column(db.String(20), default="-") class JobInfo(db.Model): __tablename__ = "job_infos" id = db.Column(db.Integer, primary_key=True) snapshot_id = db.Column(db.Integer, db.ForeignKey("snapshots.id"), nullable=False, index=True) job_index = db.Column(db.Integer) owner = db.Column(db.String(100), default="") pid = db.Column(db.String(200), default="") t_pod = db.Column(db.String(200), default="") design = db.Column(db.String(200), default="") elap_time = db.Column(db.String(20), default="") reserved_key = db.Column(db.String(100), default="") class PodInfo(db.Model): __tablename__ = "pod_infos" id = db.Column(db.Integer, primary_key=True) snapshot_id = db.Column(db.Integer, db.ForeignKey("snapshots.id"), nullable=False, index=True) rack = db.Column(db.Integer) all_hdsb = db.Column(db.Text, default="") # JSON: {"5": "T4", "30": "T5"} all_t_pods = db.Column(db.Text, default="") # JSON: {"5": "T0, T1, ..."} available_hdsb = db.Column(db.Text, default="") available_t_pods = db.Column(db.Text, default="") locked_hdsb = db.Column(db.Text, default="") locked_t_pods = db.Column(db.Text, default="") reserved_hdsb = db.Column(db.Text, default="") reserved_t_pods = db.Column(db.Text, default="") unavailable_hdsb = db.Column(db.Text, default="") unavailable_t_pods = db.Column(db.Text, default="") # ═══════════════════════════════════════════════════════════════════════════ # Parser — test_server -short 输出解析 # ═══════════════════════════════════════════════════════════════════════════ def parse_test_server_output(raw_text: str) -> dict: """ 解析 `test_server -short` 命令输出,返回结构化数据。 返回结构: { "system": {emulator, hardware, configmgr, system_status}, "racks": [ {"rack": 0, "clusters": [ {"cluster": 0, "ccd": "ONLINE", "boards": [ {"ld": 0, "status": "ONLINE", "domains": ["-","-","-","-","-","-","-","-"]} ]} ]} ], "jobs": [{"index", "owner", "pid", "t_pod", "design", "elap_time", "reserved_key"}], "pods": [{"rack", "all_hdsb", "all_t_pods", "available_hdsb", ...}], } """ result = { "system": {}, "racks": [], "jobs": [], "pods": [], } # ─── 系统信息 ────────────────────────────────────────────── # 注意: Configmgr 值可能与 "System Status" 拼接在一起没有空格 emu_match = re.search(r"Emulator:\s*(\S+)", raw_text) # Hardware 名称含空格 (如 "Palladium Z1"), 匹配到 Configmgr 为止 hw_match = re.search(r"Hardware:\s*(.+?)\s+Configmgr:", raw_text) # 非贪婪匹配直到 "System Status" 或行尾 cfg_match = re.search(r"Configmgr:\s*(.+?)(?:System\s*Status|$)", raw_text) status_match = re.search(r"System\s*Status:\s*(\w+)", raw_text) result["system"] = { "emulator": emu_match.group(1).strip() if emu_match else "", "hardware": hw_match.group(1).strip() if hw_match else "", "configmgr": cfg_match.group(1).strip() if cfg_match else "", "system_status": status_match.group(1).strip() if status_match else "", } # ─── Rack / Cluster / Board 结构 ─────────────────────────── lines = raw_text.split("\n") current_rack = None # LD 行: 序号 + 状态(ONLINE/OFFLINE) + 8个逻辑板值 board_re = re.compile( r"^\s*(\d+)\s+(ONLINE|OFFLINE)\s+" r"(\S+)\s+(\S+)\s+(\S+)\s+(\S+)\s+" r"(\S+)\s+(\S+)\s+(\S+)\s+(\S+)\s*$" ) for line in lines: # Rack 头: "Rack 0 has 2 clusters" rack_match = re.match(r"Rack\s+(\d+)\s+has\s+(\d+)\s+clusters", line) if rack_match: current_rack = int(rack_match.group(1)) result["racks"].append({ "rack": current_rack, "clusters": [], }) continue # Cluster 头: "Cluster 0 has 6 boards CCD: ONLINE" cluster_match = re.match( r"Cluster\s+(\d+)\s+has\s+(\d+)\s+boards\s+CCD:\s*(\w+)", line ) if cluster_match and result["racks"]: current_cluster = int(cluster_match.group(1)) ccd = cluster_match.group(3) result["racks"][-1]["clusters"].append({ "cluster": current_cluster, "ccd": ccd, "boards": [], }) continue # LD 板行 board_match = board_re.match(line) if board_match and result["racks"] and result["racks"][-1]["clusters"]: ld_idx = int(board_match.group(1)) status = board_match.group(2) domains = [board_match.group(i) for i in range(3, 11)] result["racks"][-1]["clusters"][-1]["boards"].append({ "ld": ld_idx, "status": status, "domains": domains, }) continue # ─── 作业信息 ────────────────────────────────────────────── job_section = re.search( r"Job Information:.*?(?=Target Pod Information:|Exit code:|$)", raw_text, re.DOTALL, ) if job_section: job_lines = job_section.group(0).split("\n") for jline in job_lines: jline = jline.strip() if not jline or jline.startswith("Job Information") or jline.startswith("Index"): continue # 分割后从两端取值: 前3个是 index/owner/pid, 后3个是 design/elaptime/reservedkey # 中间的是 T-Pod (可能含空格, 如 "-- --") parts = jline.split() if len(parts) >= 7: idx = int(parts[0]) owner = parts[1] pid = parts[2] reserved_key = parts[-1] elap_time = parts[-2] design = parts[-3] t_pod = " ".join(parts[3:-3]) result["jobs"].append({ "index": idx, "owner": owner, "pid": pid, "t_pod": t_pod, "design": design, "elap_time": elap_time, "reserved_key": reserved_key, }) # ─── T-Pod 信息 ──────────────────────────────────────────── pod_section = re.search( r"Target Pod Information:.*?(?=Exit code:|$)", raw_text, re.DOTALL, ) if pod_section: pod_text = pod_section.group(0) # 按 "Rack N:" 分割 rack_pod_re = re.compile(r"Rack\s+(\d+):(.*?)(?=Rack\s+\d+:|Exit code:|$)", re.DOTALL) for m in rack_pod_re.finditer(pod_text): rack_num = int(m.group(1)) rack_text = m.group(2) all_hdsb = {} all_t_pods = {} fields = { "available_hdsb": "", "available_t_pods": "", "locked_hdsb": "", "locked_t_pods": "", "reserved_hdsb": "", "reserved_t_pods": "", "unavailable_hdsb": "", "unavailable_t_pods": "", } for rline in rack_text.split("\n"): rline = rline.strip() if not rline: continue # "All length 5 HDSB: T4" / "All length 30 HDSB: T5" all_h_match = re.match(r"All length (\d+) HDSB:\s*(.+)", rline) if all_h_match: all_hdsb[all_h_match.group(1)] = all_h_match.group(2).strip() continue # "All length 5 T-Pods: T0, T1, ..." all_t_match = re.match(r"All length (\d+) T-Pods:\s*(.+)", rline) if all_t_match: all_t_pods[all_t_match.group(1)] = all_t_match.group(2).strip() continue # 标准字段 — 用 startswith 精确匹配避免 "Available" vs "Unavailable" 混淆 for field_key in fields: label = field_key.replace("_", " ").title() # 修正大小写和连字符: "Hdsb"→"HDSB", "T Pods"→"T-Pods" label = label.replace("Hdsb", "HDSB").replace("T Pods", "T-Pods") if rline.startswith(label + ":"): val = rline.split(":", 1)[1].strip() fields[field_key] = val break result["pods"].append({ "rack": rack_num, "all_hdsb": json.dumps(all_hdsb), "all_t_pods": json.dumps(all_t_pods), **fields, }) return result def classify_domain(raw_value: str) -> str: """将域原始值分类为状态常量。""" v = raw_value.strip() if v == "-": return DOMAIN_AVAILABLE if v == "X": return DOMAIN_NOT_EXIST if v == "D": return DOMAIN_DISABLED if v.isdigit(): return DOMAIN_DOWNLOADED return DOMAIN_AVAILABLE def calculate_stats(parsed: dict) -> dict: """从解析结果计算汇总统计。""" total_boards = online = offline = used = 0 total_dom = avail = dl = dis = ne = 0 for rack in parsed["racks"]: for cluster in rack["clusters"]: for board in cluster["boards"]: total_boards += 1 if board["status"] == "ONLINE": online += 1 else: offline += 1 board_used = False for d in board["domains"]: state = classify_domain(d) total_dom += 1 if state == DOMAIN_AVAILABLE: avail += 1 elif state == DOMAIN_DOWNLOADED: dl += 1 board_used = True elif state == DOMAIN_DISABLED: dis += 1 elif state == DOMAIN_NOT_EXIST: ne += 1 if board_used: used += 1 util = round(used / total_boards * 100, 1) if total_boards > 0 else 0.0 return { "total_boards": total_boards, "online_boards": online, "offline_boards": offline, "used_boards": used, "total_domains": total_dom, "available_domains": avail, "downloaded_domains": dl, "disabled_domains": dis, "not_exist_domains": ne, "active_jobs": len(parsed["jobs"]), "utilization": util, } # ═══════════════════════════════════════════════════════════════════════════ # Template Filters # ═══════════════════════════════════════════════════════════════════════════ @app.template_filter("domain_state") def domain_state_filter(raw_value): return classify_domain(raw_value) @app.template_filter("domain_color") def domain_color_filter(raw_value): return DOMAIN_COLORS.get(classify_domain(raw_value), "#888") @app.template_filter("domain_label") def domain_label_filter(raw_value): return DOMAIN_LABELS.get(classify_domain(raw_value), "未知") @app.template_filter("fmt_time") def fmt_time_filter(dt): if not dt: return "—" if dt.tzinfo is None: dt = dt.replace(tzinfo=timezone.utc) local = dt.astimezone() return local.strftime("%Y-%m-%d %H:%M:%S") @app.template_filter("fmt_time_short") def fmt_time_short_filter(dt): if not dt: return "—" if dt.tzinfo is None: dt = dt.replace(tzinfo=timezone.utc) local = dt.astimezone() return local.strftime("%H:%M:%S") @app.template_filter("from_json") def from_json_filter(s): if not s: return {} try: return json.loads(s) except (json.JSONDecodeError, TypeError): return {} # ═══════════════════════════════════════════════════════════════════════════ # Routes # ═══════════════════════════════════════════════════════════════════════════ @app.route("/") def dashboard(): """主仪表板 — 显示最新快照数据。""" snapshot = ScanSnapshot.query.order_by(ScanSnapshot.timestamp.desc()).first() stats = None racks_data = [] jobs_data = [] pods_data = [] history = [] if snapshot: stats = { "total_boards": snapshot.total_boards, "online_boards": snapshot.online_boards, "offline_boards": snapshot.offline_boards, "total_domains": snapshot.total_domains, "available_domains": snapshot.available_domains, "downloaded_domains": snapshot.downloaded_domains, "disabled_domains": snapshot.disabled_domains, "not_exist_domains": snapshot.not_exist_domains, "active_jobs": snapshot.active_jobs, "used_boards": snapshot.used_boards, "utilization": snapshot.utilization, } # 重建 racks 结构 boards = BoardStatus.query.filter_by(snapshot_id=snapshot.id).order_by( BoardStatus.rack, BoardStatus.cluster, BoardStatus.ld_index ).all() rack_map = {} for b in boards: if b.rack not in rack_map: rack_map[b.rack] = {} if b.cluster not in rack_map[b.rack]: rack_map[b.rack][b.cluster] = { "cluster": b.cluster, "ccd": b.ccd or "UNKNOWN", "boards": [], } rack_map[b.rack][b.cluster]["boards"].append({ "ld": b.ld_index, "status": b.status, "domains": [b.d0, b.d1, b.d2, b.d3, b.d4, b.d5, b.d6, b.d7], }) for rack_num in sorted(rack_map): racks_data.append({ "rack": rack_num, "clusters": [rack_map[rack_num][c] for c in sorted(rack_map[rack_num])], }) jobs_data = JobInfo.query.filter_by(snapshot_id=snapshot.id).all() pods_data = PodInfo.query.filter_by(snapshot_id=snapshot.id).order_by(PodInfo.rack).all() # 历史趋势 (最近50条) history = ScanSnapshot.query.order_by(ScanSnapshot.timestamp.desc()).limit(50).all() history = list(reversed(history)) # 时间正序 return render_template( "dashboard.html", snapshot=snapshot, stats=stats, racks=racks_data, jobs=jobs_data, pods=pods_data, history=history, ) @app.route("/history") def history(): """历史快照列表。""" page = request.args.get("page", 1, type=int) per_page = 30 pagination = ScanSnapshot.query.order_by(ScanSnapshot.timestamp.desc()).paginate( page=page, per_page=per_page, error_out=False ) return render_template("history.html", pagination=pagination) @app.route("/snapshot/") def snapshot_detail(sid): """单条快照详情。""" snapshot = db.get_or_404(ScanSnapshot, sid) stats = { "total_boards": snapshot.total_boards, "online_boards": snapshot.online_boards, "offline_boards": snapshot.offline_boards, "total_domains": snapshot.total_domains, "available_domains": snapshot.available_domains, "downloaded_domains": snapshot.downloaded_domains, "disabled_domains": snapshot.disabled_domains, "not_exist_domains": snapshot.not_exist_domains, "active_jobs": snapshot.active_jobs, "used_boards": snapshot.used_boards, "utilization": snapshot.utilization, } boards = BoardStatus.query.filter_by(snapshot_id=sid).order_by( BoardStatus.rack, BoardStatus.cluster, BoardStatus.ld_index ).all() rack_map = {} for b in boards: if b.rack not in rack_map: rack_map[b.rack] = {} if b.cluster not in rack_map[b.rack]: rack_map[b.rack][b.cluster] = { "cluster": b.cluster, "ccd": b.ccd or "UNKNOWN", "boards": [], } rack_map[b.rack][b.cluster]["boards"].append({ "ld": b.ld_index, "status": b.status, "domains": [b.d0, b.d1, b.d2, b.d3, b.d4, b.d5, b.d6, b.d7], }) racks_data = [] for rack_num in sorted(rack_map): racks_data.append({ "rack": rack_num, "clusters": [rack_map[rack_num][c] for c in sorted(rack_map[rack_num])], }) jobs_data = JobInfo.query.filter_by(snapshot_id=sid).all() pods_data = PodInfo.query.filter_by(snapshot_id=sid).order_by(PodInfo.rack).all() return render_template( "snapshot.html", snapshot=snapshot, stats=stats, racks=racks_data, jobs=jobs_data, pods=pods_data, ) @app.route("/paste", methods=["POST"]) def paste_data(): """手动粘贴 test_server -short 输出。""" raw_output = request.form.get("raw_output", "").strip() if not raw_output: return jsonify({"error": "无输出内容"}), 400 sid = store_snapshot(raw_output, source="paste") return redirect(url_for("dashboard")) @app.route("/api/collect", methods=["POST"]) def api_collect(): """API 端点 — 接收 test_server -short 原始输出。""" data = request.get_json(silent=True) if not data or "raw_output" not in data: return jsonify({"error": "缺少 raw_output 字段"}), 400 raw_output = data["raw_output"] source = data.get("source", "api") sid = store_snapshot(raw_output, source=source) snapshot = db.session.get(ScanSnapshot, sid) return jsonify({ "ok": True, "snapshot_id": sid, "stats": { "total_boards": snapshot.total_boards, "online_boards": snapshot.online_boards, "offline_boards": snapshot.offline_boards, "total_domains": snapshot.total_domains, "available_domains": snapshot.available_domains, "downloaded_domains": snapshot.downloaded_domains, "disabled_domains": snapshot.disabled_domains, "not_exist_domains": snapshot.not_exist_domains, "active_jobs": snapshot.active_jobs, "used_boards": snapshot.used_boards, "utilization": snapshot.utilization, }, }) @app.route("/api/latest") def api_latest(): """返回最新快照 JSON。""" snapshot = ScanSnapshot.query.order_by(ScanSnapshot.timestamp.desc()).first() if not snapshot: return jsonify({"error": "无数据"}), 404 boards = BoardStatus.query.filter_by(snapshot_id=snapshot.id).all() jobs = JobInfo.query.filter_by(snapshot_id=snapshot.id).all() pods = PodInfo.query.filter_by(snapshot_id=snapshot.id).all() return jsonify({ "snapshot": { "id": snapshot.id, "timestamp": snapshot.timestamp.isoformat() if snapshot.timestamp else None, "emulator": snapshot.emulator, "hardware": snapshot.hardware, "configmgr": snapshot.configmgr, "system_status": snapshot.system_status, "source": snapshot.source, }, "stats": { "total_boards": snapshot.total_boards, "online_boards": snapshot.online_boards, "offline_boards": snapshot.offline_boards, "total_domains": snapshot.total_domains, "available_domains": snapshot.available_domains, "downloaded_domains": snapshot.downloaded_domains, "disabled_domains": snapshot.disabled_domains, "not_exist_domains": snapshot.not_exist_domains, "active_jobs": snapshot.active_jobs, "used_boards": snapshot.used_boards, "utilization": snapshot.utilization, }, "boards": [{"rack": b.rack, "cluster": b.cluster, "ld": b.ld_index, "status": b.status, "ccd": b.ccd, "domains": [b.d0, b.d1, b.d2, b.d3, b.d4, b.d5, b.d6, b.d7]} for b in boards], "jobs": [{"index": j.job_index, "owner": j.owner, "pid": j.pid, "t_pod": j.t_pod, "design": j.design, "elap_time": j.elap_time, "reserved_key": j.reserved_key} for j in jobs], "pods": [{"rack": p.rack, "all_hdsb": json.loads(p.all_hdsb) if p.all_hdsb else {}, "all_t_pods": json.loads(p.all_t_pods) if p.all_t_pods else {}, "available_hdsb": p.available_hdsb, "available_t_pods": p.available_t_pods, "locked_hdsb": p.locked_hdsb, "locked_t_pods": p.locked_t_pods, "reserved_hdsb": p.reserved_hdsb, "reserved_t_pods": p.reserved_t_pods, "unavailable_hdsb": p.unavailable_hdsb, "unavailable_t_pods": p.unavailable_t_pods} for p in pods], }) @app.route("/api/history") def api_history(): """历史趋势 JSON (用于图表)。""" snapshots = ScanSnapshot.query.order_by(ScanSnapshot.timestamp.desc()).limit(100).all() snapshots = list(reversed(snapshots)) return jsonify({ "data": [{ "timestamp": s.timestamp.isoformat() if s.timestamp else None, "utilization": s.utilization, "online_boards": s.online_boards, "offline_boards": s.offline_boards, "downloaded_domains": s.downloaded_domains, "available_domains": s.available_domains, "active_jobs": s.active_jobs, } for s in snapshots] }) @app.route("/api/snapshot//delete", methods=["POST"]) def delete_snapshot(sid): """删除指定快照。""" snapshot = db.get_or_404(ScanSnapshot, sid) db.session.delete(snapshot) db.session.commit() return redirect(url_for("history")) @app.route("/api/cleanup", methods=["POST"]) def cleanup_old(): """清理 N 天前的快照。""" days = request.args.get("days", 30, type=int) cutoff = datetime.now(timezone.utc) - timedelta(days=days) old = ScanSnapshot.query.filter(ScanSnapshot.timestamp < cutoff).all() count = len(old) for s in old: db.session.delete(s) db.session.commit() return jsonify({"ok": True, "deleted": count}) # ═══════════════════════════════════════════════════════════════════════════ # Email Routes (延迟导入 mailer, 避免在 mailer 缺失时整个 app 起不来) # ═══════════════════════════════════════════════════════════════════════════ def _get_mailer(): try: import mailer return mailer except ImportError as e: return None @app.route("/api/email/status") def email_status(): """返回当前邮件配置状态 (不泄露密码)。""" ml = _get_mailer() if not ml: return jsonify({"ok": False, "error": "mailer.py 不存在或导入失败"}), 500 cfg = ml.get_config() errs = ml.config_errors(cfg) return jsonify({ "ok": True, "enabled": cfg["enabled"], "configured": len(errs) == 0, "config_errors": errs, "smtp_host": cfg["smtp_host"], "smtp_port": cfg["smtp_port"], "smtp_use_tls": cfg["smtp_use_tls"], "smtp_use_ssl": cfg["smtp_use_ssl"], "from_addr": cfg["from_addr"], "from_name": cfg["from_name"], "recipients": cfg["to_addrs"], "send_at": f"{cfg['send_hour']:02d}:{cfg['send_minute']:02d}", "monitor_base_url": cfg["monitor_base_url"], }) @app.route("/api/email/preview") def email_preview(): """渲染报告到文件 (不发送), 返回文件路径供下载预览。""" import tempfile ml = _get_mailer() if not ml: return jsonify({"ok": False, "error": "mailer.py 不可用"}), 500 fd, path = tempfile.mkstemp(prefix="palladium-mail-preview-", suffix=".eml", dir="/tmp") os.close(fd) ml.write_preview(path) with open(path, "r", encoding="utf-8") as f: content = f.read() try: os.unlink(path) except OSError: pass return Response(content, mimetype="text/plain; charset=utf-8", headers={"Content-Disposition": 'inline; filename="mail-preview.txt"'}) @app.route("/api/email/send", methods=["POST"]) def email_send(): """ 手动触发一次发送 (用于测试)。 安全: 除非 ALLOW_EMAIL_TRIGGER=true, 否则禁止 (防止误触发给真实收件人发邮件)。 """ ml = _get_mailer() if not ml: return jsonify({"ok": False, "error": "mailer.py 不可用"}), 500 if os.environ.get("ALLOW_EMAIL_TRIGGER", "").strip().lower() not in ("1", "true", "yes"): return jsonify({ "ok": False, "error": "未授权: 设置环境变量 ALLOW_EMAIL_TRIGGER=true 后才能手动发送" }), 403 result = ml.send_report() status = 200 if result["ok"] else 500 return jsonify(result), status # ═══════════════════════════════════════════════════════════════════════════ # Email Scheduler (后台 daemon thread, 每天指定时间触发) # ═══════════════════════════════════════════════════════════════════════════ _scheduler_started = False _scheduler_lock = threading.Lock() def _seconds_until_next(hour: int, minute: int) -> int: """距离下一次 hh:mm 还有多少秒 (0 ~ 86400)。""" now = datetime.now() target = now.replace(hour=hour, minute=minute, second=0, microsecond=0) if target <= now: target += timedelta(days=1) return int((target - now).total_seconds()) def _scheduler_loop(): """后台线程: 每天到点调一次 send_report(), 错误也不退出, 1 分钟后自愈。""" import mailer print("[scheduler] 邮件调度线程已启动", flush=True) while True: try: cfg = mailer.get_config() if not cfg["enabled"]: # 关闭, 每小时看一眼 time.sleep(3600) continue wait = _seconds_until_next(cfg["send_hour"], cfg["send_minute"]) print(f"[scheduler] 下次发送: {wait}s 后 " f"({cfg['send_hour']:02d}:{cfg['send_minute']:02d})", flush=True) time.sleep(wait) # 到点了, 触发一次 print("[scheduler] 触发每日报告发送...", flush=True) result = mailer.send_report(cfg) if result["ok"]: print(f"[scheduler] 发送成功: {result['subject']} -> {result['recipients']}", flush=True) else: print(f"[scheduler] 发送失败: {result['error']}", flush=True) except Exception as e: # 出错不退出, 1 分钟后重试 print(f"[scheduler] 异常: {type(e).__name__}: {e}", flush=True) time.sleep(60) def start_email_scheduler(): """启动后台调度线程 (只启动一次)。""" global _scheduler_started with _scheduler_lock: if _scheduler_started: return # 延迟导入, 避免 mailer 缺失时启动失败 ml = _get_mailer() if ml is None: print("[scheduler] mailer.py 不可用, 调度器未启动", flush=True) return # 检查配置: 至少 SMTP_HOST 设置了才启动 cfg = ml.get_config() if not cfg["smtp_host"]: print("[scheduler] SMTP_HOST 未设置, 调度器未启动 " "(设了之后 gunicorn 重启生效)", flush=True) return t = threading.Thread(target=_scheduler_loop, name="email-scheduler", daemon=True) t.start() _scheduler_started = True # ═══════════════════════════════════════════════════════════════════════════ # Export — 历史快照导出 (CSV / JSON) # ═══════════════════════════════════════════════════════════════════════════ EXPORT_FIELDS = [ ("id", "快照 ID"), ("timestamp_local", "时间(本地)"), ("timestamp_utc", "时间(UTC)"), ("emulator", "仿真器"), ("hardware", "硬件"), ("configmgr", "ConfigMgr"), ("system_status", "系统状态"), ("source", "数据来源"), ("total_boards", "板总数"), ("online_boards", "在线板"), ("offline_boards", "离线板"), ("used_boards", "使用板"), ("utilization", "利用率%"), ("total_domains", "域总数"), ("available_domains", "可用域"), ("downloaded_domains", "已加载域"), ("disabled_domains", "已禁用域"), ("not_exist_domains", "不存在域"), ("active_jobs", "活跃作业"), ] def _build_export_query(): """ 解析导出过滤参数, 返回构建好的 query。 参数: - days: 限制最近 N 天 (0/缺省=全部) - since: ISO 日期 (YYYY-MM-DD 或 YYYY-MM-DDTHH:MM:SS) - until: 同上 - source: 模糊匹配 source 字段 (如 "collect.py" 或 "scmp03") - limit: 最多返回 N 条 (默认 10000, 上限 100000) """ q = ScanSnapshot.query # ── 时间范围 ── days = request.args.get("days", type=int) if days and days > 0: cutoff = datetime.now(timezone.utc) - timedelta(days=days) q = q.filter(ScanSnapshot.timestamp >= cutoff) def _parse_dt(s: str): """接受 'YYYY-MM-DD' 或完整 ISO, 返回带 tz 的 datetime。""" if not s: return None s = s.strip() for fmt in ("%Y-%m-%dT%H:%M:%S", "%Y-%m-%d %H:%M:%S", "%Y-%m-%d"): try: dt = datetime.strptime(s, fmt) if dt.tzinfo is None: dt = dt.replace(tzinfo=timezone.utc) return dt except ValueError: continue return None since = _parse_dt(request.args.get("since", "")) if since is not None: q = q.filter(ScanSnapshot.timestamp >= since) until = _parse_dt(request.args.get("until", "")) if until is not None: q = q.filter(ScanSnapshot.timestamp <= until) # ── 来源过滤 ── source = request.args.get("source", "").strip() if source: q = q.filter(ScanSnapshot.source.like(f"%{source}%")) # ── 排序 + limit ── limit = min(request.args.get("limit", 10000, type=int), 100000) q = q.order_by(ScanSnapshot.timestamp.asc()).limit(limit) return q def _snapshot_to_row(snap: ScanSnapshot) -> dict: """把 ORM 对象转成可序列化的 dict (供 CSV/JSON 共享)。""" ts = snap.timestamp if ts and ts.tzinfo is None: ts = ts.replace(tzinfo=timezone.utc) return { "id": snap.id, "timestamp_utc": ts.isoformat() if ts else "", "timestamp_local": ts.astimezone().strftime("%Y-%m-%d %H:%M:%S") if ts else "", "emulator": snap.emulator or "", "hardware": snap.hardware or "", "configmgr": snap.configmgr or "", "system_status": snap.system_status or "", "source": snap.source or "", "total_boards": snap.total_boards, "online_boards": snap.online_boards, "offline_boards": snap.offline_boards, "used_boards": snap.used_boards, "utilization": snap.utilization, "total_domains": snap.total_domains, "available_domains": snap.available_domains, "downloaded_domains": snap.downloaded_domains, "disabled_domains": snap.disabled_domains, "not_exist_domains": snap.not_exist_domains, "active_jobs": snap.active_jobs, } @app.route("/api/export") def api_export(): """ 导出历史快照。 Query: - format: csv | json (默认 csv) - days: 限制最近 N 天 - since: 起始时间 (YYYY-MM-DD 或 ISO) - until: 截止时间 - source: 模糊匹配 source 字段 - limit: 最多 N 条 (默认 10000, 上限 100000) CSV 文件名: palladium_snapshots_YYYYMMDD_HHMMSS.csv JSON 文件名: palladium_snapshots_YYYYMMDD_HHMMSS.json """ q = _build_export_query() rows = [_snapshot_to_row(s) for s in q.all()] fmt = request.args.get("format", "csv").lower() ts_str = datetime.now().strftime("%Y%m%d_%H%M%S") if fmt == "json": body = { "exported_at": datetime.now(timezone.utc).isoformat(), "count": len(rows), "rows": rows, } resp = jsonify(body) resp.headers["Content-Disposition"] = ( f'attachment; filename="palladium_snapshots_{ts_str}.json"' ) return resp # ── CSV (默认) ── buf = io.StringIO() # 写 BOM 让 Excel 直接识别 UTF-8 (否则中文乱码) buf.write("\ufeff") writer = csv.writer(buf) writer.writerow([label for _, label in EXPORT_FIELDS]) for r in rows: writer.writerow([r.get(key, "") for key, _ in EXPORT_FIELDS]) csv_data = buf.getvalue() resp = Response(csv_data, mimetype="text/csv") resp.headers["Content-Disposition"] = ( f'attachment; filename="palladium_snapshots_{ts_str}.csv"' ) return resp # ═══════════════════════════════════════════════════════════════════════════ # Storage Helper # ═══════════════════════════════════════════════════════════════════════════ def store_snapshot(raw_output: str, source: str = "manual") -> int: """解析、存储一条快照,返回 snapshot id。""" parsed = parse_test_server_output(raw_output) stats = calculate_stats(parsed) sysinfo = parsed["system"] snap = ScanSnapshot( emulator=sysinfo.get("emulator", ""), hardware=sysinfo.get("hardware", ""), configmgr=sysinfo.get("configmgr", ""), system_status=sysinfo.get("system_status", ""), source=source, raw_output=raw_output, total_boards=stats["total_boards"], online_boards=stats["online_boards"], offline_boards=stats["offline_boards"], total_domains=stats["total_domains"], available_domains=stats["available_domains"], downloaded_domains=stats["downloaded_domains"], disabled_domains=stats["disabled_domains"], not_exist_domains=stats["not_exist_domains"], active_jobs=stats["active_jobs"], used_boards=stats["used_boards"], utilization=stats["utilization"], ) db.session.add(snap) db.session.flush() # 获取 snap.id for rack in parsed["racks"]: for cluster in rack["clusters"]: for board in cluster["boards"]: domains = board["domains"] db.session.add(BoardStatus( snapshot_id=snap.id, rack=rack["rack"], cluster=cluster["cluster"], ld_index=board["ld"], status=board["status"], ccd=cluster["ccd"], d0=domains[0], d1=domains[1], d2=domains[2], d3=domains[3], d4=domains[4], d5=domains[5], d6=domains[6], d7=domains[7], )) for job in parsed["jobs"]: db.session.add(JobInfo( snapshot_id=snap.id, job_index=job["index"], owner=job["owner"], pid=job["pid"], t_pod=job["t_pod"], design=job["design"], elap_time=job["elap_time"], reserved_key=job["reserved_key"], )) for pod in parsed["pods"]: db.session.add(PodInfo( snapshot_id=snap.id, rack=pod["rack"], all_hdsb=pod.get("all_hdsb", ""), all_t_pods=pod.get("all_t_pods", ""), available_hdsb=pod.get("available_hdsb", ""), available_t_pods=pod.get("available_t_pods", ""), locked_hdsb=pod.get("locked_hdsb", ""), locked_t_pods=pod.get("locked_t_pods", ""), reserved_hdsb=pod.get("reserved_hdsb", ""), reserved_t_pods=pod.get("reserved_t_pods", ""), unavailable_hdsb=pod.get("unavailable_hdsb", ""), unavailable_t_pods=pod.get("unavailable_t_pods", ""), )) db.session.commit() return snap.id # ═══════════════════════════════════════════════════════════════════════════ # Seed Sample Data (首次运行) # ═══════════════════════════════════════════════════════════════════════════ SAMPLE_OUTPUT = r"""Emulator: sc01_emu Hardware: Palladium Z1 Configmgr: V21.02.102.s002System Status: PARTIAL Rack 0 has 2 clusters Cluster 0 has 6 boards CCD: ONLINE LD Status D0 D1 D2 D3 D4 D5 D6 D7 0 ONLINE - - - - - - - - 1 ONLINE - - - - - - - - 2 ONLINE - - - - - - - - 3 ONLINE - - - - - - - - 4 ONLINE - - - - - - - - 5 ONLINE - - - - - - - - Cluster 1 has 6 boards CCD: ONLINE LD Status D0 D1 D2 D3 D4 D5 D6 D7 6 ONLINE 1 1 1 1 1 1 1 1 7 ONLINE 1 1 1 1 1 1 1 1 8 ONLINE - - - - - - - - 9 OFFLINE X X X X X X X X 10 ONLINE - - - - - - - - 11 ONLINE - - - - - - - - Rack 1 has 2 clusters Cluster 3 has 6 boards CCD: ONLINE LD Status D0 D1 D2 D3 D4 D5 D6 D7 18 ONLINE - - - - - - - - 19 ONLINE - - - - - - - - 20 ONLINE - - - - - - - - 21 ONLINE - - - - - - - - 22 ONLINE - - - - - - - - 23 ONLINE - - - - - - - - Cluster 4 has 6 boards CCD: ONLINE LD Status D0 D1 D2 D3 D4 D5 D6 D7 24 ONLINE - - - - - - - - 25 ONLINE - - - - - - - - 26 ONLINE - - - - - - - - 27 ONLINE - - - - - - - - 28 ONLINE - - - - - - - - 29 ONLINE - - - - - - - - X: Domain does not exist D: Domain is disabled -: Domain is available : Domain is downloaded or reserved Job Information: Index Owner PID T-Pod Design ElapTime ReservedKey 1 pzuser101 scmp03:106144 -- -- top 00:35:39 -- Target Pod Information: Rack 0: All length 5 HDSB: T4 All length 5 T-Pods: T0, T1, T2, T3, T5, T6, T7, T8, T9 Available HDSB: T4 Available T-Pods: T0, T1, T2, T3, T5 Locked HDSB: NONE Locked T-Pods: NONE Reserved HDSB: NONE Reserved T-Pods: NONE Unavailable HDSB: NONE Unavailable T-Pods: T6, T7, T8, T9 Rack 1: All length 5 HDSB: T4 All length 30 HDSB: T5 Available HDSB: T4, T5 Available T-Pods: NONE Locked HDSB: NONE Locked T-Pods: NONE Reserved HDSB: NONE Reserved T-Pods: NONE Unavailable HDSB: NONE Unavailable T-Pods: NONE Exit code: 0 Deinitializing IB - 0 - 0 """ def seed_sample_data(): """首次运行时写入示例数据。""" if ScanSnapshot.query.first(): return store_snapshot(SAMPLE_OUTPUT, source="sample") print("[seed] 示例数据已写入") # ═══════════════════════════════════════════════════════════════════════════ # Gunicorn hook: 让 gunicorn worker 也能启动调度器 (每个 worker 各 1 个) # 生产用法: --preload 时只有 master 会触发, --workers N 时 N 个 daemon 都会启动 # daemon 内部有 _scheduler_started 锁, 实际只跑一个; 其余 sleep 等下一天的空跑无害 # ═══════════════════════════════════════════════════════════════════════════ def _gunicorn_post_fork(server, worker): # noqa: ARG001 start_email_scheduler() # ═══════════════════════════════════════════════════════════════════════════ # Main # ═══════════════════════════════════════════════════════════════════════════ if __name__ == "__main__": with app.app_context(): db.create_all() seed_sample_data() start_email_scheduler() print("════════════════════════════════════════════════════════") print(" Palladium Z1 Monitor") print(" http://0.0.0.0:5100") print("════════════════════════════════════════════════════════") app.run(host="0.0.0.0", port=5100, debug=True)