diff --git a/pipeline_service/__init__.py b/pipeline_service/__init__.py index 5fa34e4..82c7484 100644 --- a/pipeline_service/__init__.py +++ b/pipeline_service/__init__.py @@ -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" diff --git a/pipeline_service/audit_service.py b/pipeline_service/audit_service.py deleted file mode 100644 index 53ac3d5..0000000 --- a/pipeline_service/audit_service.py +++ /dev/null @@ -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} diff --git a/pipeline_service/init.py b/pipeline_service/init.py index 4a00199..0320c51 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -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", {})