ymq eaaf5fad26 feat: 审计日志全局模块(append-only,查看/备份/删除仅 owner.audit)
- audit_service.py: audit_log 写入 + is_audit_role + list/backup/delete
- sd_audit_logs 表 + owner.audit 角色 DDL
- 审计独立性:owner.superuser 也无权查看审计日志
2026-08-14 12:10:57 +08:00

158 lines
5.9 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.

"""
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 为时间戳,删除早于该时间的日志。
删除操作本身也会审计(防自删留痕)。
"""
ns = {}
wsql = ""
if before:
wsql = " WHERE created_at<${b}$"
ns["b"] = before
# 先记录删除操作(在删除前写,确保删除动作本身留痕)
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}