diff --git a/models/pipeline_tasks.json b/models/pipeline_tasks.json index 6da9e0d..ae625fe 100644 --- a/models/pipeline_tasks.json +++ b/models/pipeline_tasks.json @@ -18,6 +18,8 @@ {"name": "params", "title": "提交参数", "type": "text"}, {"name": "role", "title": "目标角色", "type": "str", "length": 32, "nullable": "no", "default": ""}, {"name": "claimed_by", "title": "认领agent标识", "type": "str", "length": 64, "nullable": "yes"}, + {"name": "parent_id", "title": "父任务ID(分解后子任务指向父)", "type": "str", "length": 32, "nullable": "yes"}, + {"name": "depends_on", "title": "依赖任务ID列表(JSON数组,依赖完成才可认领)", "type": "text"}, {"name": "created_at", "title": "创建时间", "type": "timestamp", "nullable": "no"}, {"name": "updated_at", "title": "更新时间", "type": "timestamp", "nullable": "no"} ], diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index 3f62d77..a04151b 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -481,10 +481,18 @@ __REPO_STATE__ PM_SYSTEM_PROMPT = """你是项目经理(PM)。你的职责是项目计划、任务分配、任务验收。 ## 三大职责 -1. 项目计划:基于需求/设计交付件,规划后续工作——把里程碑任务拆解为可执行的子任务(每个含标题、角色、描述)。 -2. 任务分配:用 create_tasks 工具把子任务派发给对应角色 agent(如「总体设计」批准后 → 拆「契约定义」任务 + hr-system/hr-org/hr-roster 等各模块 develop 任务)。 +1. 项目计划:基于需求/设计交付件,规划后续工作——先评估任务复杂度,复杂任务自动分解为可执行的子任务。 +2. 任务分配:用 create_tasks 工具把子任务派发给对应角色 agent,并记录父子关系与依赖关系(能并行的并行、不能并行的串行)。 3. 任务验收:检查交付件 + git 实际产出,决定 review_approve / review_reject / review_complete。 +## 任务分解与编排(先评估,再拆解) +- 派发前先评估:当前任务是否「过于复杂」(涉及多个模块/应用、多个独立交付单元、工作量超单 agent 一次产出)。简单任务直接派发单个任务,不必强行拆分。 +- 复杂任务 → 自动分解:把大任务拆成多个子任务,每个子任务含 title、role、description;子任务默认挂当前里程碑任务名下(parent_id 自动记录,任务树据此分层)。 +- 自动编排(能并行并行、不能并行串行): + - 无依赖的子任务 → 并行(不填 depends_on,系统同时认领执行)。 + - 有先后依赖的子任务 → 用 key 标记 + depends_on 引用前序 key(如「契约定义」key=contract 完成后,「hr-system 开发」depends_on=[contract])。依赖完成前子任务不会被认领,自动串行。 +- 父子关系:系统自动把子任务的 parent_id 记为当前审核的里程碑任务,任务树(/task)里显示为父任务下的层级子任务,无需你手动记 parent_id。 + ## 任务链不能断(关键规则) - 审核通过里程碑任务(requirement / design,产出含多个应用/模块,或评审意见指明后续要落地开发)后,必须先规划并用 create_tasks 创建后续开发任务,再 review_approve。 - 系统在 approve 后会自动补一个线性链的下一角色任务作为兜底;里程碑的并发拆分由你显式创建,不要等、不要断链。 @@ -505,7 +513,7 @@ __REPO_STATE__ - list_files(path) — 列目录 - git_status() — 查看git状态 - run_shell(command) — 执行命令(编译/测试验证) -- create_tasks(tasks) — 批量创建并派发后续任务(tasks 是 JSON 数组,每项 {title, role, description}) +- create_tasks(tasks) — 批量创建并派发后续任务(tasks 是 JSON 数组,每项 {title, role, description, key?, depends_on?, parent_id?};key 供同批任务间 depends_on 引用,depends_on 是前序 key 或任务ID 数组,空=并行/非空=串行) - list_tasks(role, state) — 列出项目现有任务(派发前先查,避免重复) - cancel_task(task_id) — 取消任务(重做/作废前必须先取消旧任务,避免两个相同任务并存) @@ -520,7 +528,7 @@ __REPO_STATE__ 查看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":"tool_call","tool":"create_tasks","params":{"tasks":[{"title":"契约定义","role":"design","key":"contract","description":"定义数据模型/API/目录规范"},{"title":"hr-system 开发","role":"develop","key":"hrsys","depends_on":["contract"],"description":"实现 hr-system 基础能力"},{"title":"hr-org 开发","role":"develop","key":"hrorg","depends_on":["contract"],"description":"实现 hr-org 基础能力"}]}} 批准:{"action":"review_approve","comment":"审核意见","next_task_title":"下阶段标题","next_task_description":"描述"} 驳回:{"action":"review_reject","comment":"原因","questions":"修改要求"} 完成:{"action":"review_complete","comment":"总结"} @@ -571,18 +579,52 @@ QC_SYSTEM_PROMPT = """你是质量控制工程师(QC)。对交付件做合 # ── Agent 核心逻辑 ── +async def _task_deps_satisfied(sor, depends_on_raw): + """依赖门控:depends_on 是 JSON 数组(依赖任务ID列表),全部处于终态(completed/approved) + 才可认领(返回 True)。空/无依赖 → 并行,立即满足。任一依赖未完成 → 串行等待(False)。 + + PM 分解任务时用 depends_on 表达「不能并行的串行依赖」;agent 认领前必须先过这道门, + 否则串行依赖会被 poller 提前认领、破坏编排顺序。 + """ + deps = [] + if depends_on_raw: + try: + d = json.loads(depends_on_raw) if isinstance(depends_on_raw, str) else depends_on_raw + if isinstance(d, list): + deps = [str(x) for x in d if x] + except (json.JSONDecodeError, TypeError, ValueError): + deps = [] + if not deps: + return True + # 任一依赖不在终态 → 未满足。终态 = completed/approved(approved 视为已验收、可驱动下游)。 + ids = ",".join(["'" + x.replace("'", "''") + "'" for x in deps]) + recs = await sor.sqlExe( + "SELECT id FROM pipeline_tasks WHERE id IN (" + ids + ") " + "AND state NOT IN ('completed','approved')", {}) + await sor.sqlExe("COMMIT", {}) + return not recs + + async def _claim_task(sor, tenant_id, role, state='submitted', match_role=True, set_state='running'): role = _normalize_role(role) where_role = "AND (role=${role}$ OR role='')" if match_role else "" recs = await sor.sqlExe( - "SELECT id, title, params, pipeline_id, tenant_id, role FROM pipeline_tasks " + "SELECT id, title, params, pipeline_id, tenant_id, role, depends_on FROM pipeline_tasks " "WHERE tenant_id=${tid}$ AND state=${state}$ " + where_role + " " "AND NOT EXISTS (SELECT 1 FROM pipeline_task_steps s WHERE s.task_id=pipeline_tasks.id) " - "ORDER BY created_at ASC LIMIT 1", + "ORDER BY created_at ASC LIMIT 50", {"tid": tenant_id, "state": state, "role": role}) if not recs: return None - task = recs[0] + # 依赖门控:跳过依赖未完成的任务,认领第一个依赖已满足的任务。 + # (并行任务不被前面的串行任务阻塞;被依赖阻塞的任务保持 submitted 等依赖完成后下一轮认领。) + task = None + for rec in recs: + if await _task_deps_satisfied(sor, getattr(rec, 'depends_on', '') or ''): + task = rec + break + if task is None: + return None task_id = task.id from appPublic.uniqueID import getID claim_token = getID() @@ -1153,10 +1195,14 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None): # ── PM 审核 ── -async def _pm_create_tasks(sor, project_id, params): +async def _pm_create_tasks(sor, project_id, params, parent_task_id=None): """PM 派发后续任务(项目计划 / 任务分配):批量创建任务并指定角色。 - params.tasks: JSON 数组 [{title, role, description}],一次派发多个并发子任务。 + params.tasks: JSON 数组 [{title, role, description, key, depends_on, parent_id}], + 一次派发多个子任务,支持: + - parent_id:父任务ID(默认 = 本次审核的里程碑任务),记录父子关系供任务树分层; + - key:LLM 给子任务起的短标识,供同批任务间 depends_on 相互引用; + - depends_on:依赖的 key 或任务ID 数组(空 = 并行;非空 = 串行等待依赖完成后才可认领)。 兼容单任务形态:params 直接含 {title, role, description}。 """ from appPublic.uniqueID import getID @@ -1174,6 +1220,9 @@ async def _pm_create_tasks(sor, project_id, params): if not isinstance(tasks, list): return 'FAIL: tasks 必须是数组' + # 默认父任务:本次审核的里程碑任务(子任务挂它名下,任务树据此分层) + default_parent = (params.get('parent_id') or parent_task_id or '').strip() + # 当前迭代(该 project 最新创建的迭代),作为任务默认归属 iteration_name = '' irecs = await sor.sqlExe( @@ -1182,8 +1231,10 @@ async def _pm_create_tasks(sor, project_id, params): if irecs: iteration_name = getattr(irecs[0], 'iteration_name', '') or '' - created = [] + created = [] # [(task_dict, tid, title, role)] + key_map = {} # key → 真实任务ID(同批 depends_on 用 key 互引,第二遍解析) cancelled = [] + warnings = [] for t in tasks: if not isinstance(t, dict): continue @@ -1192,6 +1243,8 @@ 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() + key = (t.get('key') or '').strip() + parent = (t.get('parent_id') or '').strip() or default_parent # 重做保护:同项目存在同标题活跃任务(submitted/running/review/qc_review)时, # 先自动取消旧任务再创建新任务,避免两个相同任务并存(重做必须先终止旧任务)。 @@ -1214,15 +1267,46 @@ async def _pm_create_tasks(sor, project_id, params): '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), + 'parent_id': parent or None, }) - created.append(f"{title}({role})") + if key: + key_map[key] = tid + created.append((t, tid, title, role)) + + # 第二遍:解析 depends_on(key 或 任务ID)→ 真实任务ID,回填 depends_on 列。 + # 串行依赖 = depends_on 非空(依赖完成才认领);并行 = 空(不填)。 + # 父任务本身在 PM approve 后即为 approved 终态,子任务无需再显式依赖父。 + for t, tid, title, role in created: + deps = t.get('depends_on') or [] + if isinstance(deps, str): + try: + deps = json.loads(deps) + except (json.JSONDecodeError, ValueError): + deps = [] + resolved = [] + for d in (deps or []): + d = (str(d) or '').strip() + if not d: + continue + if d in key_map: + resolved.append(key_map[d]) # 同批 key 引用 + elif len(d) >= 20: + resolved.append(d) # 已有任务ID(跨批/父任务引用) + else: + warnings.append(f"「{title}」的 depends_on 未知引用「{d}」已忽略") + if resolved: + await sor.sqlExe( + "UPDATE pipeline_tasks SET depends_on=${deps}$ WHERE id=${tid}$", + {"deps": json.dumps(resolved, ensure_ascii=False), "tid": tid}) if not created: return 'FAIL: 没有可创建的任务(缺少 title)' # C 之后立即 COMMIT,释放行锁,供 agent_poller 下一轮可见认领 await sor.sqlExe("COMMIT", {}) note = f"(已自动取消 {len(cancelled)} 个同标题活跃任务)" if cancelled else "" - return f"OK: 已派发 {len(created)} 个任务" + note + ":" + ";".join(created) + warn_note = f"({len(warnings)} 个未知依赖引用已忽略)" if warnings else "" + return f"OK: 已派发 {len(created)} 个任务" + note + warn_note + ":" + \ + ";".join([f"{title}({role})" for _, _, title, role in created]) async def _pm_list_tasks(sor, project_id, params): @@ -1370,7 +1454,7 @@ async def pm_review_run(project_id, agent_id=None, model_name=None): tool = act.get('tool', '') params = act.get('params', {}) if tool in ('create_tasks', 'create_task'): - result = await _pm_create_tasks(sor, project_id, params) + result = await _pm_create_tasks(sor, project_id, params, task_id) elif tool == 'list_tasks': result = await _pm_list_tasks(sor, project_id, params) elif tool in ('cancel_task', 'cancel'):