From 78b84c710e8264189c39d37466910cc41f2e4404 Mon Sep 17 00:00:00 2001 From: ymq Date: Fri, 21 Aug 2026 22:43:11 +0800 Subject: [PATCH] =?UTF-8?q?feat(sdlc):=20=E4=BC=9A=E8=AF=9Dagent=E5=8A=A0?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E6=B5=81=E8=BD=AC=E5=B7=A5=E5=85=B7=20approv?= =?UTF-8?q?e=5Ftask/complete=5Ftask/cancel=5Ftask=EF=BC=88PM=E6=9D=83?= =?UTF-8?q?=E9=99=90=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - task_capability 新增 cancel_task 原语(CAS+审计+终态守护) - sdlc_ability 加 3 个 ToolDefinition + handler,who=agent.pm 记录 PM 权限操作 - _resolve_task_id 租户隔离前缀兜底 - SDL_PROMPT 补任务流转典型场景 - _pm_cancel_task 委托新原语,去掉重复 SQL --- pipeline_service/agent_loop.py | 12 ++-- pipeline_service/sdlc_ability.py | 95 ++++++++++++++++++++++++++++- pipeline_service/task_capability.py | 31 ++++++++++ 3 files changed, 131 insertions(+), 7 deletions(-) diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 9cd8ed8..06e185e 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -1430,8 +1430,8 @@ async def _pm_list_tasks(sor, project_id, params): async def _pm_cancel_task(sor, project_id, params): """PM 取消任务(重做/作废场景必须先取消旧任务,避免两个相同任务并存)。 - 标记 state=cancelled + 清 claimed_by;正在执行中的 agent 会在下一轮心跳检测到 - cancelled 并自行中止(见 role_agent_run 的取消检测)。 + 委托 task_capability.cancel_task(CAS + 审计),正在执行中的 agent 会在下一轮心跳 + 检测到 cancelled 并自行中止(见 role_agent_run 的取消检测)。 """ task_id = (params.get('task_id') or params.get('id') or '').strip() if not task_id: @@ -1446,10 +1446,10 @@ async def _pm_cancel_task(sor, project_id, params): state = getattr(recs[0], 'state', '') if state in ('completed', 'cancelled', 'failed', 'approved'): return f'任务「{title}」已是终态 {state},无需取消' - await sor.sqlExe( - "UPDATE pipeline_tasks SET state='cancelled', claimed_by=NULL, updated_at=NOW() WHERE id=${tid}$", - {"tid": task_id}) - await sor.sqlExe("COMMIT", {}) + from .task_capability import cancel_task + ok, msg = await cancel_task(task_id, project_id, who="agent.pm") + if not ok: + return f'FAIL: {msg}' return f'OK: 已取消任务「{title}」({task_id}),正在执行的 agent 将在下一轮心跳中止' diff --git a/pipeline_service/sdlc_ability.py b/pipeline_service/sdlc_ability.py index 15870ba..8eebc8a 100644 --- a/pipeline_service/sdlc_ability.py +++ b/pipeline_service/sdlc_ability.py @@ -56,6 +56,24 @@ SDL_TOOLS = [ parameters={"task_id": "任务ID"}, category="task", ), + ToolDefinition( + name="approve_task", + description="审核通过任务(review→approved,PM 验收)", + parameters={"task_id": "任务ID", "comment": "审核意见(可选)"}, + category="task", + ), + ToolDefinition( + name="complete_task", + description="标记任务完成(approved/review→completed)", + parameters={"task_id": "任务ID"}, + category="task", + ), + ToolDefinition( + name="cancel_task", + description="取消任务(重做/作废前必须先取消旧任务,避免两个相同任务并存)", + parameters={"task_id": "任务ID", "comment": "取消原因(可选)"}, + category="task", + ), ToolDefinition( name="start_agents", description="启动角色agent执行已提交的任务", @@ -238,7 +256,8 @@ SDL_PROMPT = """你是「开发产线」的驾驶舱 agent,负责软件项目 - 测试计划 → create_test_plan 建计划、list_plans 列计划;审批/执行/完成用 approve_plan / start_plan / complete_plan(规范见 test-plan 技能)。 - 测试用例 → create_case 建用例、list_cases 列用例;执行结果用 pass_case / fail_case / skip_case / block_case(失败用例应联动 report_bug,规范见 test-case 技能)。 - Bug 管理 → report_bug 上报、list_bugs 列 Bug;确认/修复/验证/关闭/驳回/重开用 confirm_bug / start_fix / fix_bug / verify_bug / close_bug / reject_bug / reopen_bug(规范见 bug 技能)。 -- 部署环境 → configure_env 配环境、list_envs 列环境;验证用 verify_env(规范见 deploy-env 技能)。""" +- 部署环境 → configure_env 配环境、list_envs 列环境;验证用 verify_env(规范见 deploy-env 技能)。 +- 任务流转(PM 权限操作)→ 审核通过任务 approve_task(review→approved)、标记完成 complete_task(approved/review→completed)、取消任务 cancel_task(重做/作废前必须先取消旧任务)。这三个是 PM 验收级操作,仅在用户明确要求审核/验收/完成/取消任务时使用,不主动越权推进。""" # ── 产线角色集(可插拔,替代硬编码 ROLE_SPECIFICS/ROLE_ALIASES/ROLE_CHAIN) ── @@ -530,6 +549,59 @@ async def _h_reset_task_retry(sor, p, ctx): return f"OK: 任务重复计数已归零并重新执行{resume_note}" +async def _h_approve_task(sor, p, ctx): + """审核通过任务(PM 验收):review → approved。""" + task_id = (p.get("task_id") or "").strip() + if not task_id: + return "需要任务ID" + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + full = await _resolve_task_id(sor, task_id, pid) + if not full: + return f"任务不存在: {task_id}" + from .task_capability import approve_task + ok, msg = await approve_task(full, pid, who="agent.pm", + agent_id=ctx.get("user_id", ""), + comment=(p.get("comment") or "").strip()) + return f"OK: {msg}" if ok else f"ERROR: {msg}" + + +async def _h_complete_task(sor, p, ctx): + """标记任务完成(PM):approved/review → completed。""" + task_id = (p.get("task_id") or "").strip() + if not task_id: + return "需要任务ID" + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + full = await _resolve_task_id(sor, task_id, pid) + if not full: + return f"任务不存在: {task_id}" + from .task_capability import complete_task + ok, msg = await complete_task(full, pid, who="agent.pm", + agent_id=ctx.get("user_id", "")) + return f"OK: {msg}" if ok else f"ERROR: {msg}" + + +async def _h_cancel_task(sor, p, ctx): + """取消任务(PM):任意非终态 → cancelled,清 claimed_by。""" + task_id = (p.get("task_id") or "").strip() + if not task_id: + return "需要任务ID" + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + full = await _resolve_task_id(sor, task_id, pid) + if not full: + return f"任务不存在: {task_id}" + from .task_capability import cancel_task + ok, msg = await cancel_task(full, pid, who="agent.pm", + agent_id=ctx.get("user_id", ""), + comment=(p.get("comment") or "").strip()) + return f"OK: {msg}" if ok else f"ERROR: {msg}" + + async def _h_start_agents(sor, p, ctx): pid = ctx.get("project_id", "") if not pid: @@ -952,6 +1024,24 @@ async def _resolve_id(sor, table, fid): return "" +async def _resolve_task_id(sor, task_id, pid): + """任务 ID 前缀兜底(租户隔离):精确匹配失败且入参≥6位时 LIKE 前缀解析完整 ID。""" + if not task_id: + return "" + recs = await sor.sqlExe( + "SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND tenant_id=${pid}$", + {"tid": task_id, "pid": pid}) + if recs: + return getattr(recs[0], "id", "") + if 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 recs: + return getattr(recs[0], "id", "") + return "" + + async def _get_scope(sor, table, id_field, id_val, scope_field): """反查实体的 scope 列值(如 iteration_id/plan_id)。返回 scope_val 或空。""" recs = await sor.sqlExe( @@ -1593,6 +1683,9 @@ SDL_HANDLERS = { "list_tasks": _h_list_tasks, "task_detail": _h_task_detail, "reset_task_retry": _h_reset_task_retry, + "approve_task": _h_approve_task, + "complete_task": _h_complete_task, + "cancel_task": _h_cancel_task, "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 68eea6c..070f474 100644 --- a/pipeline_service/task_capability.py +++ b/pipeline_service/task_capability.py @@ -204,6 +204,37 @@ async def mark_failed(task_id, tenant_id, who=None, agent_id=None, error=None): return True, S_FAILED +async def cancel_task(task_id, tenant_id, who=None, agent_id=None, comment=None): + """取消任务(任意非终态 → cancelled,清 claimed_by)。 + + 重做/作废场景必须先取消旧任务,避免两个相同任务并存;正在执行的 agent 会在 + 下一轮心跳检测到 cancelled 后自行中止。终态(completed/cancelled/failed/approved) + 不重复取消。CAS:仅当任务仍处原非终态时取消,防并发覆盖。 + """ + if not task_id or not tenant_id: + return False, "缺少 task_id 或 tenant_id" + 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', '') + if from_state in (S_COMPLETED, 'cancelled', S_FAILED, S_APPROVED): + return False, f"任务已是终态 {from_state},无需取消" + await sor.sqlExe( + "UPDATE pipeline_tasks SET state='cancelled', claimed_by=NULL, " + "updated_at=NOW() WHERE id=${tid}$ AND tenant_id=${tn}$ AND state=${from}$", + {"tid": task_id, "tn": tenant_id, "from": from_state}) + chk = await sor.R('pipeline_tasks', {'id': task_id}) + await sor.sqlExe("COMMIT", {}) + if not chk or getattr(chk[0], 'state', '') != 'cancelled': + return False, "取消失败(CAS): 任务状态已变化" + await record_audit(tenant_id, 'pipeline_tasks', task_id, 'cancel', + from_state=from_state, to_state='cancelled', + who=who, agent_id=agent_id, detail=comment, sor=sor) + return True, 'cancelled' + + async def retry_task(task_id, tenant_id, who=None, agent_id=None): """重试:failed → submitted(retry_count+1,清 claimed_by)。""" db, dbname = _get_db()