feat(audit): 级联操作聚合 —— 一次业务点击只产生一条审计
问题:BOM 归档/启停这类接口一次请求改动整组数据,ORM 监听器逐行写审计,
审计页被刷屏。**实测问题范围远大于 BOM**(同一秒、同接口、同操作人的条数):
/api/v1/inbound/stock/draft/start-new 924 条
/api/v1/inbound/stock/stocktake/generate-missing 898 条
/api/v1/permissions/assign 381 条
/api/v1/outbound 110 条
/api/v1/bom/save 83 条
/api/v1/bom/archive 58 条
全库共 311 次「单请求 ≥20 条」的突发,累计 4 万余行。
实现:URL 白名单 + 同事务合并(方案 A 的收敛版)
· AGGREGATE_PATH_MARKERS 按**路径段前缀**匹配(不是子串):
/api/v1/bom_draft_x 不会命中 bom,/api/v1/somebom/thing 也不会。
· 命中白名单的变更在 flush 期间累加到 session.info,由 Session
after_flush 钩子合并成**一条** AuditLog。
· details = 各行的**深度公共子集** + aggregate{count, targets}。
★ 为什么用白名单而不是全局默认合并:合并会改变审计的**语义粒度**
(日报里「修改 354 条」可能变成几十条),全局改会让业务方以为日志坏了。
白名单的失败模式也更安全 —— 新接口忘了登记只是"仍然刷屏",
而不会把两个不相干的业务动作错误合并(错误合并 = 把 A 的改动记到 B 头上,
是审计里最危险的一类错)。
★ 只取公共子集,不用某一行的值代表整批 —— 那是编造。BOM 的级联批次各行
改动完全一致(实测 /bom/archive 400 条只有 3 种签名、349 条同一个),
故合并**无损**;各行不一致时 changes 只留共同字段,其余由 targets 交代
"动过哪些对象"。
★★ 钩子必须挂 after_flush,不能挂 before_commit —— 踩过的坑:
Session.commit() 的顺序是 before_commit → flush → after_flush → COMMIT,
而累加发生在 flush **期间**。挂 before_commit 时钩子跑在累加之前,
缓冲还是空的;等 flush 填满后没人再写 —— 结果是**审计整批丢失**
(实测归档请求产出 0 条日志,业务却已提交,正是最危险的"改了但没记录")。
首版就是这个错,靠真实请求打 /bom/archive 数日志条数才抓出来。
★ 为什么不用 after_request/teardown:那跑在业务事务之外,业务回滚也会留下
一条"成功"的审计 —— 假账。after_flush 仍属同一事务,聚合日志与业务改动
同生共死(实测回滚后不留日志)。
消费端:
· changes_summary 优先识别 aggregate,读作
「批量更新 58 条;是否启用:是→否;是否归档:否→是(来源:/bom/archive)」
· SNAPSHOT_VIEWS 增加 aggregate,抽屉据此渲染「变更次数 + 受影响对象表」
· 前端把「变更对比」与「快照/汇总」从二选一改为各自独立渲染 ——
聚合日志两者都有,原先的 v-else 会让汇总块显示不出来
验证(26 + 18 项断言全过):
· 真实请求 POST /api/v1/bom/archive 打一个 58 行的 BOM:
归档前 0 条 → 归档后**恰好 1 条**,count=58、targets 58 条、
共同变更无损、业务改动同时生效(58/58 行已归档)
· 回滚后不留下聚合日志(59996 → 59996)
· 路径匹配 15 个边界(含 bom_draft_x / somebom 两个反向用例)
· 非白名单接口行为不变;日报附件 490.5K/193.4K/92.3K 与 6 列结构未变;
导出口径与列表一致;权限与凭据过滤未松动;字典 181 键零缺失
· vue-tsc 与 vite build exit=0
⚠️ 仅对**改动之后**的请求生效,历史 4 万行突发数据不变。
⚠️ 验证过程在库里留下 4 条真实审计记录(id 60278/60279 等,均是本人对
SF-9000-9 V2.2 的归档/取消归档操作,业务数据已精确还原为原状)。
审计记录未删除 —— 删审计要单独决策。
This commit is contained in:
@ -115,6 +115,69 @@ IGNORE_FIELDS = {
|
||||
IGNORE_FIELD_KEYWORDS = ('embedding', 'image_data', 'photo_blob')
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# 级联操作聚合
|
||||
# =============================================================================
|
||||
# 「一次业务点击 → 一条审计记录」。
|
||||
#
|
||||
# 背景:BOM 归档/启停、权限分配这类接口一次请求会改动整组数据,ORM 监听器
|
||||
# 遂逐行写审计。实测的突发规模(同一秒、同接口、同操作人):
|
||||
#
|
||||
# /api/v1/inbound/stock/draft/start-new 924 条
|
||||
# /api/v1/inbound/stock/stocktake/generate-missing 898 条
|
||||
# /api/v1/permissions/assign 381 条
|
||||
# /api/v1/outbound 110 条
|
||||
# /api/v1/bom/save 83 条
|
||||
# /api/v1/bom/archive 58 条
|
||||
#
|
||||
# 全库共 311 次「单请求 ≥20 条」的突发。逐行审计在本该表达"一次业务动作"的
|
||||
# 地方产生了成百条重复日志,审计页被刷屏、日报也被灌水。
|
||||
#
|
||||
# ★ 为什么用**白名单**而不是全局默认合并:
|
||||
# 合并会改变审计的**语义粒度**(日报里的"修改 354 条"可能变成几十条),
|
||||
# 全局改会让业务方以为日志坏了。白名单的失败模式也更安全 ——
|
||||
# 新接口忘了登记只是"仍然刷屏",而不会把两个不相干的业务动作错误合并
|
||||
# (错误合并 = 把 A 的改动记到 B 头上,是审计里最危险的一类错)。
|
||||
#
|
||||
# ★ 匹配规则:按**路径段前缀**,不是子串。`/api/v1/bom/SF-3500-9` 命中
|
||||
# `bom`,而 `/api/v1/bom_draft_x` 不会(段名不等)。版本前缀
|
||||
# (/api、/api/v1)先剥掉,兼容 /api/... 的老路由。
|
||||
AGGREGATE_PATH_MARKERS = (
|
||||
'bom', # save/archive/status/draft/* 与 DELETE /bom/<no>
|
||||
'permissions/assign',
|
||||
'inbound/stock/draft/start-new',
|
||||
'inbound/stock/stocktake/generate-missing',
|
||||
'outbound',
|
||||
'import/execute',
|
||||
)
|
||||
|
||||
# 聚合日志里保留的受影响对象清单上限。
|
||||
# ★ 超出时**显式标注** targets_truncated 与真实总数,不静默丢弃 ——
|
||||
# target_id 逐个都进清单的话,924 条的突发会生成几十 KB 的 details,
|
||||
# 而列表接口一页 50 条会把它翻 50 倍带出去。
|
||||
AGGREGATE_TARGETS_LIMIT = 200
|
||||
|
||||
# 累加器挂在 session.info 上。
|
||||
# ★ 不能用 connection.info:连接来自池子,会跨请求残留 ——
|
||||
# 那会把上一个请求的变更算进下一个请求的日志里。
|
||||
AGGREGATE_INFO_KEY = '_audit_aggregate'
|
||||
|
||||
|
||||
def _is_aggregate_path(url):
|
||||
"""该请求路径是否命中级联聚合白名单"""
|
||||
segments = [s for s in str(url or '').split('?')[0].split('/') if s]
|
||||
# 剥掉 api 与版本前缀(/api/v1/bom/x → bom/x;/api/bom/x → bom/x)
|
||||
if segments and segments[0].lower() == 'api':
|
||||
segments = segments[1:]
|
||||
if segments and len(segments[0]) > 1 and segments[0][0] in 'vV' and segments[0][1:].isdigit():
|
||||
segments = segments[1:]
|
||||
for marker in AGGREGATE_PATH_MARKERS:
|
||||
want = marker.split('/')
|
||||
if segments[:len(want)] == want:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# 序列化
|
||||
# =============================================================================
|
||||
@ -300,6 +363,150 @@ def _get_request_user_info():
|
||||
# 核心:写入审计日志
|
||||
# =============================================================================
|
||||
|
||||
_INSERT_AUDIT_SQL = text("""
|
||||
INSERT INTO audit_logs
|
||||
(user_id, username, display_name, action, module,
|
||||
target_id, target_name, details, ip_address, method, url, created_at)
|
||||
VALUES
|
||||
(:user_id, :username, :display_name, :action, :module,
|
||||
:target_id, :target_name, cast(:details AS jsonb), :ip_address,
|
||||
:method, :url, :created_at)
|
||||
""")
|
||||
|
||||
_MISSING = object()
|
||||
|
||||
|
||||
def _common_part(a, b):
|
||||
"""
|
||||
两个 details 的**公共部分**:键都在、且值相等;dict 递归。
|
||||
|
||||
用于把一批变更合并成"所有行都成立"的那部分。
|
||||
|
||||
★ 只取公共部分,绝不用某一行的值代表整批 —— 那是**编造**:
|
||||
把 A 行的改动说成整批的改动,读者会照着它去追一个错误的结论。
|
||||
不公共的部分由 aggregate.targets 交代"动过哪些对象"。
|
||||
"""
|
||||
if isinstance(a, dict) and isinstance(b, dict):
|
||||
out = {}
|
||||
for key in a.keys() & b.keys():
|
||||
v = _common_part(a[key], b[key])
|
||||
if v is not _MISSING:
|
||||
out[key] = v
|
||||
# 空 dict 是"没有任何公共字段",留着会渲染成空的变更列表
|
||||
return out if out else _MISSING
|
||||
if a == b:
|
||||
return a
|
||||
return _MISSING
|
||||
|
||||
|
||||
def _accumulate_aggregate(target, action, module, target_id, target_name,
|
||||
details):
|
||||
"""
|
||||
把这一次变更累加到当前 session 的聚合缓冲里(不写库)。
|
||||
|
||||
返回 True 表示已接管(调用方不要再逐行 INSERT)。
|
||||
|
||||
★ 缓冲挂在 session.info:session 是请求级的,随请求销毁;
|
||||
连接是池化的,挂在 connection.info 上会跨请求残留。
|
||||
"""
|
||||
try:
|
||||
from sqlalchemy.orm import object_session
|
||||
session = object_session(target)
|
||||
if session is None:
|
||||
return False
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
user = _get_request_user_info()
|
||||
key = (module, action, user.get('url') or '')
|
||||
|
||||
buf = session.info.setdefault(AGGREGATE_INFO_KEY, {})
|
||||
slot = buf.get(key)
|
||||
if slot is None:
|
||||
slot = buf[key] = {
|
||||
'module': module,
|
||||
'action': action,
|
||||
'user': user,
|
||||
'count': 0,
|
||||
'targets': [],
|
||||
'seen': set(),
|
||||
'first_id': target_id,
|
||||
'first_name': target_name,
|
||||
'common': None, # None = 还没有任何一行参与比较
|
||||
}
|
||||
slot['count'] += 1
|
||||
if target_id not in slot['seen']:
|
||||
slot['seen'].add(target_id)
|
||||
slot['targets'].append({'id': target_id, 'name': target_name})
|
||||
slot['common'] = (details if slot['common'] is None
|
||||
else _common_part(slot['common'], details))
|
||||
return True
|
||||
|
||||
|
||||
def _flush_aggregated_audit(session, flush_context=None):
|
||||
"""
|
||||
Session after_flush 钩子:把聚合缓冲写成**一条** AuditLog。
|
||||
|
||||
★★ 钩子必须挂在 after_flush,不能挂 before_commit —— 这是个踩过的坑:
|
||||
Session.commit() 的顺序是「before_commit → flush → after_flush → COMMIT」,
|
||||
而累加**发生在 flush 期间**(ORM 的 before_update 回调)。挂在
|
||||
before_commit 上,钩子跑的时候缓冲还是空的,等 flush 把它填满后
|
||||
再没人来写 —— 结果是**审计整批丢失**(实测归档请求产出 0 条日志,
|
||||
业务却已提交,正是最危险的"改了但没记录")。
|
||||
|
||||
★ 为什么不用 after_request / teardown:它们跑在业务事务之外,
|
||||
业务回滚也会留下一条"成功"的审计 —— 那是假账。
|
||||
after_flush 仍在同一事务内,故聚合日志与业务改动**同生共死**。
|
||||
|
||||
★ 幂等:一次 commit 可能触发多次 flush(自动 flush)。取出即清空,
|
||||
后续 flush 若有新变更会再写一条 —— 那对应另一个原子单元,是对的。
|
||||
|
||||
★ 没有缓冲 = 不是白名单请求 / 无变更,直接返回(绝大多数请求走这里)。
|
||||
"""
|
||||
slots = session.info.pop(AGGREGATE_INFO_KEY, None)
|
||||
if not slots:
|
||||
return
|
||||
try:
|
||||
conn = session.connection()
|
||||
for slot in slots.values():
|
||||
common = slot['common'] if isinstance(slot['common'], dict) else {}
|
||||
targets_total = len(slot['targets'])
|
||||
targets = slot['targets'][:AGGREGATE_TARGETS_LIMIT]
|
||||
agg = {
|
||||
'count': slot['count'],
|
||||
'targets': targets,
|
||||
}
|
||||
# 只在「变更次数 ≠ 对象数」时才带 targets_total:多数批次两者相等
|
||||
# (一次变更一行),多一个恒等于 count 的字段只会让抽屉多一行噪声
|
||||
if targets_total != slot['count']:
|
||||
agg['targets_total'] = targets_total
|
||||
if targets_total > len(targets):
|
||||
agg['targets_truncated'] = True
|
||||
details = dict(common)
|
||||
details['aggregate'] = agg
|
||||
user = slot['user']
|
||||
conn.execute(_INSERT_AUDIT_SQL, {
|
||||
'user_id': user.get('user_id'),
|
||||
'username': user.get('username') or 'system',
|
||||
'display_name': user.get('display_name') or '',
|
||||
'action': slot['action'],
|
||||
'module': slot['module'],
|
||||
'target_id': slot['first_id'],
|
||||
'target_name': (slot['first_name'] or '')[:200],
|
||||
'details': json.dumps(details, cls=_AuditJSONEncoder),
|
||||
'ip_address': user.get('ip'),
|
||||
'method': user.get('method'),
|
||||
'url': user.get('url'),
|
||||
'created_at': beijing_time(),
|
||||
})
|
||||
except Exception as e:
|
||||
# 与逐行写入一致:审计失败不影响业务,但必须留下痕迹
|
||||
try:
|
||||
current_app.logger.error(f"Audit aggregate flush failed: {e}")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _create_audit_log(connection, mapper, target, action, details):
|
||||
"""
|
||||
使用事件回调传入的 connection 直接 INSERT。
|
||||
@ -320,6 +527,14 @@ def _create_audit_log(connection, mapper, target, action, details):
|
||||
user = _get_request_user_info()
|
||||
module = _get_module_name(mapper)
|
||||
|
||||
# ★ 级联聚合白名单:命中则只累加,不逐行写库。
|
||||
# 真正的写入发生在 Session before_commit(见 _flush_aggregated_audit),
|
||||
# 一次业务点击最终只留一条日志。
|
||||
if _is_aggregate_path(user.get('url')):
|
||||
if _accumulate_aggregate(target, action, module, target_id,
|
||||
target_name, details):
|
||||
return
|
||||
|
||||
sql = text("""
|
||||
INSERT INTO audit_logs
|
||||
(user_id, username, display_name, action, module,
|
||||
@ -457,10 +672,34 @@ def after_insert_listener(mapper, connection, target):
|
||||
_BOUND_TABLES = set()
|
||||
|
||||
|
||||
_AGGREGATE_HOOK_BOUND = False
|
||||
|
||||
|
||||
def _bind_aggregate_hook():
|
||||
"""
|
||||
注册 Session after_flush 钩子(幂等)。
|
||||
|
||||
★ 用 after_flush 而不是 before_commit / after_commit,理由见
|
||||
_flush_aggregated_audit 的说明(前者太早、缓冲还空着;后者已在
|
||||
业务事务之外)。
|
||||
"""
|
||||
global _AGGREGATE_HOOK_BOUND
|
||||
if _AGGREGATE_HOOK_BOUND:
|
||||
return
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
event.listen(Session, 'after_flush', _flush_aggregated_audit)
|
||||
_AGGREGATE_HOOK_BOUND = True
|
||||
|
||||
|
||||
def register_audit_listeners(db):
|
||||
"""
|
||||
按白名单向已完成映射的模型注册事件监听器,返回本次**新增**绑定的模型数。
|
||||
|
||||
★ 务必要连带注册级联聚合的 Session 钩子(见 _bind_aggregate_hook):
|
||||
白名单接口的变更只累加、不逐行写库,钩子缺席就永远不落地 ——
|
||||
那是**丢审计**,比刷屏严重得多。故绑在这里而不是靠 import 副作用。
|
||||
|
||||
★ 为什么需要"惰性补绑"(见 ensure_audit_listeners):
|
||||
本项目大量模型是在**函数体内延迟导入**的(31 处,例如
|
||||
app/api/v1/scrap.py 内部 `from app.models.scrap_approval import ScrapApproval`),
|
||||
@ -476,6 +715,7 @@ def register_audit_listeners(db):
|
||||
写日志必然抛 AttributeError 并被静默吞掉。
|
||||
现改为按 **表名** 从 db.metadata 取模型,不依赖 app.models 的导出。
|
||||
"""
|
||||
_bind_aggregate_hook()
|
||||
count = 0
|
||||
for tablename in list(WHITELIST_TABLES):
|
||||
if tablename in _BOUND_TABLES:
|
||||
@ -508,6 +748,7 @@ def ensure_audit_listeners():
|
||||
|
||||
调用开销:一次集合差集运算;已全部绑定后立即返回。
|
||||
"""
|
||||
_bind_aggregate_hook()
|
||||
if _BOUND_TABLES >= WHITELIST_TABLES:
|
||||
return 0
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user