diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py index da4f3773b..57ef114d2 100644 --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -978,7 +978,7 @@ def _render_video( Args: Returns: - (output_path, render_duration, cover_candidates) + (output_path, render_duration, cover_candidates, voiceover_path) """ if not downloaded_videos: raise RuntimeError(f"素材下载结果为空: task_id={task_id}") @@ -1206,13 +1206,6 @@ def _upload_and_record( # ── Celery Task ────────────────────────────────────────────────────────────── -@celery_app.task( - bind=True, - name="worker.generate_video", - max_retries=2, - soft_time_limit=600, # 10 分钟软超时 - time_limit=660, # 11 分钟硬超时 -) def _sync_task_config_to_plan(source_edit_plan_id: str, task_info: dict, db) -> str | None: """将 GenerationTask 的配置同步到 EditPlan.config,返回配音本地路径(如果有)。 @@ -1295,7 +1288,7 @@ def _render_from_edit_plan( """从 EditPlan 数据库记录直接渲染(不再内存重建clips)。 Returns: - (output_path, render_duration, cover_candidates) + (output_path, render_duration, cover_candidates, voiceover_path) """ from video_processing.render_adapter import RenderAdapter from worker_app.db import SessionLocal @@ -1335,13 +1328,18 @@ def _render_from_edit_plan( output_path = result.output_path cover_candidates = getattr(result, "cover_candidates", None) - return output_path, result.duration, cover_candidates + return output_path, result.duration, cover_candidates, voiceover_path finally: db.close() - # 清理临时配音文件 - # voiceover_path 在外部作用域,这里不直接引用 +@celery_app.task( + bind=True, + name="worker.generate_video", + max_retries=2, + soft_time_limit=600, # 10 分钟软超时 + time_limit=660, # 11 分钟硬超时 +) def generate_video(self, task_id: str) -> dict: """生成视频任务 — 使用 UnifiedRenderService 统一渲染。 @@ -1439,7 +1437,7 @@ def generate_video(self, task_id: str) -> dict: gen_task.append_log("渲染模式", "从草稿数据渲染(与预览一致)") _flush_logs(task_id, gen_task) - output_path, render_duration, cover_candidates = _render_from_edit_plan( + output_path, render_duration, cover_candidates, voiceover_tmp_path = _render_from_edit_plan( task_id=task_id, source_edit_plan_id=source_edit_plan_id, task_info=task_info, @@ -1587,6 +1585,13 @@ def generate_video(self, task_id: str) -> dict: file_size, ) + # 清理临时配音文件 + if voiceover_tmp_path: + try: + Path(voiceover_tmp_path).unlink(missing_ok=True) + except OSError: + logger.warning("[task_id=%s] 清理临时配音文件失败: %s", task_id, voiceover_tmp_path) + return { "status": "completed", "task_id": task_id, diff --git a/tests/unit/test_worker_generate_video_task_binding.py b/tests/unit/test_worker_generate_video_task_binding.py new file mode 100644 index 000000000..30c6dbdc8 --- /dev/null +++ b/tests/unit/test_worker_generate_video_task_binding.py @@ -0,0 +1,40 @@ +"""Regression test: ensure worker.generate_video Celery task is bound to the +real generate_video function, not a helper introduced above it. + +Context (P0 incident 2026-08-23): a refactor inserted helper function +_sync_task_config_to_plan directly under the @celery_app.task decorator, +so Celery registered the helper as "worker.generate_video". Calling the +task with a single task_id raised TypeError and every generation job +failed immediately. This test pins the decorator target. +""" + +from __future__ import annotations + +import inspect + + +def test_generate_video_task_registered_under_expected_name(): + from worker_app.tasks.generation import generate_video + + # Celery task object exposes its registered name + assert generate_video.name == "worker.generate_video" + + +def test_generate_video_task_signature_has_task_id(): + from worker_app.tasks.generation import generate_video + + # The underlying callable must accept (self, task_id) for bind=True tasks + sig = inspect.signature(generate_video.run) + assert "task_id" in sig.parameters + # The first positional arg after self must be task_id + params = list(sig.parameters) + assert params[0] == "task_id" + + +def test_sync_task_config_to_plan_is_plain_function(): + """Helper must NOT be registered as a Celery task.""" + from worker_app.tasks.generation import _sync_task_config_to_plan + + assert not hasattr(_sync_task_config_to_plan, "run"), ( + "_sync_task_config_to_plan must be a plain function, not a Celery task" + )