文本日志只能 grep,"这台设备这次运行为什么失败"翻起来很费劲。新增一张 **结构化**的步骤明细表,把"哪一步、什么类型、哪个选择器、结果、耗时"落库。 - core/models.py:新增 task_step_log(run_id/job/设备/step_path/selector/ result/detail/duration_ms),索引 run_id、created_at、(serial,created_at)、 (job_id,created_at); - core/step_log.py(新):**专用写线程 + 有界队列**批量落库——步骤执行是热路径, 任务线程只 put_nowait(实测 9000 次入队 31ms),队列满丢弃并计数,绝不阻塞; 另有保留期清理(默认 14 天,每日 04:13 + 每次启动); - tasks/generic/task.py:_exec_one 记一条(异常=error / handler 返回 False=miss / 未知类型=unknown / 概率未触发=skip);_exec_steps 维护路径栈得到 "2.1.3" 这样的嵌套位置;**单次运行封顶 2000 条**——forever 循环任务否则会写爆表; - 运行上下文 ctx(run_id/job_id/job_name/device_name)由 TaskManager 生成, 经 create_worker(serial, params, ctx=None) 传入 worker(扩展点向后兼容); - 接口:/api/step_logs(过滤+分页)、/runs(按运行归组)、/filters(下拉选项)、 /download(CSV,带 BOM); - 前端:日志页拆成「文件日志 / 步骤明细」子分栏 + static/admin/steplog.js。 红线:新表自动进备份覆盖清单(SUMMARY_TABLES 由元数据派生),已补 TABLE_LABELS 中文标签,导出实测 14 张表、coverage_missing 为空。 文档:DATA_MODEL §2.8/§1、API §10.2、ARCHITECTURE §1.1/§2.2/§3.1、 DEVELOPMENT §5.2 与红线表、DEPLOY §5.2、TASK_DEV §8.3、README。
255 lines
11 KiB
Python
255 lines
11 KiB
Python
"""管理域 API:用户管理/日志查看。"""
|
||
import time
|
||
from io import BytesIO
|
||
|
||
from flask import Blueprint, jsonify, request, send_file
|
||
from flask_login import current_user
|
||
from flask import session
|
||
|
||
from core import logger
|
||
from core.logger import get_logger
|
||
from core.models import db, User
|
||
from web import context
|
||
from web.auth import (admin_required, perm_required, _validate_perms,
|
||
PERM_LOGS)
|
||
|
||
_log = get_logger("web")
|
||
bp = Blueprint("admin", __name__)
|
||
|
||
@bp.route("/api/users")
|
||
@admin_required
|
||
def api_users_list():
|
||
users = [{"id": u.id, "username": u.username,
|
||
"is_admin": u.is_admin, "perms": u.get_perms()} for u in User.query.all()]
|
||
return jsonify({"ok": True, "users": users})
|
||
|
||
@bp.route("/api/users", methods=["POST"])
|
||
@admin_required
|
||
def api_users_create():
|
||
data = request.json or {}
|
||
username = (data.get("username") or "").strip()
|
||
password = data.get("password", "")
|
||
if not username or not password:
|
||
return jsonify({"ok": False, "error": "用户名和密码不能为空"}), 400
|
||
if User.query.filter_by(username=username).first():
|
||
return jsonify({"ok": False, "error": "用户名已存在"}), 400
|
||
# 默认普通用户(最小权限);管理员由 is_admin 决定,权限位照常保存(降级后生效)
|
||
u = User(username=username, is_admin=bool(data.get("is_admin", False)))
|
||
u.set_perms(_validate_perms(data.get("perms") or []))
|
||
u.set_password(password)
|
||
db.session.add(u)
|
||
db.session.commit()
|
||
_log.info(f"创建用户 {username} (admin={u.is_admin}, perms={u.get_perms()})")
|
||
return jsonify({"ok": True, "msg": "用户已创建"})
|
||
|
||
@bp.route("/api/users/<int:uid>", methods=["PUT"])
|
||
@admin_required
|
||
def api_users_update(uid):
|
||
u = User.query.get(uid)
|
||
if not u:
|
||
return jsonify({"ok": False, "error": "用户不存在"}), 404
|
||
data = request.json or {}
|
||
if "password" in data and data["password"]:
|
||
u.set_password(data["password"])
|
||
if "is_admin" in data:
|
||
new_admin = bool(data["is_admin"])
|
||
if u.is_admin and not new_admin and User.query.filter_by(is_admin=True).count() <= 1:
|
||
return jsonify({"ok": False, "error": "不能取消最后一个管理员"}), 400
|
||
u.is_admin = new_admin
|
||
if "perms" in data:
|
||
u.set_perms(_validate_perms(data["perms"]))
|
||
db.session.commit()
|
||
_log.info(f"更新用户 {u.username} (admin={u.is_admin}, perms={u.get_perms()})")
|
||
return jsonify({"ok": True, "msg": "用户已更新"})
|
||
|
||
@bp.route("/api/users/<int:uid>", methods=["DELETE"])
|
||
@admin_required
|
||
def api_users_delete(uid):
|
||
u = User.query.get(uid)
|
||
if not u:
|
||
return jsonify({"ok": False, "error": "用户不存在"}), 404
|
||
if u.username == "admin":
|
||
return jsonify({"ok": False, "error": "不能删除默认管理员"}), 400
|
||
if u.id == current_user.id:
|
||
return jsonify({"ok": False, "error": "不能删除当前登录用户"}), 400
|
||
if u.is_admin and User.query.filter_by(is_admin=True).count() <= 1:
|
||
return jsonify({"ok": False, "error": "不能删除最后一个管理员"}), 400
|
||
db.session.delete(u)
|
||
db.session.commit()
|
||
_log.info(f"删除用户 {u.username}")
|
||
return jsonify({"ok": True, "msg": "用户已删除"})
|
||
|
||
|
||
# ================== API:日志查看 ==================
|
||
#
|
||
# 过滤/白名单都在 core.logger 里做(读取侧与写入侧共用同一份文件清单,
|
||
# 见 core/logger.py「日志读取」一节);这里只负责取参数 + 组装响应。
|
||
|
||
|
||
def _log_filters():
|
||
"""从 query string 取过滤条件(/api/logs 与 /api/logs/download 共用)。"""
|
||
return dict(keyword=request.args.get("q", ""),
|
||
min_level=request.args.get("level", ""),
|
||
since=request.args.get("since", ""),
|
||
until=request.args.get("until", ""))
|
||
|
||
|
||
@bp.route("/api/logs")
|
||
@perm_required(PERM_LOGS)
|
||
def api_logs():
|
||
"""读日志:可按关键字 / 最低级别 / 时间范围过滤,返回结构化行。
|
||
|
||
参数:file 文件名(白名单)、q 关键字、level 最低级别(ERROR 含以上…)、
|
||
since/until 时间、lines 返回条数(取**最近** N 条命中)。
|
||
"""
|
||
files = logger.list_log_files()
|
||
names = [f["name"] for f in files]
|
||
# 默认仍落在 core.log(接口按 mtime 倒序返回,names[0] 会随当前哪个模块在
|
||
# 写而漂移,作为"打开日志页时看哪个"不稳定)
|
||
default = "core.log" if "core.log" in names else (names[0] if names else "")
|
||
current = request.args.get("file") or default
|
||
try:
|
||
lines = int(request.args.get("lines", 300))
|
||
except (TypeError, ValueError):
|
||
lines = 300
|
||
res = logger.query_log(current, limit=lines, **_log_filters())
|
||
return jsonify({"ok": not res["error"], "error": res["error"], "file": current,
|
||
"files": files, "rows": res["rows"], "matched": res["matched"],
|
||
"scanned": res["scanned"], "truncated": res["truncated"]})
|
||
|
||
|
||
@bp.route("/api/logs/download")
|
||
@perm_required(PERM_LOGS)
|
||
def api_logs_download():
|
||
"""下载日志文件(带同样的过滤条件;无条件时就是整个文件)。
|
||
|
||
用 BytesIO 发送而不是 send_file(路径):Windows 上流式发送时文件句柄可能
|
||
到 close 仍未释放,而日志文件正被日志线程持续写入,按路径发容易踩锁。
|
||
"""
|
||
name = request.args.get("file", "")
|
||
text, err = logger.read_log_text(name, **_log_filters())
|
||
if err:
|
||
return jsonify({"ok": False, "error": err}), 404
|
||
data = text.encode("utf-8")
|
||
fname = f"{name}_{time.strftime('%Y%m%d_%H%M%S')}.txt"
|
||
_log.info("下载日志: %s(%d 字节)by %s", name, len(data),
|
||
getattr(current_user, "username", ""))
|
||
return send_file(BytesIO(data), as_attachment=True, download_name=fname,
|
||
mimetype="text/plain; charset=utf-8")
|
||
|
||
|
||
# ================== API:任务步骤明细 ==================
|
||
#
|
||
# 数据源是 task_step_log 表(结构化),不是 logs/task.log 文本——所以能按设备/
|
||
# 任务/时间/结果过滤。写入侧见 core/step_log.py(异步写线程 + 保留期清理)。
|
||
|
||
def _step_filters():
|
||
"""从 query string 取步骤明细的过滤条件(列表与下载共用)。"""
|
||
return dict(serial=request.args.get("serial", ""),
|
||
job_id=request.args.get("job_id", ""),
|
||
result=request.args.get("result", ""),
|
||
run_id=request.args.get("run_id", ""),
|
||
keyword=request.args.get("q", ""),
|
||
since=logger.norm_ts(request.args.get("since", "")),
|
||
until=logger.norm_ts(request.args.get("until", ""), end=True))
|
||
|
||
|
||
def _step_to_dict(r):
|
||
return {"id": r.id, "run_id": r.run_id, "job_id": r.job_id,
|
||
"job_name": r.job_name, "serial": r.serial,
|
||
"device_name": r.device_name, "step_path": r.step_path,
|
||
"step_label": r.step_label, "step_type": r.step_type,
|
||
"selector": r.selector, "result": r.result, "detail": r.detail,
|
||
"duration_ms": r.duration_ms, "created_at": r.created_at}
|
||
|
||
|
||
@bp.route("/api/step_logs")
|
||
@perm_required(PERM_LOGS)
|
||
def api_step_logs():
|
||
"""任务步骤明细(分页,最近的在前)。"""
|
||
from core import step_log
|
||
try:
|
||
limit = int(request.args.get("limit", 200))
|
||
offset = int(request.args.get("offset", 0))
|
||
except (TypeError, ValueError):
|
||
limit, offset = 200, 0
|
||
rows, total = step_log.query(limit=limit, offset=offset, **_step_filters())
|
||
return jsonify({"ok": True, "rows": [_step_to_dict(r) for r in rows],
|
||
"total": total, "limit": limit, "offset": offset,
|
||
"stats": step_log.stats(),
|
||
"keep_days": step_log.KEEP_DAYS,
|
||
"max_rows_per_run": step_log.MAX_ROWS_PER_RUN})
|
||
|
||
|
||
@bp.route("/api/step_logs/runs")
|
||
@perm_required(PERM_LOGS)
|
||
def api_step_logs_runs():
|
||
"""按 run_id 归组的一次运行概览(哪次运行、几步、失败几步)。"""
|
||
from core import step_log
|
||
f = _step_filters()
|
||
rows = step_log.runs(serial=f["serial"], job_id=f["job_id"],
|
||
since=f["since"], until=f["until"],
|
||
limit=request.args.get("limit", 100))
|
||
return jsonify({"ok": True, "runs": [
|
||
{"run_id": r.run_id, "job_id": r.job_id, "job_name": r.job_name,
|
||
"serial": r.serial, "device_name": r.device_name,
|
||
"started_at": r.started_at, "ended_at": r.ended_at,
|
||
"steps": int(r.steps or 0), "failures": int(r.failures or 0)}
|
||
for r in rows]})
|
||
|
||
|
||
@bp.route("/api/step_logs/filters")
|
||
@perm_required(PERM_LOGS)
|
||
def api_step_logs_filters():
|
||
"""筛选下拉的可选项:**只在明细里出现过的**设备与任务(按数据自洽,不依赖
|
||
设备池/任务表,也不要求调用方另有 task/device 权限)。"""
|
||
from core import step_log
|
||
from core.models import TaskStepLog, db
|
||
devices = (db.session.query(TaskStepLog.serial, TaskStepLog.device_name)
|
||
.filter(TaskStepLog.serial != "").distinct().limit(500).all())
|
||
jobs = (db.session.query(TaskStepLog.job_id, TaskStepLog.job_name)
|
||
.filter(TaskStepLog.job_id != "").distinct().limit(500).all())
|
||
return jsonify({"ok": True,
|
||
"devices": sorted(({"serial": s, "device_name": n}
|
||
for s, n in devices),
|
||
key=lambda x: x["device_name"] or x["serial"]),
|
||
"jobs": sorted(({"job_id": i, "job_name": n} for i, n in jobs),
|
||
key=lambda x: x["job_name"] or x["job_id"]),
|
||
"results": ["ok", "miss", "error", "unknown", "skip", "cap"],
|
||
"keep_days": step_log.KEEP_DAYS,
|
||
"max_rows_per_run": step_log.MAX_ROWS_PER_RUN})
|
||
|
||
|
||
@bp.route("/api/step_logs/download")
|
||
@perm_required(PERM_LOGS)
|
||
def api_step_logs_download():
|
||
"""导出步骤明细为 CSV(带同样的过滤条件)。
|
||
|
||
加 UTF-8 BOM:不加的话 Excel 打开中文是乱码(这是给运维看的表,
|
||
大概率会被 Excel 打开)。
|
||
"""
|
||
import csv
|
||
from io import StringIO
|
||
from core import step_log
|
||
f = _step_filters()
|
||
# 导出上限:和列表页共用同一次查询,最多 2 万行(再多请缩小时间范围)
|
||
rows, total = step_log.query(limit=20000, offset=0, **f)
|
||
buf = StringIO()
|
||
w = csv.writer(buf)
|
||
w.writerow(["时间", "设备", "设备名", "任务", "运行ID", "步骤路径", "步骤",
|
||
"类型", "选择器", "结果", "耗时(ms)", "详情"])
|
||
for r in reversed(rows): # 导出的时间序与页面相反:文件里按正序更好读
|
||
w.writerow([r.created_at, r.serial, r.device_name, r.job_name,
|
||
r.run_id, r.step_path, r.step_label, r.step_type,
|
||
r.selector, r.result, r.duration_ms, r.detail])
|
||
data = buf.getvalue().encode("utf-8-sig") # utf-8-sig = UTF-8 带 BOM(Excel 中文不乱码)
|
||
fname = f"step_log_{time.strftime('%Y%m%d_%H%M%S')}.csv"
|
||
_log.info("导出步骤明细: %d 行(命中 %d)by %s", len(rows), total,
|
||
getattr(current_user, "username", ""))
|
||
return send_file(BytesIO(data), as_attachment=True, download_name=fname,
|
||
mimetype="text/csv; charset=utf-8")
|
||
|
||
|
||
# ================== API:运行控制 ==================
|
||
|