Files
track/backend/app/services/audit_service.py
duxingchen 39697ca3ad feat(audit): 每日活动打点表 —— 修正日活「上线/下线时间」口径
【问题】
上线/下线时间此前取自审计记录(写操作)时间,当天只翻看、没做写操作的人
会被整条漏掉;而取登录时间更错(Refresh Token 有效期 7 天,用户不必每天登录,
会出现「登录次数 0 却操作 35 次」的自相矛盾报表)。

【方案 C+:一天一人一行的小状态表】
- 新增 user_daily_seen(user_id, day, first_seen_at, last_seen_at),
  迁移 k1l2m3n4o5p6(紧接 j1k2l3m4n5o6)
- 打点挂在 RequestContextMiddleware —— 它是最外层,能覆盖【所有】请求,
  含不被审计的普通 GET。挂审计中间件没用:那里只记写操作,正是漏人的原因
- 节流:进程内缓存,同一用户 2 分钟内只落盘一次,把"每请求一次写库"
  压到"每人每 2 分钟一次";代价是末次活动时间最多落后 2 分钟
- UPSERT on_conflict_do_update 只刷 last_seen_at,first_seen_at 保持当天首次值
- 打点失败全部吞掉并记日志,绝不影响业务请求

【为什么不复用 audit_logs】
「末次活动」是需要不断 UPDATE 的状态,而审计流水必须只增不改 ——
能改的审计记录等于没有审计价值。写进审计表还会让表随访问量线性膨胀。

【查询合并】
get_daily_usage 改为「活动表 ∪ 审计表」并集:只有活动记录的(只看不操作)
和只有审计记录的(本表上线前的历史数据)都会出现。
上线/下线时间取两者的【最早 / 最晚】而非"活动表优先"——打点有 2 分钟节流、
跨零点或写库失败时可能晚于当天第一次写操作,取 min/max 后结果永不劣于任一来源。

零前端改动、零移动端发版:打点在服务端,PC 与移动端同一套口径、同一张表。
2026-09-21 13:05:35 +08:00

453 lines
17 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 time
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.core.time_utils import get_beijing_time
from app.models.audit_log import AuditLog
from app.models.user_daily_seen import UserDailySeen
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
# ============================================================
# 每日活动打点(日活报表的「上线时间 / 下线时间」来源)
# ============================================================
# 同一用户两次落盘之间的最小间隔(秒)。
#
# 打点挂在「每个请求」上,但不希望每个请求都写一次数据库 —— 那会把
# user_daily_seen 变成热点。这里用进程内缓存做节流:同一用户 2 分钟内
# 只落盘一次。代价是「末次活动时间」最多落后真实值 2 分钟,
# 对"日活统计"这个精度要求完全够用。
#
# 多 worker 部署时每个进程各持一份缓存,实际写库频率最多放大到 worker 数倍
# 4 worker × 每人每 2 分钟 1 次),依然可忽略。
_TOUCH_INTERVAL_S = 120.0
_touch_cache: dict[str, float] = {}
# 缓存只增不减会缓慢泄漏(键是 user_id量级 = 用户数,实际很小)。
# 超过阈值就整体清空 —— 代价只是多写几次库,换来内存有界。
_TOUCH_CACHE_MAX = 5000
async def touch_daily_seen(user_id: str | None) -> None:
"""记录「该用户此刻活动过」。首次 INSERT、其后只刷新 last_seen_at。
唯一的消费方是日活报表的上线/下线时间(见 get_daily_usage
刻意不写进 audit_logs那是只增不改的审计流水而本表是需要不断
UPDATE 的状态(详见 UserDailySeen 模型注释)。
任何异常都吞掉 —— 活动打点失败绝不能影响业务请求本身。
"""
if not user_id:
return
now_mono = time.monotonic()
last = _touch_cache.get(user_id)
if last is not None and now_mono - last < _TOUCH_INTERVAL_S:
return # 节流窗口内,跳过
if len(_touch_cache) > _TOUCH_CACHE_MAX:
_touch_cache.clear()
# 先占位再写库:同一用户的并发请求不会同时打进来
_touch_cache[user_id] = now_mono
try:
from sqlalchemy.dialects.postgresql import insert as pg_insert
now = get_beijing_time()
day = now.date() # 北京时间自然日(与报表分日口径一致)
async with AsyncSessionLocal() as db:
await db.execute(
pg_insert(UserDailySeen)
.values(user_id=user_id, day=day, first_seen_at=now, last_seen_at=now)
# 冲突时只刷新 last_seen_atfirst_seen_at 保持当天首次值不变
.on_conflict_do_update(
index_elements=["user_id", "day"],
set_={"last_seen_at": now},
)
)
await db.commit()
except Exception: # noqa: BLE001 —— 打点失败不影响业务
logger.exception("记录每日活动失败(已忽略)")
# ============================================================
# 日活 / 使用统计
# ============================================================
# 成功 = 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不是"上线次数"
· 操作频次 = 当天该用户的全部审计记录数(代表系统使用深度)
· 登录/登出次数 = 成功登录 / 成功登出数
· 上线/下线时间 = 当天**首次 / 末次活动**(优先取 user_daily_seen
⚠️ 上线/下线时间【不能】取登录/登出时间。
Access/Refresh Token 有效期内refresh 7 天)用户无需重新登录,
于是"周一登录、周二到周日继续用"会导致周二~周日:
登录次数=0、登录时间=空,但操作次数却是几十 —— 报表自相矛盾。
⚠️ 也不能只取审计表的写操作时间:审计中间件只记写操作,普通 GET 不入账,
当天只翻看、没做写操作的人会被整条漏掉。
故上线/下线时间优先取 user_daily_seen挂在每个请求上打点
仅对本表上线前的历史数据回退到审计表的写操作时间。
为什么用 `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()
# ── 活动表:当天首次/末次活动(覆盖"只翻看不操作"的人)──
# day 列是北京时间 DATE与上面的 day_col 口径一致,可直接按 (user_id, day) 对齐。
# 取 start.date() ~ end.date()end 是次日 00:00 的半开上界,故用 <)。
seen_rows = (
await db.execute(
select(
UserDailySeen.user_id, UserDailySeen.day,
UserDailySeen.first_seen_at, UserDailySeen.last_seen_at,
).where(
UserDailySeen.day >= start.date(),
UserDailySeen.day < end.date(),
)
)
).all()
seen = {
(s.user_id, s.day.strftime("%Y-%m-%d")): (s.first_seen_at, s.last_seen_at)
for s in seen_rows
}
audit = {
(r.user_id, r.day.strftime("%Y-%m-%d") if hasattr(r.day, "strftime") else str(r.day)): r
for r in rows
}
# ── 合并 ──
# 并集:只有审计记录的人(本表上线前的历史数据)和只有活动记录的人
# (当天只翻看、没做写操作)都要出现,各自缺的部分留空/计 0。
items: list[dict] = []
for key in set(audit) | set(seen):
user_id, day = key
a = audit.get(key)
first_seen, last_seen = seen.get(key, (None, None))
# 取「两者的最早/最晚」,而不是简单地"活动表优先"
# 活动表靠请求触发且有 2 分钟节流,极端情况(跨零点被节流、
# 打点写库失败被吞掉)可能晚于当天第一次写操作。
# 取 min/max 后,结果永远不会比任一来源更差,也不需要为兜底写分支逻辑。
audit_first = a.first_active_at if a else None
audit_last = a.last_active_at if a else None
first_candidates = [t for t in (first_seen, audit_first) if t is not None]
last_candidates = [t for t in (last_seen, audit_last) if t is not None]
items.append({
"day": day,
"user_id": user_id,
# 姓名字段只有审计记录里有(活动表为了轻量刻意不冗余存)
"display_name": a.display_name if a else None,
"role": a.role if a else None,
"login_count": (a.login_count or 0) if a else 0,
"logout_count": (a.logout_count or 0) if a else 0,
"op_count": (a.op_count or 0) if a else 0,
"first_active_at": min(first_candidates) if first_candidates else None,
"last_active_at": max(last_candidates) if last_candidates else None,
})
# 与 SQL 里的排序保持一致:日期倒序 → 操作次数倒序
items.sort(key=lambda x: (x["day"], x["op_count"]), reverse=True)
return items