diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 8ef1082..48d2e87 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -415,7 +415,16 @@ __REPO_STATE__ - test:测试覆盖充分、发现问题记录完整 - deploy:部署配置完整、可一键部署""" -PM_SYSTEM_PROMPT = """你是项目经理(PM)。审核交付件并推动项目前进。 +PM_SYSTEM_PROMPT = """你是项目经理(PM)。你的职责是项目计划、任务分配、任务验收。 + +## 三大职责 +1. 项目计划:基于需求/设计交付件,规划后续工作——把里程碑任务拆解为可执行的子任务(每个含标题、角色、描述)。 +2. 任务分配:用 create_tasks 工具把子任务派发给对应角色 agent(如「总体设计」批准后 → 拆「契约定义」任务 + hr-system/hr-org/hr-roster 等各模块 develop 任务)。 +3. 任务验收:检查交付件 + git 实际产出,决定 review_approve / review_reject / review_complete。 + +## 任务链不能断(关键规则) +- 审核通过里程碑任务(requirement / design,产出含多个应用/模块,或评审意见指明后续要落地开发)后,必须先规划并用 create_tasks 创建后续开发任务,再 review_approve。 +- 系统在 approve 后会自动补一个线性链的下一角色任务作为兜底;里程碑的并发拆分由你显式创建,不要等、不要断链。 ## 项目信息 工作目录:__WORKSPACE__ @@ -433,15 +442,20 @@ __REPO_STATE__ - list_files(path) — 列目录 - git_status() — 查看git状态 - run_shell(command) — 执行命令(编译/测试验证) +- create_tasks(tasks) — 批量创建并派发后续任务(tasks 是 JSON 数组,每项 {title, role, description}) +- list_tasks(role, state) — 列出项目现有任务(派发前先查,避免重复) ## 审核流程 -1. 先用工具检查代码是否实际产出(git_status/list_files/read_file) +1. 先用工具检查代码/交付件是否实际产出(git_status / list_files / read_file) 2. 如需编译验证,用 run_shell -3. 最后给出决策 +3. 若审核将通过、且本任务产出需要拆成多个后续任务,先 list_tasks 查重,再 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":"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":"修改要求"} 完成:{"action":"review_complete","comment":"总结"} @@ -945,6 +959,89 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None): # ── PM 审核 ── +async def _pm_create_tasks(sor, project_id, params): + """PM 派发后续任务(项目计划 / 任务分配):批量创建任务并指定角色。 + + params.tasks: JSON 数组 [{title, role, description}],一次派发多个并发子任务。 + 兼容单任务形态:params 直接含 {title, role, description}。 + """ + from appPublic.uniqueID import getID + tasks = params.get('tasks') or [] + if isinstance(tasks, str): + try: + tasks = json.loads(tasks) + except (json.JSONDecodeError, ValueError): + return 'FAIL: tasks 必须是 JSON 数组' + if not tasks: + if params.get('title'): + tasks = [{'title': params.get('title'), 'role': params.get('role'), 'description': params.get('description')}] + else: + return 'FAIL: 需要 tasks 数组或 title' + if not isinstance(tasks, list): + return 'FAIL: tasks 必须是数组' + + # 当前迭代(该 project 最新创建的迭代),作为任务默认归属 + iteration_name = '' + irecs = await sor.sqlExe( + "SELECT iteration_name FROM sd_iterations WHERE project_id=${pid}$ ORDER BY created_at DESC LIMIT 1", + {"pid": project_id}) + if irecs: + iteration_name = getattr(irecs[0], 'iteration_name', '') or '' + + created = [] + for t in tasks: + if not isinstance(t, dict): + continue + title = (t.get('title') or '').strip() + if not title: + continue + role = _normalize_role(t.get('role') or 'agent.develop') + desc = (t.get('description') or '').strip() + tparams = {"description": desc, "pm_assigned": True} + if iteration_name: + tparams["iteration_id"] = iteration_name + tid = getID() + await sor.C('pipeline_tasks', { + 'id': tid, 'tenant_id': project_id, 'pipeline_id': 'role_task', + 'owner_id': 'pm', 'title': title, 'state': 'submitted', + 'role': role, 'params': json.dumps(tparams, ensure_ascii=False), + }) + created.append(f"{title}({role})") + + if not created: + return 'FAIL: 没有可创建的任务(缺少 title)' + # C 之后立即 COMMIT,释放行锁,供 agent_poller 下一轮可见认领 + await sor.sqlExe("COMMIT", {}) + return f"OK: 已派发 {len(created)} 个任务:" + ";".join(created) + + +async def _pm_list_tasks(sor, project_id, params): + """PM 查看项目现有任务(派发前查重)。""" + role = (params.get('role') or '').strip() + state = (params.get('state') or '').strip() + sql = "SELECT id, title, role, state FROM pipeline_tasks WHERE tenant_id=${pid}$" + p = {"pid": project_id} + if role: + sql += " AND role=${role}$" + p["role"] = _normalize_role(role) + if state: + sql += " AND state=${state}$" + p["state"] = state + sql += " ORDER BY created_at ASC LIMIT 40" + recs = await sor.sqlExe(sql, p) + # 纯 SELECT 后 COMMIT,避免持有 MDL 锁阻塞后续 DDL + await sor.sqlExe("COMMIT", {}) + if not recs: + return "暂无任务" + icons = {"submitted": "⏳", "running": "🔄", "review": "👀", + "approved": "✅", "completed": "✔️", "failed": "❌", "waiting": "⏸️"} + lines = ["| 状态 | 任务 | 角色 |", "|------|------|------|"] + for r in recs: + st = getattr(r, 'state', '') + lines.append(f"| {icons.get(st, st)} | {getattr(r, 'title', '')} | {getattr(r, 'role', '')} |") + return "\n".join(lines) + + async def pm_review_run(project_id, agent_id=None, model_name=None): db = _get_db() async with db.sqlorContext("pipeline") as sor: @@ -997,12 +1094,14 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): {"tid": task_id}) await sor.sqlExe("COMMIT", {}) # auto-inject 兜底:最后两轮强制要求给出 review_* 决策,禁止再 tool_call, - # 否则 deepseek 会一直读文件/git_status 耗尽 5 轮 → 审核超时 → 打回重跑 → 死循环 + # 否则 deepseek 会一直读文件/git_status 耗尽 5 轮 → 审核超时 → 打回重跑 → 死循环。 + # 例外:允许 create_tasks(里程碑拆后续任务的「项目计划/任务分配」职责),派发后必须立即决策。 if turn >= 3: msgs.append({"role": "user", "content": - "已检查足够信息。现在必须立即输出最终决策 JSON," - "只能是 review_approve / review_reject / review_complete 三者之一," - "禁止再调用任何工具。"}) + "已检查足够信息。现在必须立即收尾:" + "若本任务审批通过后需拆分为多个后续任务,现在用 create_tasks 一次性派发;" + "随后必须立即输出 review_approve / review_reject / review_complete 三者之一," + "禁止再调用其它工具。"}) try: raw = await llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id) except Exception as e: @@ -1025,7 +1124,12 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): elif act.get('action') == 'tool_call': tool = act.get('tool', '') params = act.get('params', {}) - result = await _exec_agent_tool(tool, params, workspace_dir) + if tool in ('create_tasks', 'create_task'): + result = await _pm_create_tasks(sor, project_id, params) + elif tool == 'list_tasks': + result = await _pm_list_tasks(sor, project_id, params) + else: + result = await _exec_agent_tool(tool, params, workspace_dir) msgs.append({"role": "assistant", "content": raw}) msgs.append({"role": "user", "content": f"工具 {tool} 结果:\n{result}"}) else: