diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index ff03ec8..3f62d77 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -507,17 +507,19 @@ __REPO_STATE__ - run_shell(command) — 执行命令(编译/测试验证) - create_tasks(tasks) — 批量创建并派发后续任务(tasks 是 JSON 数组,每项 {title, role, description}) - list_tasks(role, state) — 列出项目现有任务(派发前先查,避免重复) +- cancel_task(task_id) — 取消任务(重做/作废前必须先取消旧任务,避免两个相同任务并存) ## 审核流程 1. 先用工具检查代码/交付件是否实际产出(git_status / list_files / read_file) 2. 如需编译验证,用 run_shell -3. 若审核将通过、且本任务产出需要拆成多个后续任务,先 list_tasks 查重,再 create_tasks 派发 +3. 若审核将通过、且本任务产出需要拆成多个后续任务,先 list_tasks 查重,再 create_tasks 派发;发现需重做的旧任务仍在活跃,先 cancel_task 取消,再 create_tasks 4. 最后给出决策 ## 输出格式(每次一个JSON) 查看文件:{"action":"tool_call","tool":"read_file","params":{"path":"相对路径"}} 查看git:{"action":"tool_call","tool":"git_status","params":{}} 查现有任务:{"action":"tool_call","tool":"list_tasks","params":{"role":"develop","state":"submitted"}} +取消旧任务:{"action":"tool_call","tool":"cancel_task","params":{"task_id":"旧任务ID"}} 派发后续任务:{"action":"tool_call","tool":"create_tasks","params":{"tasks":[{"title":"契约定义","role":"design","description":"定义数据模型/API/目录规范"},{"title":"hr-system 开发","role":"develop","description":"实现 hr-system 基础能力"}]}} 批准:{"action":"review_approve","comment":"审核意见","next_task_title":"下阶段标题","next_task_description":"描述"} 驳回:{"action":"review_reject","comment":"原因","questions":"修改要求"} @@ -1007,6 +1009,12 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None): "UPDATE pipeline_tasks SET updated_at=NOW() WHERE id=${tid}$ AND state='running'", {"tid": task_id}) await sor.sqlExe("COMMIT", {}) + # 取消检测:任务被 cancel_task 标记 cancelled 后立即中止,避免与重派的新任务重复执行 + _st = await sor.sqlExe("SELECT state FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id}) + await sor.sqlExe("COMMIT", {}) + if _st and getattr(_st[0], 'state', '') == 'cancelled': + logger.info(f"task cancelled mid-run: {task_id}") + return {"status": "cancelled", "task_id": task_id} try: resp = await llm_call_msgs_native(msgs, tools=tools_schema, model=model_name, temperature=0.4, org_id=org_id) except Exception as e: @@ -1175,6 +1183,7 @@ async def _pm_create_tasks(sor, project_id, params): iteration_name = getattr(irecs[0], 'iteration_name', '') or '' created = [] + cancelled = [] for t in tasks: if not isinstance(t, dict): continue @@ -1183,6 +1192,20 @@ async def _pm_create_tasks(sor, project_id, params): continue role = _normalize_role(t.get('role') or 'agent.develop') desc = (t.get('description') or '').strip() + + # 重做保护:同项目存在同标题活跃任务(submitted/running/review/qc_review)时, + # 先自动取消旧任务再创建新任务,避免两个相同任务并存(重做必须先终止旧任务)。 + dup = await sor.sqlExe( + "SELECT id FROM pipeline_tasks WHERE tenant_id=${pid}$ AND title=${title}$ " + "AND state IN ('submitted','running','review','qc_review')", + {"pid": project_id, "title": title}) + for d in (dup or []): + did = getattr(d, 'id', '') + await sor.sqlExe( + "UPDATE pipeline_tasks SET state='cancelled', claimed_by=NULL, updated_at=NOW() WHERE id=${did}$", + {"did": did}) + cancelled.append(did) + tparams = {"description": desc, "pm_assigned": True} if iteration_name: tparams["iteration_id"] = iteration_name @@ -1198,7 +1221,8 @@ async def _pm_create_tasks(sor, project_id, params): return 'FAIL: 没有可创建的任务(缺少 title)' # C 之后立即 COMMIT,释放行锁,供 agent_poller 下一轮可见认领 await sor.sqlExe("COMMIT", {}) - return f"OK: 已派发 {len(created)} 个任务:" + ";".join(created) + note = f"(已自动取消 {len(cancelled)} 个同标题活跃任务)" if cancelled else "" + return f"OK: 已派发 {len(created)} 个任务" + note + ":" + ";".join(created) async def _pm_list_tasks(sor, project_id, params): @@ -1228,6 +1252,32 @@ async def _pm_list_tasks(sor, project_id, params): return "\n".join(lines) +async def _pm_cancel_task(sor, project_id, params): + """PM 取消任务(重做/作废场景必须先取消旧任务,避免两个相同任务并存)。 + + 标记 state=cancelled + 清 claimed_by;正在执行中的 agent 会在下一轮心跳检测到 + cancelled 并自行中止(见 role_agent_run 的取消检测)。 + """ + task_id = (params.get('task_id') or params.get('id') or '').strip() + if not task_id: + return 'FAIL: 需要 task_id' + recs = await sor.sqlExe( + "SELECT id, title, state FROM pipeline_tasks WHERE id=${tid}$ AND tenant_id=${pid}$", + {"tid": task_id, "pid": project_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return 'FAIL: 任务不存在或不属于当前项目' + title = getattr(recs[0], 'title', '') + 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", {}) + return f'OK: 已取消任务「{title}」({task_id}),正在执行的 agent 将在下一轮心跳中止' + + async def pm_review_run(project_id, agent_id=None, model_name=None): db = _get_db() async with db.sqlorContext("pipeline") as sor: @@ -1279,6 +1329,12 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): "UPDATE pipeline_tasks SET updated_at=NOW() WHERE id=${tid}$ AND state='review'", {"tid": task_id}) await sor.sqlExe("COMMIT", {}) + # 取消检测:任务被 cancel_task 标记 cancelled 后立即中止审核 + _st = await sor.sqlExe("SELECT state FROM pipeline_tasks WHERE id=${tid}$", {"tid": task_id}) + await sor.sqlExe("COMMIT", {}) + if _st and getattr(_st[0], 'state', '') == 'cancelled': + logger.info(f"pm review cancelled mid-run: {task_id}") + return {"status": "cancelled", "task_id": task_id} # auto-inject 兜底:最后两轮强制要求给出 review_* 决策,禁止再 tool_call, # 否则 deepseek 会一直读文件/git_status 耗尽 5 轮 → 审核超时 → 打回重跑 → 死循环。 # 例外:允许 create_tasks(里程碑拆后续任务的「项目计划/任务分配」职责),派发后必须立即决策。 @@ -1317,6 +1373,8 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): result = await _pm_create_tasks(sor, project_id, params) elif tool == 'list_tasks': result = await _pm_list_tasks(sor, project_id, params) + elif tool in ('cancel_task', 'cancel'): + result = await _pm_cancel_task(sor, project_id, params) else: result = await _exec_agent_tool(tool, params, workspace_dir) msgs.append({"role": "assistant", "content": raw})