feat(pm): 固化 PM 项目计划/任务分配/任务验收职责,新增 create_tasks 批量派发后续任务,修复里程碑审批后链断
This commit is contained in:
parent
7ec0fb611f
commit
dfa441b09b
@ -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:
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user