from fastapi import APIRouter, HTTPException, status from app.deps import AuthDep from app.schemas.tasks import OcrTaskRequest, TaskStatusResponse from app.services.task_tracker import task_tracker from app.workers.tasks import ocr_pipeline router = APIRouter(prefix="/tasks") @router.post("/ocr", response_model=TaskStatusResponse) async def enqueue_ocr_task(payload: OcrTaskRequest, auth: AuthDep) -> TaskStatusResponse: """记录任务并投递 Celery,阶段 0 返回占位任务。""" task = task_tracker.create_task(user_id=auth.user_id, document_id=payload.document_id, task_type="ocr") ocr_pipeline.delay( task_id=task.task_id, document_id=payload.document_id, file_url=str(payload.file_url), user_id=auth.user_id, ) return task @router.get("/{task_id}", response_model=TaskStatusResponse) async def get_task_status(task_id: str, auth: AuthDep) -> TaskStatusResponse: task = task_tracker.get_task(task_id=task_id, user_id=auth.user_id) if not task: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Task not found") return task