fix: 渲染产物临时目录不在render_plan中提前清理,改由调用方上传后清理 #1484

Merged
auto-approve-bot merged 3 commits from fix/temp-dir-cleanup-race into develop 2026-08-24 21:21:00 +08:00
3 changed files with 110 additions and 36 deletions
@@ -84,6 +84,7 @@ class RenderAdapterResult:
cover_candidates: list[dict] | None = (
None # 封面候选帧 [{"image_url": "...", "frame_time": 5.0, "storage_key": "..."}]
)
temp_dir: str | None = None # 渲染临时目录,成功时由调用方清理,失败时由 finally 清理
def __post_init__(self):
if self.rendered_clip_ids is None:
@@ -200,7 +201,7 @@ class RenderAdapter:
self._report_progress(progress_cb, 35.0, "准备 BGM 音频")
# 3~6. 统一渲染核心流程(BGM + ASR + 渲染 + 缩略图 + 上传)
return self._do_render(
result = self._do_render(
plan=plan,
clips=ready_clips,
asset_path_map=asset_path_map,
@@ -212,6 +213,11 @@ class RenderAdapter:
failed_clip_ids=failed_clip_ids,
voiceover_audio_path=voiceover_audio_path,
)
# 成功时将临时目录所有权转移给调用方,阻止 finally 清理
if result.success and temp_dir:
result.temp_dir = temp_dir
temp_dir = None # 阻止 finally 块清理
return result
except subprocess.CalledProcessError as exc:
stderr_text = (exc.stderr or "").strip()
+47 -35
View File
@@ -535,11 +535,11 @@ def _render_from_edit_plan(
task_id: str,
source_edit_plan_id: str,
task_info: dict,
) -> tuple[Path, float, list[dict] | None]:
) -> tuple[Path, float, list[dict] | None, str | None, str | None]:
"""从 EditPlan 数据库记录直接渲染(不再内存重建clips)。
Returns:
(output_path, render_duration, cover_candidates, voiceover_path)
(output_path, render_duration, cover_candidates, voiceover_path, temp_dir)
"""
from video_processing.render_adapter import RenderAdapter
from worker_app.db import SessionLocal
@@ -578,8 +578,9 @@ def _render_from_edit_plan(
output_path = result.output_path
cover_candidates = getattr(result, "cover_candidates", None)
render_temp_dir = getattr(result, "temp_dir", None)
return output_path, result.duration, cover_candidates, voiceover_path
return output_path, result.duration, cover_candidates, voiceover_path, render_temp_dir
finally:
db.close()
@@ -678,6 +679,7 @@ def generate_video(self, task_id: str) -> dict:
source_edit_plan_id = task_info.get("source_edit_plan_id", "")
if source_edit_plan_id:
voiceover_tmp_path: str | None = None
render_temp_dir: str | None = None
try:
logger.info(
"[task_id=%s] 使用 EditPlan 数据库路径渲染: plan_id=%s",
@@ -690,40 +692,50 @@ def generate_video(self, task_id: str) -> dict:
gen_task.append_log("渲染模式", "从草稿数据渲染(与预览一致)")
_flush_logs(task_id, gen_task)
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,
)
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,
output_path, render_duration, cover_candidates, voiceover_tmp_path, render_temp_dir = (
_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)
)
# 从这里开始,render_temp_dir 已赋值,必须确保异常时也能清理
try:
if gen_task:
gen_task.append_log("渲染", f"渲染完成, 时长={render_duration:.1f}s")
_flush_logs(task_id, gen_task)
_update_task_progress(task_id, 95, "上传完成")
_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, "上传完成")
finally:
# 清理渲染临时目录(无论后续步骤成功与否都清理)
if render_temp_dir:
import shutil
shutil.rmtree(render_temp_dir, ignore_errors=True)
logger.info("[task_id=%s] 渲染临时目录已清理: %s", task_id, render_temp_dir)
# ── 4.5 封面帧持久化 ────────────────────────────────────────────
try:
+56
View File
@@ -0,0 +1,56 @@
"""回归测试:渲染产物临时目录不在 render_plan 中提前清理。
根因:render_plan 的 finally 块在返回前清理了临时目录,
但 generation.py 还需要访问其中的文件进行 OSS 上传。
修复:将清理责任交给调用方(generation.py),render_plan 只在失败时清理。
"""
import ast
def test_render_adapter_result_has_temp_dir_field():
"""RenderAdapterResult 包含 temp_dir 字段"""
with open("apps/worker/video_processing/render_adapter.py") as f:
source = f.read()
tree = ast.parse(source)
for node in ast.walk(tree):
if isinstance(node, ast.ClassDef) and node.name == "RenderAdapterResult":
for item in node.body:
if isinstance(item, ast.AnnAssign) and isinstance(item.target, ast.Name):
if item.target.id == "temp_dir":
return
raise AssertionError("RenderAdapterResult 缺少 temp_dir 字段")
def test_render_plan_does_not_cleanup_on_success():
"""render_plan 成功时不在 finally 中清理临时目录(通过将 temp_dir 置为 None"""
with open("apps/worker/video_processing/render_adapter.py") as f:
source = f.read()
# 成功路径必须将 temp_dir 置为 None,以阻止 finally 清理
assert "temp_dir = None" in source, "render_plan 成功时应将 temp_dir 置为 None 以阻止 finally 清理"
def test_render_plan_passes_temp_dir_to_result():
"""render_plan 将 temp_dir 传递给返回结果"""
with open("apps/worker/video_processing/render_adapter.py") as f:
source = f.read()
assert "result.temp_dir = temp_dir" in source, "render_plan 应将 temp_dir 设置到 result 上"
def test_generation_cleans_up_temp_dir():
"""generation.py 在上传完成后清理临时目录"""
with open("apps/worker/worker_app/tasks/generation.py") as f:
source = f.read()
# 验证 _render_from_edit_plan 返回 temp_dir
assert "render_temp_dir" in source, "generation.py 应接收 render_temp_dir"
# 验证有清理逻辑(shutil.rmtree(render_temp_dir...
assert "rmtree(render_temp_dir" in source, "generation.py 应清理 render_temp_dir"
# 验证清理发生在上传之后(通过查找顺序)
upload_pos = source.find("_upload_and_record")
cleanup_pos = source.find("rmtree(render_temp_dir")
assert upload_pos > 0 and cleanup_pos > upload_pos, "清理临时目录应在 _upload_and_record 之后执行"