Files
track/backend/app/services/audit_service.py
duxingchen 1fea30b03b feat(audit): 日活使用统计 + 双 CSV 导出(审计/上下线,列可自定义)
【日活统计:GET /audit/daily-usage】
按【北京时间自然日 × 操作人】聚合,单个 GROUP BY 完成(count(*) FILTER),
不用窗口函数。指标:上线/下线时间、操作次数、登录/登出次数。

⚠️ 上线/下线时间取【当天首次/末次活动】,刻意不取登录时间:
   Refresh Token 有效期 7 天,用户不必每天重新登录。按登录算会出现
   「登录次数 0、上线时间空,但操作次数 35」——报表自相矛盾。
   时间一律按 +08:00 分日与渲染,否则早班(00:00~08:00)操作会掉到前一天。

【CSV 导出:两个端点 + 列自定义】
- /audit/logs/export        整体审计导出,筛选维度与列表页完全一致
- /audit/daily-usage/export 上下线导出,每人一行
- 列清单由后端统一维护并经 /audit/options 下发(log_export_columns /
  usage_export_columns),前端不硬编码表头,避免两端漂移
- ⚠️ 响应带 UTF-8 BOM:Excel 靠它识别编码,否则中文表头全乱码
- 单次上限 5 万行,超出经 X-Export-Truncated 头告知前端明确提示
  (静默截断比报错更危险)
- list_audit_logs 与 export_audit_logs 共用 _log_filters,
  保证「看到的」与「导出的」永远是同一批数据

【前端】
- 审计页页头新增「人员统计」「导出 CSV」两个按钮,现有表格与筛选零改动
- 人员统计走抽屉(AuditUsagePanel):日期范围+快捷键、日活表格、
  上下线次数彩色标签、北京时间渲染、导出前弹列勾选面板
- ExportColumnsModal 为两处导出共用,默认全选
2026-09-21 12:06:43 +08:00

340 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""审计服务 — 写入与检索
写入方案的取舍(与 MOM/KCGL 不同,理由如下)
--------------------------------------------------
MOM 用 SQLAlchemy event listener + **同事务**写入:优点是全自动、业务代码零改动;
缺点是业务事务回滚时审计记录一起被回滚掉 —— 而失败/被拒的操作恰恰是最需要
留痕的(比如越权尝试、参数错误导致的 4xx
Track 改为:响应生成后,用**独立 session** 写入审计。
- 业务回滚不影响审计,失败操作照样留痕
- 审计写入失败也不影响业务(全包裹 try/except仅记日志
- 代价:审计与业务不是原子提交,极端情况(响应后进程立即被 kill可能丢一条。
对内部系统的操作审计,这个取舍划算。
"""
from __future__ import annotations
import logging
import uuid
from datetime import datetime
from sqlalchemy import and_, func, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.database import AsyncSessionLocal
from app.models.audit_log import AuditLog
logger = logging.getLogger("track.audit")
# 绝不落库的敏感字段名(命中即替换为 ***
# 登录请求体含明文密码,一旦进审计表就成了长期泄露面
_SENSITIVE_KEYS = frozenset(
{"password", "passwd", "pwd", "token", "access_token", "refresh_token",
"secret", "api_key", "authorization", "password_hash"}
)
# 模块 / 动作 的中文标签(前端下拉与列表展示用)
MODULE_LABELS: dict[str, str] = {
"auth": "认证登录",
"product": "产品管理",
"task": "任务流转",
"order": "订单管理",
"record": "任务记录",
"print": "标签打印",
"material": "物料",
"user": "用户",
"notification": "消息通知",
"upload": "文件上传",
"dashboard": "看板统计",
"analytics": "效能分析",
"screen": "数据大屏",
"holiday": "节假日配置",
"app": "App版本",
"external": "外部系统对接",
"audit": "审计日志",
"other": "其它",
}
ACTION_LABELS: dict[str, str] = {
"create": "新增",
"update": "修改",
"delete": "删除",
"read": "查询",
"export": "导出",
"login": "登录",
"logout": "登出",
"refresh": "刷新令牌",
"print": "打印",
"upload": "上传",
"finalize": "收口",
"receive": "接收",
"transfer": "转交",
"reject": "驳回",
"recall": "撤回",
"spawn": "派发",
"end": "结束分支",
"complete": "完结",
}
def sanitize_details(details: dict | None) -> dict | None:
"""递归剔除敏感字段,避免密码/令牌落库"""
if not details:
return details
def _clean(value):
if isinstance(value, dict):
return {
k: ("***" if str(k).lower() in _SENSITIVE_KEYS else _clean(v))
for k, v in value.items()
}
if isinstance(value, list):
return [_clean(v) for v in value]
return value
return _clean(details)
async def record_audit(
*,
action: str,
module: str,
user_id: str | None = None,
display_name: str | None = None,
role: str | None = None,
target_type: str | None = None,
target_id: str | None = None,
target_name: str | None = None,
details: dict | None = None,
ip_address: str | None = None,
user_agent: str | None = None,
method: str | None = None,
url: str | None = None,
status_code: int | None = None,
error_message: str | None = None,
request_id: str | None = None,
) -> None:
"""写入一条审计记录。**绝不抛异常**:审计失败不能影响业务。"""
try:
async with AsyncSessionLocal() as session:
session.add(
AuditLog(
id=uuid.uuid4(),
user_id=user_id,
display_name=display_name,
role=role,
action=action,
module=module,
target_type=target_type,
target_id=str(target_id) if target_id is not None else None,
target_name=target_name,
details=sanitize_details(details),
ip_address=ip_address,
user_agent=user_agent[:500] if user_agent else None,
method=method,
url=url[:500] if url else None,
status_code=status_code,
error_message=error_message,
request_id=request_id,
)
)
await session.commit()
except Exception:
# 用 exception 级别但吞掉异常:保证调用方业务流程不受影响
logger.exception(
"审计写入失败(已忽略,不影响业务)",
extra={"extra_fields": {"action": action, "module": module, "url": url}},
)
async def list_audit_logs(
db: AsyncSession,
*,
user_id: str | None = None,
module: str | None = None,
action: str | None = None,
target_id: str | None = None,
request_id: str | None = None,
status_code: int | None = None,
start: datetime | None = None,
end: datetime | None = None,
skip: int = 0,
limit: int = 50,
) -> tuple[list[AuditLog], int]:
"""审计日志检索(按时间倒序)。返回 (当前页, 真实总数)。
真实总数走独立 COUNT —— 前端分页器依赖它,不能用 len(当前页)。
"""
filters = _log_filters(
user_id=user_id, module=module, action=action, target_id=target_id,
request_id=request_id, status_code=status_code, start=start, end=end,
)
total = await db.scalar(
select(func.count()).select_from(AuditLog).where(*filters)
) or 0
rows = (
await db.execute(
select(AuditLog)
.where(*filters)
.order_by(AuditLog.created_at.desc())
.offset(skip)
.limit(limit)
)
).scalars().all()
return list(rows), total
# ============================================================
# 导出
# ============================================================
# 单次导出的行数上限。审计表只增不减,全量导出迟早会撑爆内存与浏览器,
# 故设硬上限;超出时向上层返回 truncated=True由前端明确提示「已截断」——
# 静默截断会让使用者以为导全了,比报错更危险。
EXPORT_MAX_ROWS = 50000
def _log_filters(
*,
user_id: str | None = None,
module: str | None = None,
action: str | None = None,
target_id: str | None = None,
request_id: str | None = None,
status_code: int | None = None,
start: datetime | None = None,
end: datetime | None = None,
) -> list:
"""审计日志的筛选条件 —— list_audit_logs 与 export_audit_logs 共用。
抽出来的唯一目的:保证「列表看到的」和「导出出去的」永远是同一批数据。
两处各写一份迟早会漂移,而导出与列表不一致是最让人不信任的那种 bug。
"""
filters = []
if user_id:
filters.append(AuditLog.user_id.ilike(f"%{user_id}%"))
if module:
filters.append(AuditLog.module == module)
if action:
filters.append(AuditLog.action == action)
if target_id:
filters.append(AuditLog.target_id == target_id)
if request_id:
filters.append(AuditLog.request_id == request_id)
if status_code is not None:
filters.append(AuditLog.status_code == status_code)
if start:
filters.append(AuditLog.created_at >= start)
if end:
filters.append(AuditLog.created_at <= end)
return filters
async def export_audit_logs(
db: AsyncSession, *, limit: int = EXPORT_MAX_ROWS, **kwargs,
) -> tuple[list[AuditLog], bool]:
"""导出用:按筛选条件取全部记录(不分页)。返回 (rows, truncated)。
多取一行来判断是否被截断 —— 比再跑一次 COUNT 便宜。
"""
rows = (
await db.execute(
select(AuditLog)
.where(*_log_filters(**kwargs))
.order_by(AuditLog.created_at.desc())
.limit(limit + 1)
)
).scalars().all()
truncated = len(rows) > limit
return list(rows[:limit]), truncated
# ============================================================
# 日活 / 使用统计
# ============================================================
# 成功 = 2xx/3xx。登录失败401也要留痕但不应计入"上线次数"。
_OK_STATUS_UPPER = 400
async def get_daily_usage(
db: AsyncSession, *, start: datetime, end: datetime,
) -> list[dict]:
"""按【北京时间自然日 × 操作人】聚合用量 —— 日活报表的数据源。
start/end 为半开区间 [start, end),调用方按北京时间日界传入。
全部指标由**一个 GROUP BY 查询**算出,不用窗口函数:
· 登录/登出次数 = 成功登录 / 成功登出数(最终凭证是 login_count不是"上线次数"
· 操作频次 = 当天该用户的全部审计记录数(代表系统使用深度)
· 上线/下线时间 = 当天**首次 / 末次活动**时间(任意审计记录)
⚠️ 上线/下线时间【不能】取登录/登出时间。
Access/Refresh Token 有效期内refresh 7 天)用户无需重新登录,
于是"周一登录、周二到周日继续用"会导致周二~周日:
登录次数=0、登录时间=空,但操作次数却是几十 —— 报表自相矛盾。
改用活动口径后,工人当天的第一次操作就是真实上线时间。
(审计中间件是全站的,移动端的接收/转交/完工同样入账,故自动覆盖移动端,
无需前端上报心跳。)
为什么用 `count(*) FILTER (WHERE ...)`:分组内一次扫描同时算出多个条件计数,
比多次子查询或 UNION 简单得多且语义一眼可读。Postgres 原生支持。
⚠️ 按【北京时间】分日created_at 是 timestamptz实存 UTC
直接按 UTC 分日会让 00:00~08:00 的早班操作掉到前一天。
"""
day_col = func.date(func.timezone("Asia/Shanghai", AuditLog.created_at))
login_ok = and_(
AuditLog.action == "login", AuditLog.status_code < _OK_STATUS_UPPER,
)
logout_ok = and_(
AuditLog.action == "logout", AuditLog.status_code < _OK_STATUS_UPPER,
)
stmt = (
select(
day_col.label("day"),
AuditLog.user_id.label("user_id"),
# 同一用户的 display_name / role 是一致的,取 max 只是为了
# 在 GROUP BY 下拿到一个非空代表值(避免再套一层 DISTINCT ON
func.max(AuditLog.display_name).label("display_name"),
func.max(AuditLog.role).label("role"),
func.count().filter(login_ok).label("login_count"),
func.count().filter(logout_ok).label("logout_count"),
func.count().label("op_count"),
# 上线/下线时间取「任意记录」的首末,而不是登录/登出的首末(原因见 docstring
func.min(AuditLog.created_at).label("first_active_at"),
func.max(AuditLog.created_at).label("last_active_at"),
)
.where(
AuditLog.created_at >= start,
AuditLog.created_at < end,
# 只统计"人"未认证请求如登录前的探测、refresh没有操作人
# 混进来会让"日活人数"虚高。若要排查匿名异常流量,走日志列表页按
# 结果/来源 IP 过滤更合适。
AuditLog.user_id.isnot(None),
)
.group_by(day_col, AuditLog.user_id)
.order_by(day_col.desc(), func.count().desc())
)
rows = (await db.execute(stmt)).all()
return [
{
"day": r.day.strftime("%Y-%m-%d") if hasattr(r.day, "strftime") else str(r.day),
"user_id": r.user_id,
"display_name": r.display_name,
"role": r.role,
"login_count": r.login_count or 0,
"logout_count": r.logout_count or 0,
"op_count": r.op_count or 0,
"first_active_at": r.first_active_at,
"last_active_at": r.last_active_at,
}
for r in rows
]