diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 6c76a3d..4d3efa4 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -1695,8 +1695,8 @@ def _classify_failure(last_error: str) -> str: async def handle_failed_task(task_id: str, project_id: str) -> dict: """处理一个失败任务:判定重跑或报故障。由 failed poller 周期调用。 - - retry_count < 3 且失败原因判定为瞬时 → 重跑(retry_count+1, 回 submitted) - - retry_count >= 3 或判定为永久错误 → 报故障(建问题通知用户,任务置 waiting) + - retry_count < 最大重复数(task_max_retry,默认3) 且失败原因判定为瞬时 → 重跑(retry_count+1, 回 submitted) + - retry_count >= 最大重复数 或判定为永久错误 → 暂停任务链(pause_project) + 报故障(建问题通知用户,任务置 waiting) """ db = _get_db() 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 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 await sor.sqlExe("COMMIT", {}) - # 已重试 3 次仍未成功 → 报故障 - if retry_count >= 3: + # 已重试 max_retry 次仍未成功 → 报故障 + if retry_count >= max_retry: decision = "fault" else: 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) 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 reporter = "agent.main_agent" if role == "agent.pm" else "agent.pm" qid = await raise_problem( "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, 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) - 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): diff --git a/pipeline_service/sdlc_ability.py b/pipeline_service/sdlc_ability.py index c03b13c..68c455e 100644 --- a/pipeline_service/sdlc_ability.py +++ b/pipeline_service/sdlc_ability.py @@ -50,6 +50,12 @@ SDL_TOOLS = [ parameters={"task_id": "任务ID"}, category="task", ), + ToolDefinition( + name="reset_task_retry", + description="人工处理后恢复任务:重复计数清零+重新执行+恢复项目推进(任务重复超限抛故障给人工后使用)", + parameters={"task_id": "任务ID"}, + category="task", + ), ToolDefinition( name="start_agents", description="启动角色agent执行已提交的任务", @@ -217,6 +223,7 @@ SDL_PROMPT = """你是「开发产线」的驾驶舱 agent,负责软件项目 - 用户报告异常/故障 → 先用 diagnose_project 定位根因,再用 task_detail 查详情。 - 仓库管理 → add_repo / list_repos / clone_repo。 - 发现卡点(任务卡死、审核超时、失败、僵尸 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 沿冒泡路径转下一个处理方)。 - 功能/需求管理 → 用户提新需求时 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 技能)。 @@ -458,6 +465,40 @@ async def _h_task_detail(sor, p, ctx): }, 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): pid = ctx.get("project_id", "") if not pid: @@ -1497,6 +1538,7 @@ SDL_HANDLERS = { "create_task": _h_create_task, "list_tasks": _h_list_tasks, "task_detail": _h_task_detail, + "reset_task_retry": _h_reset_task_retry, "start_agents": _h_start_agents, "diagnose_project": _h_diagnose_project, "list_deliverables": _h_list_deliverables, diff --git a/pipeline_service/task_capability.py b/pipeline_service/task_capability.py index 51c7ae5..aa1ea03 100644 --- a/pipeline_service/task_capability.py +++ b/pipeline_service/task_capability.py @@ -183,6 +183,32 @@ async def retry_task(task_id, tenant_id, who=None, agent_id=None): 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): """挂起等回答:running → waiting(清 claimed_by)。""" return await _transition(task_id, tenant_id, S_RUNNING, S_WAITING, 'suspend', diff --git a/pipeline_service/workspace.py b/pipeline_service/workspace.py index 7456a95..83bf811 100644 --- a/pipeline_service/workspace.py +++ b/pipeline_service/workspace.py @@ -60,6 +60,16 @@ async def get_max_concurrent_agents(sor): 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): """读当前项目的 workspace 目录(统一 workspace_*.dspy 的重复逻辑)。