【扩大采集:核心业务数据的「查看详情」】
新增 _TRACKED_READ_PREFIXES(notifications / tasks / orders /
products / records)—— products 是补的:GET /products/scan/{sn}(扫码查询)
是整个车间最高频的读操作,不采它等于没采"活跃度"。
_should_audit 的 GET 分支改为三级判断:
敏感读(export/download/print) → 受跟踪前缀且非裸列表 → 否则不采
用 _is_bare_list 跳过「拉整个列表」:
· 列表接口被前端高频轮询(消息、任务列表尤其明显),逐条留痕会让
audit_logs 迅速膨胀,真正有价值的操作反而被淹没;
· 只有「查看详情」(/tasks/{id}) 才代表用户真的点开了某条业务数据。
判定用「去掉末尾斜杠后是否恰好等于某前缀」,天然排除查询串。
【修正文案】
- "read": "查询" → "查看详情"(前者易被误解成"随便搜了一下")
- 新增 "mark_read": "标为已读",并在 _SEGMENT_ACTION 补 "read" 映射 ——
否则 PUT /notifications/{id}/read 会回退到 _METHOD_ACTION(PUT→update),
把"点开一条通知"记成"修改了某样东西"
- "refresh": "刷新令牌" → "上线"(token 2 小时一换,业务上视作一次上线)
⚠️ 两点须知:
1. notifications / orders 目录下【只有列表路由】,而列表按规则不采集,
故这两个模块不会产生查看记录 —— 后端没有"查看单条消息"的接口。
且消息列表被前端轮询,采集它反而会造成日志爆炸,跳过是正确的。
2. products 被纳入后,工人每扫一次码就多一条记录(50 人 × 每天数百次
≈ 上万条/天)。若嫌吵,去掉该前缀一行即可。
注:读接口鉴权已在上一提交补齐,故这些"查看详情"记录能正确挂上操作人。
459 lines
18 KiB
Python
459 lines
18 KiB
Python
"""审计服务 — 写入与检索
|
||
|
||
写入方案的取舍(与 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": "删除",
|
||
# 只用于被采集的 GET(核心业务详情 / 敏感读)。
|
||
# 叫「查看详情」而不是「查询」:前者说明用户确实点开了某条业务数据,
|
||
# 后者容易被误解成"随便搜了一下"。
|
||
"read": "查看详情",
|
||
"export": "导出",
|
||
"login": "登录",
|
||
"logout": "登出",
|
||
# 刷新令牌 = 用户重新开始使用系统(token 2 小时一换,7 天免登录),
|
||
# 业务上视作一次「上线」,比"刷新令牌"这种技术词更贴近车间口径
|
||
"refresh": "上线",
|
||
"print": "打印",
|
||
"upload": "上传",
|
||
"finalize": "收口",
|
||
"receive": "接收",
|
||
"transfer": "转交",
|
||
"reject": "驳回",
|
||
"recall": "撤回",
|
||
"spawn": "派发",
|
||
"end": "结束分支",
|
||
"complete": "完结",
|
||
"mark_read": "标为已读",
|
||
}
|
||
|
||
|
||
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_at,first_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
|