feat(sdlc): 会话agent加任务流转工具 approve_task/complete_task/cancel_task(PM权限)

- task_capability 新增 cancel_task 原语(CAS+审计+终态守护)
- sdlc_ability 加 3 个 ToolDefinition + handler,who=agent.pm 记录 PM 权限操作
- _resolve_task_id 租户隔离前缀兜底
- SDL_PROMPT 补任务流转典型场景
- _pm_cancel_task 委托新原语,去掉重复 SQL
This commit is contained in:
ymq 2026-08-21 22:43:11 +08:00
parent 6dbb574862
commit 78b84c710e
3 changed files with 131 additions and 7 deletions

View File

@ -1430,8 +1430,8 @@ async def _pm_list_tasks(sor, project_id, params):
async def _pm_cancel_task(sor, project_id, params):
"""PM 取消任务(重做/作废场景必须先取消旧任务,避免两个相同任务并存)。
标记 state=cancelled + claimed_by正在执行中的 agent 会在下一轮心跳检测到
cancelled 并自行中止 role_agent_run 的取消检测
委托 task_capability.cancel_taskCAS + 审计正在执行中的 agent 会在下一轮心跳
检测到 cancelled 并自行中止 role_agent_run 的取消检测
"""
task_id = (params.get('task_id') or params.get('id') or '').strip()
if not task_id:
@ -1446,10 +1446,10 @@ async def _pm_cancel_task(sor, project_id, params):
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", {})
from .task_capability import cancel_task
ok, msg = await cancel_task(task_id, project_id, who="agent.pm")
if not ok:
return f'FAIL: {msg}'
return f'OK: 已取消任务「{title}」({task_id}),正在执行的 agent 将在下一轮心跳中止'

View File

@ -56,6 +56,24 @@ SDL_TOOLS = [
parameters={"task_id": "任务ID"},
category="task",
),
ToolDefinition(
name="approve_task",
description="审核通过任务review→approvedPM 验收)",
parameters={"task_id": "任务ID", "comment": "审核意见(可选)"},
category="task",
),
ToolDefinition(
name="complete_task",
description="标记任务完成approved/review→completed",
parameters={"task_id": "任务ID"},
category="task",
),
ToolDefinition(
name="cancel_task",
description="取消任务(重做/作废前必须先取消旧任务,避免两个相同任务并存)",
parameters={"task_id": "任务ID", "comment": "取消原因(可选)"},
category="task",
),
ToolDefinition(
name="start_agents",
description="启动角色agent执行已提交的任务",
@ -238,7 +256,8 @@ SDL_PROMPT = """你是「开发产线」的驾驶舱 agent负责软件项目
- 测试计划 create_test_plan 建计划list_plans 列计划审批/执行/完成用 approve_plan / start_plan / complete_plan规范见 test-plan 技能
- 测试用例 create_case 建用例list_cases 列用例执行结果用 pass_case / fail_case / skip_case / block_case失败用例应联动 report_bug规范见 test-case 技能
- Bug 管理 report_bug 上报list_bugs Bug确认/修复/验证/关闭/驳回/重开用 confirm_bug / start_fix / fix_bug / verify_bug / close_bug / reject_bug / reopen_bug规范见 bug 技能
- 部署环境 configure_env 配环境list_envs 列环境验证用 verify_env规范见 deploy-env 技能"""
- 部署环境 configure_env 配环境list_envs 列环境验证用 verify_env规范见 deploy-env 技能
- 任务流转PM 权限操作 审核通过任务 approve_taskreviewapproved标记完成 complete_taskapproved/reviewcompleted取消任务 cancel_task重做/作废前必须先取消旧任务这三个是 PM 验收级操作仅在用户明确要求审核/验收/完成/取消任务时使用不主动越权推进"""
# ── 产线角色集(可插拔,替代硬编码 ROLE_SPECIFICS/ROLE_ALIASES/ROLE_CHAIN ──
@ -530,6 +549,59 @@ async def _h_reset_task_retry(sor, p, ctx):
return f"OK: 任务重复计数已归零并重新执行{resume_note}"
async def _h_approve_task(sor, p, ctx):
"""审核通过任务PM 验收review → approved。"""
task_id = (p.get("task_id") or "").strip()
if not task_id:
return "需要任务ID"
pid = ctx.get("project_id", "")
if not pid:
return "请先切换到项目"
full = await _resolve_task_id(sor, task_id, pid)
if not full:
return f"任务不存在: {task_id}"
from .task_capability import approve_task
ok, msg = await approve_task(full, pid, who="agent.pm",
agent_id=ctx.get("user_id", ""),
comment=(p.get("comment") or "").strip())
return f"OK: {msg}" if ok else f"ERROR: {msg}"
async def _h_complete_task(sor, p, ctx):
"""标记任务完成PMapproved/review → completed。"""
task_id = (p.get("task_id") or "").strip()
if not task_id:
return "需要任务ID"
pid = ctx.get("project_id", "")
if not pid:
return "请先切换到项目"
full = await _resolve_task_id(sor, task_id, pid)
if not full:
return f"任务不存在: {task_id}"
from .task_capability import complete_task
ok, msg = await complete_task(full, pid, who="agent.pm",
agent_id=ctx.get("user_id", ""))
return f"OK: {msg}" if ok else f"ERROR: {msg}"
async def _h_cancel_task(sor, p, ctx):
"""取消任务PM任意非终态 → cancelled清 claimed_by。"""
task_id = (p.get("task_id") or "").strip()
if not task_id:
return "需要任务ID"
pid = ctx.get("project_id", "")
if not pid:
return "请先切换到项目"
full = await _resolve_task_id(sor, task_id, pid)
if not full:
return f"任务不存在: {task_id}"
from .task_capability import cancel_task
ok, msg = await cancel_task(full, pid, who="agent.pm",
agent_id=ctx.get("user_id", ""),
comment=(p.get("comment") or "").strip())
return f"OK: {msg}" if ok else f"ERROR: {msg}"
async def _h_start_agents(sor, p, ctx):
pid = ctx.get("project_id", "")
if not pid:
@ -952,6 +1024,24 @@ async def _resolve_id(sor, table, fid):
return ""
async def _resolve_task_id(sor, task_id, pid):
"""任务 ID 前缀兜底租户隔离精确匹配失败且入参≥6位时 LIKE 前缀解析完整 ID。"""
if not task_id:
return ""
recs = await sor.sqlExe(
"SELECT id FROM pipeline_tasks WHERE id=${tid}$ AND tenant_id=${pid}$",
{"tid": task_id, "pid": pid})
if recs:
return getattr(recs[0], "id", "")
if len(task_id) >= 6:
recs = await sor.sqlExe(
"SELECT id FROM pipeline_tasks WHERE id LIKE ${prefix}$ AND tenant_id=${pid}$ LIMIT 1",
{"prefix": task_id + "%", "pid": pid})
if recs:
return getattr(recs[0], "id", "")
return ""
async def _get_scope(sor, table, id_field, id_val, scope_field):
"""反查实体的 scope 列值(如 iteration_id/plan_id。返回 scope_val 或空。"""
recs = await sor.sqlExe(
@ -1593,6 +1683,9 @@ SDL_HANDLERS = {
"list_tasks": _h_list_tasks,
"task_detail": _h_task_detail,
"reset_task_retry": _h_reset_task_retry,
"approve_task": _h_approve_task,
"complete_task": _h_complete_task,
"cancel_task": _h_cancel_task,
"start_agents": _h_start_agents,
"diagnose_project": _h_diagnose_project,
"list_deliverables": _h_list_deliverables,

View File

@ -204,6 +204,37 @@ async def mark_failed(task_id, tenant_id, who=None, agent_id=None, error=None):
return True, S_FAILED
async def cancel_task(task_id, tenant_id, who=None, agent_id=None, comment=None):
"""取消任务(任意非终态 → cancelled清 claimed_by
重做/作废场景必须先取消旧任务避免两个相同任务并存正在执行的 agent 会在
下一轮心跳检测到 cancelled 后自行中止终态(completed/cancelled/failed/approved)
不重复取消CAS仅当任务仍处原非终态时取消防并发覆盖
"""
if not task_id or not tenant_id:
return False, "缺少 task_id 或 tenant_id"
db, dbname = _get_db()
async with db.sqlorContext(dbname) as sor:
recs = await sor.R('pipeline_tasks', {'id': task_id})
if not recs:
return False, "任务不存在"
from_state = getattr(recs[0], 'state', '')
if from_state in (S_COMPLETED, 'cancelled', S_FAILED, S_APPROVED):
return False, f"任务已是终态 {from_state},无需取消"
await sor.sqlExe(
"UPDATE pipeline_tasks SET state='cancelled', claimed_by=NULL, "
"updated_at=NOW() WHERE id=${tid}$ AND tenant_id=${tn}$ AND state=${from}$",
{"tid": task_id, "tn": tenant_id, "from": from_state})
chk = await sor.R('pipeline_tasks', {'id': task_id})
await sor.sqlExe("COMMIT", {})
if not chk or getattr(chk[0], 'state', '') != 'cancelled':
return False, "取消失败(CAS): 任务状态已变化"
await record_audit(tenant_id, 'pipeline_tasks', task_id, 'cancel',
from_state=from_state, to_state='cancelled',
who=who, agent_id=agent_id, detail=comment, sor=sor)
return True, 'cancelled'
async def retry_task(task_id, tenant_id, who=None, agent_id=None):
"""重试failed → submittedretry_count+1清 claimed_by"""
db, dbname = _get_db()