fix(worker): P0 修复@celery_app.task装饰器错位导致所有生成任务崩溃 #1470

Merged
auto-approve-bot merged 3 commits from fix/p0-celery-task-decorator-misplacement into develop 2026-08-23 17:16:00 +08:00
2 changed files with 220 additions and 173 deletions
+181 -173
View File
@@ -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 统一渲染。
@@ -1428,173 +1426,183 @@ def generate_video(self, task_id: str) -> dict:
# ── 新路径:有 source_edit_plan_id 时直接从数据库 EditPlan 渲染 ──
source_edit_plan_id = task_info.get("source_edit_plan_id", "")
if source_edit_plan_id:
logger.info(
"[task_id=%s] 使用 EditPlan 数据库路径渲染: plan_id=%s",
task_id,
source_edit_plan_id,
)
_update_task_progress(task_id, 30, "加载草稿数据")
if gen_task:
gen_task.append_log("渲染模式", "从草稿数据渲染(与预览一致)")
_flush_logs(task_id, gen_task)
output_path, render_duration, cover_candidates = _render_from_edit_plan(
task_id=task_id,
source_edit_plan_id=source_edit_plan_id,
task_info=task_info,
)
if gen_task:
gen_task.append_log("渲染", f"渲染完成, 时长={render_duration:.1f}s")
_flush_logs(task_id, gen_task)
_update_task_progress(task_id, 80, "渲染完成")
# ── 4. 上传 OSS + 查重记录 ───────────────────────────────
_update_task_progress(task_id, 85, "开始上传")
file_url, duration, file_size, video_count = _upload_and_record(
task_id=task_id,
output_path=output_path,
project_id=project_id,
batch_id=batch_id,
editing_mode=editing_mode,
user_id=user_id,
video_name=task_info.get("video_title", ""),
)
if gen_task:
gen_task.append_log(
"OSS上传",
f"上传成功, 大小={file_size}",
file_size=file_size,
file_url=file_url,
voiceover_tmp_path: str | None = None
try:
logger.info(
"[task_id=%s] 使用 EditPlan 数据库路径渲染: plan_id=%s",
task_id,
source_edit_plan_id,
)
_flush_logs(task_id, gen_task)
_update_task_progress(task_id, 30, "加载草稿数据")
_update_task_progress(task_id, 95, "上传完成")
if gen_task:
gen_task.append_log("渲染模式", "从草稿数据渲染(与预览一致)")
_flush_logs(task_id, gen_task)
# ── 4.5 封面帧持久化 ────────────────────────────────────────────
try:
if cover_candidates:
first = cover_candidates[0]
cover_frame_url = first.get("image_url") or first.get("url") or ""
if cover_frame_url:
_cover_session = SessionLocal()
try:
from packages.adapters.sqlalchemy_impl.models import (
GenerationTaskModel,
)
_cover_model = (
_cover_session.query(GenerationTaskModel)
.filter(GenerationTaskModel.id == task_id)
.first()
)
if _cover_model:
_cover_model.cover_url = cover_frame_url
meta = dict(_cover_model.metadata or {})
meta["cover_candidates"] = cover_candidates
_cover_model.metadata = meta
_cover_session.commit()
finally:
_cover_session.close()
except Exception:
logger.warning("[task_id=%s] 封面帧持久化失败", task_id, exc_info=True)
# ── 5. 标记完成 ──────────────────────────────────────────────────
_update_task_status(task_id, "mark_completed", result_count=video_count)
# 5.1 更新标题使用次数
try:
_title_session = SessionLocal()
try:
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
from packages.adapters.sqlalchemy_impl.title_library_repository import (
SQLAlchemyTitleLibraryRepository,
)
_task_repo = SQLAlchemyGenerationTaskRepository(_title_session)
_gen_task = _task_repo.get(task_id)
if _gen_task and _gen_task.title_ids and _gen_task.created_by_user_id:
_title_repo = SQLAlchemyTitleLibraryRepository(_title_session)
for _tid in _gen_task.title_ids:
try:
_title_repo.increment_usage_count(_tid, _gen_task.created_by_user_id)
except Exception:
logger.warning(
"[task_id=%s] 更新标题使用次数失败: title_id=%s",
task_id,
_tid,
exc_info=True,
)
finally:
_title_session.close()
except Exception:
logger.warning("[task_id=%s] 更新标题使用次数异常", task_id, exc_info=True)
# 5.2 更新素材使用次数
try:
from worker_app.core.asset_usage import mark_asset_used_for_generation
_asset_session = SessionLocal()
try:
from packages.adapters.sqlalchemy_impl.asset_repository import (
SQLAlchemyAssetRepository,
)
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
_task_repo = SQLAlchemyGenerationTaskRepository(_asset_session)
_asset_repo = SQLAlchemyAssetRepository(_asset_session)
_gen_task = _task_repo.get(task_id)
if _gen_task and _gen_task.asset_ids:
for _aid in _gen_task.asset_ids:
try:
_asset = _asset_repo.get(_aid)
if _asset:
mark_asset_used_for_generation(_asset)
_asset_repo.update(_asset)
except Exception:
logger.warning(
"[task_id=%s] 更新素材使用次数失败: asset_id=%s",
task_id,
_aid,
exc_info=True,
)
finally:
_asset_session.close()
except Exception:
logger.warning("[task_id=%s] 更新素材使用次数异常", task_id, exc_info=True)
if gen_task:
gen_task.append_log(
"任务完成",
f"视频生成完成: 时长={duration:.2f}s, 大小={file_size}",
duration=round(duration, 2),
file_size=file_size,
video_count=video_count,
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,
)
_flush_logs(task_id, gen_task)
logger.info(
"[task_id=%s] [任务完成] duration=%.2fs file_size=%d (edit_plan path)",
task_id,
duration,
file_size,
)
if gen_task:
gen_task.append_log("渲染", f"渲染完成, 时长={render_duration:.1f}s")
_flush_logs(task_id, gen_task)
return {
"status": "completed",
"task_id": task_id,
"output_path": str(output_path),
"file_size": file_size,
"duration": duration,
"mode": editing_mode.value,
}
_update_task_progress(task_id, 80, "渲染完成")
# ── 4. 上传 OSS + 查重记录 ───────────────────────────────
_update_task_progress(task_id, 85, "开始上传")
file_url, duration, file_size, video_count = _upload_and_record(
task_id=task_id,
output_path=output_path,
project_id=project_id,
batch_id=batch_id,
editing_mode=editing_mode,
user_id=user_id,
video_name=task_info.get("video_title", ""),
)
if gen_task:
gen_task.append_log(
"OSS上传",
f"上传成功, 大小={file_size}",
file_size=file_size,
file_url=file_url,
)
_flush_logs(task_id, gen_task)
_update_task_progress(task_id, 95, "上传完成")
# ── 4.5 封面帧持久化 ────────────────────────────────────────────
try:
if cover_candidates:
first = cover_candidates[0]
cover_frame_url = first.get("image_url") or first.get("url") or ""
if cover_frame_url:
_cover_session = SessionLocal()
try:
from packages.adapters.sqlalchemy_impl.models import (
GenerationTaskModel,
)
_cover_model = (
_cover_session.query(GenerationTaskModel)
.filter(GenerationTaskModel.id == task_id)
.first()
)
if _cover_model:
_cover_model.cover_url = cover_frame_url
meta = dict(_cover_model.metadata or {})
meta["cover_candidates"] = cover_candidates
_cover_model.metadata = meta
_cover_session.commit()
finally:
_cover_session.close()
except Exception:
logger.warning("[task_id=%s] 封面帧持久化失败", task_id, exc_info=True)
# ── 5. 标记完成 ──────────────────────────────────────────────────
_update_task_status(task_id, "mark_completed", result_count=video_count)
# 5.1 更新标题使用次数
try:
_title_session = SessionLocal()
try:
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
from packages.adapters.sqlalchemy_impl.title_library_repository import (
SQLAlchemyTitleLibraryRepository,
)
_task_repo = SQLAlchemyGenerationTaskRepository(_title_session)
_gen_task = _task_repo.get(task_id)
if _gen_task and _gen_task.title_ids and _gen_task.created_by_user_id:
_title_repo = SQLAlchemyTitleLibraryRepository(_title_session)
for _tid in _gen_task.title_ids:
try:
_title_repo.increment_usage_count(_tid, _gen_task.created_by_user_id)
except Exception:
logger.warning(
"[task_id=%s] 更新标题使用次数失败: title_id=%s",
task_id,
_tid,
exc_info=True,
)
finally:
_title_session.close()
except Exception:
logger.warning("[task_id=%s] 更新标题使用次数异常", task_id, exc_info=True)
# 5.2 更新素材使用次数
try:
from worker_app.core.asset_usage import mark_asset_used_for_generation
_asset_session = SessionLocal()
try:
from packages.adapters.sqlalchemy_impl.asset_repository import (
SQLAlchemyAssetRepository,
)
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
SQLAlchemyGenerationTaskRepository,
)
_task_repo = SQLAlchemyGenerationTaskRepository(_asset_session)
_asset_repo = SQLAlchemyAssetRepository(_asset_session)
_gen_task = _task_repo.get(task_id)
if _gen_task and _gen_task.asset_ids:
for _aid in _gen_task.asset_ids:
try:
_asset = _asset_repo.get(_aid)
if _asset:
mark_asset_used_for_generation(_asset)
_asset_repo.update(_asset)
except Exception:
logger.warning(
"[task_id=%s] 更新素材使用次数失败: asset_id=%s",
task_id,
_aid,
exc_info=True,
)
finally:
_asset_session.close()
except Exception:
logger.warning("[task_id=%s] 更新素材使用次数异常", task_id, exc_info=True)
if gen_task:
gen_task.append_log(
"任务完成",
f"视频生成完成: 时长={duration:.2f}s, 大小={file_size}",
duration=round(duration, 2),
file_size=file_size,
video_count=video_count,
)
_flush_logs(task_id, gen_task)
logger.info(
"[task_id=%s] [任务完成] duration=%.2fs file_size=%d (edit_plan path)",
task_id,
duration,
file_size,
)
return {
"status": "completed",
"task_id": task_id,
"output_path": str(output_path),
"file_size": file_size,
"duration": duration,
"mode": editing_mode.value,
}
finally:
# 无论任务成功或失败,都清理临时配音文件,避免磁盘泄漏
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)
# DEPRECATED: 以下为旧路径,仅兼容无 source_edit_plan_id 的旧调用,后续移除
with tempfile.TemporaryDirectory(prefix="xiaoxia-generation-") as temp_dir:
@@ -0,0 +1,39 @@
"""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
# For bind=True tasks Celery binds self at call time, so run() signature
# starts directly with task_id (verified on Celery 5.x).
sig = inspect.signature(generate_video.run)
params = list(sig.parameters)
assert params[0] == "task_id", f"expected task_id as first param, got {params}"
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"