5 Commits

20 changed files with 2345 additions and 236 deletions

View File

@@ -19,6 +19,8 @@ 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
@@ -223,12 +225,16 @@ def _require_approval_or_admin(
"""
高风险操作审批旁路:
- admin 用户直接放行(返回 None
- 普通用户创建审批实例返回审批待定响应code=202 None
- 资源创建者Owner直接放行返回 None
- 其他普通用户创建审批实例返回审批待定响应code=202非 None
code=202 使前端响应拦截器走业务错误分支,弹提示并 reject
避免前端误认为删除成功。
"""
if is_admin(current_user):
return None
# 资源创建者直接放行,无需审批
if _check_owner(resource_type, resource_id, current_user.get("id")):
return None
store = get_platform_store()
instance = store.create_approval_instance({
"resource_type": resource_type,
@@ -243,6 +249,28 @@ def _require_approval_or_admin(
}
def _check_owner(resource_type: str, resource_id: str, user_id: str | None) -> bool:
"""直接查数据库判断 user_id 是否为资源的 created_by。"""
from app.core.auth import OWNER_TABLES
table_info = OWNER_TABLES.get(resource_type)
if not table_info:
return False
table, column = table_info
store = get_platform_store()
with store.connect() as conn:
row = conn.execute(f"SELECT {column} FROM {table} WHERE id=?", (resource_id,)).fetchone()
if not row:
return False
owner = row[column]
if column == "payload":
try:
import json
owner = json.loads(owner or "{}").get("created_by")
except (TypeError, ValueError):
owner = None
return owner == user_id
def _node_for_task(task: dict[str, Any]) -> dict[str, Any] | None:
return next((node for node in get_platform_store().compute_nodes() if node["id"] == task.get("compute_node_id")), None)
@@ -894,6 +922,11 @@ async def test_online_model(payload: dict[str, Any] = Body(...)) -> dict[str, An
@router.post("/model-manage")
@audit_log(
action=AuditActions.CREATE_MODEL,
target_type="model",
detail_template="创建模型: {name}",
)
async def create_model(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]:
payload.setdefault("created_by", current_user.get("id"))
try:
@@ -908,17 +941,24 @@ 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)
@router.put("/model-manage/{model_id}")
async def update_model(model_id: str, payload: dict[str, Any] = Body(...)) -> dict[str, Any]:
@audit_log(
action=AuditActions.UPDATE_MODEL,
target_type="model",
detail_template="更新模型: {model_id}",
)
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:
@@ -926,7 +966,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:
@@ -935,16 +978,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 "")
@@ -960,8 +1002,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")
@@ -1362,6 +1403,12 @@ async def dataset_list(current_user: dict = Depends(get_current_user)) -> dict[s
@router.post("/dataset-manage")
@audit_log(
action=AuditActions.CREATE_DATASET,
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"))
try:
@@ -1383,6 +1430,11 @@ async def dataset_detail(dataset_id: str, current_user: dict = Depends(get_curre
@router.put("/dataset-manage/{dataset_id}")
@audit_log(
action=AuditActions.UPDATE_DATASET,
target_type="dataset",
detail_template="更新数据集: {dataset_id}",
)
async def update_dataset(dataset_id: str, payload: dict[str, Any] = Body(...)) -> dict[str, Any]:
try:
return ok(get_platform_store().update_dataset(dataset_id, payload))
@@ -1391,6 +1443,12 @@ async def update_dataset(dataset_id: str, payload: dict[str, Any] = Body(...)) -
@router.delete("/dataset-manage/{dataset_id}")
@audit_log(
action=AuditActions.DELETE_DATASET,
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")
@@ -1425,18 +1483,25 @@ async def fine_tune_list(current_user: dict = Depends(get_current_user)) -> dict
tasks = get_platform_store().tasks()
if current_user.get("role") == "admin" or current_user.get("protected"):
return ok(tasks)
# 普通用户可见:自己创建的 + ACL 授权的
user_id = current_user.get("id")
accessible = set(filter_accessible_resource_ids("fine-tune", [t["id"] for t in tasks], current_user))
return ok([t for t in tasks if t["id"] in accessible])
result = [t for t in tasks if t.get("created_by") == user_id or t["id"] in accessible]
return ok(result)
@router.post("/fine-tune")
@audit_log(
action=AuditActions.CREATE_FINE_TUNE,
target_type="fine_tune",
detail_template="创建微调任务: {name}",
)
async def create_fine_tune(payload: dict[str, Any] = Body(...), current_user: dict = Depends(get_current_user)) -> dict[str, Any]:
payload.setdefault("created_by", current_user.get("id"))
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:
@@ -1447,6 +1512,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),
@@ -1606,6 +1672,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:
@@ -1652,6 +1719,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")
@@ -1701,8 +1769,11 @@ async def model_eval_list(current_user: dict = Depends(get_current_user)) -> dic
tasks = get_platform_store().eval_tasks()
if current_user.get("role") == "admin" or current_user.get("protected"):
return ok(tasks)
# 普通用户可见:自己创建的 + ACL 授权的
user_id = current_user.get("id")
accessible = set(filter_accessible_resource_ids("eval", [t["id"] for t in tasks], current_user))
return ok([t for t in tasks if t["id"] in accessible])
result = [t for t in tasks if t.get("created_by") == user_id or t["id"] in accessible]
return ok(result)
@router.get("/model-eval/{task_id}")
@@ -1733,10 +1804,12 @@ 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()
# 1. Create eval task record
payload.setdefault("created_by", current_user.get("id"))
task = store.create_eval_task({**payload, "status": "pending"})
# 2. Resolve model path (supports both regular models and trained models)
@@ -1938,19 +2011,18 @@ async def model_compare_list(current_user: dict = Depends(get_current_user)) ->
tasks = get_platform_store().compare_tasks()
if is_admin(current_user):
return ok(tasks)
# 普通用户可见:自己创建的 + ACL 授权的
user_id = current_user.get("id")
accessible = filter_accessible_resource_ids_batch("compare", [item["id"] for item in tasks], current_user)
return ok([item for item in tasks if item["id"] in accessible])
result = [item for item in tasks if item.get("created_by") == user_id or item["id"] in accessible]
return ok(result)
@router.post("/model-compare")
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"]})
@@ -2008,17 +2080,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)
@@ -2081,6 +2171,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]:
"""异步派发模型加载到算力节点,立即返回。
@@ -2091,8 +2182,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:
@@ -2299,8 +2395,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:
@@ -2433,8 +2530,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)

199
backend/app/core/audit.py Normal file
View File

@@ -0,0 +1,199 @@
"""
审计日志装饰器模块
提供 @audit_log 装饰器,用于自动记录关键业务操作的审计日志。
使用示例:
from app.core.audit import audit_log
@audit_log(action="create_dataset", target_type="dataset")
async def create_dataset(request: Request, ...):
...
"""
from __future__ import annotations
import asyncio
import functools
import time
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
logger = get_logger("app.audit")
F = TypeVar("F", bound=Callable[..., Any])
def audit_log(
action: str,
target_type: str = "",
*,
detail_template: str = "",
extract_target_id: Optional[Callable[[Any], str]] = None,
) -> Callable[[F], F]:
"""
审计日志装饰器
Args:
action: 操作类型,如 create_dataset、update_model 等
target_type: 目标资源类型,如 dataset、model 等
detail_template: 日志详情模板(支持 format 参数)
extract_target_id: 从返回值中提取目标 ID 的函数
Returns:
装饰后的函数
"""
def decorator(func: F) -> F:
if asyncio.iscoroutinefunction(func):
@functools.wraps(func)
async def async_wrapper(*args, **kwargs):
started_at = time.perf_counter()
trace_id = request_id_var.get("-")
try:
result = await func(*args, **kwargs)
elapsed_ms = (time.perf_counter() - started_at) * 1000
target_id = _extract_target_id(result, kwargs, extract_target_id)
detail = _build_detail(detail_template, kwargs)
_record_audit(
action=action,
target_type=target_type,
target_id=target_id,
detail=detail,
trace_id=trace_id,
duration_ms=elapsed_ms,
)
return result
except Exception:
logger.error(
"审计日志记录失败 action=%s", action, exc_info=True
)
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("-")
try:
result = func(*args, **kwargs)
elapsed_ms = (time.perf_counter() - started_at) * 1000
target_id = _extract_target_id(result, kwargs, extract_target_id)
detail = _build_detail(detail_template, kwargs)
_record_audit(
action=action,
target_type=target_type,
target_id=target_id,
detail=detail,
trace_id=trace_id,
duration_ms=elapsed_ms,
)
return result
except Exception:
logger.error(
"审计日志记录失败 action=%s", action, exc_info=True
)
raise
return sync_wrapper # type: ignore
return decorator
def _extract_target_id(
result: Any, kwargs: dict, extractor: Optional[Callable[[Any], str]]
) -> Optional[str]:
"""从返回值或 kwargs 中提取目标 ID"""
if extractor:
try:
return extractor(result)
except Exception:
pass
if isinstance(result, dict):
return result.get("id")
# 尝试从路径参数中提取
for key in ("dataset_id", "model_id", "task_id", "resource_id"):
val = kwargs.get(key)
if val:
return str(val)
return None
def _build_detail(template: str, kwargs: dict) -> str:
"""构建审计详情"""
if not template:
return ""
try:
return template.format(**kwargs)
except (KeyError, IndexError):
return template
def _record_audit(
action: str,
target_type: str,
target_id: Optional[str],
detail: str,
trace_id: str,
duration_ms: float,
) -> None:
"""通过已有的 record_audit 方法写入审计日志"""
try:
from app.db.platform_store import get_platform_store
store = get_platform_store()
store.record_audit(
action=action,
target_type=target_type or None,
target_id=target_id,
detail=f"{detail} trace_id={trace_id} duration_ms={duration_ms:.1f}" if detail else f"trace_id={trace_id} duration_ms={duration_ms:.1f}",
)
except Exception:
logger.error("写入审计日志失败 action=%s", action, exc_info=True)
# ==================== 预定义的审计操作常量 ====================
class AuditActions:
"""预定义的审计操作类型"""
# 数据集操作
CREATE_DATASET = "create_dataset"
UPDATE_DATASET = "update_dataset"
DELETE_DATASET = "delete_dataset"
# 模型操作
CREATE_MODEL = "create_model"
UPDATE_MODEL = "update_model"
DELETE_MODEL = "delete_model"
# 微调任务
CREATE_FINE_TUNE = "create_fine_tune"
UPDATE_FINE_TUNE = "update_fine_tune"
DELETE_FINE_TUNE = "delete_fine_tune"
# 推理任务
CREATE_INFERENCE = "create_inference"
UPDATE_INFERENCE = "update_inference"
DELETE_INFERENCE = "delete_inference"
# 用户管理
CREATE_USER = "create_user"
UPDATE_USER = "update_user"
DELETE_USER = "delete_user"
# 租户管理
CREATE_TENANT = "create_tenant"
UPDATE_TENANT = "update_tenant"
DELETE_TENANT = "delete_tenant"
# 权限授权
GRANT_ACL = "grant_acl"
REVOKE_ACL = "revoke_acl"
# 系统配置
UPDATE_CONFIG = "update_config"

View File

@@ -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]
# 如果列是 payloadJSON需要解析后提取 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

View File

@@ -4,11 +4,12 @@ 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
import time
from typing import Any
from typing import Any, Callable, Optional
from uuid import uuid4
from fastapi import FastAPI, Request
@@ -17,6 +18,74 @@ from app.core.config import Settings, get_settings
request_id_var: ContextVar[str] = ContextVar("request_id", default="-")
# ==================== 敏感数据脱敏规则 ====================
SENSITIVE_PATTERNS: dict[str, Callable | str] = {
"token": "***",
"password": "***",
"access_token": "***",
"refresh_token": "***",
"secret_key": "***",
"authorization": "***",
"bearer": "***",
"api_key": "***",
"private_key": "***",
}
def mask_value(key: str, value: Any) -> str:
"""对单个值进行脱敏处理"""
if value is None:
return ""
str_val = str(value)
handler = SENSITIVE_PATTERNS.get(key)
if callable(handler):
return handler(str_val)
elif isinstance(handler, str):
# 支持正则替换模式,如 r"1\d{3}\d{4}"
try:
return re.sub(handler, "***", str_val)
except re.error:
return "***"
return handler
def mask_sensitive_dict(data: dict) -> dict:
"""递归脱敏字典中的敏感字段"""
if not data or not isinstance(data, dict):
return data
result = {}
for key, value in data.items():
result[key] = mask_value(key, value)
return result
def mask_sensitive_string(text: str) -> str:
"""从文本中脱敏常见敏感信息"""
if not text:
return text
patterns = [
(r'Bearer\s+[A-Za-z0-9\-._]+', '***'),
(r'token\s*[:=]\s*', '***'),
(r'password\s*[:=]\s*', '***'),
(r'secret[_-]?key\s*[:=]', '***'),
(r'api[-_]?key\s*[:=]', '***'),
(r'private[_-]?key\s*[:=]', '***'),
(r'\d{11}', r'\d{3}\*\d{4}'), # 手机号/身份证
(r'1[3-9]\d{9}', r'1\*{3}\*{4}'), # 手机号
]
for pattern, replacement in patterns:
try:
text = re.sub(pattern, replacement, text, flags=re.IGNORECASE)
except re.error:
pass
return text
# ==================== RequestId Filter ====================
class RequestIdFilter(logging.Filter):
def filter(self, record: LogRecord) -> bool:
@@ -24,8 +93,30 @@ class RequestIdFilter(logging.Filter):
return True
# ==================== Enhanced JSON Formatter ====================
class JsonLogFormatter(logging.Formatter):
"""Format one JSON object per line for ELK/Filebeat collection."""
"""
增强的 JSON 日志格式化器,支持结构化字段输出。
输出示例:
{
"@timestamp": "2026-08-17T18:30:00.123Z",
"level": "INFO",
"logger": "dataset.router",
"message": "数据集创建成功",
"module": "dataset.router",
"function": "create_dataset",
"file": "dataset/router.py",
"line": 45,
"process": 12345,
"thread": "MainThread",
"request_id": "req-abc123",
"user_id": "u_admin",
"client_ip": "192.168.1.100",
"extra": {...}
}
"""
def format(self, record: LogRecord) -> str:
payload: dict[str, Any] = {
@@ -44,13 +135,26 @@ class JsonLogFormatter(logging.Formatter):
"thread_name": record.threadName,
"request_id": getattr(record, "request_id", "-"),
}
# 从 record 中提取额外字段(通过 extra 参数传入)
for attr in ("user_id", "client_ip", "target_type", "target_id",
"duration_ms", "status_code", "error"):
val = getattr(record, attr, None)
if val is not None:
payload[attr] = val
# 处理异常信息
if record.exc_info:
payload["exception"] = self.formatException(record.exc_info)
if record.stack_info:
payload["stack"] = self.formatStack(record.stack_info)
return json.dumps(payload, ensure_ascii=False, separators=(",", ":"))
# ==================== DateSizeRotatingFileHandler ====================
# (保持不变,已有实现)
class DateSizeRotatingFileHandler(Handler):
"""Rotate log files by date and size while keeping date in every file name."""
@@ -160,6 +264,157 @@ class DateSizeRotatingFileHandler(Handler):
path.unlink(missing_ok=True)
# ==================== Structured Logger 封装 ====================
class StructuredLogger:
"""
结构化日志记录器,提供统一的日志接口。
使用方式:
logger = get_structured_logger('dataset.router')
logger.info('创建数据集', dataset_id='ds_123')
"""
def __init__(self, name: str, module: str = ""):
self.logger = logging.getLogger(name)
self.name = name
self.module = module
@property
def trace_id(self) -> str:
return request_id_var.get("-")
def info(self, message: str, **extra: Any) -> None:
self._log("INFO", message, **extra)
def warning(self, message: str, **extra: Any) -> None:
self._log("WARNING", message, **extra)
def error(self, message: str, **extra: Any) -> None:
self._log("ERROR", message, **extra)
def debug(self, message: str, **extra: Any) -> None:
self._log("DEBUG", message, **extra)
def _log(self, level: str, message: str, **extra: Any) -> None:
"""统一日志记录方法"""
log_entry: dict[str, Any] = {
"timestamp": datetime.utcnow().isoformat(),
"level": level,
"logger": self.name,
"module": self.module,
"message": message,
"trace_id": self.trace_id,
"extra": extra,
}
self.logger.log(getattr(logging, level, logging.INFO), json.dumps(log_entry, ensure_ascii=False, default=str))
def get_structured_logger(name: str, module: str = "") -> StructuredLogger:
"""获取结构化日志记录器"""
return StructuredLogger(name, module)
# ==================== 快捷函数 ====================
def get_logger(name: str) -> logging.Logger:
"""获取标准 Python logger"""
return logging.getLogger(name)
def set_request_id(request_id: str) -> None:
"""设置当前请求的追踪 ID"""
request_id_var.set(request_id)
def setup_request_logging(app: FastAPI) -> None:
"""配置 FastAPI 请求日志中间件"""
logger = get_logger("app.access")
@app.middleware("http")
async def request_logging_middleware(request: Request, call_next): # type: ignore[no-untyped-def]
request_id = request.headers.get("X-Request-ID") or str(uuid4())
token = request_id_var.set(request_id)
started_at = time.perf_counter()
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:
log_method = logger.debug
if response.status_code >= 400:
log_method = logger.warning
log_method(
"request completed method=%s path=%s status_code=%s duration_ms=%.2f client=%s",
request.method,
request.url.path,
response.status_code,
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:
elapsed_ms = (time.perf_counter() - started_at) * 1000
logger.exception(
"request failed method=%s path=%s duration_ms=%.2f client=%s",
request.method,
request.url.path,
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)
# ==================== 配置函数 ====================
def configure_logging(settings: Settings | None = None) -> None:
settings = settings or get_settings()
@@ -209,57 +464,5 @@ def configure_logging(settings: Settings | None = None) -> None:
logger.handlers.clear()
logger.propagate = True
logging.getLogger("uvicorn.access").setLevel(logging.WARNING)
logging.getLogger("psycopg.pool").setLevel(logging.ERROR)
def get_logger(name: str) -> logging.Logger:
return logging.getLogger(name)
def set_request_id(request_id: str) -> None:
request_id_var.set(request_id)
def setup_request_logging(app: FastAPI) -> None:
logger = get_logger("app.access")
@app.middleware("http")
async def request_logging_middleware(request: Request, call_next): # type: ignore[no-untyped-def]
request_id = request.headers.get("X-Request-ID") or str(uuid4())
token = request_id_var.set(request_id)
started_at = time.perf_counter()
try:
response = await call_next(request)
elapsed_ms = (time.perf_counter() - started_at) * 1000
# Docker/frontend probes and polling endpoints are intentionally
# quiet at INFO; failures remain visible at WARNING/ERROR.
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:
log_method = logger.debug
if response.status_code >= 400:
log_method = logger.warning
log_method(
"request completed method=%s path=%s status_code=%s duration_ms=%.2f client=%s",
request.method,
request.url.path,
response.status_code,
elapsed_ms,
request.client.host if request.client else "-",
)
response.headers["X-Request-ID"] = request_id
return response
except Exception:
elapsed_ms = (time.perf_counter() - started_at) * 1000
logger.exception(
"request failed method=%s path=%s duration_ms=%.2f client=%s",
request.method,
request.url.path,
elapsed_ms,
request.client.host if request.client else "-",
)
raise
finally:
request_id_var.reset(token)
logging.getLogger("uvicorn.access").setLevel(logging.WARNING)
logging.getLogger("psycopg.pool").setLevel(logging.ERROR)

382
backend/app/core/op_log.py Normal file
View File

@@ -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)

View File

@@ -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",
@@ -555,6 +555,37 @@ class PlatformStore:
extra_path = schema_dir / extra
if extra_path.exists():
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 # 列不存在时忽略
# 修复历史数据:为 eval_tasks 表回填 created_by从 payload JSON 中提取)
try:
conn.execute("""
UPDATE eval_tasks SET created_by = payload::json->>'created_by'
WHERE (created_by IS NULL OR created_by = '')
AND payload IS NOT NULL
AND payload::json->>'created_by' IS NOT 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(
@@ -1285,7 +1316,13 @@ class PlatformStore:
def create_user(self, payload: dict[str, Any]) -> dict[str, Any]:
user_id = new_id("u")
permissions = payload.get("permissions") or (ALL_PERMISSIONS if payload.get("role") == "admin" else ["dashboard"])
role = payload.get("role", "user")
if role == "admin":
permissions = ALL_PERMISSIONS
else:
# 普通用户:默认拥有所有业务权限,仅排除 user-settings 和 compute
role = "user"
permissions = [p for p in ALL_PERMISSIONS if p not in ("user-settings", "compute")]
with self.connect() as conn:
conn.execute(
"""
@@ -1298,7 +1335,7 @@ class PlatformStore:
payload["username"],
hash_password(payload.get("password", "platform123")),
payload.get("display_name") or payload["username"],
payload.get("role", "viewer"),
role,
payload.get("status", "active"),
json_dumps(permissions),
utcnow(),
@@ -1318,8 +1355,8 @@ class PlatformStore:
# 管理员权限不可更改,必须是全部
perms = ALL_PERMISSIONS
else:
# 非 admin 用户不能拥有 user-settings 权限
perms = [p for p in (perms or []) if p != "user-settings"]
# 非 admin 用户不能拥有 user-settings 和 compute 权限
perms = [p for p in (perms or []) if p not in ("user-settings", "compute")]
payload = {**payload, "permissions": perms}
values = {
"role": payload.get("role", row["role"]),
@@ -2004,6 +2041,7 @@ class PlatformStore:
"process_id": None,
"train_duration": "",
"create_time": now,
"created_by": payload.get("created_by"),
}
with self.connect() as conn:
conn.execute(
@@ -2502,8 +2540,8 @@ class PlatformStore:
if data.get("metric") == "custom":
data["metric"] = data["metric_label"]
conn.execute(
"INSERT INTO eval_tasks (id, name, payload, status, create_time) VALUES (?, ?, ?, ?, ?)",
(task_id, name, json_dumps(data), status, now),
"INSERT INTO eval_tasks (id, name, payload, status, create_time, created_by) VALUES (?, ?, ?, ?, ?, ?)",
(task_id, name, json_dumps(data), status, now, data.get("created_by")),
)
return self.eval_task(task_id)

View File

@@ -51,10 +51,18 @@ CREATE TABLE IF NOT EXISTS audit_logs (
time TEXT
);
CREATE INDEX IF NOT EXISTS idx_audit_tenant ON audit_logs(tenant_id);
CREATE INDEX IF NOT EXISTS idx_audit_project ON audit_logs(project_id);
CREATE INDEX IF NOT EXISTS idx_audit_action ON audit_logs(action);
CREATE INDEX IF NOT EXISTS idx_audit_time ON audit_logs(time);
-- 幂等升级 audit_logs 表:新增字段(已存在则跳过)
ALTER TABLE audit_logs ADD COLUMN IF NOT EXISTS trace_id TEXT;
ALTER TABLE audit_logs ADD COLUMN IF NOT EXISTS request_method TEXT;
ALTER TABLE audit_logs ADD COLUMN IF NOT EXISTS request_path TEXT;
ALTER TABLE audit_logs ADD COLUMN IF NOT EXISTS status_code INTEGER;
ALTER TABLE audit_logs ADD COLUMN IF NOT EXISTS duration_ms REAL;
ALTER TABLE audit_logs ADD COLUMN IF NOT EXISTS extra JSONB;
-- 索引(已存在则跳过)
CREATE INDEX IF NOT EXISTS idx_audit_trace ON audit_logs(trace_id);
CREATE INDEX IF NOT EXISTS idx_audit_actor_time ON audit_logs(actor_id, time);
CREATE INDEX IF NOT EXISTS idx_audit_target ON audit_logs(target_type, target_id);
CREATE TABLE IF NOT EXISTS retention_policies (
id TEXT PRIMARY KEY,
@@ -66,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);

View File

@@ -9,7 +9,8 @@ from fastapi import APIRouter, Body, Depends, File, UploadFile
from fastapi.responses import FileResponse
from app.api.v1.endpoints.platform import ok, fail
from app.core.auth import get_current_user
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
@@ -63,18 +64,33 @@ def list_tasks(
) -> dict[str, Any]:
store = get_platform_store()
with store.connect() as conn:
rows = conn.execute(
"SELECT * FROM data_convert_tasks WHERE deleted_at IS NULL "
"ORDER BY create_time DESC LIMIT %s OFFSET %s",
(page_size, (page - 1) * page_size),
).fetchall()
total = conn.execute(
"SELECT COUNT(*) FROM data_convert_tasks WHERE deleted_at IS NULL"
).fetchone()[0]
if is_admin(current_user):
# 管理员可见全部
rows = conn.execute(
"SELECT * FROM data_convert_tasks WHERE deleted_at IS NULL "
"ORDER BY create_time DESC LIMIT %s OFFSET %s",
(page_size, (page - 1) * page_size),
).fetchall()
total = conn.execute(
"SELECT COUNT(*) FROM data_convert_tasks WHERE deleted_at IS NULL"
).fetchone()[0]
else:
# 普通用户只能看到自己创建的
user_id = current_user.get("id")
rows = conn.execute(
"SELECT * FROM data_convert_tasks WHERE deleted_at IS NULL AND created_by=%s "
"ORDER BY create_time DESC LIMIT %s OFFSET %s",
(user_id, page_size, (page - 1) * page_size),
).fetchall()
total = conn.execute(
"SELECT COUNT(*) FROM data_convert_tasks WHERE deleted_at IS NULL AND created_by=%s",
(user_id,)
).fetchone()[0]
return ok({"items": [dict(r) for r in rows], "total": total})
@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),
@@ -85,12 +101,13 @@ def create_task(
task_id = new_id("dct")
output_filename = _safe_output_filename(payload.get("output_filename"))
description = str(payload.get("description") or "").strip()
user_id = current_user.get("id")
store = get_platform_store()
with store.connect() as conn:
conn.execute(
"INSERT INTO data_convert_tasks (id, name, description, output_filename) "
"VALUES (%s, %s, %s, %s)",
(task_id, name, description, output_filename),
"INSERT INTO data_convert_tasks (id, name, description, output_filename, created_by) "
"VALUES (%s, %s, %s, %s, %s)",
(task_id, name, description, output_filename, user_id),
)
# 创建目录
_input_dir(task_id).mkdir(parents=True, exist_ok=True)
@@ -118,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(...),
@@ -190,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:
@@ -211,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),
@@ -317,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:
@@ -325,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),

View File

@@ -6,6 +6,7 @@ from typing import Any
from app.api.v1.endpoints.platform import ok, fail
from app.db.platform_store import get_platform_store
from app.core.auth import get_current_user, has_resource_access, is_admin
from app.core.audit import audit_log, AuditActions
router = APIRouter(prefix="/resources", tags=["resource"])
@@ -25,6 +26,11 @@ def get_acl(resource_type: str, resource_id: str, current_user: dict = Depends(g
@router.put("/{resource_type}/{resource_id}/acl")
@audit_log(
action=AuditActions.GRANT_ACL,
target_type="",
detail_template="设置资源授权: {resource_type}/{resource_id}",
)
def set_acl(
resource_type: str,
resource_id: str,

View File

@@ -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}

View File

@@ -0,0 +1,412 @@
# 生产级日志系统设计方案
> 版本v1.0
> 日期2026-08-17
> 状态:待评审
---
## 一、现状分析
### 1.1 当前日志架构
```
┌─────────────┐
│ FastAPI │ ← 请求入口
└──────┬──────┘
┌─────────────┐
│ Logging │ ← Python logging 模块
│ Middleware │
└──────┬──────┘
├──────────────────┬──────────────────┐
▼ ▼
┌─────────────┐ ┌─────────────┐
│ Console │ │ File │ ← 输出目标
│ (开发环境) │ │ (JSON格式) │
└─────────────┘ └─────────────┘
┌─────────────┐
│ audit_logs │ ← 审计日志表
│ (PostgreSQL) │
└─────────────┘
```
### 1.2 现有组件
| 组件 | 文件路径 | 功能 |
|------|----------|------|
| `logging.py` | `backend/app/core/` | 日志配置、JSON 格式化、按日期/大小轮转 |
| `platform_store.py` | `backend/app/db/` | `record_audit()` 审计日志写入 |
| `002_governance.sql` | `backend/app/db/sql/` | `audit_logs` 表结构 |
### 1.3 存在的问题
| 问题 | 影响 | 严重程度 |
|------|------|----------|
| **无结构化日志分级** | DEBUG/INFO/WARNING/ERROR 全部混在一起,无法按级别过滤查看 | 🔴 高 |
| **无请求链路追踪** | 一个请求从进入到返回经过哪些服务/函数,无法串联 | 🔴 高 |
| **审计日志与业务耦合** | 各模块手动调用 `record_audit()`,容易遗漏 | 🟡 中 |
| **无敏感数据脱敏** | 用户 token、密码等可能明文记录 | 🔴 高 |
| **无日志聚合查询** | 无法按用户/时间范围/操作类型快速检索 | 🟡 中 |
| **无告警通知** | 系统异常无法主动推送通知 | 🟡 中 |
| **日志文件无归档策略** | 只有简单的过期删除,无压缩归档 | 🟢 低 |
---
## 二、设计目标
### 2.1 核心原则
1. **结构化** - 日志有固定 schema便于机器解析和查询
2. **可追溯** - 每个请求有唯一 ID可串联完整调用链路
3. **分级输出** - 不同环境输出不同级别,生产环境不输出 DEBUG
4. **安全合规** - 敏感数据自动脱敏token、密码、手机号等
5. **高性能** - 日志写入不影响业务接口性能(异步写入)
6. **可观测** - 支持快速检索、统计、告警
### 2.2 日志分级标准
| 级别 | 使用场景 | 示例 | 生产环境 |
|------|----------|------|:--------:|
| **DEBUG** | 开发调试 | 变量值、SQL 语句、完整堆栈 | ❌ 不输出 |
| **INFO** | 正常流程记录 | 任务创建成功、用户登录 | ✅ 记录 |
| **WARNING** | 可恢复异常 | 重试操作、参数校验失败、资源不足 | ✅ 记录 |
| **ERROR** | 需要人工介入 | 数据库连接失败、第三方 API 超时 | ✅ 记录 + 告警 |
| **CRITICAL** | 系统不可用 | 磁盘满、主节点宕机 | ✅ 记录 + 立即告警 |
---
## 三、技术方案
### 3.1 整体架构
```
┌─────────────────────────────────────────────────────────────────────┐
│ 应用层 (Application Layer) │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ 数据集管理 │ │ 微调训练 │ │ 模型推理 │ │ 用户认证 │ ... │
│ └─────┬────┘ └─────┬────┘ └─────┬────┘ └─────┬────┘ │
│ │ │ │ │ │
│ └────────────┴───────────┴──────────┘ │
│ ▼ │
│ ┌──────────────┐ │
│ │ Structured │ ← 结构化日志中间件 │
│ │ Logger │ │
│ └──────┬───────┘ │
│ │ │
│ ┌────────────┬────────────┬─────────────┐ │
│ ▼ ▼ ▼ │ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │
│ │ Console │ │ File │ │ 审计DB │ │ 告警 │ │
│ │ (开发) │ │ (JSON) │ │ (PG) │ │(可选) │ │
│ └──────────┘ └──────────┘ └──────────┘ └────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘
┌─────────────────────────────────────────────────────────────┐
│ 可观测层 (Observability) │
├─────────────────────────────────────────────────────────────┤
│ ┌───────────┐ ┌───────────┐ ┌───────────┐ │
│ │ Grafana │ │ Kibana │ │ PagerDuty │ ... │
│ │ (查询) │ │ (分析) │ │ (告警) │ │
│ └───────────┘ └───────────┘ └───────────┘ │
└─────────────────────────────────────────────────────────────┘
```
### 3.2 日志 Schema 设计
#### 3.2.1 应用日志 (app.log)
```json
{
"timestamp": "2026-08-17T10:30:00.000Z",
"level": "INFO",
"trace_id": "req-abc123",
"parent_span_id": "span-xyz789", // OpenTelemetry Span
"request": {
"method": "POST",
"path": "/dataset-manage",
"client_ip": "192.168.1.100",
"user_agent": "Mozilla/5.0...",
"user_id": "u_admin"
},
"module": "dataset.router",
"function": "create_dataset",
"message": "数据集创建成功",
"extra": {
"dataset_id": "ds_abc123",
"dataset_name": "训练数据"
},
"duration_ms": 125,
"status_code": 200,
"error": null
}
```
#### 3.2.2 审计日志 (audit_logs 表)
```sql
-- 已有表结构(保持不变)
CREATE TABLE IF NOT EXISTS audit_logs (
id TEXT PRIMARY KEY,
tenant_id TEXT,
project_id TEXT,
actor_id TEXT, -- 操作人
action TEXT, -- 操作类型: create/delete/update/acl.set/login...
target_type TEXT, -- 资源类型: dataset/model/fine-tune/user...
target_id TEXT, -- 资源 ID
detail TEXT, -- 详细信息 JSON
client_ip TEXT, -- 客户端 IP
time TEXT, -- 操作时间
-- 新增字段
trace_id TEXT, -- 关联应用日志的请求追踪 ID
request_method TEXT, -- HTTP 方法
request_path TEXT, -- 请求路径
status_code INTEGER, -- 响应状态码
duration_ms REAL, -- 耗时(ms)
extra JSONB -- 扩展信息
);
-- 新增索引
CREATE INDEX IF NOT EXISTS idx_audit_trace ON audit_logs(trace_id);
CREATE INDEX IF NOT EXISTS idx_audit_actor_time ON audit_logs(actor_id, time);
```
### 3.3 日志中间件设计
```python
# backend/app/core/logging.py 新增
class StructuredLogger:
"""结构化日志记录器"""
def __init__(self, name: str):
self.logger = logging.getLogger(name)
self.trace_id = context_var.get("trace_id")
def info(self, msg: str, **kwargs):
self._log("INFO", msg, **kwargs)
def warning(self, msg: str, **kwargs):
self._log("WARNING", msg, **kwargs)
def error(self, msg: str, **kwargs):
self._log("ERROR", msg, **kwargs)
def _log(self, level: str, msg: str,
user_id: str = None,
target_type: str = None,
target_id: str = None,
duration_ms: float = None,
status_code: int = None,
error: Exception = None,
**extra):
"""统一日志记录方法"""
log_entry = {
"timestamp": datetime.utcnow().isoformat(),
"level": level,
"trace_id": self.trace_id.get(),
"request": {
"user_id": user_id or current_user_id(),
"client_ip": client_ip(),
# ...
},
"module": calling_module,
"message": msg,
"target": {
"type": target_type,
"id": target_id,
},
"extra": extra,
"duration_ms": duration_ms,
"error": format_exception(error) if error else None,
}
# 1. 写入控制台/文件
self.logger.log(level, json.dumps(log_entry))
# 2. 异步写入审计表(如果需要)
if level in ("WARNING", "ERROR", "CRITICAL"):
async_write_audit(log_entry)
```
### 3.4 装饰器模式(推荐)
使用 Python 裁饰器自动记录,避免手动调用:
```python
# backend/app/core/log_decorator.py
def audit_log(action: str, target_type: str = ""):
"""审计日志装饰器"""
def decorator(func):
@wraps(func)
async def wrapper(*args, **kwargs):
result = await func(*args, **kwargs)
# 自动记录审计日志
record_audit(
action=action,
target_type=target_type,
target_id=kwargs.get('id') or result.get('id'),
detail=f"params={kwargs}"
)
return result
return wrapper
return decorator
# 使用示例
@audit_log("dataset.create", "dataset")
async def create_dataset(...):
# 业务逻辑
pass
```
### 3.5 敏感数据脱敏规则
```python
# backend/app/core/masking.py
SENSITIVE_FIELDS = {
"token": "***",
"password": "***",
"phone": lambda x: f"{x[:3]}****{x[-4:]}",
"email": lambda x: x[0] + "***" + x.split("@")[1] if "@" in x else "***",
"id_card": lambda x: f"{x[:6]}********{x[-4:]}",
}
def mask_sensitive(data: dict) -> dict:
"""递归脱敏字典中的敏感字段"""
for key, value in data.items():
if key in SENSITIVE_FIELDS:
data[key] = SENSITIVE_FIELDS[key](value) if callable(SENSITIVE_FIELDS[key]) else "***"
elif isinstance(value, dict):
mask_sensitive(value)
return data
```
---
## 四、实施计划
### 4.1 Phase 1基础增强1-2 天)
- [ ] **P1-1** 升级 `JsonLogFormatter`,增加 `trace_id` 字段
- [ ] **P1-2** 新增 `StructuredLogger` 封装类
- [ ] **P1-3** 统一所有模块的日志格式为 JSON
- [ ] **P1-4** 实现 `mask_sensitive()` 脱敏函数
- [ ] **P1-5** 审计日志表新增 `trace_id``duration_ms` 字段
### 4.2 Phase 2自动化2-3 天)
- [ ] **P2-1** 编写 `@audit_log` 装饰器
- [ ] **P2-2** 为关键业务接口添加装饰器:
- 数据集 CRUD
- 模型 CRUD
- 微调任务创建/删除
- 用户登录/登出
- ACL 授权变更
- [ ] **P2-3** 实现日志异步写入队列(避免影响性能)
### 4.3 Phase 3可观测性3-5 天)
- [ ] **P3-1** 集成 ELK Stack 或 Loki可选
- [ ] **P3-2** 编写 Grafana 仪表板:
- 请求量趋势图
- 错误率统计
- 慢接口 TOP10
- 用户操作审计面板
- [ ] [ ] **P3-3** 实现告警规则(错误率超阈值触发)
---
## 五、配置示例
### 5.1 日志配置 (settings)
```yaml
# config.yaml 或 .env
LOGGING:
level: INFO # 生产环境用 INFO开发用 DEBUG
dir: ./logs
file_prefix: app
max_bytes: 50MB # 单文件最大 50MB
retention_days: 30 # 保留 30 天
error_prefix: error # 错误日志单独文件
json: true # JSON 格式输出
AUDIT:
enabled: true
auto_record: true # 是否自动记录(通过装饰器)
sensitive_mask: true # 启用敏感数据脱敏
```
### 5.2 日志输出示例
**控制台输出(开发环境):**
```
2026-08-17 18:30:00.123 | INFO | pid=12345 | MainThread | req=req-abc | dataset.router:create_dataset | dataset/router.py:45 | 数据集创建成功 {"dataset_id":"ds_abc"}
```
**文件输出JSON 格式):**
```json
{"@timestamp":"2026-08-17T18:30:00.123Z","level":"INFO","logger":"dataset.router","message":"数据集创建成功","module":"dataset.router","function":"create_dataset","file":"dataset/router.py","line":45,"process":12345,"thread":"MainThread","request_id":"req-abc","extra":{"dataset_id":"ds_abc"}}
```
**审计日志查询 SQL**
```sql
-- 查询某用户最近7天的所有操作
SELECT time, action, target_type, target_id, detail, client_ip
FROM audit_logs
WHERE actor_id = 'u_admin'
AND time >= now() - interval '7 days'
ORDER BY time DESC;
-- 查询某资源的授权变更历史
SELECT * FROM audit_logs
WHERE action LIKE '%acl%'
AND target_id = 'ds_abc123'
ORDER BY time DESC;
```
---
## 六、附录
### A. 日志关键字段说明
| 字段 | 类型 | 说明 | 示例 |
|------|------|------|------|
| `trace_id` | string | 请求唯一标识,用于串联一次请求的所有日志 | `req-uuid-1234` |
| `parent_span_id` | string | 父 Span ID用于分布式追踪 | `span-parent-5678` |
| `actor_id` | string | 操作人用户 ID | `u_admin` |
| `action` | string | 操作动作 | `dataset.create`, `model.delete`, `login.success` |
| `target_type` | string | 操作的资源类型 | `dataset`, `trained_model`, `user` |
| `target_id` | string | 资源 ID | `ds_abc123` |
| `detail` | string/json | 操作详情 | `{"name": "训练数据", "type": "train"}` |
| `client_ip` | string | 客户端 IP | `192.168.1.100` |
| `duration_ms` | real | 接口耗时(ms) | `125.5` |
| `status_code` | int | HTTP 状态码 | `200`, `404`, `500` |
### B. 推荐的 Python 日志库对比
| 库 | 特点 | 适用场景 |
|-----|------|---------|
| `structlog` | 结构化日志,高性能 | 推荐 ✅ |
| `loguru` | 简单易用,自动配置 | 小型项目 |
| `logging` | Python 标准库 | 当前已使用 |
### C. 参考链接
- [Python logging cookbook](https://docs.python.org/3/howto/logging.html)
- [ELK Stack 官方文档](https://www.elastic.co/guide/index.html)
- [OpenTelemetry 规范](https://opentelemetry.io/docs/)

View File

@@ -1,22 +0,0 @@
4.2.3 评估工作台
负责训练后模型的质量评测和问题诊断。这是平台闭环的核心环节。
评估流程图
1.评估方式
评估默认使用本地部署的 LLM 作为评审模型,不依赖外部 API。评审模型对每条测试数据从四个子维度打分核心事实正确性、信息完整性、无幻觉、格式合规性可选。汇总为三档判定正确、部分正确、错误。
用户可以选择额外启动人工复核——平台按错误类型分层抽样建议五十到两百条题目,用户在界面上逐条确认或修改 LLM 评审的判定。人工复核的结果用于校准 LLM 评审——一致率超过百分之八十时 LLM 评审结果标记为"可信",低于百分之六十时标记为"以人工为准"。
平台同时提供自动指标作为参考——如关键字段匹配率。这些指标不单独作为判定依据,仅作为快速参考。
2.错误诊断
评审模型在完成评分后,额外输出一个错误分类标签。标签从五种固定类型中选择:混淆(模型回答的值像是另一个实体的属性值)、不完整(事实正确但缺少部分信息)、格式偏差(语义正确但措辞与预期不符)、幻觉(回答中存在标准答案没有的内容)、其他(不属于以上任何类型)。
这种分类方式简单可落地——它是一个固定枚举的分类任务评审模型的prompt 中已包含每种类型的定义和判别示例,不需要额外的自然语言聚类或机器学习算法。
3.评估报告
评估完成后自动生成报告,分为四个部分:总览面板:整体得分和各维度通过率。如果做了人工复核,展示 LLM 评审和人工判定的一致率及可信度标记。
错误分类面板:按五种分类标签分组的错误列表,每组展示数量和占比。点击展开可查看具体错误样例(问题、标准答案、模型回答、评审模型的原因描述)。
修复建议面板:根据错误分类的统计分布,自动生成方向性建议。如"混淆"类错误占比最高时,建议检查训练数据中指令相似但答案不同的样本对;"不完整"类占比最高时,建议统一同类问题的答案详略标准。某类错误的绝对数量不足百分之三时不单独给建议,"其他"类占比最高时提示用户人工分析错误样例。
迭代对比面板:如果存在上轮评估记录,展示两轮各分类标签的数量变化。标注"混淆类错误从35条降到14条下降60%,修复可能生效"。
4.建议有效性的验证
平台在每次评估完成后将分类标签的统计数据(每种标签的数量和占比)存入评估记录。下一轮评估的对应数据与之对比,计算差值和变化百分比。某类错误数量下降超过百分之二十,标注"调整可能生效";变化不足百分之十,标注"调整可能未生效或生效不显著";反向上升则标注"建议检查本次调整方向是否正确"。

View File

@@ -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<OperationLogStats>('/system/operation-logs/stats', params)
/** 获取操作日志中出现的模块列表 */
export const getOperationLogModules = () =>
get<{ value: string; label: string }[]>('/system/operation-logs/modules')

View File

@@ -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' },
],

View File

@@ -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<string, PermissionCode> = {
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',

View File

@@ -0,0 +1,437 @@
<script setup lang="ts">
import { onMounted, reactive, ref, computed } from 'vue'
import { ElMessage } from 'element-plus'
import {
getOperationLogs,
getOperationLogStats,
type OperationLog,
type OperationLogQuery,
type OperationLogStats,
} from '@/api/modules/operation-log'
const loading = ref(false)
const statsLoading = ref(false)
const logs = ref<OperationLog[]>([])
const total = ref(0)
const stats = ref<OperationLogStats | null>(null)
// 默认筛选:只看失败
const query = reactive<OperationLogQuery>({
user_id: '',
module: '',
action: '',
status: 'failure',
keyword: '',
start_time: '',
end_time: '',
limit: 50,
offset: 0,
})
const timeRange = ref<[string, string] | null>(null)
// 模块选项
const moduleOptions = [
{ value: 'fine-tune', label: '模型训练' },
{ value: 'model-eval', label: '模型评测' },
{ value: 'model-inference', label: '模型推理' },
{ value: 'model-manage', label: '模型管理' },
{ value: 'dataset', label: '数据集' },
{ value: 'data-process', label: '数据处理' },
{ value: 'data-convert', label: '数据类型转换' },
{ value: 'compute', label: '算力节点' },
{ value: 'system', label: '系统' },
]
// 动作选项
const actionOptions = [
{ value: 'create', label: '创建' },
{ value: 'start', label: '启动' },
{ value: 'stop', label: '停止' },
{ value: 'delete', label: '删除' },
{ value: 'upload', label: '上传' },
{ value: 'convert', label: '转换' },
{ value: 'merge', label: '合并' },
{ value: 'download', label: '导出' },
{ value: 'import', label: '导入' },
{ value: 'publish', label: '发布' },
{ value: 'request', label: '系统请求' },
]
// 状态选项
const statusOptions = [
{ value: '', label: '全部' },
{ value: 'failure', label: '仅失败' },
{ value: 'success', label: '仅成功' },
]
// 详情弹窗
const detailVisible = ref(false)
const detailLog = ref<OperationLog | null>(null)
// 是否只看失败
const onlyFailures = computed(() => query.status === 'failure')
function applyTimeRange() {
if (timeRange.value && timeRange.value.length === 2) {
query.start_time = timeRange.value[0]
query.end_time = timeRange.value[1]
} else {
query.start_time = ''
query.end_time = ''
}
}
async function load() {
loading.value = true
try {
const params: OperationLogQuery = { ...query }
Object.keys(params).forEach((k) => {
if (params[k as keyof OperationLogQuery] === '' || params[k as keyof OperationLogQuery] === null) {
delete params[k as keyof OperationLogQuery]
}
})
const res = await getOperationLogs(params)
logs.value = res.items
total.value = res.total
} catch {
ElMessage.error('加载操作日志失败')
} finally {
loading.value = false
}
}
async function loadStats() {
statsLoading.value = true
try {
const params: { start_time?: string; end_time?: string } = {}
if (query.start_time) params.start_time = query.start_time
if (query.end_time) params.end_time = query.end_time
stats.value = await getOperationLogStats(params)
} catch {
// 静默失败
} finally {
statsLoading.value = false
}
}
function handleSearch() {
query.offset = 0
load()
loadStats()
}
function handlePageChange(page: number) {
query.offset = (page - 1) * (query.limit || 50)
load()
}
function showDetail(row: OperationLog) {
detailLog.value = row
detailVisible.value = true
}
function statusTagType(status: string) {
return status === 'success' ? 'success' : 'danger'
}
function moduleLabel(module?: string) {
return moduleOptions.find((m) => m.value === module)?.label || module || '-'
}
function actionLabel(action?: string) {
return actionOptions.find((a) => a.value === action)?.label || action || '-'
}
function toggleOnlyFailures() {
query.status = onlyFailures.value ? '' : 'failure'
handleSearch()
}
onMounted(() => {
load()
loadStats()
})
</script>
<template>
<div class="page">
<div class="page-header">
<h2 class="page-title">系统操作日志</h2>
<el-button :type="onlyFailures ? 'danger' : 'default'" @click="toggleOnlyFailures">
{{ onlyFailures ? '只看失败 ' : '显示全部' }}
</el-button>
</div>
<!-- 统计卡片 -->
<el-row :gutter="12" class="stats-row" v-loading="statsLoading">
<el-col :span="4">
<el-card class="stat-card" shadow="hover">
<div class="stat-value">{{ stats?.total ?? '-' }}</div>
<div class="stat-label">总操作数</div>
</el-card>
</el-col>
<el-col :span="4">
<el-card class="stat-card" shadow="hover">
<div class="stat-value" style="color: #67c23a">{{ stats?.success ?? '-' }}</div>
<div class="stat-label">成功</div>
</el-card>
</el-col>
<el-col :span="4">
<el-card class="stat-card" shadow="hover">
<div class="stat-value" style="color: #f56c6c">{{ stats?.failure ?? '-' }}</div>
<div class="stat-label">失败</div>
</el-card>
</el-col>
<el-col :span="4">
<el-card class="stat-card" shadow="hover">
<div class="stat-value" :style="{ color: (stats?.failure_rate ?? 0) > 5 ? '#f56c6c' : '#909399' }">
{{ stats?.failure_rate ?? '-' }}%
</div>
<div class="stat-label">失败率</div>
</el-card>
</el-col>
<el-col :span="8">
<el-card class="stat-card" shadow="hover">
<div class="stat-label" style="margin-bottom: 4px">各模块失败分布</div>
<div class="stat-modules">
<el-tag
v-for="m in stats?.module_failures || []"
:key="m.module"
type="danger"
size="small"
style="margin-right: 4px; margin-bottom: 4px"
>
{{ moduleOptions.find((opt) => opt.value === m.module)?.label || m.module }}: {{ m.count }}
</el-tag>
<span v-if="!stats?.module_failures?.length" style="color: #c0c4cc">暂无失败</span>
</div>
</el-card>
</el-col>
</el-row>
<!-- 筛选区域 -->
<el-card class="filter-card">
<el-form :inline="true">
<el-form-item label="状态">
<el-select v-model="query.status" placeholder="全部状态" clearable style="width: 120px" @change="handleSearch">
<el-option v-for="s in statusOptions" :key="s.value" :label="s.label" :value="s.value" />
</el-select>
</el-form-item>
<el-form-item label="模块">
<el-select v-model="query.module" placeholder="全部模块" clearable style="width: 150px" @change="handleSearch">
<el-option v-for="m in moduleOptions" :key="m.value" :label="m.label" :value="m.value" />
</el-select>
</el-form-item>
<el-form-item label="动作">
<el-select v-model="query.action" placeholder="全部动作" clearable style="width: 120px" @change="handleSearch">
<el-option v-for="a in actionOptions" :key="a.value" :label="a.label" :value="a.value" />
</el-select>
</el-form-item>
<el-form-item label="报错关键字">
<el-input
v-model="query.keyword"
placeholder="搜索报错信息/异常类型"
clearable
style="width: 220px"
@keyup.enter="handleSearch"
/>
</el-form-item>
<el-form-item label="用户">
<el-input v-model="query.user_id" placeholder="用户ID/用户名" clearable style="width: 140px" @keyup.enter="handleSearch" />
</el-form-item>
<el-form-item label="时间范围">
<el-date-picker
v-model="timeRange"
type="datetimerange"
value-format="YYYY-MM-DDTHH:mm:ss"
range-separator=""
start-placeholder="开始时间"
end-placeholder="结束时间"
clearable
style="width: 360px"
@change="handleSearch"
/>
</el-form-item>
<el-form-item>
<el-button type="primary" @click="handleSearch">查询</el-button>
</el-form-item>
</el-form>
</el-card>
<!-- 日志列表 -->
<el-table
:data="logs"
v-loading="loading"
border
stripe
class="log-table"
:row-class-name="({ row }) => row.status === 'failure' ? 'failure-row' : ''"
>
<el-table-column prop="create_time" label="时间" width="170" align="center">
<template #default="{ row }">
{{ row.create_time ? new Date(row.create_time).toLocaleString('zh-CN') : '-' }}
</template>
</el-table-column>
<el-table-column label="用户" width="110" align="center" show-overflow-tooltip>
<template #default="{ row }">{{ row.username || row.user_id || '-' }}</template>
</el-table-column>
<el-table-column label="模块" width="110" align="center">
<template #default="{ row }">{{ moduleLabel(row.module) }}</template>
</el-table-column>
<el-table-column label="动作" width="90" align="center">
<template #default="{ row }">{{ actionLabel(row.action) }}</template>
</el-table-column>
<el-table-column label="目标" min-width="140" align="center" show-overflow-tooltip>
<template #default="{ row }">{{ row.target_name || row.target_id || '-' }}</template>
</el-table-column>
<el-table-column label="状态" width="80" align="center">
<template #default="{ row }">
<el-tag :type="statusTagType(row.status)" size="small" effect="dark">
{{ row.status === 'success' ? '成功' : '失败' }}
</el-tag>
</template>
</el-table-column>
<el-table-column label="异常类型" width="160" align="center" show-overflow-tooltip>
<template #default="{ row }">
<el-tag v-if="row.error_type" type="danger" size="small" effect="plain">{{ row.error_type }}</el-tag>
<span v-else style="color: #c0c4cc">-</span>
</template>
</el-table-column>
<el-table-column label="报错信息" min-width="250" show-overflow-tooltip>
<template #default="{ row }">
<span v-if="row.error_message" class="error-text">{{ row.error_message }}</span>
<span v-else style="color: #c0c4cc">-</span>
</template>
</el-table-column>
<el-table-column label="耗时(ms)" width="90" align="center">
<template #default="{ row }">
<span :style="{ color: row.duration_ms > 3000 ? '#e6a23c' : '' }">
{{ row.duration_ms ? row.duration_ms : '-' }}
</span>
</template>
</el-table-column>
<el-table-column label="操作" width="80" align="center" fixed="right">
<template #default="{ row }">
<el-button link type="primary" size="small" @click="showDetail(row)">详情</el-button>
</template>
</el-table-column>
</el-table>
<!-- 分页 -->
<div class="pager">
<el-pagination
background
layout="total, prev, pager, next"
:total="total"
:page-size="query.limit"
:current-page="Math.floor((query.offset || 0) / (query.limit || 50)) + 1"
@current-change="handlePageChange"
/>
</div>
<!-- 详情弹窗 -->
<el-dialog v-model="detailVisible" title="操作日志详情" width="850px" top="5vh">
<el-descriptions :column="2" border v-if="detailLog">
<el-descriptions-item label="时间">
{{ detailLog.create_time ? new Date(detailLog.create_time).toLocaleString('zh-CN') : '-' }}
</el-descriptions-item>
<el-descriptions-item label="状态">
<el-tag :type="statusTagType(detailLog.status)" size="small" effect="dark">
{{ detailLog.status === 'success' ? '成功' : '失败' }}
</el-tag>
</el-descriptions-item>
<el-descriptions-item label="用户">{{ detailLog.username || detailLog.user_id || '-' }}</el-descriptions-item>
<el-descriptions-item label="IP">{{ detailLog.client_ip || '-' }}</el-descriptions-item>
<el-descriptions-item label="模块">{{ moduleLabel(detailLog.module) }}</el-descriptions-item>
<el-descriptions-item label="动作">{{ actionLabel(detailLog.action) }}</el-descriptions-item>
<el-descriptions-item label="目标类型">{{ detailLog.target_type || '-' }}</el-descriptions-item>
<el-descriptions-item label="目标ID">{{ detailLog.target_id || '-' }}</el-descriptions-item>
<el-descriptions-item label="目标名称">{{ detailLog.target_name || '-' }}</el-descriptions-item>
<el-descriptions-item label="耗时">{{ detailLog.duration_ms ? detailLog.duration_ms + ' ms' : '-' }}</el-descriptions-item>
<el-descriptions-item label="出错函数" :span="2">
<span v-if="detailLog.func_name" class="func-name">{{ detailLog.func_name }}</span>
<span v-else style="color: #c0c4cc">-</span>
</el-descriptions-item>
<el-descriptions-item label="请求方法">{{ detailLog.request_method || '-' }}</el-descriptions-item>
<el-descriptions-item label="请求路径">{{ detailLog.request_path || '-' }}</el-descriptions-item>
<el-descriptions-item label="Trace ID" :span="2">{{ detailLog.trace_id || '-' }}</el-descriptions-item>
<el-descriptions-item label="异常类型" :span="2">
<el-tag v-if="detailLog.error_type" type="danger" size="small" effect="dark">{{ detailLog.error_type }}</el-tag>
<span v-else style="color: #c0c4cc">-</span>
</el-descriptions-item>
<el-descriptions-item label="报错信息" :span="2">
<div v-if="detailLog.error_message" class="error-box">{{ detailLog.error_message }}</div>
<span v-else style="color: #c0c4cc">-</span>
</el-descriptions-item>
<el-descriptions-item label="异常堆栈 (Traceback)" :span="2">
<pre v-if="detailLog.error_traceback" class="traceback-box">{{ detailLog.error_traceback }}</pre>
<span v-else style="color: #c0c4cc">无堆栈信息</span>
</el-descriptions-item>
<el-descriptions-item label="操作详情" :span="2">
<pre v-if="detailLog.detail" class="detail-box">{{ detailLog.detail }}</pre>
<span v-else style="color: #c0c4cc">-</span>
</el-descriptions-item>
</el-descriptions>
</el-dialog>
</div>
</template>
<style scoped lang="scss">
.page { padding: 16px; }
.page-header { display: flex; align-items: center; justify-content: space-between; margin-bottom: 16px; }
.page-title { margin: 0; font-size: 18px; }
.stats-row { margin-bottom: 16px; }
.stat-card {
text-align: center;
.stat-value { font-size: 24px; font-weight: bold; }
.stat-label { font-size: 12px; color: #909399; }
.stat-modules { text-align: left; min-height: 40px; }
}
.filter-card { margin-bottom: 16px; }
.log-table { margin-top: 8px; }
.pager { margin-top: 12px; text-align: right; }
.error-text { color: #f56c6c; font-weight: 500; }
.func-name {
font-family: 'Courier New', monospace;
font-size: 12px;
color: #e6a23c;
background: #fdf6ec;
padding: 2px 6px;
border-radius: 3px;
}
.error-box {
color: #f56c6c;
background: #fef0f0;
padding: 8px 12px;
border-radius: 4px;
font-size: 13px;
word-break: break-all;
}
.traceback-box {
background: #2d2d2d;
color: #f48771;
padding: 12px;
border-radius: 4px;
font-size: 12px;
font-family: 'Courier New', monospace;
white-space: pre-wrap;
word-break: break-all;
max-height: 300px;
overflow-y: auto;
}
.detail-box {
background: #f5f7fa;
padding: 8px 12px;
border-radius: 4px;
font-size: 12px;
white-space: pre-wrap;
word-break: break-all;
}
:deep(.failure-row) {
background-color: #fef0f0 !important;
}
:deep(.failure-row:hover > td) {
background-color: #fde2e2 !important;
}
</style>

View File

@@ -2,7 +2,7 @@
import { onMounted, ref } from 'vue'
import { ElMessage, ElMessageBox } from 'element-plus'
import { Plus, Delete, Refresh } from '@element-plus/icons-vue'
import type { TagProps, UploadRequestOptions } from 'element-plus'
import type { TagProps, UploadRequestOptions, UploadFile } from 'element-plus'
import PageCard from '@/components/PageCard.vue'
import {
getDataConvertTasks,
@@ -16,6 +16,8 @@ const loading = ref(false)
const tasks = ref<DataConvertTask[]>([])
const showCreate = ref(false)
const form = ref({ name: '', outputName: 'converted-data' })
const fileList = ref<UploadFile[]>([])
const creating = ref(false)
async function load() {
loading.value = true
@@ -27,21 +29,91 @@ async function load() {
}
}
// 文件上传前的校验(仅校验文件格式)
function beforeUpload(file: UploadFile) {
// 检查文件类型
const isJson = file.name.endsWith('.json') || file.raw?.type === 'application/json'
if (!isJson) {
ElMessage.error('只能上传 .json 格式的文件')
return false
}
return true
}
// 手动点击"创建并上传"
async function submitCreate() {
if (!form.value.name) {
// 校验任务名称
if (!form.value.name || !form.value.name.trim()) {
ElMessage.warning('请填写任务名称')
return
}
await createDataConvertTask({
name: form.value.name,
output_filename: form.value.outputName + '.jsonl',
})
ElMessage.success('任务创建成功')
showCreate.value = false
form.value = { name: '', outputName: 'converted-data' }
load()
// 检查名称是否重复
const exists = tasks.value.some((t) => t.name === form.value.name.trim())
if (exists) {
ElMessage.error(`数据集管理中已存在名为「${form.value.name}」的任务,请换一个名称`)
return
}
// 检查是否选择了文件
if (!fileList.value || fileList.value.length === 0) {
ElMessage.warning('请选择要上传的 JSON 文件')
return
}
creating.value = true
try {
// 1. 创建转换任务
const task = await createDataConvertTask({
name: form.value.name.trim(),
output_filename: form.value.outputName.trim() + '.jsonl',
})
// 2. 获取任务 ID
const taskId = (task as any)?.id || task?.id || (task as any)?.data?.id
if (!taskId) {
throw new Error('创建任务失败,服务端未返回任务 ID')
}
// 3. 上传文件到刚创建的任务
const file = fileList.value[0].raw
if (!file) {
throw new Error('文件信息丢失,请重新选择文件')
}
const res = await uploadSourceFiles(taskId, [file])
const data = (res as any)?.data || res
if (data?.auto_converted) {
ElMessage.success(
`创建成功!文件已上传并自动转换完成(输入 ${data.input_count} 条 / 输出 ${data.output_count} 条),结果已导入数据集`,
)
} else if (data?.error) {
ElMessage.error(`文件上传成功但转换失败:${data.error}`)
} else {
ElMessage.warning('文件已上传,等待后台转换处理...')
}
// 关闭弹窗并刷新列表
showCreate.value = false
resetForm()
load()
} catch (e: any) {
// 根据错误类型给出更清晰的提示
const msg = e?.message || e?.toString() || '未知错误'
if (msg.includes('409') || msg.includes('conflict') || msg.includes('已存在') || msg.includes('duplicate')) {
ElMessage.error(`数据集管理中已存在名为「${form.value.name}」的任务,请换一个名称`)
} else if (msg.includes('400')) {
ElMessage.error(`参数错误:${msg}`)
} else if (msg.includes('403') || msg.includes('权限')) {
ElMessage.error(`没有权限执行此操作:${msg}`)
} else {
ElMessage.error(`操作失败:${msg}`)
}
} finally {
creating.value = false
}
}
// 自定义上传(表格中的上传按钮仍使用此方法)
async function customUpload(options: UploadRequestOptions) {
const taskId = options.data?.taskId as string
if (!taskId) {
@@ -59,11 +131,22 @@ async function customUpload(options: UploadRequestOptions) {
ElMessage.warning('上传完成,但转换失败:' + (data?.error || '未知错误'))
}
load()
} catch {
ElMessage.error('上传失败')
} catch (e: any) {
ElMessage.error('上传失败' + (e?.message || '未知错误'))
}
}
// 重置表单
function resetForm() {
form.value = { name: '', outputName: 'converted-data' }
fileList.value = []
}
// 移除已选文件
function handleRemoveFile(file: UploadFile) {
fileList.value = fileList.value.filter((f) => f.uid !== file.uid)
}
async function handleDelete(task: DataConvertTask) {
try {
await ElMessageBox.confirm(
@@ -136,21 +219,42 @@ onMounted(load)
</el-table-column>
</el-table>
<!-- 新建任务弹窗 -->
<el-dialog v-model="showCreate" title="新建转换任务" width="480px">
<!-- 新建任务弹窗一步到位填写信息 + 上传文件 -->
<el-dialog v-model="showCreate" title="新建转换任务" width="520px" :close-on-click-modal="false">
<el-form label-width="100px">
<el-form-item label="任务名称" required>
<el-input v-model="form.name" placeholder="请输入任务名称" />
<el-input v-model="form.name" placeholder="请输入任务名称" clearable />
</el-form-item>
<el-form-item label="输出文件名">
<el-input v-model="form.outputName" placeholder="converted-data">
<el-input v-model="form.outputName" placeholder="converted-data" clearable>
<template #append>.jsonl</template>
</el-input>
</el-form-item>
<el-form-item label="上传文件" required>
<el-upload
ref="uploadRef"
v-model:file-list="fileList"
:auto-upload="false"
:limit="1"
accept=".json"
:before-upload="beforeUpload"
:on-remove="handleRemoveFile"
:disabled="creating"
drag
>
<el-icon class="el-icon--upload"><Plus /></el-icon>
<div class="el-upload__text"> JSON 文件拖到此处<em>点击上传</em></div>
<template #tip>
<div class="el-upload__tip">仅支持 .json 格式文件点击"创建并上传"按钮后自动创建任务并上传文件</div>
</template>
</el-upload>
</el-form-item>
</el-form>
<template #footer>
<el-button @click="showCreate = false">取消</el-button>
<el-button type="primary" @click="submitCreate">创建</el-button>
<el-button @click="showCreate = false" :disabled="creating">取消</el-button>
<el-button type="primary" :loading="creating" :disabled="!form.name || fileList.length === 0" @click="submitCreate">
{{ creating ? '处理中...' : '创建并上传' }}
</el-button>
</template>
</el-dialog>
</PageCard>
@@ -162,4 +266,10 @@ onMounted(load)
gap: 10px;
margin-bottom: 16px;
}
:deep(.el-upload-dragger) {
width: 100%;
.el-upload__text {
padding: 20px 0;
}
}
</style>

View File

@@ -1,4 +1,4 @@
<script setup lang="ts">
<script setup lang="ts">
import { ref, reactive, computed, onMounted } from 'vue'
import { useRouter } from 'vue-router'
import { ElMessage, type FormInstance, type FormRules } from 'element-plus'
@@ -272,9 +272,14 @@ async function handleSubmit() {
await startFineTune({ ...payload, task_id: taskId })
ElMessage.success('训练任务已创建并启动')
} catch {
await updateFineTune(taskId, { status: 'failed' })
try {
await updateFineTune(taskId, { status: 'failed' })
} catch {
// 忽略状态更新失败
}
ElMessage.error('任务已创建,但训练启动失败')
}
// 无论启动成功还是失败,都跳转到训练列表
router.push('/fine-tune')
} catch {
ElMessage.error('训练任务创建失败,请稍后重试')

View File

@@ -3,7 +3,7 @@ import { reactive, ref } from 'vue'
import { useRouter } from 'vue-router'
import { ElMessage } from 'element-plus'
import { createUser } from '@/api/modules/system'
import type { CreateUserPayload, PermissionCode } from '@/types'
import type { CreateUserPayload } from '@/types'
const router = useRouter()
const submitting = ref(false)
@@ -12,26 +12,11 @@ const form = reactive<CreateUserPayload>({
username: '',
display_name: '',
password: 'platform123',
role: 'viewer',
role: 'user',
status: 'active',
permissions: ['dashboard'],
permissions: [],
})
const permissionOptions: PermissionCode[] = [
'dashboard',
'fine-tune',
'model-eval',
'model-inference',
'model-manage',
'dataset',
'data-process',
'data-convert',
'compute',
'hardware',
'logs',
'user-settings',
]
async function submit() {
submitting.value = true
try {
@@ -49,19 +34,18 @@ async function submit() {
<h1>创建用户</h1>
<el-form :model="form" label-width="110px" class="user-form">
<el-form-item label="账号">
<el-input v-model="form.username" />
<el-input v-model="form.username" placeholder="登录用户名" />
</el-form-item>
<el-form-item label="显示名称">
<el-input v-model="form.display_name" />
<el-input v-model="form.display_name" placeholder="如:张三" />
</el-form-item>
<el-form-item label="初始密码">
<el-input v-model="form.password" type="password" show-password />
</el-form-item>
<el-form-item label="角色">
<el-select v-model="form.role">
<el-select v-model="form.role" style="width: 100%">
<el-option label="管理员" value="admin" />
<el-option label="操作员" value="operator" />
<el-option label="观察员" value="viewer" />
<el-option label="普通用户" value="user" />
</el-select>
</el-form-item>
<el-form-item label="状态">
@@ -70,10 +54,13 @@ async function submit() {
<el-radio value="disabled">禁用</el-radio>
</el-radio-group>
</el-form-item>
<el-form-item label="页面权限">
<el-checkbox-group v-model="form.permissions">
<el-checkbox v-for="item in permissionOptions" :key="item" :value="item">{{ item }}</el-checkbox>
</el-checkbox-group>
<el-form-item label="权限说明">
<el-alert type="info" :closable="false" show-icon>
<template #title>
<span v-if="form.role === 'admin'">管理员拥有全部权限包括用户管理平台治理算力节点</span>
<span v-else>普通用户可见服务看板模型服务数据治理其他工具平台性能查看日志数据集和微调模型仅创建者和被授权用户可见</span>
</template>
</el-alert>
</el-form-item>
<el-form-item>
<el-button @click="router.back()">返回</el-button>
@@ -89,6 +76,6 @@ async function submit() {
}
.user-form {
max-width: 760px;
max-width: 640px;
}
</style>

View File

@@ -2,7 +2,6 @@
import { onMounted, reactive, ref } from 'vue'
import { ElMessage, ElMessageBox } from 'element-plus'
import {
changeMyPassword,
deleteUser,
getUsers,
resetUserPassword,
@@ -81,39 +80,11 @@ async function confirmResetPwd() {
}
}
// ---------- 用户自改密码 ----------
const myPwdDialog = reactive({ visible: false, oldPassword: '', newPassword: '', saving: false })
function openChangeMyPwd() {
myPwdDialog.oldPassword = ''
myPwdDialog.newPassword = ''
myPwdDialog.visible = true
}
async function confirmChangeMyPwd() {
if (!myPwdDialog.oldPassword.trim() || !myPwdDialog.newPassword.trim()) {
ElMessage.warning('请填写旧密码和新密码')
return
}
if (myPwdDialog.newPassword.length < 6) {
ElMessage.warning('新密码至少 6 位')
return
}
myPwdDialog.saving = true
try {
await changeMyPassword(myPwdDialog.oldPassword.trim(), myPwdDialog.newPassword.trim())
ElMessage.success('密码修改成功')
myPwdDialog.visible = false
} catch {
ElMessage.error('密码修改失败,请检查旧密码是否正确')
} finally {
myPwdDialog.saving = false
}
}
// ---------- 删除 ----------
async function removeUser(row: SystemUser) {
try {
await ElMessageBox.confirm(
`确定删除用户 ${row.display_name}${row.username} 吗?该操作不可恢复。`,
`确定删除用户 "${row.display_name}${row.username}" 吗?该操作不可恢复。`,
'删除用户',
{ type: 'warning', confirmButtonText: '删除', cancelButtonText: '取消' },
)
@@ -139,7 +110,6 @@ async function removeUser(row: SystemUser) {
<p>管理平台账号角色状态与登录密码</p>
</div>
<div>
<el-button @click="openChangeMyPwd">修改密码</el-button>
<el-button type="primary" @click="$router.push('/user-settings/create')">创建用户</el-button>
</div>
</header>
@@ -147,7 +117,13 @@ async function removeUser(row: SystemUser) {
<el-table :data="users" border>
<el-table-column prop="username" label="账号" min-width="140" />
<el-table-column prop="display_name" label="显示名称" min-width="160" />
<el-table-column prop="role" label="角色" width="120" />
<el-table-column prop="role" label="角色" width="120">
<template #default="{ row }">
<el-tag :type="asSystemUser(row).role === 'admin' ? 'danger' : 'info'" size="small">
{{ asSystemUser(row).role === 'admin' ? '管理员' : '普通用户' }}
</el-tag>
</template>
</el-table-column>
<el-table-column label="状态" width="130">
<template #default="{ row }">
<el-tag :type="statusTagType(asSystemUser(row).status)" size="small">{{ statusLabel(asSystemUser(row).status) }}</el-tag>
@@ -189,22 +165,6 @@ async function removeUser(row: SystemUser) {
<el-button type="primary" :loading="pwdDialog.saving" @click="confirmResetPwd">确定重置</el-button>
</template>
</el-dialog>
<!-- 修改自己的密码 -->
<el-dialog v-model="myPwdDialog.visible" title="修改密码" width="420px">
<el-form label-width="80px">
<el-form-item label="旧密码">
<el-input v-model="myPwdDialog.oldPassword" placeholder="请输入当前密码" show-password />
</el-form-item>
<el-form-item label="新密码">
<el-input v-model="myPwdDialog.newPassword" placeholder="至少 6 位" show-password />
</el-form-item>
</el-form>
<template #footer>
<el-button @click="myPwdDialog.visible = false">取消</el-button>
<el-button type="primary" :loading="myPwdDialog.saving" @click="confirmChangeMyPwd">确认修改</el-button>
</template>
</el-dialog>
</section>
</template>