From 0b87a72a416ca0a633fe9fd8103be062a0c68f30 Mon Sep 17 00:00:00 2001 From: yueli Date: Mon, 28 Sep 2026 11:07:10 +0800 Subject: [PATCH] =?UTF-8?q?feat(audit):=20=E7=BA=A7=E8=81=94=E6=93=8D?= =?UTF-8?q?=E4=BD=9C=E8=81=9A=E5=90=88=20=E2=80=94=E2=80=94=20=E4=B8=80?= =?UTF-8?q?=E6=AC=A1=E4=B8=9A=E5=8A=A1=E7=82=B9=E5=87=BB=E5=8F=AA=E4=BA=A7?= =?UTF-8?q?=E7=94=9F=E4=B8=80=E6=9D=A1=E5=AE=A1=E8=AE=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题: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 的归档/取消归档操作,业务数据已精确还原为原状)。 审计记录未删除 —— 删审计要单独决策。 --- inventory-backend/app/core/audit_listener.py | 241 ++++++++++++++++++ .../app/services/audit_export_service.py | 39 +++ inventory-backend/app/utils/audit_labels.py | 18 ++ inventory-web/src/views/system/AuditLog.vue | 6 +- 4 files changed, 302 insertions(+), 2 deletions(-) diff --git a/inventory-backend/app/core/audit_listener.py b/inventory-backend/app/core/audit_listener.py index 498e904..be5edae 100644 --- a/inventory-backend/app/core/audit_listener.py +++ b/inventory-backend/app/core/audit_listener.py @@ -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/ + '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: diff --git a/inventory-backend/app/services/audit_export_service.py b/inventory-backend/app/services/audit_export_service.py index d609bb2..3b0b45d 100644 --- a/inventory-backend/app/services/audit_export_service.py +++ b/inventory-backend/app/services/audit_export_service.py @@ -38,6 +38,7 @@ from app.utils.audit_labels import ( BOOLEAN_FIELDS, CHANGES_KEY, PERSON_NAME_FIELDS, + SNAPSHOT_AGGREGATE, SNAPSHOT_CREATED, SNAPSHOT_DELETED, USER_ID_FIELDS, @@ -1057,6 +1058,37 @@ def _child_text(child): f"#{child['id']}" if child.get('id') else '') +def _aggregate_summary(row, agg, changes, limit): + """ + 级联聚合日志的摘要,形如: + 批量更新 58 条:是否启用:是→否;是否归档:否→是(来源:/bom/archive) + + ★ 聚合行必须走这条分支,不能落到下面按快照/字段名罗列的逻辑: + 它的 details 是 {共同变更 + aggregate},按字段名罗列会得到 + 「新增(BOM管理):count、targets」这种毫无意义的串。 + + changes: 该行的共同变更(各行都成立的字段)。**可能为空** —— + 各行改动不一致时不敢用某一行的值代表整批,那属于编造。 + """ + verb = {'CREATE': '批量新增', 'UPDATE': '批量更新', + 'DELETE': '批量删除'}.get(canon_action(row.action), '批量操作') + parts = [f"{verb} {agg.get('count', 0)} 条"] + if changes: + parts.append(';'.join( + f"{lb}:{fmt_value(o)}→{fmt_value(n)}" for lb, o, n in changes)) + # 受影响对象数与变更次数是两回事(前者去重):336 次变更可能只涉及 21 个对象 + total = agg.get('targets_total') + if total and total != agg.get('count'): + parts.append(f"涉及 {total} 个对象") + if agg.get('targets_truncated'): + parts.append(f"对象清单仅列前 {len(agg.get('targets') or [])} 个") + + body = truncate(':'.join(parts[:1]) + (';' + ';'.join(parts[1:]) + if len(parts) > 1 else ''), limit) + source = url_source(row.url) + return f"{body}(来源:/{source})" if source else body + + def changes_summary(row, ref_maps, limit=200, resolved=None, material=None, child=None): """ @@ -1069,6 +1101,13 @@ def changes_summary(row, ref_maps, limit=200, resolved=None, material=None, 台账摘要和明细长表,调用方预计算一次传进来即可 —— 否则每行 要跑两遍 resolve_value(4 万行的导出上是可感知的浪费)。 """ + # ★ 级联聚合行:优先于下面所有分支(它的 details 形状与快照/变更都不同) + details = row.details if isinstance(row.details, dict) else {} + agg = details.get(SNAPSHOT_AGGREGATE) + if isinstance(agg, dict): + common = resolved_changes(row.module, changes_of(row), ref_maps) + return _aggregate_summary(row, agg, common, limit) + if canon_action(row.action) == 'UPDATE': if resolved is None: resolved = resolved_changes(row.module, changes_of(row), ref_maps) diff --git a/inventory-backend/app/utils/audit_labels.py b/inventory-backend/app/utils/audit_labels.py index 9fffc23..79ba9d9 100644 --- a/inventory-backend/app/utils/audit_labels.py +++ b/inventory-backend/app/utils/audit_labels.py @@ -333,6 +333,18 @@ FIELD_LABELS = { 'product_image': '产品图片', 'product_image_remark': '产品图片备注', 'purchase_link': '采购链接', + + # --- 级联聚合日志的汇总字段 --- + # ★ 这些是 audit_listener 写聚合日志时自造的元字段(见其 AGGREGATE_PATH_MARKERS), + # 不是任何业务表上的列。名字都很通用(count/targets), + # 目前业务快照里没有同名键;将来若出现,要改成按 details 结构区分而非按字段名。 + 'count': '变更次数', + 'targets': '受影响对象', + 'targets_total': '对象总数', + 'targets_truncated': '对象清单已截断', + # 早期版本的聚合载荷里带过 action(与日志自身的 action 重复,已不再写入), + # 但历史行里还在,留个标签免得它在详情里裸露英文 + 'action': '操作类型', } # ============================================================================= @@ -398,10 +410,16 @@ SNAPSHOT_DELETED = 'deleted_snapshot' SNAPSHOT_PAYLOAD = 'payload' CHANGES_KEY = 'changes' +# 聚合键:级联接口的批量变更被合并成一条日志时,受影响对象清单存在这里 +# (见 audit_listener 的级联聚合)。它不是"快照",但走同一套分层渲染, +# 故一并列在 VIEWS 里 —— 抽屉里会把它渲染成一张对象清单表。 +SNAPSHOT_AGGREGATE = 'aggregate' + SNAPSHOT_VIEWS = ( (SNAPSHOT_CREATED, '新增数据快照'), (SNAPSHOT_DELETED, '删除前数据快照'), (SNAPSHOT_PAYLOAD, '业务数据'), + (SNAPSHOT_AGGREGATE, '批量汇总'), ) diff --git a/inventory-web/src/views/system/AuditLog.vue b/inventory-web/src/views/system/AuditLog.vue index 0257df1..1fc35d2 100644 --- a/inventory-web/src/views/system/AuditLog.vue +++ b/inventory-web/src/views/system/AuditLog.vue @@ -264,8 +264,10 @@ - -