From 04c3c1412ce4d89201762e93773d1cf138d8b0cf Mon Sep 17 00:00:00 2001 From: wangjiming Date: Tue, 18 Aug 2026 14:49:12 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9=E6=99=AE=E9=80=9A=E7=94=A8?= =?UTF-8?q?=E6=88=B7=E7=9A=84=E6=95=B0=E6=8D=AE=E7=B1=BB=E5=9E=8B=E8=BD=AC?= =?UTF-8?q?=E6=8D=A2=E5=9C=A8=E6=95=B0=E6=8D=AE=E9=9B=86=E7=AE=A1=E7=90=86?= =?UTF-8?q?=E7=9C=8B=E4=B8=8D=E8=A7=81=E7=9A=84=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/app/api/v1/endpoints/platform.py | 89 ++-- backend/app/core/auth.py | 35 +- backend/app/core/logging.py | 46 +- backend/app/core/op_log.py | 382 +++++++++++++++ backend/app/db/platform_store.py | 21 +- backend/app/db/sql/002_governance.sql | 35 ++ backend/app/modules/data_convert/router.py | 7 + backend/app/modules/system/router.py | 140 ++++++ frontend/src/api/modules/operation-log.ts | 58 +++ frontend/src/components/AppSidebar.vue | 1 + frontend/src/router/index.ts | 7 + frontend/src/views/audit/OperationLogView.vue | 437 ++++++++++++++++++ 12 files changed, 1221 insertions(+), 37 deletions(-) create mode 100644 backend/app/core/op_log.py create mode 100644 frontend/src/api/modules/operation-log.ts create mode 100644 frontend/src/views/audit/OperationLogView.vue diff --git a/backend/app/api/v1/endpoints/platform.py b/backend/app/api/v1/endpoints/platform.py index 8fdfb09..7ead5c1 100644 --- a/backend/app/api/v1/endpoints/platform.py +++ b/backend/app/api/v1/endpoints/platform.py @@ -17,6 +17,7 @@ import httpx from app.core.auth import filter_accessible_resource_ids, filter_accessible_resource_ids_batch, get_current_user, has_resource_access, is_admin from app.core.config import get_settings from app.core.audit import audit_log, AuditActions +from app.core.op_log import op_log, OpModule, OpAction from app.db.platform_store import get_platform_store from app.modules.compute_gateway.client import ComputeNodeClient from app.modules.compute_gateway.sync import fetch_eval_result_content, poll_compute_jobs_once @@ -904,12 +905,11 @@ async def create_model(payload: dict[str, Any] = Body(...), current_user: dict = @router.get("/model-manage/{model_id}") async def model_detail(model_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: + # 基座模型(配置模型)是平台共享资源,所有登录用户均可查看 try: model = get_platform_store().model(model_id) except KeyError: raise fail(404, "model not found") - if not has_resource_access("model", model_id, current_user, "read"): - raise fail(403, "no permission to access this model") return ok(model) @@ -919,7 +919,10 @@ async def model_detail(model_id: str, current_user: dict = Depends(get_current_u target_type="model", detail_template="更新模型: {model_id}", ) -async def update_model(model_id: str, payload: dict[str, Any] = Body(...)) -> dict[str, Any]: +async def update_model(model_id: str, payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: + # 基座模型(配置模型)只有管理员可以编辑 + if not is_admin(current_user): + raise fail(403, "只有管理员可以修改模型配置") try: return ok(get_platform_store().update_model(model_id, payload)) except KeyError: @@ -927,7 +930,10 @@ async def update_model(model_id: str, payload: dict[str, Any] = Body(...)) -> di @router.put("/model-manage/{model_id}/purpose") -async def update_model_purpose(model_id: str, payload: dict[str, Any] = Body(...)) -> dict[str, Any]: +async def update_model_purpose(model_id: str, payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: + # 基座模型(配置模型)只有管理员可以修改用途 + if not is_admin(current_user): + raise fail(403, "只有管理员可以修改模型用途") try: return ok(get_platform_store().update_model(model_id, {"purpose": payload.get("purpose", "training")})) except KeyError: @@ -936,16 +942,15 @@ async def update_model_purpose(model_id: str, payload: dict[str, Any] = Body(... @router.delete("/model-manage/{model_id}") async def delete_model(model_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: - if not has_resource_access("model", model_id, current_user, "delete"): - raise fail(403, "no permission to delete this model") - pending = _require_approval_or_admin("model", model_id, current_user, f"删除模型 {model_id}") - if pending: - return pending + # 基座模型(配置模型)只有管理员可以删除 + if not is_admin(current_user): + raise fail(403, "只有管理员可以删除模型配置") get_platform_store().delete_model(model_id) return ok({"deleted": model_id}) @router.post("/model-manage/merge") +@op_log(module=OpModule.MODEL_MANAGE, action=OpAction.MERGE, target_type="trained_model", target_name_param="trained_model_id") async def merge_model(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: store = get_platform_store() trained_model_id = str(payload.get("trained_model_id") or payload.get("model_id") or payload.get("model_name") or "") @@ -961,8 +966,7 @@ async def merge_model(payload: dict[str, Any] = Body(...), current_user: dict = raise fail(404, "trained model not found") if not has_resource_access("trained_model", trained_model["id"], current_user, "execute"): raise fail(403, "no permission to merge this trained model") - if payload.get("base_model_id") and not has_resource_access("model", str(payload["base_model_id"]), current_user, "execute"): - raise fail(403, "no permission to use merge base model") + # 基座模型(配置模型)是平台共享资源,不需要 ACL 授权即可使用 base_model_path = payload.get("base_model_path") or (trained_model and trained_model.get("base_model_path")) adapter_path = ( payload.get("adapter_path") @@ -1282,6 +1286,7 @@ async def dataset_list(current_user: dict = Depends(get_current_user)) -> dict[s target_type="dataset", detail_template="创建数据集: {name}", ) +@op_log(module=OpModule.DATASET, action=OpAction.CREATE, target_type="dataset", target_name_param="name") async def create_dataset(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: payload.setdefault("created_by", current_user.get("id")) dataset = get_platform_store().create_dataset(payload) @@ -1318,6 +1323,7 @@ async def update_dataset(dataset_id: str, payload: dict[str, Any] = Body(...)) - target_type="dataset", detail_template="删除数据集: {dataset_id}", ) +@op_log(module=OpModule.DATASET, action=OpAction.DELETE, target_type="dataset", target_name_param="dataset_id") async def delete_dataset(dataset_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: if not has_resource_access("dataset", dataset_id, current_user, "delete"): raise fail(403, "no permission to delete this dataset") @@ -1367,8 +1373,7 @@ async def create_fine_tune(payload: dict[str, Any] = Body(...), current_user: di if not is_admin(current_user): model_id = str(payload.get("base_model") or payload.get("base_model_id") or "") dataset_id = str(payload.get("train_dataset_id") or "") - if model_id and not has_resource_access("model", model_id, current_user, "execute"): - raise fail(403, "no permission to use this base model") + # 基座模型(配置模型)是平台共享资源,不需要 ACL 授权 if dataset_id and not has_resource_access("dataset", dataset_id, current_user, "execute"): raise fail(403, "no permission to use this dataset") try: @@ -1379,6 +1384,7 @@ async def create_fine_tune(payload: dict[str, Any] = Body(...), current_user: di @router.post("/fine-tune/start") +@op_log(module=OpModule.FINE_TUNE, action=OpAction.START, target_type="fine_tune", target_name_param="name", detail_params=["task_id", "base_model", "train_dataset_id"]) async def start_fine_tune( payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user), @@ -1538,6 +1544,7 @@ async def update_fine_tune(task_id: str, payload: dict[str, Any] = Body(...), cu @router.post("/fine-tune/stop/{task_id}") +@op_log(module=OpModule.FINE_TUNE, action=OpAction.STOP, target_type="fine_tune", target_name_param="task_id") async def stop_fine_tune(task_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: store = get_platform_store() try: @@ -1584,6 +1591,7 @@ async def retry_fine_tune(task_id: str, payload: dict[str, Any] | None = Body(de @router.delete("/fine-tune/{task_id}") +@op_log(module=OpModule.FINE_TUNE, action=OpAction.DELETE, target_type="fine_tune", target_name_param="task_id") async def delete_fine_tune(task_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: if not has_resource_access("fine-tune", task_id, current_user, "delete"): raise fail(403, "no permission to delete this task") @@ -1665,6 +1673,7 @@ async def model_eval_detail(task_id: str, current_user: dict = Depends(get_curre @router.post("/model-eval/start") +@op_log(module=OpModule.MODEL_EVAL, action=OpAction.START, target_type="eval_task", target_name_param="name", detail_params=["model_id", "dataset_id"]) async def model_eval_start(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: """Start an evaluation task: submit eval job to compute node.""" store = get_platform_store() @@ -1878,11 +1887,7 @@ async def model_compare_list(current_user: dict = Depends(get_current_user)) -> async def model_compare_create(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: payload.setdefault("created_by", current_user.get("id")) model_ids = payload.get("model_ids") or payload.get("models") or [] - if not is_admin(current_user): - for model_id in model_ids: - if isinstance(model_id, dict): model_id = model_id.get("id") or model_id.get("model_id") - if model_id and not has_resource_access("model", str(model_id), current_user, "execute"): - raise fail(403, "no permission to use inference model") + # 基座模型(配置模型)是平台共享资源,不需要 ACL 授权即可用于推理 task = get_platform_store().create_compare_task(payload) return ok({"id": task["id"]}) @@ -1940,17 +1945,35 @@ async def _unload_from_compute_node(store: Any, task: dict[str, Any] | None = No @router.delete("/model-compare/{task_id}") +@op_log(module=OpModule.INFERENCE, action=OpAction.DELETE, target_type="inference", target_name_param="task_id") async def model_compare_delete(task_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: # 先删记录(快),再 best-effort 释放算力节点上的模型——删除绝不被卸载阻塞 try: task = get_platform_store().compare_task(task_id) except KeyError: raise fail(404, "compare task not found") - if not has_resource_access("compare", task_id, current_user, "delete"): - raise fail(403, "no permission to delete inference task") - pending = _require_approval_or_admin("compare", task_id, current_user, f"删除推理任务 {task_id}") - if pending: - return pending + # 先判断是否为任务创建者本人:如果是,直接允许删除(不需要 ACL 授权也不需要审批) + is_owner = False + payload_obj = task + if isinstance(payload_obj, dict): + is_owner = (payload_obj.get("created_by") == current_user.get("id")) + else: + try: + import json + payload_str = str(payload_obj.get("payload", "{}") or "{}") + payload_obj = json.loads(payload_str) if payload_str.startswith("{") else {} + is_owner = (payload_obj.get("created_by") == current_user.get("id")) + except (TypeError, ValueError): + pass + + # 管理员或任务创建者:直接删除 + # 其他用户:需要 ACL delete 权限 + 审批流程 + if not is_admin(current_user) and not is_owner: + if not has_resource_access("compare", task_id, current_user, "delete"): + raise fail(403, "no permission to delete inference task") + pending = _require_approval_or_admin("compare", task_id, current_user, f"删除推理任务 {task_id}") + if pending: + return pending get_platform_store().delete_compare_task(task_id) try: await _unload_from_compute_node(get_platform_store(), task=task) @@ -2013,6 +2036,7 @@ def _invalidate_superseded_models(store: Any, task_id: str, loaded_models: list[ @router.post("/model-compare/{task_id}/load") +@op_log(module=OpModule.INFERENCE, action=OpAction.START, target_type="inference", target_name_param="task_id") async def model_compare_load(task_id: str, current_user: dict = Depends(get_current_user)) -> dict[str, Any]: """异步派发模型加载到算力节点,立即返回。 @@ -2023,8 +2047,13 @@ async def model_compare_load(task_id: str, current_user: dict = Depends(get_curr try: store = get_platform_store() task = store.compare_task(task_id) - if not has_resource_access("compare", task_id, current_user, "execute"): - raise fail(403, "no permission to load inference task") + # 先判断是否为任务创建者本人或管理员:如果是,直接允许操作 + is_owner = False + if isinstance(task, dict): + is_owner = (task.get("created_by") == current_user.get("id")) + if not is_admin(current_user) and not is_owner: + if not has_resource_access("compare", task_id, current_user, "execute"): + raise fail(403, "no permission to load inference task") models = task.get("models") or [] if isinstance(models, str): try: @@ -2231,8 +2260,9 @@ async def model_chat_local_status() -> dict[str, Any]: @router.post("/model-chat/trained/preload") async def model_chat_trained_preload(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]: resource_id = str(payload.get("trained_model_id") or payload.get("model_id") or payload.get("resource_id") or "") - if resource_id and not has_resource_access("trained_model", resource_id, current_user, "execute") and not has_resource_access("model", resource_id, current_user, "execute"): - raise fail(403, "no permission to load this model") + # 训练模型需要 ACL 授权;基座模型(配置模型)是平台共享资源,不需要 ACL + if resource_id and not has_resource_access("trained_model", resource_id, current_user, "execute"): + raise fail(403, "no permission to load this trained model") """Load a trained model (base + adapter) on the compute node for inference.""" model_path = (payload.get("model_name_or_path") or "").strip() if not model_path: @@ -2365,8 +2395,9 @@ async def archive_node_files( if not is_admin(current_user): model_id = str(payload.get("model_id") or "") dataset_id = str(payload.get("dataset_id") or "") - if not model_id or not has_resource_access("model", model_id, current_user, "execute"): - raise fail(403, "no permission to evaluate this model") + # 基座模型(配置模型)是平台共享资源,不需要 ACL 授权 + if not model_id: + raise fail(400, "model_id is required") if not dataset_id or not has_resource_access("dataset", dataset_id, current_user, "execute"): raise fail(403, "no permission to evaluate this dataset") node = next((item for item in store.compute_nodes() if item["id"] == node_id), None) diff --git a/backend/app/core/auth.py b/backend/app/core/auth.py index 155135e..403e75e 100644 --- a/backend/app/core/auth.py +++ b/backend/app/core/auth.py @@ -180,9 +180,22 @@ def filter_accessible_resource_ids( accessible = {r["resource_id"] for r in rows} if resource_type in OWNER_TABLES: table, column = OWNER_TABLES[resource_type] - with store.connect() as conn: - owned = conn.execute(f"SELECT id FROM {table} WHERE {column}=?", (user_id,)).fetchall() - accessible.update(row["id"] for row in owned) + if column == "payload": + # payload 是 JSON 字符串,需要查出后解析 created_by + with store.connect() as conn: + owned = conn.execute(f"SELECT id, {column} FROM {table}").fetchall() + for row in owned: + try: + import json + payload = json.loads(row[column] or "{}") + if payload.get("created_by") == user_id: + accessible.add(row["id"]) + except (TypeError, ValueError): + pass + else: + with store.connect() as conn: + owned = conn.execute(f"SELECT id FROM {table} WHERE {column}=?", (user_id,)).fetchall() + accessible.update(row["id"] for row in owned) return [rid for rid in all_ids if rid in accessible] @@ -210,8 +223,18 @@ def filter_accessible_resource_ids_batch( table, column = table_info with store.connect() as conn: owned = conn.execute( - f"SELECT id FROM {table} WHERE id IN ({placeholders}) AND {column}=?", - (*resource_ids, user["id"]), + f"SELECT id, {column} FROM {table} WHERE id IN ({placeholders})", + (*resource_ids,), ).fetchall() - accessible.update(row["id"] for row in owned) + for row in owned: + owner = row[column] + # 如果列是 payload(JSON),需要解析后提取 created_by + if column == "payload": + try: + import json + owner = json.loads(owner or "{}").get("created_by") + except (TypeError, ValueError): + owner = None + if owner == user["id"]: + accessible.add(row["id"]) return accessible diff --git a/backend/app/core/logging.py b/backend/app/core/logging.py index e519375..51453f3 100644 --- a/backend/app/core/logging.py +++ b/backend/app/core/logging.py @@ -4,6 +4,7 @@ from contextvars import ContextVar from datetime import date, datetime, timedelta import json import logging +import sys from logging import Handler, LogRecord from pathlib import Path import re @@ -338,7 +339,7 @@ def setup_request_logging(app: FastAPI) -> None: try: response = await call_next(request) elapsed_ms = (time.perf_counter() - started_at) * 1000 - + noisy_paths = ("/health", "/system-info", "/compute/jobs/", "/model-eval/", "/model-compare/") log_method = logger.debug if request.method == "GET" and response.status_code < 400 else logger.info if any(request.url.path.endswith(path) or path in request.url.path for path in noisy_paths) and response.status_code < 400: @@ -353,6 +354,27 @@ def setup_request_logging(app: FastAPI) -> None: elapsed_ms, request.client.host if request.client else "-", ) + + # 5xx 系统错误自动写入操作日志(未被 @op_log 覆盖的系统级异常) + if response.status_code >= 500: + try: + from app.core.op_log import log_operation, OpModule, OpStatus + log_operation( + module=OpModule.SYSTEM, + action="request", + target_type="api", + target_name=request.url.path, + status=OpStatus.FAILURE, + error_message=f"HTTP {response.status_code} - 系统内部错误", + error_type="HTTPError", + detail=f'{{"method":"{request.method}","path":"{request.url.path}","status":{response.status_code}}}', + func_name="request_logging_middleware", + request=request, + duration_ms=elapsed_ms, + ) + except Exception: + pass # 日志写入失败不影响主流程 + response.headers["X-Request-ID"] = request_id return response except Exception: @@ -364,6 +386,28 @@ def setup_request_logging(app: FastAPI) -> None: elapsed_ms, request.client.host if request.client else "-", ) + + # 未被捕获的异常,写入操作日志 + try: + import traceback as _tb + from app.core.op_log import log_operation, OpModule, OpStatus + log_operation( + module=OpModule.SYSTEM, + action="request", + target_type="api", + target_name=request.url.path, + status=OpStatus.FAILURE, + error_message=str(sys.exc_info()[1])[:1000] if sys.exc_info()[1] else "未知异常", + error_type=type(sys.exc_info()[1]).__name__ if sys.exc_info()[1] else "UnknownError", + error_traceback="".join(_tb.format_exception(*sys.exc_info()))[:5000], + func_name="request_logging_middleware", + detail=f'{{"method":"{request.method}","path":"{request.url.path}"}}', + request=request, + duration_ms=elapsed_ms, + ) + except Exception: + pass # 日志写入失败不影响主流程 + raise finally: request_id_var.reset(token) diff --git a/backend/app/core/op_log.py b/backend/app/core/op_log.py new file mode 100644 index 0000000..2020356 --- /dev/null +++ b/backend/app/core/op_log.py @@ -0,0 +1,382 @@ +""" +操作日志工具模块 + +提供 @op_log 装饰器和 log_operation 函数,用于记录用户在各业务模块的详细操作。 +自动捕获成功/失败状态、完整报错堆栈、操作耗时等。 + +核心设计: + - 失败操作必须清晰记录完整异常堆栈(traceback) + - 记录异常类型(如 RuntimeError / ValueError / ConnectionError) + - 记录具体出错的函数名和文件位置,方便定位 bug + - 记录 HTTP 状态码,方便区分用户错误(4xx)和系统错误(5xx) + +使用示例: + from app.core.op_log import op_log, OpModule, OpAction + + @router.post("/inference/start") + @op_log(module=OpModule.INFERENCE, action=OpAction.START, target_type="inference") + async def start_inference(...): + ... +""" + +from __future__ import annotations + +import asyncio +import functools +import json +import time +import traceback +from datetime import datetime, timezone +from typing import Any, Callable, Optional, TypeVar + +from fastapi import Request + +from app.core.logging import get_logger, request_id_var +from app.db.platform_store import get_platform_store, new_id, utcnow + +logger = get_logger("app.op_log") + +F = TypeVar("F", bound=Callable[..., Any]) + + +class OpModule: + """业务模块常量""" + FINE_TUNE = "fine-tune" # 模型训练 + MODEL_EVAL = "model-eval" # 模型评测 + INFERENCE = "model-inference" # 模型推理 + MODEL_MANAGE = "model-manage" # 模型管理 + DATASET = "dataset" # 数据集 + DATA_PROCESS = "data-process" # 数据处理 + DATA_CONVERT = "data-convert" # 数据类型转换 + COMPUTE = "compute" # 算力节点 + SYSTEM = "system" # 系统 + + +class OpAction: + """操作动作常量""" + CREATE = "create" + UPDATE = "update" + DELETE = "delete" + START = "start" + STOP = "stop" + UPLOAD = "upload" + DOWNLOAD = "download" + CONVERT = "convert" + MERGE = "merge" + IMPORT = "import" + LOGIN = "login" + LOGOUT = "logout" + PUBLISH = "publish" + RETRY = "retry" + + +class OpStatus: + """操作状态常量""" + SUCCESS = "success" + FAILURE = "failure" + + +def op_log( + module: str, + action: str, + target_type: str = "", + *, + target_name_param: str = "name", + detail_params: Optional[list[str]] = None, +) -> Callable[[F], F]: + """ + 操作日志装饰器 + + 自动记录: + - 谁在什么时间操作了什么 + - 成功还是失败 + - 失败时记录完整异常堆栈(traceback)、异常类型、异常消息 + - 出错的函数名和文件位置,方便定位 bug + - 操作耗时(ms) + - 客户端 IP、请求路径 + + Args: + module: 业务模块(OpModule 常量) + action: 操作动作(OpAction 常量) + target_type: 资源类型 + target_name_param: 从 kwargs 中提取目标名称的参数名 + detail_params: 需要记录到 detail 的参数名列表 + """ + def decorator(func: F) -> F: + func_name = f"{func.__module__}.{func.__qualname__}" + + if asyncio.iscoroutinefunction(func): + @functools.wraps(func) + async def async_wrapper(*args, **kwargs): + started_at = time.perf_counter() + trace_id = request_id_var.get("-") + + user = _extract_user(args, kwargs) + request = _extract_request(args) + + target_name = _get_param(kwargs, target_name_param, "") + target_id = _get_param(kwargs, "task_id", "") or _get_param(kwargs, "dataset_id", "") or _get_param(kwargs, "model_id", "") + + detail_dict = _build_detail(detail_params, kwargs) + detail_str = json.dumps(detail_dict, ensure_ascii=False) if detail_dict else "" + + try: + result = await func(*args, **kwargs) + elapsed_ms = (time.perf_counter() - started_at) * 1000 + + if not target_id and isinstance(result, dict): + target_id = str(result.get("id", "")) + + _write_log( + module=module, + action=action, + target_type=target_type, + target_id=str(target_id) if target_id else None, + target_name=str(target_name) if target_name else None, + status=OpStatus.SUCCESS, + error_message="", + error_type="", + error_traceback="", + func_name=func_name, + detail=detail_str, + user=user, + request=request, + trace_id=trace_id, + duration_ms=elapsed_ms, + ) + return result + + except Exception as exc: + elapsed_ms = (time.perf_counter() - started_at) * 1000 + + # 捕获完整异常堆栈 + tb_lines = traceback.format_exception(type(exc), exc, exc.__traceback__) + full_traceback = "".join(tb_lines) + error_msg = str(exc)[:1000] + error_type = type(exc).__name__ + + _write_log( + module=module, + action=action, + target_type=target_type, + target_id=str(target_id) if target_id else None, + target_name=str(target_name) if target_name else None, + status=OpStatus.FAILURE, + error_message=error_msg, + error_type=error_type, + error_traceback=full_traceback, + func_name=func_name, + detail=detail_str, + user=user, + request=request, + trace_id=trace_id, + duration_ms=elapsed_ms, + ) + raise + + return async_wrapper # type: ignore + else: + @functools.wraps(func) + def sync_wrapper(*args, **kwargs): + started_at = time.perf_counter() + trace_id = request_id_var.get("-") + + user = _extract_user(args, kwargs) + request = _extract_request(args) + + target_name = _get_param(kwargs, target_name_param, "") + target_id = _get_param(kwargs, "task_id", "") or _get_param(kwargs, "dataset_id", "") or _get_param(kwargs, "model_id", "") + + detail_dict = _build_detail(detail_params, kwargs) + detail_str = json.dumps(detail_dict, ensure_ascii=False) if detail_dict else "" + + try: + result = func(*args, **kwargs) + elapsed_ms = (time.perf_counter() - started_at) * 1000 + + if not target_id and isinstance(result, dict): + target_id = str(result.get("id", "")) + + _write_log( + module=module, + action=action, + target_type=target_type, + target_id=str(target_id) if target_id else None, + target_name=str(target_name) if target_name else None, + status=OpStatus.SUCCESS, + error_message="", + error_type="", + error_traceback="", + func_name=func_name, + detail=detail_str, + user=user, + request=request, + trace_id=trace_id, + duration_ms=elapsed_ms, + ) + return result + + except Exception as exc: + elapsed_ms = (time.perf_counter() - started_at) * 1000 + + tb_lines = traceback.format_exception(type(exc), exc, exc.__traceback__) + full_traceback = "".join(tb_lines) + error_msg = str(exc)[:1000] + error_type = type(exc).__name__ + + _write_log( + module=module, + action=action, + target_type=target_type, + target_id=str(target_id) if target_id else None, + target_name=str(target_name) if target_name else None, + status=OpStatus.FAILURE, + error_message=error_msg, + error_type=error_type, + error_traceback=full_traceback, + func_name=func_name, + detail=detail_str, + user=user, + request=request, + trace_id=trace_id, + duration_ms=elapsed_ms, + ) + raise + + return sync_wrapper # type: ignore + + return decorator + + +def log_operation( + *, + module: str, + action: str, + target_type: str = "", + target_id: str = "", + target_name: str = "", + status: str = OpStatus.SUCCESS, + error_message: str = "", + error_type: str = "", + error_traceback: str = "", + detail: str = "", + func_name: str = "", + user: Optional[dict] = None, + request: Optional[Request] = None, + duration_ms: float = 0, +) -> None: + """手动记录操作日志(不方便用装饰器时使用)""" + trace_id = request_id_var.get("-") + _write_log( + module=module, + action=action, + target_type=target_type, + target_id=target_id or None, + target_name=target_name or None, + status=status, + error_message=error_message, + error_type=error_type, + error_traceback=error_traceback, + func_name=func_name, + detail=detail, + user=user, + request=request, + trace_id=trace_id, + duration_ms=duration_ms, + ) + + +def _build_detail(detail_params: Optional[list[str]], kwargs: dict) -> dict: + """从 kwargs 中提取需要记录的参数""" + detail_dict = {} + if detail_params: + for p in detail_params: + val = kwargs.get(p) + if val is not None: + detail_dict[p] = str(val)[:200] + return detail_dict + + +def _extract_user(args: tuple, kwargs: dict) -> Optional[dict]: + """从函数参数中提取 current_user dict""" + for arg in args: + if isinstance(arg, dict) and "id" in arg and "username" in arg: + return arg + for v in kwargs.values(): + if isinstance(v, dict) and "id" in v and "username" in v: + return v + return None + + +def _extract_request(args: tuple) -> Optional[Request]: + """从函数参数中提取 Request 对象""" + for arg in args: + if isinstance(arg, Request): + return arg + return None + + +def _get_param(kwargs: dict, key: str, default: str = "") -> str: + """安全获取参数值""" + val = kwargs.get(key, default) + if val is None: + return default + return str(val) + + +def _write_log( + module: str, + action: str, + target_type: str, + target_id: Optional[str], + target_name: Optional[str], + status: str, + error_message: str, + error_type: str, + error_traceback: str, + func_name: str, + detail: str, + user: Optional[dict], + request: Optional[Request], + trace_id: str, + duration_ms: float, +) -> None: + """写入操作日志到数据库""" + try: + store = get_platform_store() + log_id = new_id("op") + user_id = user.get("id") if user else None + username = user.get("username") if user else None + client_ip = None + req_method = None + req_path = None + if request: + client_ip = request.client.host if request.client else None + req_method = request.method + req_path = request.url.path + + with store.connect() as conn: + conn.execute( + """ + INSERT INTO operation_logs + (id, user_id, username, module, action, target_type, target_id, + target_name, status, error_message, error_type, error_traceback, + func_name, detail, client_ip, + request_method, request_path, trace_id, duration_ms, create_time) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + log_id, user_id, username, module, action, + target_type or None, target_id, target_name, + status, + error_message[:1000] if error_message else None, + error_type or None, + error_traceback[:5000] if error_traceback else None, + func_name or None, + detail[:2000] if detail else None, + client_ip, req_method, req_path, trace_id, + round(duration_ms, 2), utcnow(), + ), + ) + except Exception: + logger.error("写入操作日志失败 module=%s action=%s", module, action, exc_info=True) diff --git a/backend/app/db/platform_store.py b/backend/app/db/platform_store.py index 34526f6..ad5f2ff 100644 --- a/backend/app/db/platform_store.py +++ b/backend/app/db/platform_store.py @@ -534,7 +534,7 @@ class PlatformStore: }, ) for table in ("models", "datasets", "eval_tasks"): - self._ensure_columns(conn, table, {"deleted_at": "TEXT", "deleted_by": "TEXT", "tenant_id": "TEXT", "project_id": "TEXT"}) + self._ensure_columns(conn, table, {"deleted_at": "TEXT", "deleted_by": "TEXT", "tenant_id": "TEXT", "project_id": "TEXT", "created_by": "TEXT"}) self._ensure_columns( conn, "resource_replicas", @@ -557,6 +557,25 @@ class PlatformStore: conn.executescript(extra_path.read_text(encoding="utf-8")) # data_convert_tasks 表补充 created_by 字段(用于数据隔离) self._ensure_columns(conn, "data_convert_tasks", {"created_by": "TEXT"}) + # 修复历史数据:将 data_convert_tasks.created_by 回填到关联的 datasets 记录 + try: + conn.execute(""" + UPDATE datasets SET created_by = dct.created_by + FROM data_convert_tasks dct + WHERE datasets.task_id = dct.id + AND datasets.source = 'upload' + AND (datasets.created_by IS NULL OR datasets.created_by = '') + AND dct.created_by IS NOT NULL + AND dct.deleted_at IS NULL + """) + except Exception: + pass # 列不存在时忽略 + # 清理超过 7 天的操作日志 + try: + cutoff = (datetime.now(timezone.utc) - timedelta(days=7)).isoformat() + conn.execute("DELETE FROM operation_logs WHERE create_time < %s", (cutoff,)) + except Exception: + pass # 表不存在时忽略,下次启动会建表 def _column_names(self, conn: PgConnection, table_name: str) -> set[str]: columns = conn.execute( diff --git a/backend/app/db/sql/002_governance.sql b/backend/app/db/sql/002_governance.sql index 60bc239..fa991da 100644 --- a/backend/app/db/sql/002_governance.sql +++ b/backend/app/db/sql/002_governance.sql @@ -74,3 +74,38 @@ CREATE TABLE IF NOT EXISTS retention_policies ( create_by TEXT, updated_at TEXT ); + +-- ===================== 操作日志表 ===================== +-- 记录用户在每个业务模块的详细操作(成功/失败、报错信息等) +CREATE TABLE IF NOT EXISTS operation_logs ( + id TEXT PRIMARY KEY, + user_id TEXT, -- 操作用户 ID + username TEXT, -- 操作用户名(冗余,便于查询) + module TEXT, -- 业务模块:fine-tune/model-eval/model-inference/dataset/data-process/data-convert/model-manage + action TEXT, -- 具体动作:create/start/stop/delete/upload/convert 等 + target_type TEXT, -- 资源类型 + target_id TEXT, -- 资源 ID + target_name TEXT, -- 资源名称(便于阅读) + status TEXT NOT NULL, -- success / failure + error_message TEXT, -- 失败时的报错信息 + detail TEXT, -- 操作详情 JSON(参数摘要) + client_ip TEXT, -- 客户端 IP + request_method TEXT, -- HTTP 方法 + request_path TEXT, -- 请求路径 + trace_id TEXT, -- 链路追踪 ID + duration_ms REAL, -- 耗时(ms) + create_time TEXT -- 操作时间 +); + +-- 幂等升级 operation_logs 表:新增字段(已存在则跳过) +ALTER TABLE operation_logs ADD COLUMN IF NOT EXISTS error_type TEXT; -- 异常类型:RuntimeError / ValueError / ConnectionError +ALTER TABLE operation_logs ADD COLUMN IF NOT EXISTS error_traceback TEXT; -- 完整异常堆栈 +ALTER TABLE operation_logs ADD COLUMN IF NOT EXISTS func_name TEXT; -- 出错的函数名 + +-- 索引 +CREATE INDEX IF NOT EXISTS idx_op_user_time ON operation_logs(user_id, create_time DESC); +CREATE INDEX IF NOT EXISTS idx_op_module_time ON operation_logs(module, create_time DESC); +CREATE INDEX IF NOT EXISTS idx_op_status ON operation_logs(status, create_time DESC); +CREATE INDEX IF NOT EXISTS idx_op_action ON operation_logs(action); +CREATE INDEX IF NOT EXISTS idx_op_create_time ON operation_logs(create_time DESC); +CREATE INDEX IF NOT EXISTS idx_op_error_type ON operation_logs(error_type); diff --git a/backend/app/modules/data_convert/router.py b/backend/app/modules/data_convert/router.py index 472c19a..4ed11e0 100644 --- a/backend/app/modules/data_convert/router.py +++ b/backend/app/modules/data_convert/router.py @@ -10,6 +10,7 @@ from fastapi.responses import FileResponse from app.api.v1.endpoints.platform import ok, fail from app.core.auth import get_current_user, is_admin +from app.core.op_log import op_log, OpModule, OpAction from app.db.platform_store import get_platform_store, new_id @@ -89,6 +90,7 @@ def list_tasks( @router.post("") +@op_log(module=OpModule.DATA_CONVERT, action=OpAction.CREATE, target_type="convert_task", target_name_param="name") def create_task( payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user), @@ -133,6 +135,7 @@ def get_task( @router.post("/{task_id}/source-files") +@op_log(module=OpModule.DATA_CONVERT, action=OpAction.UPLOAD, target_type="convert_task", target_name_param="task_id") async def upload_source_files( task_id: str, files: list[UploadFile] = File(...), @@ -205,6 +208,7 @@ async def upload_source_files( "size": f"{size_bytes} B", "count": output_count, "description": f"由数据类型转换任务 {task_id} 自动导入", + "created_by": task.get("created_by") or current_user.get("id"), }) dataset_id = dataset["id"] with store.connect() as conn: @@ -226,6 +230,7 @@ async def upload_source_files( @router.post("/{task_id}/run") +@op_log(module=OpModule.DATA_CONVERT, action=OpAction.CONVERT, target_type="convert_task", target_name_param="task_id") def run_convert( task_id: str, current_user: dict = Depends(get_current_user), @@ -332,6 +337,7 @@ def import_as_dataset( "size": f"{size_bytes} B", "count": task["output_count"], "description": description, + "created_by": task.get("created_by") or (current_user.get("id") if current_user else None), }) dataset_id = dataset["id"] with store.connect() as conn: @@ -340,6 +346,7 @@ def import_as_dataset( @router.delete("/{task_id}") +@op_log(module=OpModule.DATA_CONVERT, action=OpAction.DELETE, target_type="convert_task", target_name_param="task_id") def delete_task( task_id: str, current_user: dict = Depends(get_current_user), diff --git a/backend/app/modules/system/router.py b/backend/app/modules/system/router.py index 410bae7..7ebd5aa 100644 --- a/backend/app/modules/system/router.py +++ b/backend/app/modules/system/router.py @@ -122,3 +122,143 @@ def audit_logs_export( media_type="text/csv", headers={"Content-Disposition": "attachment; filename=audit_logs.csv"}, ) + + +# ===================== 操作日志 ===================== + +@router.get("/operation-logs") +def operation_logs( + user_id: str | None = Query(default=None, description="按用户 ID 筛选"), + module: str | None = Query(default=None, description="按模块筛选: fine-tune/model-eval/model-inference/dataset/data-convert/model-manage"), + action: str | None = Query(default=None, description="按动作筛选: create/start/stop/delete/upload/convert/merge"), + status: str | None = Query(default=None, description="按状态筛选: success/failure(不传则查全部)"), + keyword: str | None = Query(default=None, description="关键字搜索报错信息(error_message)"), + start_time: str | None = Query(default=None, description="ISO8601 起始时间"), + end_time: str | None = Query(default=None, description="ISO8601 结束时间"), + limit: int = Query(default=50, ge=1, le=200), + offset: int = Query(default=0, ge=0), + current_user: dict = Depends(get_current_user), +) -> dict: + """操作日志查询:按用户/模块/动作/状态/关键字/时间范围分页过滤。""" + if not is_admin(current_user): + from app.api.v1.endpoints.platform import fail + raise fail(403, "admin permission required") + store = get_platform_store() + conditions = [] + params: list = [] + if user_id: + conditions.append("user_id = %s") + params.append(user_id) + if module: + conditions.append("module = %s") + params.append(module) + if action: + conditions.append("action = %s") + params.append(action) + if status: + conditions.append("status = %s") + params.append(status) + if keyword: + conditions.append("(error_message ILIKE %s OR error_type ILIKE %s)") + params.append(f"%{keyword}%") + params.append(f"%{keyword}%") + if start_time: + conditions.append("create_time >= %s") + params.append(start_time) + if end_time: + conditions.append("create_time <= %s") + params.append(end_time) + where = " WHERE " + " AND ".join(conditions) if conditions else "" + with store.connect() as conn: + rows = conn.execute( + f"SELECT * FROM operation_logs{where} ORDER BY create_time DESC LIMIT %s OFFSET %s", + tuple(params + [limit, offset]), + ).fetchall() + total = conn.execute(f"SELECT COUNT(*) FROM operation_logs{where}", tuple(params)).fetchone()[0] + return {"code": 0, "message": "ok", "data": {"items": [dict(r) for r in rows], "total": total}} + + +@router.get("/operation-logs/stats") +def operation_logs_stats( + start_time: str | None = Query(default=None, description="ISO8601 起始时间"), + end_time: str | None = Query(default=None, description="ISO8601 结束时间"), + current_user: dict = Depends(get_current_user), +) -> dict: + """操作日志统计:总操作数、成功数、失败数、失败率、各模块失败分布、最近错误列表。""" + if not is_admin(current_user): + from app.api.v1.endpoints.platform import fail + raise fail(403, "admin permission required") + store = get_platform_store() + conditions = [] + params: list = [] + if start_time: + conditions.append("create_time >= %s") + params.append(start_time) + if end_time: + conditions.append("create_time <= %s") + params.append(end_time) + where = " WHERE " + " AND ".join(conditions) if conditions else "" + failure_where = where + " AND status = 'failure'" if where else " WHERE status = 'failure'" + + with store.connect() as conn: + # 总计 + row = conn.execute( + f"SELECT status, COUNT(*) as cnt FROM operation_logs{where} GROUP BY status", tuple(params) + ).fetchall() + total_count = 0 + success_count = 0 + failure_count = 0 + for r in row: + total_count += r["cnt"] + if r["status"] == "success": + success_count = r["cnt"] + elif r["status"] == "failure": + failure_count = r["cnt"] + failure_rate = round(failure_count / total_count * 100, 2) if total_count > 0 else 0 + + # 各模块失败数 + module_stats = conn.execute( + f"SELECT module, COUNT(*) as cnt FROM operation_logs{failure_where} GROUP BY module ORDER BY cnt DESC", + tuple(params), + ).fetchall() + + # 各异常类型分布 + error_type_stats = conn.execute( + f"SELECT error_type, COUNT(*) as cnt FROM operation_logs{failure_where} AND error_type IS NOT NULL GROUP BY error_type ORDER BY cnt DESC LIMIT 10", + tuple(params), + ).fetchall() + + # 最近的 10 条错误 + recent_errors = conn.execute( + f"SELECT * FROM operation_logs{failure_where} ORDER BY create_time DESC LIMIT 10", + tuple(params), + ).fetchall() + + return { + "code": 0, + "message": "ok", + "data": { + "total": total_count, + "success": success_count, + "failure": failure_count, + "failure_rate": failure_rate, + "module_failures": [{"module": r["module"], "count": r["cnt"]} for r in module_stats], + "error_types": [{"type": r["error_type"], "count": r["cnt"]} for r in error_type_stats], + "recent_errors": [dict(r) for r in recent_errors], + }, + } + + +@router.get("/operation-logs/modules") +def operation_log_modules(current_user: dict = Depends(get_current_user)) -> dict: + """返回操作日志中出现的模块列表(用于筛选下拉框)。""" + if not is_admin(current_user): + from app.api.v1.endpoints.platform import fail + raise fail(403, "admin permission required") + store = get_platform_store() + with store.connect() as conn: + rows = conn.execute( + "SELECT DISTINCT module FROM operation_logs WHERE module IS NOT NULL ORDER BY module" + ).fetchall() + modules = [{"value": r["module"], "label": r["module"]} for r in rows] + return {"code": 0, "message": "ok", "data": modules} diff --git a/frontend/src/api/modules/operation-log.ts b/frontend/src/api/modules/operation-log.ts new file mode 100644 index 0000000..de0c0c5 --- /dev/null +++ b/frontend/src/api/modules/operation-log.ts @@ -0,0 +1,58 @@ +import { get } from '../request' + +export interface OperationLog { + id: string + user_id?: string + username?: string + module?: string + action?: string + target_type?: string + target_id?: string + target_name?: string + status: string + error_message?: string + error_type?: string + error_traceback?: string + func_name?: string + detail?: string + client_ip?: string + request_method?: string + request_path?: string + trace_id?: string + duration_ms?: number + create_time?: string +} + +export interface OperationLogQuery { + user_id?: string + module?: string + action?: string + status?: string + keyword?: string + start_time?: string + end_time?: string + limit?: number + offset?: number +} + +export interface OperationLogStats { + total: number + success: number + failure: number + failure_rate: number + module_failures: { module: string; count: number }[] + error_types: { type: string; count: number }[] + recent_errors: OperationLog[] +} + +/** 操作日志查询 */ +export const getOperationLogs = (query: OperationLogQuery = {}) => + get<{ items: OperationLog[]; total: number }>('/system/operation-logs', query) + +/** 操作日志统计 */ +export const getOperationLogStats = (params?: { start_time?: string; end_time?: string }) => + get('/system/operation-logs/stats', params) + +/** 获取操作日志中出现的模块列表 */ +export const getOperationLogModules = () => + get<{ value: string; label: string }[]>('/system/operation-logs/modules') diff --git a/frontend/src/components/AppSidebar.vue b/frontend/src/components/AppSidebar.vue index e2a8de9..7d188fb 100644 --- a/frontend/src/components/AppSidebar.vue +++ b/frontend/src/components/AppSidebar.vue @@ -81,6 +81,7 @@ const menuGroups: MenuGroup[] = [ { key: 'projects', label: '项目空间', icon: 'fa-folder', to: '/projects', permission: 'user-settings' }, { key: 'resource-acl', label: '资源授权', icon: 'fa-key', to: '/resource-acl', permission: 'user-settings' }, { key: 'audit-logs', label: '审计日志', icon: 'fa-history', to: '/audit-logs', permission: 'user-settings' }, + { key: 'operation-logs', label: '操作日志', icon: 'fa-list', to: '/operation-logs', permission: 'user-settings' }, { key: 'approval-templates', label: '审批模板', icon: 'fa-list-alt', to: '/approval-templates', permission: 'user-settings' }, { key: 'approval-instances', label: '审批中心', icon: 'fa-check-square', to: '/approval-instances', permission: 'user-settings' }, ], diff --git a/frontend/src/router/index.ts b/frontend/src/router/index.ts index 82e14d3..710987a 100644 --- a/frontend/src/router/index.ts +++ b/frontend/src/router/index.ts @@ -62,6 +62,12 @@ const routes: RouteRecordRaw[] = [ component: () => import('@/views/audit/AuditLogView.vue'), meta: { title: '审计日志', permission: 'user-settings' }, }, + { + path: 'operation-logs', + name: 'operation-logs', + component: () => import('@/views/audit/OperationLogView.vue'), + meta: { title: '操作日志', permission: 'user-settings' }, + }, { path: 'approval-templates', name: 'approval-templates', @@ -355,6 +361,7 @@ const permissionBySegment: Record = { tenants: 'user-settings', projects: 'user-settings', 'audit-logs': 'user-settings', + 'operation-logs': 'user-settings', 'approval-templates': 'user-settings', 'approval-instances': 'user-settings', 'resource-acl': 'user-settings', diff --git a/frontend/src/views/audit/OperationLogView.vue b/frontend/src/views/audit/OperationLogView.vue new file mode 100644 index 0000000..4a0485e --- /dev/null +++ b/frontend/src/views/audit/OperationLogView.vue @@ -0,0 +1,437 @@ + + + + +