Skip to content

54|后端接口实现(下)

上半部分实现了用户认证和资产查询/创建。下半部分完成资产更新/删除、任务管理、事件日志,以及审计功能。

一、资产更新与删除

python
# routers/assets.py

@router.put("/{asset_id}", response_model=AssetOut)
async def update_asset(
    asset_id: int,
    data: AssetUpdate,
    db: Session = Depends(get_db),
    user: User = Depends(get_current_user),
):
    asset = db.get(Asset, asset_id)
    if not asset:
        raise HTTPException(status_code=404, detail="资产不存在")
    if user.role != "admin" and asset.owner_id != user.id:
        raise HTTPException(status_code=403, detail="无权限")

    for field, value in data.dict(exclude_unset=True).items():
        setattr(asset, field, value)

    db.commit()
    db.refresh(asset)
    return asset

@router.delete("/{asset_id}")
async def delete_asset(
    asset_id: int,
    db: Session = Depends(get_db),
    user: User = Depends(require_admin),
):
    asset = db.get(Asset, asset_id)
    if not asset:
        raise HTTPException(status_code=404, detail="资产不存在")

    db.delete(asset)
    db.commit()
    return {"deleted": asset_id}

data.dict(exclude_unset=True) 只返回客户端传了的字段,实现部分更新效果(即使路由装饰器是 PUT)。

二、任务管理接口

python
# routers/tasks.py
from fastapi import APIRouter, Depends, HTTPException, BackgroundTasks
from sqlalchemy.orm import Session
from sqlalchemy import select

from app.database import get_db
from app.models import Task, Asset
from app.schemas import TaskCreate, TaskOut
from app.routers.users import get_current_user

router = APIRouter(prefix="/api/tasks", tags=["tasks"])

@router.get("/", response_model=list[TaskOut])
async def list_tasks(
    asset_id: int | None = None,
    db: Session = Depends(get_db),
    user: User = Depends(get_current_user),
):
    stmt = select(Task)
    if asset_id:
        stmt = stmt.where(Task.asset_id == asset_id)
    if user.role != "admin":
        stmt = stmt.join(Asset).where(Asset.owner_id == user.id)

    return db.execute(stmt).scalars().all()

@router.post("/", response_model=TaskOut, status_code=201)
async def create_task(
    data: TaskCreate,
    background_tasks: BackgroundTasks,
    db: Session = Depends(get_db),
    user: User = Depends(get_current_user),
):
    asset = db.get(Asset, data.asset_id)
    if not asset:
        raise HTTPException(status_code=404, detail="资产不存在")
    if user.role != "admin" and asset.owner_id != user.id:
        raise HTTPException(status_code=403, detail="无权限")

    task = Task(
        asset_id=data.asset_id,
        action=data.action,
        status="pending",
        created_by=user.id,
    )
    db.add(task)
    db.commit()
    db.refresh(task)

    background_tasks.add_task(execute_task, task.id)
    return task

def execute_task(task_id: int):
    """模拟任务执行"""
    import time
    from app.database import SessionLocal

    db = SessionLocal()
    try:
        task = db.get(Task, task_id)
        task.status = "running"
        db.commit()

        time.sleep(3)   # 模拟执行

        task.status = "done"
        task.result = "执行成功"
        task.completed_at = datetime.utcnow()
        db.commit()
    except Exception as e:
        task.status = "failed"
        task.result = str(e)
        db.commit()
    finally:
        db.close()

三、事件日志接口

python
# routers/events.py
from fastapi import APIRouter, Depends
from sqlalchemy.orm import Session
from sqlalchemy import select

from app.database import get_db
from app.models import Event
from app.routers.users import require_admin

router = APIRouter(prefix="/api/events", tags=["events"])

@router.get("/")
async def list_events(
    skip: int = 0,
    limit: int = 50,
    db: Session = Depends(get_db),
    user: User = Depends(require_admin),
):
    stmt = select(Event).order_by(Event.created_at.desc()).offset(skip).limit(limit)
    return db.execute(stmt).scalars().all()

四、审计日志中间件

python
# app/audit.py
from fastapi import Request
from sqlalchemy.orm import Session
from app.models import Event

async def log_event(
    db: Session,
    user_id: int,
    action: str,
    target_type: str,
    target_id: int,
    detail: str = "",
):
    event = Event(
        user_id=user_id,
        action=action,
        target_type=target_type,
        target_id=target_id,
        detail=detail,
    )
    db.add(event)
    db.commit()

在资产删除等敏感操作中调用:

python
@router.delete("/{asset_id}")
async def delete_asset(
    asset_id: int,
    db: Session = Depends(get_db),
    user: User = Depends(require_admin),
):
    asset = db.get(Asset, asset_id)
    ...

    await log_event(
        db=db,
        user_id=user.id,
        action="delete_asset",
        target_type="asset",
        target_id=asset_id,
        detail=f"删除资产 {asset.hostname}",
    )
    return {"deleted": asset_id}

五、常见错误

后台任务中直接使用路由的 db session

python
# 错误:请求结束后 session 已关闭
background_tasks.add_task(execute_task, task.id, db)

# 正确:后台任务自己创建 session
def execute_task(task_id: int):
    db = SessionLocal()
    try:
        ...
    finally:
        db.close()

事件日志写入失败影响主业务

python
# 错误:日志写入异常导致资产删除回滚
try:
    db.delete(asset)
    db.commit()
    log_event(...)   # 如果这里出错,前面已经 commit 了,无法回滚
except:
    db.rollback()

# 正确:日志失败不阻塞主业务,或用后台任务写日志
db.delete(asset)
db.commit()
background_tasks.add_task(log_event_async, ...)