feat: add pipeline_restart() - restart completed/failed/cancelled tasks with new version
This commit is contained in:
parent
99fd9753b2
commit
23047eed95
@ -264,7 +264,49 @@ async def pipeline_cancel(tenant_id, task_id):
|
|||||||
result["message"] = "任务已取消"
|
result["message"] = "任务已取消"
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
result["message"] = str(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():
|
def pipeline_handlers():
|
||||||
@ -302,6 +344,7 @@ def load_pipeline_service():
|
|||||||
env.pipeline_pause = pipeline_pause
|
env.pipeline_pause = pipeline_pause
|
||||||
env.pipeline_resume = pipeline_resume
|
env.pipeline_resume = pipeline_resume
|
||||||
env.pipeline_cancel = pipeline_cancel
|
env.pipeline_cancel = pipeline_cancel
|
||||||
|
env.pipeline_restart = pipeline_restart
|
||||||
|
|
||||||
# Handler management
|
# Handler management
|
||||||
env.pipeline_register_handler = register_handler
|
env.pipeline_register_handler = register_handler
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user