8935196fcd
Deploy / Deploy Staging (push) Failing after 138h4m33s
CI/CD Pipeline / Frontend Lint (push) Failing after 138h4m39s
CI/CD Pipeline / Validate Code Quality And Tests (push) Failing after 138h4m39s
Deploy / Staging E2E Tests (push) Failing after 1796h40m0s
Deploy / Production Browser E2E (push) Failing after 1796h40m1s
Deploy / Build Production Runtime Images (push) Failing after 1796h40m6s
Deploy / Deploy Production (push) Failing after 1797h11m28s
333 lines
12 KiB
Python
Executable File
333 lines
12 KiB
Python
Executable File
"""Job API 路由 — Phase 8 任务 2.10.
|
||
|
||
提供统一异步任务管理 RESTful 接口:
|
||
- POST /api/v1/jobs 创建任务
|
||
- GET /api/v1/jobs/{job_id} 任务详情
|
||
- GET /api/v1/projects/{project_id}/jobs 项目任务列表
|
||
- GET /api/v1/projects/{project_id}/jobs/stats 任务统计
|
||
- PUT /api/v1/jobs/{job_id}/progress 更新进度
|
||
- POST /api/v1/jobs/{job_id}/complete 标记完成
|
||
- POST /api/v1/jobs/{job_id}/fail 标记失败
|
||
- POST /api/v1/jobs/{job_id}/retry 重试任务
|
||
- POST /api/v1/jobs/{job_id}/cancel 取消任务
|
||
- POST /api/v1/jobs/{job_id}/submit 提交执行
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from typing import Any
|
||
|
||
from app.auth import AuthenticatedUser, get_current_user
|
||
from app.core.celery_app import celery_app
|
||
from app.dependencies import get_db_session, get_job_repository, get_project_repository
|
||
from app.schemas.job import (
|
||
CompleteJobRequest,
|
||
CreateJobRequest,
|
||
FailJobRequest,
|
||
JobResponse,
|
||
JobStatisticsResponse,
|
||
ListJobsResponse,
|
||
UpdateProgressRequest,
|
||
job_to_response,
|
||
)
|
||
from fastapi import APIRouter, Depends, HTTPException, Query, status
|
||
|
||
from packages.application.jobs import (
|
||
CancelJobUseCase,
|
||
CompleteJobCommand,
|
||
CompleteJobUseCase,
|
||
CreateJobCommand,
|
||
CreateJobUseCase,
|
||
FailJobCommand,
|
||
FailJobUseCase,
|
||
GetJobStatisticsUseCase,
|
||
GetJobUseCase,
|
||
ListJobsUseCase,
|
||
RetryJobUseCase,
|
||
SubmitJobUseCase,
|
||
UpdateJobProgressCommand,
|
||
UpdateJobProgressUseCase,
|
||
)
|
||
from packages.domain.job import JobType
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter()
|
||
|
||
# 任务类型 → Celery task name 映射
|
||
_JOB_TYPE_TO_CELERY_TASK: dict[str, str] = {
|
||
JobType.VIDEO_COMPOSE: "worker.compose_video",
|
||
JobType.RENDER_EDIT_PLAN: "worker.render_edit_plan",
|
||
JobType.ASSET_INGEST: "worker.ingest_asset",
|
||
JobType.CLASSIFICATION: "worker.classify_asset",
|
||
JobType.VOICE_EXTRACTION: "worker.extract_voice",
|
||
JobType.GENERATION: "worker.generate_video",
|
||
}
|
||
|
||
|
||
def _check_project_access(project_id: str, user_id: str, project_repository) -> None:
|
||
"""检查用户是否有项目访问权限。"""
|
||
project = project_repository.find_by_id(project_id)
|
||
if project is None:
|
||
raise HTTPException(status_code=404, detail=f"Project {project_id} not found")
|
||
if not project.can_access(user_id):
|
||
raise HTTPException(status_code=403, detail="Access denied to project")
|
||
|
||
|
||
# ── 创建任务 ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.post("/jobs", response_model=JobResponse, status_code=status.HTTP_201_CREATED)
|
||
def create_job(
|
||
request: CreateJobRequest,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
project_repository: Any = Depends(get_project_repository),
|
||
) -> JobResponse:
|
||
"""创建异步任务。
|
||
|
||
创建后任务处于 pending 状态,需要调用 /submit 提交执行。
|
||
"""
|
||
_check_project_access(request.project_id, authenticated_user.user.id, project_repository)
|
||
|
||
# 校验 job_type
|
||
try:
|
||
JobType(request.job_type)
|
||
except ValueError:
|
||
raise HTTPException(
|
||
status_code=400,
|
||
detail=f"不支持的任务类型: {request.job_type}," f"可选值: {[t.value for t in JobType]}",
|
||
)
|
||
|
||
use_case = CreateJobUseCase(job_repo)
|
||
job = use_case.execute(
|
||
CreateJobCommand(
|
||
project_id=request.project_id,
|
||
job_type=request.job_type,
|
||
payload=request.payload,
|
||
source_id=request.source_id,
|
||
created_by_user_id=authenticated_user.user.id,
|
||
max_retries=request.max_retries,
|
||
)
|
||
)
|
||
|
||
return job_to_response(job)
|
||
|
||
|
||
# ── 提交执行 ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.post("/jobs/{job_id}/submit", response_model=JobResponse)
|
||
def submit_job(
|
||
job_id: str,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""提交任务执行。
|
||
|
||
将任务状态从 pending 切换为 running,并 dispatch Celery 异步任务。
|
||
"""
|
||
# 权限检查:先获取任务并验证权限,再执行状态变更
|
||
job = job_repo.get(job_id)
|
||
if job is None:
|
||
raise HTTPException(status_code=404, detail=f"Job {job_id} not found")
|
||
if job.created_by_user_id and job.created_by_user_id != authenticated_user.user.id:
|
||
raise HTTPException(status_code=403, detail="Access denied to this job")
|
||
|
||
use_case = SubmitJobUseCase(job_repo)
|
||
|
||
try:
|
||
job = use_case.execute(job_id)
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
|
||
# Dispatch Celery 任务
|
||
celery_task_name = _JOB_TYPE_TO_CELERY_TASK.get(job.job_type.value)
|
||
if celery_task_name:
|
||
result = celery_app.send_task(celery_task_name, args=[job.id], kwargs=job.payload)
|
||
job.celery_task_id = result.id
|
||
job_repo.update(job)
|
||
logger.info("已提交 Celery 任务: job_id=%s celery_task_id=%s", job.id, result.id)
|
||
|
||
return job_to_response(job)
|
||
|
||
|
||
# ── 查询接口 ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.get("/jobs/{job_id}", response_model=JobResponse)
|
||
def get_job(
|
||
job_id: str,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""获取任务详情。"""
|
||
use_case = GetJobUseCase(job_repo)
|
||
job = use_case.execute(job_id)
|
||
if job is None:
|
||
raise HTTPException(status_code=404, detail=f"Job {job_id} not found")
|
||
return job_to_response(job)
|
||
|
||
|
||
@router.get("/projects/{project_id}/jobs", response_model=ListJobsResponse)
|
||
def list_project_jobs(
|
||
project_id: str,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
project_repository: Any = Depends(get_project_repository),
|
||
job_type: str | None = Query(default=None, description="按任务类型过滤"),
|
||
status_filter: str | None = Query(default=None, alias="status", description="按状态过滤"),
|
||
limit: int = Query(default=50, ge=1, le=200),
|
||
offset: int = Query(default=0, ge=0),
|
||
) -> ListJobsResponse:
|
||
"""获取项目下的任务列表。"""
|
||
_check_project_access(project_id, authenticated_user.user.id, project_repository)
|
||
|
||
use_case = ListJobsUseCase(job_repo)
|
||
jobs = use_case.execute(
|
||
project_id=project_id,
|
||
job_type=job_type,
|
||
status=status_filter,
|
||
limit=limit,
|
||
offset=offset,
|
||
)
|
||
items = [job_to_response(j) for j in jobs]
|
||
return ListJobsResponse(items=items, total=len(items))
|
||
|
||
|
||
@router.get("/projects/{project_id}/jobs/stats", response_model=JobStatisticsResponse)
|
||
def get_job_statistics(
|
||
project_id: str,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
project_repository: Any = Depends(get_project_repository),
|
||
) -> JobStatisticsResponse:
|
||
"""获取项目任务统计摘要。"""
|
||
_check_project_access(project_id, authenticated_user.user.id, project_repository)
|
||
|
||
use_case = GetJobStatisticsUseCase(job_repo)
|
||
stats = use_case.execute(project_id)
|
||
return JobStatisticsResponse(**stats)
|
||
|
||
|
||
# ── 进度更新 ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.put("/jobs/{job_id}/progress", response_model=JobResponse)
|
||
def update_job_progress(
|
||
job_id: str,
|
||
request: UpdateProgressRequest,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""更新任务进度。"""
|
||
use_case = UpdateJobProgressUseCase(job_repo)
|
||
|
||
try:
|
||
job = use_case.execute(
|
||
UpdateJobProgressCommand(
|
||
job_id=job_id,
|
||
progress=request.progress,
|
||
current_stage=request.current_stage,
|
||
)
|
||
)
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
|
||
return job_to_response(job)
|
||
|
||
|
||
# ── 完成 / 失败 ────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.post("/jobs/{job_id}/complete", response_model=JobResponse)
|
||
def complete_job(
|
||
job_id: str,
|
||
request: CompleteJobRequest,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""标记任务完成。"""
|
||
use_case = CompleteJobUseCase(job_repo)
|
||
|
||
try:
|
||
job = use_case.execute(CompleteJobCommand(job_id=job_id, result=request.result))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
|
||
return job_to_response(job)
|
||
|
||
|
||
@router.post("/jobs/{job_id}/fail", response_model=JobResponse)
|
||
def fail_job(
|
||
job_id: str,
|
||
request: FailJobRequest,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""标记任务失败。"""
|
||
use_case = FailJobUseCase(job_repo)
|
||
|
||
try:
|
||
job = use_case.execute(FailJobCommand(job_id=job_id, error_message=request.error_message))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
|
||
return job_to_response(job)
|
||
|
||
|
||
# ── 重试 / 取消 ────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@router.post("/jobs/{job_id}/retry", response_model=JobResponse)
|
||
def retry_job(
|
||
job_id: str,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""重试失败任务。
|
||
|
||
将任务重置为 pending,retry_count + 1,但不自动 dispatch。
|
||
需要再次调用 /submit 提交执行。
|
||
"""
|
||
# 权限检查:先获取任务并验证权限,再执行状态变更
|
||
job = job_repo.get(job_id)
|
||
if job is None:
|
||
raise HTTPException(status_code=404, detail=f"Job {job_id} not found")
|
||
if job.created_by_user_id and job.created_by_user_id != authenticated_user.user.id:
|
||
raise HTTPException(status_code=403, detail="Access denied to this job")
|
||
|
||
use_case = RetryJobUseCase(job_repo)
|
||
|
||
try:
|
||
job = use_case.execute(job_id)
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
|
||
return job_to_response(job)
|
||
|
||
|
||
@router.post("/jobs/{job_id}/cancel", response_model=JobResponse)
|
||
def cancel_job(
|
||
job_id: str,
|
||
authenticated_user: AuthenticatedUser = Depends(get_current_user),
|
||
job_repo: Any = Depends(get_job_repository),
|
||
) -> JobResponse:
|
||
"""取消任务。"""
|
||
# 权限检查:先获取任务并验证权限,再执行状态变更
|
||
job = job_repo.get(job_id)
|
||
if job is None:
|
||
raise HTTPException(status_code=404, detail=f"Job {job_id} not found")
|
||
if job.created_by_user_id and job.created_by_user_id != authenticated_user.user.id:
|
||
raise HTTPException(status_code=403, detail="Access denied to this job")
|
||
|
||
use_case = CancelJobUseCase(job_repo)
|
||
|
||
try:
|
||
job = use_case.execute(job_id)
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
|
||
return job_to_response(job)
|