diff --git a/pipeline_service/__init__.py b/pipeline_service/__init__.py index 82c7484..5fa34e4 100644 --- a/pipeline_service/__init__.py +++ b/pipeline_service/__init__.py @@ -54,4 +54,13 @@ 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" diff --git a/pipeline_service/audit_service.py b/pipeline_service/audit_service.py new file mode 100644 index 0000000..a12fc03 --- /dev/null +++ b/pipeline_service/audit_service.py @@ -0,0 +1,157 @@ +""" +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} diff --git a/pipeline_service/init.py b/pipeline_service/init.py index 0320c51..4a00199 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -660,6 +660,21 @@ 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", {})