From 43de087282adbf43cedcca31a596b21d37be1d70 Mon Sep 17 00:00:00 2001 From: Xiaoxia AI Date: Wed, 24 Jun 2026 08:39:32 +0800 Subject: [PATCH] feat(tasks): add project task center --- apps/api/app/api/router.py | 5 + apps/api/app/api/routes/task_center.py | 105 ++++++++++++++++++ apps/api/app/schemas/task_center.py | 22 ++++ apps/web/e2e/core-generation.spec.ts | 6 + apps/web/src/api/tasks.ts | 22 ++++ apps/web/src/pages/workspace/ProjectTasks.tsx | 82 ++++++++++++++ .../src/pages/workspace/WorkspaceDetail.tsx | 3 + apps/web/src/router/index.tsx | 4 + .../sqlalchemy_impl/ingest_job_repository.py | 4 + tests/unit/test_task_center.py | 25 +++++ 10 files changed, 278 insertions(+) create mode 100644 apps/api/app/api/routes/task_center.py create mode 100644 apps/api/app/schemas/task_center.py create mode 100644 apps/web/src/api/tasks.ts create mode 100644 apps/web/src/pages/workspace/ProjectTasks.tsx create mode 100644 tests/unit/test_task_center.py diff --git a/apps/api/app/api/router.py b/apps/api/app/api/router.py index 50347f8ae..a29372ec9 100644 --- a/apps/api/app/api/router.py +++ b/apps/api/app/api/router.py @@ -10,6 +10,7 @@ from app.api.routes.ingest_jobs import router as ingest_jobs_router from app.api.routes.project_management import router as project_management_router from app.api.routes.project_titles import router as project_titles_router from app.api.routes.projects import router as projects_router +from app.api.routes.task_center import router as task_center_router from app.api.routes.upload import router as upload_router from app.api.routes.workspaces import router as workspaces_router from fastapi import APIRouter @@ -35,6 +36,10 @@ api_router.include_router( project_titles_router, tags=["标题库"], ) +api_router.include_router( + task_center_router, + tags=["任务中心"], +) api_router.include_router( asset_diagnosis_router, tags=["素材诊断"], diff --git a/apps/api/app/api/routes/task_center.py b/apps/api/app/api/routes/task_center.py new file mode 100644 index 000000000..d821acc0b --- /dev/null +++ b/apps/api/app/api/routes/task_center.py @@ -0,0 +1,105 @@ +from typing import Any + +from app.api.routes.permissions import require_workspace_member +from app.auth import AuthenticatedUser, get_current_user +from app.dependencies import ( + get_generation_task_repository, + get_ingest_job_repository, + get_project_repository, + get_workspace_member_repository, +) +from app.schemas.task_center import ListProjectTasksResponse, ProjectTaskResponse +from fastapi import APIRouter, Depends, HTTPException + +from packages.ports.workspace_member_repository import WorkspaceMemberRepository + +router = APIRouter() + + +def _humanize_task_error(error_message: str) -> str: + raw = (error_message or "").strip() + if not raw: + return "" + lower = raw.lower() + if "ffmpeg" in lower or "ffprobe" in lower or "invalid data" in lower or "moov atom" in lower: + return "视频素材格式无法识别,请重新导出为常见 MP4/H.264 后再试。" + if "oss" in lower or "bucket" in lower or "storage" in lower: + return "素材存储服务读取或写入失败,请稍后重试或联系小虾检查 OSS。" + if "not found" in lower or "no such file" in lower: + return "任务依赖的素材或文件不存在,请确认素材仍在项目中。" + return f"任务失败:{raw}" + + +def _generation_step(task) -> str: + if task.status.value == "pending": + return "等待 Worker 执行" + if task.status.value == "running": + return "正在生成成片" + if task.status.value == "completed": + return "生成完成" + if task.status.value == "failed": + return "生成失败" + return task.status.value + + +def _ingest_step(job) -> str: + if job.status.value == "pending": + return "等待导入" + if job.status.value == "processing": + return "正在分析素材" + if job.status.value == "completed": + return "导入完成" + if job.status.value == "failed": + return "导入失败" + return job.status.value + + +@router.get("/projects/{project_id}/tasks", response_model=ListProjectTasksResponse) +def list_project_tasks( + project_id: str, + authenticated_user: AuthenticatedUser = Depends(get_current_user), + project_repository: Any = Depends(get_project_repository), + ingest_job_repository: Any = Depends(get_ingest_job_repository), + generation_task_repository: Any = Depends(get_generation_task_repository), + workspace_member_repository: WorkspaceMemberRepository = Depends(get_workspace_member_repository), +) -> ListProjectTasksResponse: + project = project_repository.find_by_id(project_id) + if project is None: + raise HTTPException(status_code=404, detail="Project not found") + require_workspace_member(project.workspace_id, authenticated_user, workspace_member_repository) + + items: list[ProjectTaskResponse] = [] + for job in ingest_job_repository.list_by_project(project_id): + items.append(ProjectTaskResponse( + id=f"ingest:{job.id}", + task_type="ingest", + workspace_id=job.workspace_id, + project_id=job.project_id, + status=job.status.value, + progress=100.0 if job.status.value == "completed" else 0.0, + current_step=_ingest_step(job), + error_message=job.error_message, + user_message=_humanize_task_error(job.error_message), + retryable=job.status.value == "failed", + source_id=job.id, + created_at=job.created_at, + updated_at=job.updated_at, + )) + for task in generation_task_repository.list_by_project(project_id): + items.append(ProjectTaskResponse( + id=f"generation:{task.id}", + task_type="generation", + workspace_id=task.workspace_id, + project_id=task.project_id, + status=task.status.value, + progress=task.progress, + current_step=_generation_step(task), + error_message=task.error_message, + user_message=_humanize_task_error(task.error_message), + retryable=task.status.value == "failed", + source_id=task.id, + created_at=task.created_at, + updated_at=task.completed_at or task.started_at or task.created_at, + )) + items.sort(key=lambda item: item.updated_at or item.created_at or "", reverse=True) + return ListProjectTasksResponse(items=items) diff --git a/apps/api/app/schemas/task_center.py b/apps/api/app/schemas/task_center.py new file mode 100644 index 000000000..44841e449 --- /dev/null +++ b/apps/api/app/schemas/task_center.py @@ -0,0 +1,22 @@ +from datetime import datetime +from pydantic import BaseModel, Field + + +class ProjectTaskResponse(BaseModel): + id: str + task_type: str + workspace_id: str + project_id: str + status: str + progress: float + current_step: str + error_message: str = "" + user_message: str = "" + retryable: bool = False + source_id: str = "" + created_at: datetime | None = None + updated_at: datetime | None = None + + +class ListProjectTasksResponse(BaseModel): + items: list[ProjectTaskResponse] = Field(default_factory=list) diff --git a/apps/web/e2e/core-generation.spec.ts b/apps/web/e2e/core-generation.spec.ts index 900f74656..dc600edf2 100644 --- a/apps/web/e2e/core-generation.spec.ts +++ b/apps/web/e2e/core-generation.spec.ts @@ -188,5 +188,11 @@ test.describe('Core generation and download flow', () => { expect(titleAfterGeneration.status(), await titleAfterGeneration.text()).toBe(200); const titlesData = (await titleAfterGeneration.json()) as { items: ProjectTitleResponse[] }; expect(titlesData.items.find((item) => item.id === titleData.id)?.usage_count).toBe(1); + + await page.goto(`/projects/${projectData.id}/tasks`); + await expect(page.getByText('项目任务中心')).toBeVisible({ timeout: 20_000 }); + await expect(page.getByText('视频生成')).toBeVisible({ timeout: 20_000 }); + await expect(page.getByText('已完成')).toBeVisible({ timeout: 20_000 }); + await expect(page.getByText(createdTask.id)).toBeVisible({ timeout: 20_000 }); }); }); diff --git a/apps/web/src/api/tasks.ts b/apps/web/src/api/tasks.ts new file mode 100644 index 000000000..fc8fe8ff4 --- /dev/null +++ b/apps/web/src/api/tasks.ts @@ -0,0 +1,22 @@ +import apiClient from './client'; + +export interface ProjectTaskItem { + id: string; + task_type: 'ingest' | 'generation' | string; + workspace_id: string; + project_id: string; + status: string; + progress: number; + current_step: string; + error_message: string; + user_message: string; + retryable: boolean; + source_id: string; + created_at?: string | null; + updated_at?: string | null; +} + +export const getProjectTasks = async (projectId: string): Promise => { + const response = await apiClient.get(`/projects/${projectId}/tasks`); + return response.data.items; +}; diff --git a/apps/web/src/pages/workspace/ProjectTasks.tsx b/apps/web/src/pages/workspace/ProjectTasks.tsx new file mode 100644 index 000000000..384d30cf0 --- /dev/null +++ b/apps/web/src/pages/workspace/ProjectTasks.tsx @@ -0,0 +1,82 @@ +import React from 'react'; +import { Alert, Button, Card, Empty, List, Progress, Space, Tag, Typography } from 'antd'; +import { useQuery } from '@tanstack/react-query'; +import { useParams } from 'react-router-dom'; +import { getProjectTasks } from '@/api/tasks'; + +const statusMap: Record = { + pending: { label: '排队中', color: 'default' }, + processing: { label: '处理中', color: 'blue' }, + running: { label: '运行中', color: 'blue' }, + completed: { label: '已完成', color: 'green' }, + failed: { label: '失败', color: 'red' }, + cancelled: { label: '已取消', color: 'default' }, +}; + +const taskTypeLabels: Record = { + ingest: '素材导入', + generation: '视频生成', +}; + +const ProjectTasks: React.FC = () => { + const { id } = useParams<{ id: string }>(); + const projectId = id || ''; + + const tasksQuery = useQuery({ + queryKey: ['project-tasks', projectId], + queryFn: () => getProjectTasks(projectId), + enabled: !!projectId, + refetchInterval: 5000, + }); + + return ( +
+ tasksQuery.refetch()}>刷新}> + {tasksQuery.isError && } + {tasksQuery.data?.some((task) => task.retryable) && ( + + )} + }} + renderItem={(task) => { + const status = statusMap[task.status] || { label: task.status, color: 'default' }; + return ( + + + {taskTypeLabels[task.task_type] || task.task_type} + {status.label} + {task.current_step} + {task.retryable ? 可重试(暂未开放) : null} + + } + description={ + + + {task.user_message ? {task.user_message} : null} + + 任务 ID:{task.source_id} + + + } + /> + + ); + }} + /> + +
+ ); +}; + +export const Component = ProjectTasks; +export default ProjectTasks; diff --git a/apps/web/src/pages/workspace/WorkspaceDetail.tsx b/apps/web/src/pages/workspace/WorkspaceDetail.tsx index 44836ae2f..6ac291b06 100644 --- a/apps/web/src/pages/workspace/WorkspaceDetail.tsx +++ b/apps/web/src/pages/workspace/WorkspaceDetail.tsx @@ -100,6 +100,9 @@ const WorkspaceDetail: React.FC = () => { , + , , diff --git a/apps/web/src/router/index.tsx b/apps/web/src/router/index.tsx index fb936b0e4..ad830de96 100644 --- a/apps/web/src/router/index.tsx +++ b/apps/web/src/router/index.tsx @@ -72,6 +72,10 @@ export const router = createBrowserRouter([ path: 'projects/:id/generation', lazy: () => import('@/pages/workspace/ProjectGeneration'), }, + { + path: 'projects/:id/tasks', + lazy: () => import('@/pages/workspace/ProjectTasks'), + }, { path: 'projects/:id/results', lazy: () => import('@/pages/workspace/ProjectResults'), diff --git a/packages/adapters/sqlalchemy_impl/ingest_job_repository.py b/packages/adapters/sqlalchemy_impl/ingest_job_repository.py index 141231f48..81d3db3ec 100644 --- a/packages/adapters/sqlalchemy_impl/ingest_job_repository.py +++ b/packages/adapters/sqlalchemy_impl/ingest_job_repository.py @@ -42,6 +42,10 @@ class SQLAlchemyIngestJobRepository: updated_at=model.updated_at, ) + def list_by_project(self, project_id: str) -> list[IngestJob]: + models = self.session.query(IngestJobModel).filter(IngestJobModel.project_id == project_id).all() + return [self.get(model.id) for model in models if self.get(model.id) is not None] + def update(self, job: IngestJob) -> IngestJob: model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job.id).first() if model is None: diff --git a/tests/unit/test_task_center.py b/tests/unit/test_task_center.py new file mode 100644 index 000000000..e9f2adaca --- /dev/null +++ b/tests/unit/test_task_center.py @@ -0,0 +1,25 @@ +import sys +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "apps" / "api")) + +from app.api.routes.task_center import _humanize_task_error, _ingest_step +from packages.domain import IngestJob, IngestJobStatus + + +def test_humanize_task_error_for_media_failures(): + assert "MP4" in _humanize_task_error("ffmpeg invalid data found") + assert "OSS" in _humanize_task_error("oss bucket read failed") + assert "不存在" in _humanize_task_error("file not found") + + +def test_ingest_step_is_user_readable(): + job = IngestJob.create( + workspace_id="workspace-1", + project_id="project-1", + library_id="library-1", + storage_key="uploads/video.mp4", + ) + job.status = IngestJobStatus.COMPLETED + + assert _ingest_step(job) == "导入完成"