diff --git a/pipeline_service/task_capability.py b/pipeline_service/task_capability.py index 58aae0a..64b4b0d 100644 --- a/pipeline_service/task_capability.py +++ b/pipeline_service/task_capability.py @@ -292,9 +292,16 @@ async def retry_task(task_id, tenant_id, who=None, agent_id=None): async def reset_task_retry(task_id, tenant_id, who=None, agent_id=None): - """人工处理后恢复:retry_count 归零,waiting/failed → submitted 重新执行。 + """人工处理后恢复:retry_count 归零,waiting/failed/submitted → submitted 重新执行。 任务重复超限后任务链暂停、故障抛给人工;人工处理完毕调用本函数把该任务重复数清零并重新执行。 + + state IN 含 submitted(2026-09-16 实测 bug):标准解链三步 resolve_problem + (resume_task=True) → reset_task_retry → resume_project 中,resolve 已把 + waiting→submitted,本函数旧 CAS 只认 waiting/failed → UPDATE 0 行;而成功 + 判定只看「终态==submitted」→ 空转也报 True,retry_count 从未清零、陈旧 + last_error 残留,任务带着满重试计数重新认领,一次退回即再触顶 fault。 + running 态仍排除(防中途拽回 submitted 打断在跑协程)。 """ db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: @@ -302,15 +309,20 @@ async def reset_task_retry(task_id, tenant_id, who=None, agent_id=None): if not recs: return False, "任务不存在" from_state = getattr(recs[0], 'state', '') + if from_state not in ('waiting', 'failed', 'submitted'): + return False, f"恢复失败(任务状态 {from_state} 不在 waiting/failed/submitted)" await sor.sqlExe( "UPDATE pipeline_tasks SET retry_count=0, state='submitted', " "claimed_by=NULL, last_error=NULL, updated_at=NOW() " - "WHERE id=${tid}$ AND tenant_id=${tn}$ AND state IN ('waiting','failed')", + "WHERE id=${tid}$ AND tenant_id=${tn}$ AND state IN ('waiting','failed','submitted')", {"tid": task_id, "tn": tenant_id}) chk = await sor.R('pipeline_tasks', {'id': task_id}) await sor.sqlExe("COMMIT", {}) - if not chk or getattr(chk[0], 'state', '') != S_SUBMITTED: - return False, "恢复失败(任务不在 waiting/failed 状态)" + # 成功判定=UPDATE 真实生效(state 与 retry_count 双核对),不只看终态—— + # 终态可能本来就是 submitted(空转假成功,见 docstring 实测事故)。 + if (not chk or getattr(chk[0], 'state', '') != S_SUBMITTED + or int(getattr(chk[0], 'retry_count', -1) or -1) != 0): + return False, "恢复失败(计数未清零或状态竞态,请重试)" await record_audit(tenant_id, 'pipeline_tasks', task_id, 'reset_retry', from_state=from_state, to_state=S_SUBMITTED, who=who, agent_id=agent_id, sor=sor)