fix(task): reset_task_retry假成功空转——解链三步曲(resolve_problem resume_task=True先把waiting→submitted,再reset_task_retry)下旧CAS只认waiting/failed→UPDATE 0行,成功判定却只看终态==submitted→空转报True:retry_count未清零+陈旧last_error残留,任务带满计数重新认领一次退回即再触顶fault(2026-09-16 pbls M1b实测17:22 fault解链后18:08 retry_count又=3实锤)。修:state IN加submitted(running仍排除防打断在跑协程)+成功判定改state与retry_count双核对+入参态前置校验
This commit is contained in:
parent
5f98cd79a7
commit
f7a3fd364a
@ -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)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user