Appearance
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, ...)