Files
yu-zhi-ran/platform/backend/app/api/tasks.py
T
Yuzhiran Dev 233e23016c feat: 内容数据迁移至数据库,合规审查全链路打通
- 文章 HTML 存储从文件系统迁移至 articles 表,删除 releases 目录
- 合规审查从 DB 读取 HTML,审查结果写回 DB,通过后自动推进至待发布
- 新增 todayCount 筛选按钮,与系统概览统计数据一致
- 全屏预览修复:提升 z-index 超过侧边栏,添加退出全屏/关闭按钮
- 统一 '优化' → '审查' 命名,消除前后端术语不一致
- 调度器创作完成后自动触发审查(生成 → 审查 → 待发布)
- 清理旧备份/调试文件、过期大纲和研究笔记
2026-05-13 17:33:56 +08:00

237 lines
6.6 KiB
Python

import uuid
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.orm import Session
from typing import List, Optional
from ..database import get_db
from ..models import ContentTask, Topic
from ..schemas import ContentTaskCreate, ContentTaskResponse
from .auth import get_current_user
router = APIRouter(prefix="/api/tasks", tags=["tasks"])
@router.get("", response_model=List[ContentTaskResponse])
def list_tasks(
status: Optional[str] = None,
topic_id: Optional[str] = None,
stage: Optional[str] = None,
limit: int = Query(50, le=200),
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
query = db.query(ContentTask)
if status:
query = query.filter(ContentTask.status == status)
if topic_id:
query = query.filter(ContentTask.topic_id == topic_id)
if stage:
query = query.filter(ContentTask.stage == stage)
return query.order_by(ContentTask.created_at.desc()).limit(limit).all()
@router.get("/active")
def get_active_tasks(
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
return db.query(ContentTask).filter(
ContentTask.status == "running"
).order_by(ContentTask.started_at.desc()).all()
@router.post("", response_model=ContentTaskResponse)
def create_task(
data: ContentTaskCreate,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
if data.topic_id:
topic = db.query(Topic).filter(Topic.id == data.topic_id).first()
if not topic:
raise HTTPException(status_code=404, detail="选题不存在")
task_id = f"task_{uuid.uuid4().hex[:16]}"
task = ContentTask(
task_id=task_id,
topic_id=data.topic_id,
stage=data.stage,
status="pending",
created_by=data.created_by or current_user.username
)
db.add(task)
db.commit()
db.refresh(task)
return task
@router.get("/{task_id}", response_model=ContentTaskResponse)
def get_task(
task_id: str,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
task = db.query(ContentTask).filter(ContentTask.task_id == task_id).first()
if not task:
raise HTTPException(status_code=404, detail="任务不存在")
return task
@router.put("/{task_id}/start")
def start_task(
task_id: str,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
task = db.query(ContentTask).filter(ContentTask.task_id == task_id).first()
if not task:
raise HTTPException(status_code=404, detail="任务不存在")
from datetime import datetime, timezone
task.status = "running"
task.started_at = datetime.now(timezone.utc)
task.message = "任务已启动"
task.progress = 0
db.commit()
db.refresh(task)
return task
@router.put("/{task_id}/progress")
def update_progress(
task_id: str,
progress: int = Query(..., ge=0, le=100),
message: Optional[str] = None,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
task = db.query(ContentTask).filter(ContentTask.task_id == task_id).first()
if not task:
raise HTTPException(status_code=404, detail="任务不存在")
task.progress = progress
if message:
task.message = message
db.commit()
db.refresh(task)
return task
@router.put("/{task_id}/complete")
def complete_task(
task_id: str,
result_data: dict = None,
message: Optional[str] = None,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
task = db.query(ContentTask).filter(ContentTask.task_id == task_id).first()
if not task:
raise HTTPException(status_code=404, detail="任务不存在")
from datetime import datetime, timezone
finished = datetime.now(timezone.utc)
task.status = "completed"
task.finished_at = finished
task.progress = 100
if message:
task.message = message
if result_data:
task.result_data = result_data
if task.started_at:
task.duration = int((finished - task.started_at).total_seconds())
db.commit()
db.refresh(task)
return task
@router.put("/{task_id}/fail")
def fail_task(
task_id: str,
error_msg: str = Query(...),
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
task = db.query(ContentTask).filter(ContentTask.task_id == task_id).first()
if not task:
raise HTTPException(status_code=404, detail="任务不存在")
from datetime import datetime, timezone
finished = datetime.now(timezone.utc)
task.status = "failed"
task.finished_at = finished
task.error_msg = error_msg
if task.started_at:
task.duration = int((finished - task.started_at).total_seconds())
db.commit()
db.refresh(task)
return task
@router.delete("/{task_id}")
def cancel_task(
task_id: str,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
task = db.query(ContentTask).filter(ContentTask.task_id == task_id).first()
if not task:
raise HTTPException(status_code=404, detail="任务不存在")
task.status = "cancelled"
db.commit()
return {"ok": True}
@router.post("/run-creator")
def run_creator_task(
topic_id: Optional[str] = None,
db: Session = Depends(get_db),
current_user=Depends(get_current_user)
):
import threading
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
task_id = f"task_{uuid.uuid4().hex[:16]}"
task = ContentTask(
task_id=task_id,
topic_id=topic_id,
stage="creator",
status="running",
started_at=now,
created_by=current_user.username
)
db.add(task)
db.commit()
db.refresh(task)
from ..core.generator import run_creator
def _run():
try:
result = run_creator(topic_id)
finished = datetime.now(timezone.utc)
task.status = "completed"
task.finished_at = finished
task.progress = 100
task.message = "创作完成"
task.result_data = result or {}
if task.started_at:
task.duration = int((finished - task.started_at).total_seconds())
db.commit()
except Exception as e:
finished = datetime.now(timezone.utc)
task.status = "failed"
task.finished_at = finished
task.error_msg = str(e)
if task.started_at:
task.duration = int((finished - task.started_at).total_seconds())
db.commit()
thread = threading.Thread(target=_run)
thread.start()
return task