diff --git a/apps/api/app/api/routes/generation_tasks.py b/apps/api/app/api/routes/generation_tasks.py index 870b5ee17..1c8cfb9a1 100644 --- a/apps/api/app/api/routes/generation_tasks.py +++ b/apps/api/app/api/routes/generation_tasks.py @@ -1,5 +1,6 @@ from fastapi import APIRouter, Depends, HTTPException +from app.core.celery_app import celery_app from app.dependencies import get_generation_task_repository, get_generated_video_repository from app.schemas.generation_task import CreateGenerationTaskRequest, GenerationTaskResponse from app.schemas.generated_video import GeneratedVideoResponse, ListGeneratedVideosResponse @@ -30,6 +31,7 @@ def create_generation_task( created_by_user_id=request.created_by_user_id, ) ) + celery_app.send_task("worker.generate_video", args=[task.id]) return GenerationTaskResponse( id=task.id, workspace_id=task.workspace_id, diff --git a/docs/PHASE7-PROGRESS.md b/docs/PHASE7-PROGRESS.md index 9fc38e4b7..58035b778 100644 --- a/docs/PHASE7-PROGRESS.md +++ b/docs/PHASE7-PROGRESS.md @@ -2,7 +2,7 @@ **Phase**: Phase 7 - 核心视频剪辑业务 **状态**: 🔄 进行中 -**最后更新**: 2026-06-18 19:56 GMT+8 +**最后更新**: 2026-06-18 20:11 GMT+8 --- @@ -48,6 +48,8 @@ - [x] ClassificationJob 流程第一轮打通 - [x] GenerationTask 主线骨架已落地 - [x] GeneratedVideo 主线骨架已落地 +- [x] 生成任务创建后自动触发 worker +- [x] 生成结果最小闭环测试已落地 - [ ] 前端主链路联调完成 ### 4. 本轮已完成的具体验证 @@ -56,7 +58,7 @@ - [x] `tests/integration/test_upload_pipeline.py` 通过 - [x] `tests/integration/test_classification_pipeline.py` 通过 - [x] `tests/integration/test_projects.py` 通过 -- [x] `tests/integration/test_generation_pipeline.py` 通过 +- [x] `tests/integration/test_generation_pipeline.py` 通过(含生成结果最小闭环) - [x] 素材与生成主线相关目录编译检查通过 --- @@ -94,7 +96,7 @@ - [x] 上传 / Asset 创建链路第一轮对齐 - [x] ClassificationJob 与 Asset 元数据更新链路第一轮打通 - [x] GenerationTask / GeneratedVideo 主线骨架已补齐 -- [ ] API / Adapter / Worker 的生成结果流继续联调 +- [x] API / Adapter / Worker 的生成结果流第一轮联调 - [ ] 测试补齐与回归验证 ### Step 3:保留 Agent 体系设计,等待 runtime 修复后再恢复实跑 @@ -105,7 +107,7 @@ ## 五、当前判断 -**当前 Phase 7 已完成素材前半主链打通,并补出生成链主线骨架,整体仍保持在既定规则内推进。** +**当前 Phase 7 已完成素材前半主链打通,并把生成链推进到“最小可运行闭环”,整体仍保持在既定规则内推进。** 当前执行策略是: - 暂停 Agent 实跑 @@ -123,7 +125,7 @@ - 答:Phase 7 - 核心视频剪辑业务 2. **当前 Phase 主要在做什么?** - - 答:已切入 Phase 7 第一批业务开发,当前已完成素材前半主链收口,并补出生成链主线骨架,正在进入提交与 CI/CD 验证阶段 + - 答:已切入 Phase 7 第一批业务开发,当前已完成素材前半主链收口,并把生成链推进到最小可运行闭环,正在持续提交与 CI/CD 验证 3. **当前最重要的阻塞点是什么?** - 答:OpenClaw 子 Agent runtime 暂不稳定,因此暂停 Agent 实跑;另外数据库字段命名仍有历史包袱,但已通过映射兼容,不阻断主线开发 @@ -150,4 +152,4 @@ --- -**状态结论**:Phase 7 未跑偏,已暂停 Agent 实跑并切回主会话直开;当前素材前半主链已打通,生成链主线骨架已落地,现进入提交与 CI/CD 验证阶段。 +**状态结论**:Phase 7 未跑偏,已暂停 Agent 实跑并切回主会话直开;当前素材前半主链已打通,生成链已进入最小可运行闭环,现继续通过提交与 CI/CD 验证推进。 diff --git a/tests/integration/test_generation_pipeline.py b/tests/integration/test_generation_pipeline.py index d511b0a4b..ef40ecf1c 100644 --- a/tests/integration/test_generation_pipeline.py +++ b/tests/integration/test_generation_pipeline.py @@ -1,5 +1,7 @@ +from datetime import datetime, timezone + from packages.application import CreateGenerationTaskCommand, CreateGenerationTaskUseCase -from packages.domain import GenerationTaskStatus +from packages.domain import GeneratedVideo, GenerationTaskStatus class DummyGenerationTaskRepository: @@ -21,6 +23,57 @@ class DummyGenerationTaskRepository: return task +class DummyGeneratedVideoRepository: + def __init__(self): + self.items = {} + + def create(self, video): + self.items[video.id] = video + return video + + def get(self, video_id): + return self.items.get(video_id) + + def list_by_project(self, project_id): + return [video for video in self.items.values() if video.project_id == project_id] + + def list_by_generation_task(self, generation_task_id): + return [video for video in self.items.values() if video.generation_task_id == generation_task_id] + + +def simulate_generate_video(task_id: str, task_repo: DummyGenerationTaskRepository, video_repo: DummyGeneratedVideoRepository) -> dict: + task = task_repo.get(task_id) + if task is None: + return {"status": "failed", "error": "task not found"} + + task.status = GenerationTaskStatus.RUNNING + task.progress = 20.0 + task.started_at = task.started_at or datetime.now(timezone.utc) + task_repo.update(task) + + video = GeneratedVideo.create( + workspace_id=task.workspace_id, + project_id=task.project_id, + generation_task_id=task.id, + name=f"{task.id}.mp4", + file_url=f"https://example.invalid/generated/{task.id}.mp4", + file_size=2048, + duration=12.5, + width=1920, + height=1080, + fps=25.0, + ) + video_repo.create(video) + + task.status = GenerationTaskStatus.COMPLETED + task.progress = 100.0 + task.result_count = 1 + task.completed_at = datetime.now(timezone.utc) + task_repo.update(task) + + return {"status": "completed", "task_id": task.id, "video_id": video.id} + + def test_create_generation_task_smoke(): repo = DummyGenerationTaskRepository() use_case = CreateGenerationTaskUseCase(repo) @@ -39,3 +92,32 @@ def test_create_generation_task_smoke(): assert task.asset_library_id == "lib-1" assert task.status == GenerationTaskStatus.PENDING assert repo.get(task.id) is not None + + +def test_generation_pipeline_smoke(): + task_repo = DummyGenerationTaskRepository() + video_repo = DummyGeneratedVideoRepository() + use_case = CreateGenerationTaskUseCase(task_repo) + task = use_case.execute( + CreateGenerationTaskCommand( + workspace_id="ws-1", + project_id="proj-1", + asset_library_id="lib-1", + strategy_id="str-1", + voice_library_id="voice-1", + created_by_user_id="user-1", + ) + ) + + result = simulate_generate_video(task.id, task_repo, video_repo) + + assert result["status"] == "completed" + updated_task = task_repo.get(task.id) + assert updated_task is not None + assert updated_task.status == GenerationTaskStatus.COMPLETED + assert updated_task.result_count == 1 + + videos = video_repo.list_by_generation_task(task.id) + assert len(videos) == 1 + assert videos[0].project_id == "proj-1" + assert videos[0].generation_task_id == task.id