diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py index c28eb29b4..3b92242bf 100644 --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -95,6 +95,42 @@ def _update_task_status(task_id: str, status_action: str, **kwargs) -> bool: return False + +def _update_task_progress(task_id: str, progress: float, stage: str = "") -> bool: + """更新 GenerationTask 进度(独立 session,异常不向外抛出)。 + + Args: + task_id: 任务 ID + progress: 进度值(0-100) + stage: 阶段描述(仅用于日志) + + Returns: + True 表示更新成功 + """ + try: + from packages.adapters.sqlalchemy_impl.models import GenerationTaskModel + + session = SessionLocal() + try: + model = session.query(GenerationTaskModel).filter( + GenerationTaskModel.id == task_id + ).first() + if model: + model.progress = progress + session.commit() + if stage: + logger.info( + "GenerationTask 进度更新: task_id=%s progress=%.0f%% stage=%s", + task_id, progress, stage, + ) + return True + return False + finally: + session.close() + except Exception as e: + logger.error("更新任务进度异常: task_id=%s progress=%s error=%s", task_id, progress, e) + return False + # ── 日志持久化辅助 ──────────────────────────────────────────────────────────── @@ -1283,6 +1319,7 @@ def generate_video(self, task_id: str) -> dict: # 标记任务为 running _update_task_status(task_id, "mark_processing") + _update_task_progress(task_id, 10, "任务启动") try: editing_mode = EditingMode(mode) if (mode := task_info["mode"]) else EditingMode.ONE_TAKE @@ -1316,7 +1353,10 @@ def generate_video(self, task_id: str) -> dict: ) _flush_logs(task_id, gen_task) + _update_task_progress(task_id, 30, "素材下载完成") + # ── 3. 渲染 + 混音 ─────────────────────────────────────────────── + _update_task_progress(task_id, 40, "开始渲染") output_path, render_duration = _render_video( task_id=task_id, downloaded_videos=downloaded_videos, @@ -1336,7 +1376,10 @@ def generate_video(self, task_id: str) -> dict: 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, @@ -1356,6 +1399,8 @@ def generate_video(self, task_id: str) -> dict: ) _flush_logs(task_id, gen_task) + _update_task_progress(task_id, 95, "上传完成") + # ── 5. 标记完成 ────────────────────────────────────────────────── _update_task_status(task_id, "mark_completed", result_count=video_count)