feat: 任务最大重复数(task_max_retry)——超限暂停任务链+抛故障给人工,人工处理后reset_task_retry清零重执行

This commit is contained in:
ymq 2026-08-19 07:29:08 +08:00
parent 981b58b0d2
commit 536f51434f
4 changed files with 97 additions and 8 deletions

View File

@ -1695,8 +1695,8 @@ def _classify_failure(last_error: str) -> str:
async def handle_failed_task(task_id: str, project_id: str) -> dict: async def handle_failed_task(task_id: str, project_id: str) -> dict:
"""处理一个失败任务:判定重跑或报故障。由 failed poller 周期调用。 """处理一个失败任务:判定重跑或报故障。由 failed poller 周期调用。
- retry_count < 3 且失败原因判定为瞬时 → 重跑(retry_count+1, 回 submitted) - retry_count < 最大重复数(task_max_retry,默认3) 且失败原因判定为瞬时 → 重跑(retry_count+1, 回 submitted)
- retry_count >= 3 或判定为永久错误 → 报故障(建问题通知用户,任务置 waiting) - retry_count >= 最大重复数 或判定为永久错误 → 暂停任务链(pause_project) + 报故障(建问题通知用户,任务置 waiting)
""" """
db = _get_db() db = _get_db()
async with db.sqlorContext("pipeline") as sor: async with db.sqlorContext("pipeline") as sor:
@ -1717,11 +1717,15 @@ async def handle_failed_task(task_id: str, project_id: str) -> dict:
retry_count = 0 retry_count = 0
last_error = getattr(t, "last_error", "") or "" last_error = getattr(t, "last_error", "") or ""
# 最大重复数:appbase params 表 task_max_retry(默认 3),超限暂停任务链 + 抛故障给人工
from .workspace import get_max_task_retry
max_retry = await get_max_task_retry(sor)
# 释放 SELECT 的元数据锁,避免决策/建问题期间长时间持有 MDL # 释放 SELECT 的元数据锁,避免决策/建问题期间长时间持有 MDL
await sor.sqlExe("COMMIT", {}) await sor.sqlExe("COMMIT", {})
# 已重试 3 次仍未成功 → 报故障 # 已重试 max_retry 次仍未成功 → 报故障
if retry_count >= 3: if retry_count >= max_retry:
decision = "fault" decision = "fault"
else: else:
decision = _classify_failure(last_error) decision = _classify_failure(last_error)
@ -1732,17 +1736,24 @@ async def handle_failed_task(task_id: str, project_id: str) -> dict:
logger.info("failed task retry: task=%s role=%s attempt=%d ok=%s", task_id, role, retry_count + 1, ok) logger.info("failed task retry: task=%s role=%s attempt=%d ok=%s", task_id, role, retry_count + 1, ok)
return {"status": "retry", "task_id": task_id, "attempt": retry_count + 1} return {"status": "retry", "task_id": task_id, "attempt": retry_count + 1}
# 报故障:建问题通知用户,任务置 waiting 等待人工介入 # 报故障:① 暂停任务链(pause_project,PM 不再推进)② 建问题通知用户,任务置 waiting 等人工介入
try:
from .project_capability import pause_project
pok, pmsg = await pause_project(project_id, who="agent.pm")
logger.info("failed task pause project: task=%s pause=%s msg=%s", task_id, pok, pmsg)
except Exception as e:
logger.warning("failed task pause_project error: %s", e)
from .communication import raise_problem from .communication import raise_problem
reporter = "agent.main_agent" if role == "agent.pm" else "agent.pm" reporter = "agent.main_agent" if role == "agent.pm" else "agent.pm"
qid = await raise_problem( qid = await raise_problem(
"fault_report", "fault_report",
f"任务「{title}」已失败 {retry_count} 次,自动重试无法解决,请人工介入。失败原因:{last_error or '未知'}", f"任务「{title}」重复 {retry_count} 次已达上限({max_retry}),任务链已暂停,请人工介入处理。失败原因:{last_error or '未知'}",
reporter, tenant_id=project_id, task_id=task_id, reporter, tenant_id=project_id, task_id=task_id,
first_handler_role="agent.main_agent", first_handler_role="agent.main_agent",
context={"fault": True, "last_error": last_error, "retry_count": retry_count}) context={"fault": True, "last_error": last_error, "retry_count": retry_count,
"max_retry": max_retry, "chain_paused": True})
logger.info("failed task fault: task=%s role=%s qid=%s", task_id, role, qid) logger.info("failed task fault: task=%s role=%s qid=%s", task_id, role, qid)
return {"status": "fault", "task_id": task_id, "question_id": qid} return {"status": "fault", "task_id": task_id, "question_id": qid, "chain_paused": True}
async def role_agent_loop(project_id, role, agent_id=None, model_name=None, max_iterations=10): async def role_agent_loop(project_id, role, agent_id=None, model_name=None, max_iterations=10):

View File

@ -50,6 +50,12 @@ SDL_TOOLS = [
parameters={"task_id": "任务ID"}, parameters={"task_id": "任务ID"},
category="task", category="task",
), ),
ToolDefinition(
name="reset_task_retry",
description="人工处理后恢复任务:重复计数清零+重新执行+恢复项目推进(任务重复超限抛故障给人工后使用)",
parameters={"task_id": "任务ID"},
category="task",
),
ToolDefinition( ToolDefinition(
name="start_agents", name="start_agents",
description="启动角色agent执行已提交的任务", description="启动角色agent执行已提交的任务",
@ -217,6 +223,7 @@ SDL_PROMPT = """你是「开发产线」的驾驶舱 agent,负责软件项目
- 用户报告异常/故障 → 先用 diagnose_project 定位根因,再用 task_detail 查详情。 - 用户报告异常/故障 → 先用 diagnose_project 定位根因,再用 task_detail 查详情。
- 仓库管理 → add_repo / list_repos / clone_repo。 - 仓库管理 → add_repo / list_repos / clone_repo。
- 发现卡点(任务卡死、审核超时、失败、僵尸 claimed_by)时,主动定位根因并推动修复。 - 发现卡点(任务卡死、审核超时、失败、僵尸 claimed_by)时,主动定位根因并推动修复。
- 任务重复超限(retry_count 达 task_max_retry 上限)→ 系统已暂停任务链并抛 fault_report 故障;你(主 agent)收到故障后先冒泡给用户(ask_user/answer_question),用户处理完毕调 reset_task_retry 恢复该任务(重复计数清零+重新执行+恢复项目推进)。
- 发现「待回答问题」> 0 或角色提问时,加载 team-communication 技能按冒泡链处理(list_questions 列问题 → 能答则 answer_question → 不能答 escalate_question 沿冒泡路径转下一个处理方)。 - 发现「待回答问题」> 0 或角色提问时,加载 team-communication 技能按冒泡链处理(list_questions 列问题 → 能答则 answer_question → 不能答 escalate_question 沿冒泡路径转下一个处理方)。
- 功能/需求管理 → 用户提新需求时 propose_feature;查功能清单 list_features;需求评审/验收时 approve_feature / reject_feature / verify_feature(状态机规范见 feature 技能)。 - 功能/需求管理 → 用户提新需求时 propose_feature;查功能清单 list_features;需求评审/验收时 approve_feature / reject_feature / verify_feature(状态机规范见 feature 技能)。
- 项目管理 → list_projects 列项目;启动/完成/归档/重开项目用 start_project / complete_project / archive_project / reopen_project;暂停/恢复推进用 pause_project / resume_project(仅用户明确指令「暂停推进」时才 pause,默认必须推进,状态机规范见 project 技能)。 - 项目管理 → list_projects 列项目;启动/完成/归档/重开项目用 start_project / complete_project / archive_project / reopen_project;暂停/恢复推进用 pause_project / resume_project(仅用户明确指令「暂停推进」时才 pause,默认必须推进,状态机规范见 project 技能)。
@ -458,6 +465,40 @@ async def _h_task_detail(sor, p, ctx):
}, ensure_ascii=False, indent=2) }, ensure_ascii=False, indent=2)
async def _h_reset_task_retry(sor, p, ctx):
"""人工处理后恢复:任务重复计数清零 + 重新执行 + 恢复项目推进(任务链)。
任务重复超限后任务链暂停、故障抛给人工;人工处理完毕调用本工具恢复该任务并继续执行。
"""
task_id = (p.get("task_id") or "").strip()
if not task_id:
return "需要任务ID"
pid = ctx.get("project_id", "")
if not pid:
return "请先切换到项目"
# 前缀解析(LLM 可能传截断 ID)
recs = await sor.sqlExe(
"SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND tenant_id=${pid}$",
{"tid": task_id, "pid": pid})
if not recs and len(task_id) >= 6:
recs = await sor.sqlExe(
"SELECT id FROM pipeline_tasks WHERE id LIKE ${prefix}$ AND tenant_id=${pid}$ LIMIT 1",
{"prefix": task_id + "%", "pid": pid})
if not recs:
return f"任务不存在: {task_id}"
task_id = getattr(recs[0], 'id', '') or task_id
from .task_capability import reset_task_retry
ok, msg = await reset_task_retry(task_id, pid, who="agent.main_agent")
if not ok:
return f"ERROR: {msg}"
# 恢复项目推进(若因故障被 pause_project 暂停)
from .project_capability import resume_project
rok, rmsg = await resume_project(pid, who="agent.main_agent")
resume_note = ",项目已恢复推进" if rok else f"(项目恢复: {rmsg})"
return f"OK: 任务重复计数已归零并重新执行{resume_note}"
async def _h_start_agents(sor, p, ctx): async def _h_start_agents(sor, p, ctx):
pid = ctx.get("project_id", "") pid = ctx.get("project_id", "")
if not pid: if not pid:
@ -1497,6 +1538,7 @@ SDL_HANDLERS = {
"create_task": _h_create_task, "create_task": _h_create_task,
"list_tasks": _h_list_tasks, "list_tasks": _h_list_tasks,
"task_detail": _h_task_detail, "task_detail": _h_task_detail,
"reset_task_retry": _h_reset_task_retry,
"start_agents": _h_start_agents, "start_agents": _h_start_agents,
"diagnose_project": _h_diagnose_project, "diagnose_project": _h_diagnose_project,
"list_deliverables": _h_list_deliverables, "list_deliverables": _h_list_deliverables,

View File

@ -183,6 +183,32 @@ async def retry_task(task_id, tenant_id, who=None, agent_id=None):
return True, S_SUBMITTED return True, S_SUBMITTED
async def reset_task_retry(task_id, tenant_id, who=None, agent_id=None):
"""人工处理后恢复:retry_count 归零,waiting/failed → submitted 重新执行。
任务重复超限后任务链暂停、故障抛给人工;人工处理完毕调用本函数把该任务重复数清零并重新执行。
"""
db, dbname = _get_db()
async with db.sqlorContext(dbname) as sor:
recs = await sor.R('pipeline_tasks', {'id': task_id})
if not recs:
return False, "任务不存在"
from_state = getattr(recs[0], 'state', '')
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')",
{"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 状态)"
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)
return True, S_SUBMITTED
async def suspend_task(task_id, tenant_id, who=None, agent_id=None): async def suspend_task(task_id, tenant_id, who=None, agent_id=None):
"""挂起等回答:running → waiting(清 claimed_by)。""" """挂起等回答:running → waiting(清 claimed_by)。"""
return await _transition(task_id, tenant_id, S_RUNNING, S_WAITING, 'suspend', return await _transition(task_id, tenant_id, S_RUNNING, S_WAITING, 'suspend',

View File

@ -60,6 +60,16 @@ async def get_max_concurrent_agents(sor):
return 3 return 3
async def get_max_task_retry(sor):
"""读任务最大重复数(appbase params 表 task_max_retry,默认 3)。超限即暂停任务链、抛故障给人工。"""
val = await get_param(sor, 'task_max_retry', '3')
try:
n = int(float(str(val)))
return max(1, n)
except (ValueError, TypeError):
return 3
async def get_workspace_dir(sor, uid): async def get_workspace_dir(sor, uid):
"""读当前项目的 workspace 目录(统一 workspace_*.dspy 的重复逻辑)。 """读当前项目的 workspace 目录(统一 workspace_*.dspy 的重复逻辑)。