refactor: 审计逻辑迁移到独立模块 app_audit
- 移除 pipeline_service/audit_service.py(迁移到 app_audit) - __init__.py 移除 audit import,init.py 移除 sd_audit_logs/owner.audit DDL
This commit is contained in:
parent
c7bb13c17d
commit
92ee88d1f4
@ -54,13 +54,4 @@ from .work_env import (
|
||||
get_org_pubkey,
|
||||
)
|
||||
|
||||
# 审计日志全局模块(append-only,查看/备份/删除仅 owner.audit)
|
||||
from .audit_service import (
|
||||
audit_log,
|
||||
is_audit_role,
|
||||
list_audit_logs,
|
||||
backup_audit_logs,
|
||||
delete_audit_logs,
|
||||
)
|
||||
|
||||
__version__ = "3.5.0"
|
||||
|
||||
@ -1,160 +0,0 @@
|
||||
"""
|
||||
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}
|
||||
@ -660,21 +660,6 @@ def load_pipeline_service():
|
||||
"PRIMARY KEY (id), UNIQUE KEY uk_owner (owner_type, owner_id)"
|
||||
") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4", {})
|
||||
_debug("sd_work_envs table ready")
|
||||
# 审计日志表(append-only,全局)
|
||||
await sor.sqlExe(
|
||||
"CREATE TABLE IF NOT EXISTS sd_audit_logs ("
|
||||
"id varchar(32) NOT NULL, user_id varchar(32), username varchar(100),"
|
||||
"action varchar(50) NOT NULL, target varchar(200), detail text,"
|
||||
"result varchar(10), client_ip varchar(64),"
|
||||
"created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,"
|
||||
"PRIMARY KEY (id), KEY idx_user (user_id), KEY idx_action (action),"
|
||||
"KEY idx_created (created_at)"
|
||||
") ENGINE=InnoDB DEFAULT CHARSET=utf8mb4", {})
|
||||
_debug("sd_audit_logs table ready")
|
||||
# owner.audit 角色(审计独立角色)
|
||||
await sor.sqlExe(
|
||||
"INSERT IGNORE INTO role (id, orgtypeid, name) VALUES ('owner.audit', 'owner', 'audit')", {})
|
||||
_debug("owner.audit role ready")
|
||||
# agent_config columns
|
||||
try:
|
||||
await sor.sqlExe("ALTER TABLE pipelines ADD COLUMN agent_config text", {})
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user