Files
palladium-monitor/app.py
T

1192 lines
48 KiB
Python

#!/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/<int:sid>")
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/<int:sid>/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
<num>: 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)