From 23047eed95a03aec1afe9522e95362186b8438c6 Mon Sep 17 00:00:00 2001 From: yumoqing Date: Wed, 8 Jul 2026 21:03:20 +0800 Subject: [PATCH] feat: add pipeline_restart() - restart completed/failed/cancelled tasks with new version --- pipeline_service/init.py | 45 +++++++++++++++++++++++++++++++++++++++- 1 file changed, 44 insertions(+), 1 deletion(-) diff --git a/pipeline_service/init.py b/pipeline_service/init.py index c5e4b09..40738b0 100644 --- a/pipeline_service/init.py +++ b/pipeline_service/init.py @@ -264,7 +264,49 @@ async def pipeline_cancel(tenant_id, task_id): result["message"] = "任务已取消" except Exception as e: result["message"] = str(e) - return json.dumps(result, ensure_ascii=False) + return json.dumps(result, ensure_ascii=False, default=str) + + +async def pipeline_restart(tenant_id, task_id): + """重启已完成/失败/取消的任务 — 重置所有步骤为pending,新版本继续。""" + result = {"success": False} + try: + task = await get_task(tenant_id, task_id) + if not task: + result["message"] = "任务不存在" + return json.dumps(result, ensure_ascii=False) + + if is_running(task_id): + result["message"] = "任务正在执行中,请先暂停" + return json.dumps(result, ensure_ascii=False) + + # Stop any lingering executor + await stop_task(task_id) + + # Get step names to reset + steps = await get_task_steps(task_id) + step_names = [s['step_name'] for s in steps] + + # New version + current_version = task.get("current_version", task.get("current_Version", 1)) + if isinstance(current_version, str): + current_version = int(current_version) + new_version = current_version + 1 + await update_task_version(task_id, new_version) + + # Reset all steps to pending + await reset_steps(task_id, step_names) + + # Restart execution + await update_task_state(task_id, TASK_RUNNING) + await start_task(task_id) + + result["success"] = True + result["new_version"] = new_version + result["message"] = f"任务已重新启动 v{new_version}" + except Exception as e: + result["message"] = str(e) + return json.dumps(result, ensure_ascii=False, default=str) def pipeline_handlers(): @@ -302,6 +344,7 @@ def load_pipeline_service(): env.pipeline_pause = pipeline_pause env.pipeline_resume = pipeline_resume env.pipeline_cancel = pipeline_cancel + env.pipeline_restart = pipeline_restart # Handler management env.pipeline_register_handler = register_handler