""" pipeline_service/audit_service.py - 审计日志全局模块 独立设计的全局审计模块。审计日志 append-only,查看/备份/删除权限仅赋予 `owner.audit` 角色(审计独立性:owner.superuser 也不能看,防管理员删自己的记录)。 - audit_log: 各模块在关键操作时写入审计(系统自动,非用户触发) - is_audit_role: 判断用户是否 owner.audit(API 层权限校验用) - list_audit_logs / backup_audit_logs / delete_audit_logs: 仅 owner.audit 审计事件清单(action 取值): 认证:login / login_fail / logout 权限:role_change / perm_change / user_role_change 工作环境:work_env_set / org_key_gen / remote_bwrap 部署账号:account_create / account_remove / sandbox_run 用户机构:user_create / user_disable / user_delete / org_change 审计自身:audit_delete / audit_backup """ import json import logging logger = logging.getLogger("pipeline.audit") # 合法审计动作白名单 VALID_ACTIONS = { "login", "login_fail", "logout", "role_change", "perm_change", "user_role_change", "work_env_set", "org_key_gen", "remote_bwrap", "account_create", "account_remove", "sandbox_run", "user_create", "user_disable", "user_delete", "org_change", "audit_delete", "audit_backup", } async def audit_log(sor, user_id: str, username: str, action: str, target: str = "", detail: str = "", result: str = "ok", client_ip: str = "") -> bool: """写入一条审计日志(append-only,不 update)。 审计写入失败不应阻断主流程,返回 False 由调用方决定是否处理。 """ if action not in VALID_ACTIONS: action = "unknown" try: from appPublic.uniqueID import getID data = { "id": getID(), "user_id": user_id or "", "username": username or "", "action": action, "target": (target or "")[:200], "detail": (detail or "")[:2000], "result": result or "ok", "client_ip": (client_ip or "")[:64], } await sor.C("sd_audit_logs", data) return True except Exception as e: logger.exception("audit_log failed: %s", e) return False async def is_audit_role(sor, user_id: str) -> bool: """判断用户是否 owner.audit 角色。""" if not user_id: return False try: recs = await sor.sqlExe( "SELECT 1 FROM userrole WHERE userid=${u}$ AND roleid='owner.audit' LIMIT 1", {"u": user_id}) return bool(recs) except Exception as e: logger.exception("is_audit_role failed: %s", e) return False def _row_to_dict(row) -> dict: """DB 行 → dict(created_at 等 datetime 转字符串)。""" d = {} for k in ("id", "user_id", "username", "action", "target", "detail", "result", "client_ip", "created_at"): v = getattr(row, k, None) if hasattr(v, "isoformat"): v = v.isoformat() d[k] = v return d async def list_audit_logs(sor, filters: dict = None, limit: int = 100, offset: int = 0): """查询审计日志(仅 owner.audit)。""" filters = filters or {} where = [] ns = {} if filters.get("user_id"): where.append("user_id=${u}$") ns["u"] = filters["user_id"] if filters.get("username"): where.append("username=${un}$") ns["un"] = filters["username"] if filters.get("action"): where.append("action=${a}$") ns["a"] = filters["action"] if filters.get("from"): where.append("created_at>=${f}$") ns["f"] = filters["from"] if filters.get("to"): where.append("created_at<=${t}$") ns["t"] = filters["to"] wsql = (" WHERE " + " AND ".join(where)) if where else "" try: limit = max(1, min(int(limit or 100), 1000)) except (ValueError, TypeError): limit = 100 try: offset = max(0, int(offset or 0)) except (ValueError, TypeError): offset = 0 recs = await sor.sqlExe( "SELECT * FROM sd_audit_logs" + wsql + " ORDER BY created_at DESC LIMIT " + str(limit) + " OFFSET " + str(offset), ns) return [_row_to_dict(r) for r in recs] async def backup_audit_logs(sor, user_id: str, username: str, filters: dict = None, client_ip: str = "") -> dict: """导出审计日志(仅 owner.audit),返回 JSON 内容供下载。""" # 查最近 10000 条(备份上限,避免一次拉全表拖垮内存) recs = await sor.sqlExe( "SELECT * FROM sd_audit_logs ORDER BY created_at DESC LIMIT 10000", {}) rows = [_row_to_dict(r) for r in recs] # 备份操作本身也审计 await audit_log(sor, user_id, username, "audit_backup", target="sd_audit_logs", detail="备份 " + str(len(rows)) + " 条", result="ok", client_ip=client_ip) return {"ok": True, "count": len(rows), "data": rows} async def delete_audit_logs(sor, user_id: str, username: str, before: str = "", client_ip: str = "") -> dict: """删除审计日志(仅 owner.audit)。before 为时间戳,删除早于该时间的日志。 删除操作本身也会审计(audit_delete 留痕),且 audit_delete 记录**永不删除**, 保证删除操作的证据永久保留(防自删)。 """ ns = {} # 排除 audit_delete 记录,保证删除留痕不被删 if before: wsql = " WHERE created_at<${b}$ AND action != 'audit_delete'" ns["b"] = before else: wsql = " WHERE action != 'audit_delete'" # 先记录删除操作(在删除前写,确保删除动作本身留痕) await audit_log(sor, user_id, username, "audit_delete", target="sd_audit_logs", detail="删除 before=" + (before or "全部"), result="ok", client_ip=client_ip) recs = await sor.sqlExe("SELECT COUNT(*) AS c FROM sd_audit_logs" + wsql, ns) count = getattr(recs[0], "c", 0) if recs else 0 await sor.sqlExe("DELETE FROM sd_audit_logs" + wsql, ns) return {"ok": True, "deleted": count}