feat: PM任务分解与编排——父子关系(parent_id)+依赖串行(depends_on)+任务树分层

This commit is contained in:
ymq 2026-08-19 05:53:28 +08:00
parent f542fdad8a
commit 834505fcfc
2 changed files with 99 additions and 13 deletions

View File

@ -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"}
],

View File

@ -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 一次产出简单任务直接派发单个任务不必强行拆分
- 复杂任务 自动分解把大任务拆成多个子任务每个子任务含 titleroledescription子任务默认挂当前里程碑任务名下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/approvedapproved 视为已验收、可驱动下游)。
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默认 = 本次审核的里程碑任务记录父子关系供任务树分层
- keyLLM 给子任务起的短标识供同批任务间 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_onkey 或 任务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'):