From 0a90f6031532ab31ba3b0ab887871776c217a870 Mon Sep 17 00:00:00 2001 From: ymq Date: Sat, 15 Aug 2026 09:26:01 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20SDLC=E8=83=BD=E5=8A=9B=E8=BF=81?= =?UTF-8?q?=E7=A7=BB=E5=88=B0sdlc=5Fability=E5=8F=AF=E6=8F=92=E6=8B=94?= =?UTF-8?q?=E8=83=BD=E5=8A=9B=E5=8C=85=E2=80=94=E2=80=9414=E5=B7=A5?= =?UTF-8?q?=E5=85=B7+prompt+handler=E7=8B=AC=E7=AB=8B=E5=87=BD=E6=95=B0?= =?UTF-8?q?=EF=BC=8CAgentExecutor=E6=8C=89pipeline=5Fid=E5=8A=A8=E6=80=81?= =?UTF-8?q?=E6=8C=82=E8=BD=BD=EF=BC=8C=E5=88=A0=E9=99=A4=E7=A1=AC=E7=BC=96?= =?UTF-8?q?=E7=A0=81SDLC=20handler?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pipeline_service/__init__.py | 3 + pipeline_service/agent_loop_v2.py | 470 +++----------------------- pipeline_service/sdlc_ability.py | 535 ++++++++++++++++++++++++++++++ 3 files changed, 580 insertions(+), 428 deletions(-) create mode 100644 pipeline_service/sdlc_ability.py diff --git a/pipeline_service/__init__.py b/pipeline_service/__init__.py index 82c7484..63203fb 100644 --- a/pipeline_service/__init__.py +++ b/pipeline_service/__init__.py @@ -28,6 +28,9 @@ from .state import ( # v2: Agent 执行引擎(对照 Hermes Agent) from .agent_loop_v2 import AgentExecutor, run_agent +# 产线能力包注册(import 即注册到 pipeline_core 的 PipelineAbility 注册表) +from . import sdlc_ability # noqa: F401 + # 应用部署主机逻辑账号系统(零 root 逻辑隔离 + bwrap 沙箱) from .deploy_account import ( ensure_account, diff --git a/pipeline_service/agent_loop_v2.py b/pipeline_service/agent_loop_v2.py index 8dbe92d..7101045 100644 --- a/pipeline_service/agent_loop_v2.py +++ b/pipeline_service/agent_loop_v2.py @@ -64,6 +64,7 @@ class AgentExecutor: self.project_id = project_id self.user_id = user_id self.org_id = "" # loaded from project context + self.pipeline_id = "" # loaded from project context(产线能力包 key) self.workspace_dir = workspace_dir or "/tmp/pipeline_ws" self.model_name = model_name or config.model_name @@ -292,20 +293,21 @@ class AgentExecutor: except ImportError: self._memory_store = None - # 加载 org_id + workspace_dir(从项目上下文) + # 加载 org_id + workspace_dir + pipeline_id(从项目上下文) if self.project_id: try: from sqlor.dbpools import DBPools db = DBPools() async with db.sqlorContext("pipeline") as sor: recs = await sor.sqlExe( - "SELECT org_id, workspace_dir FROM sd_projects WHERE id=${pid}$", + "SELECT org_id, workspace_dir, pipeline_id FROM sd_projects WHERE id=${pid}$", {"pid": self.project_id}) if recs: self.org_id = getattr(recs[0], 'org_id', '') or '' ws = getattr(recs[0], 'workspace_dir', '') or '' if ws: self.workspace_dir = ws + self.pipeline_id = getattr(recs[0], 'pipeline_id', '') or '' except Exception: pass @@ -507,19 +509,51 @@ class AgentExecutor: logger.error(f"Tool {tool_name} handler error: {e}") return f"ERROR: {str(e)[:300]}" - # 2. SDLC 内建工具 + # 2. 产线能力包工具(从 PipelineAbility 注册表按 pipeline_id 取 handler) + result = await self._execute_ability_tool(tool_name, params) + if result is not None: + return result + + # 3. 通用内建工具(core 内核:项目管理/终端/文件/搜索/会话/规划) result = await self._execute_sdlc_tool(tool_name, params) if result is not None: return result - # 3. ask_user + # 4. ask_user if tool_name == "ask_user": return f"QUESTION: {params.get('question', '')}" return f"未知工具: {tool_name}" + def _build_ctx(self) -> dict: + """构造传给产线能力 handler 的上下文。""" + return { + "project_id": self.project_id, + "user_id": self.user_id, + "workspace_dir": self.workspace_dir, + "model_name": self.model_name, + "config": self.config, + } + + async def _execute_ability_tool(self, tool_name: str, params: dict) -> Optional[str]: + """产线能力包工具:按 pipeline_id 从 PipelineAbility 注册表取 handler 执行。""" + try: + from pipeline_core import get_ability, DEFAULT_ABILITY_ID + + ability = get_ability(self.pipeline_id) or get_ability(DEFAULT_ABILITY_ID) + if not ability or tool_name not in ability.handlers: + return None + + from sqlor.dbpools import DBPools + db = DBPools() + async with db.sqlorContext("pipeline") as sor: + return await ability.handlers[tool_name](sor, params, self._build_ctx()) + except Exception as e: + logger.error(f"ability tool {tool_name} error: {e}") + return f"ERROR: {str(e)[:300]}" + async def _execute_sdlc_tool(self, tool_name: str, params: dict) -> Optional[str]: - """SDLC 内建工具(兼容 cockpit_chat.dspy)""" + """通用内建工具(core 内核,产线无关)。""" # 导入现有 handler try: from sqlor.dbpools import DBPools @@ -536,30 +570,13 @@ class AgentExecutor: p = params or {} pid = self.project_id - # 所有工具映射(含别名) + # 通用内建工具映射(core 内核,产线无关)。 + # 产线专属工具(create_task/diagnose_project 等)已迁移到 sdlc_ability 能力包, + # 由 _execute_ability_tool 按 pipeline_id 动态挂载。 handlers = { "switch_project": self._t_switch_project, "create_project": self._t_create_project, - "create_task": self._t_create_task, - "list_tasks": self._t_list_tasks, - "get_task_list": self._t_list_tasks, - "get_tasks": self._t_list_tasks, # 别名 - "task_detail": self._t_task_detail, - "get_task": self._t_task_detail, - "get_task_detail": self._t_task_detail, # 别名 - "start_agents": self._t_start_agents, - "diagnose_project": self._t_diagnose_project, - "list_deliverables": self._t_list_deliverables, - "view_deliverable": self._t_view_deliverable, - "get_deliverable": self._t_view_deliverable, # 别名 - "list_questions": self._t_list_questions, - "answer_question": self._t_answer_question, - "add_repo": self._t_add_repo, - "list_repos": self._t_list_repos, - "clone_repo": self._t_clone_repo, "run_command": self._t_run_command, - "check_progress": self._t_check_progress, - "add_bug": self._t_add_bug, # ── 通用工具集(Hermes CLI 能力子集)── "read_file": self._t_read_file, "write_file": self._t_write_file, @@ -678,357 +695,6 @@ class AgentExecutor: await self._persist_project(sor, pid_val) return f"OK: 已创建项目 {name}" - async def _t_create_task(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - from appPublic.uniqueID import getID - - title = p.get("title", "").strip() - role = p.get("role", "develop") - desc = p.get("description", "") - if not title: - return "需要任务标题" - - tid = getID() - await sor.C("pipeline_tasks", { - "id": tid, "tenant_id": pid, "pipeline_id": "role_task", - "owner_id": self.user_id or "user", "title": title, - "state": "submitted", "role": role, - "params": json.dumps({"description": desc}, ensure_ascii=False), - }) - return f"OK: 已创建任务 {title}(角色: {role})" - - async def _t_list_tasks(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - role = p.get("role", "") - state = p.get("state", "") - - sql = "SELECT id, title, role, state, created_at FROM pipeline_tasks WHERE tenant_id=${pid}$" - params_dict = {"pid": pid} - if role: - sql += " AND role=${role}$" - params_dict["role"] = role - if state: - sql += " AND state=${state}$" - params_dict["state"] = state - sql += " ORDER BY created_at DESC LIMIT 30" - - recs = await sor.sqlExe(sql, params_dict) - if not recs: - return "暂无任务" - - # v2.1 format: markdown table - icons = { - "submitted": "⏳", "running": "🔄", "review": "👀", - "approved": "✅", "completed": "✔️", "failed": "❌", - "waiting": "⏸️", - } - lines = ["| 状态 | 任务 | 角色 | ID |", - "|------|------|------|----|"] - for r in recs: - st = getattr(r, 'state', '?') - icon = icons.get(st, "❓") - title = (getattr(r, 'title', '') or '')[:60] - role = getattr(r, 'role', '') - tid = getattr(r, 'id', '') - lines.append(f"| {icon} {st} | {title} | {role} | {tid} |") - return "\n".join(lines) - - async def _t_task_detail(self, sor, p, pid): - tid = p.get("task_id", "") - if not tid: - return "需要任务ID" - - recs = await sor.sqlExe( - "SELECT * FROM pipeline_tasks WHERE id=${tid}$", {"tid": tid}) - if not recs and len(tid) >= 6: - # 前缀匹配兜底(兼容 8 位截断 ID) - recs = await sor.sqlExe( - "SELECT * FROM pipeline_tasks WHERE id LIKE ${prefix}$ LIMIT 1", - {"prefix": tid + "%"}) - if not recs: - return f"任务不存在: {tid}" - - r = recs[0] - return json.dumps({ - "id": getattr(r, "id", ""), - "title": getattr(r, "title", ""), - "state": getattr(r, "state", ""), - "role": getattr(r, "role", ""), - "params": getattr(r, "params", "{}"), - }, ensure_ascii=False, indent=2) - - async def _t_start_agents(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - # 找 submitted 任务 - recs = await sor.sqlExe( - "SELECT id, title, role FROM pipeline_tasks " - "WHERE tenant_id=${pid}$ AND state='submitted' ORDER BY created_at ASC LIMIT 10", - {"pid": pid}) - if not recs: - return "没有待执行的任务" - - results = [] - for r in recs: - tid = getattr(r, "id", "") - role = getattr(r, "role", "develop") - try: - from pipeline_service.agent_loop import role_agent_run - - result = await role_agent_run(pid, role) - results.append(f"- {getattr(r, 'title', '')}: {result.get('status', '?')}") - except Exception as e: - results.append(f"- {getattr(r, 'title', '')}: error={str(e)[:100]}") - - return "Agent 执行结果:\n" + "\n".join(results) - - async def _t_diagnose_project(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - # 汇总各种状态 - submitted = await sor.sqlExe( - "SELECT COUNT(*) as c FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='submitted'", - {"pid": pid}) - running = await sor.sqlExe( - "SELECT COUNT(*) as c FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='running'", - {"pid": pid}) - failed = await sor.sqlExe( - "SELECT COUNT(*) as c FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='failed'", - {"pid": pid}) - review = await sor.sqlExe( - "SELECT COUNT(*) as c FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='review'", - {"pid": pid}) - # questions 表可能没有 answered 列,用 status='pending' 兜底 - qc = 0 - try: - questions = await sor.sqlExe( - "SELECT COUNT(*) as c FROM pipeline_agent_questions WHERE tenant_id=${pid}$ AND status='pending'", - {"pid": pid}) - qc = getattr(questions[0], "c", 0) if questions else 0 - except Exception: - pass - - sc = getattr(submitted[0], "c", 0) if submitted else 0 - rc = getattr(running[0], "c", 0) if running else 0 - fc = getattr(failed[0], "c", 0) if failed else 0 - vc = getattr(review[0], "c", 0) if review else 0 - - lines = [ - "项目诊断:", - f"- 待执行: {sc}", - f"- 运行中: {rc}", - f"- 待审核: {vc}", - f"- 失败: {fc}", - f"- 待回答问题: {qc}", - ] - - # 列出失败/待执行任务完整 ID,方便 agent 直接 task_detail - if fc: - failed_rows = await sor.sqlExe( - "SELECT id, title, role FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='failed' ORDER BY created_at DESC LIMIT 10", - {"pid": pid}) - lines.append("失败任务:") - for fr in (failed_rows or []): - lines.append(f" - [{getattr(fr, 'role', '?')}] {getattr(fr, 'title', '')} (id={getattr(fr, 'id', '')})") - if sc: - submitted_rows = await sor.sqlExe( - "SELECT id, title, role FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='submitted' ORDER BY created_at ASC LIMIT 10", - {"pid": pid}) - lines.append("待执行任务:") - for sr in (submitted_rows or []): - lines.append(f" - [{getattr(sr, 'role', '?')}] {getattr(sr, 'title', '')} (id={getattr(sr, 'id', '')})") - - # 运行中任务明细:claimed_by + 心跳陈旧度,帮助 agent 直接识别僵尸任务 - # (进程崩溃/协程挂起后心跳不再更新,超时即疑似僵尸,会被 poller 自动回收重跑) - if rc: - running_rows = await sor.sqlExe( - "SELECT id, title, role, claimed_by, " - "TIMESTAMPDIFF(MINUTE, updated_at, NOW()) AS mins " - "FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='running' " - "ORDER BY updated_at ASC LIMIT 10", - {"pid": pid}) - lines.append("运行中任务:") - for rr in (running_rows or []): - rrole = getattr(rr, 'role', '?') - rtitle = getattr(rr, 'title', '') - rid = getattr(rr, 'id', '') - rcb = getattr(rr, 'claimed_by', '') or '' - try: - mins = int(getattr(rr, 'mins', 0) or 0) - except (TypeError, ValueError): - mins = 0 - if mins >= 20: - flag = f"⚠️心跳超时{mins}分钟(疑似僵尸,将被回收重跑)" - else: - flag = f"已运行{mins}分钟" - lines.append(f" - [{rrole}] {rtitle} (id={rid}) | {flag} | claimed_by={rcb[:8]}") - - return "\n".join(lines) - - async def _t_list_deliverables(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - role = p.get("role", "") - sql = "SELECT id, title, deliverable_type, review_status FROM pipeline_deliverables WHERE project_id=${pid}$" - params_dict = {"pid": pid} - sql += " ORDER BY created_at DESC LIMIT 20" - - recs = await sor.sqlExe(sql, params_dict) - if not recs: - return "暂无交付件" - - lines = [] - for r in recs: - lines.append( - f"- [{getattr(r, 'review_status', '?')}] {getattr(r, 'title', '')} " - f"({getattr(r, 'deliverable_type', '')}) id={getattr(r, 'id', '')}" - ) - return "\n".join(lines) - - async def _t_view_deliverable(self, sor, p, pid): - did = p.get("deliverable_id", "") - if not did: - return "需要交付件ID" - - recs = await sor.sqlExe( - "SELECT * FROM pipeline_deliverables WHERE id=${did}$", {"did": did}) - if not recs and len(did) >= 6: - recs = await sor.sqlExe( - "SELECT * FROM pipeline_deliverables WHERE id LIKE ${prefix}$ LIMIT 1", - {"prefix": did + "%"}) - if not recs: - return f"交付件不存在: {did}" - - r = recs[0] - return getattr(r, "content", "")[:4000] or "(空)" - - async def _t_list_questions(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - try: - recs = await sor.sqlExe( - "SELECT id, question, answer_source FROM pipeline_agent_questions " - "WHERE tenant_id=${pid}$ AND status='pending' ORDER BY created_at DESC LIMIT 10", - {"pid": pid}) - except Exception: - # 兼容无 status 列的情况 - try: - recs = await sor.sqlExe( - "SELECT id, question, answer_source FROM pipeline_agent_questions " - "WHERE tenant_id=${pid}$ ORDER BY created_at DESC LIMIT 10", - {"pid": pid}) - except Exception: - return "暂无问题数据" - if not recs: - return "没有待回答问题" - - lines = [] - for r in recs: - lines.append(f"- [{getattr(r, 'id', '')}] {getattr(r, 'question', '')[:200]}") - return "\n".join(lines) - - async def _t_answer_question(self, sor, p, pid): - qid = p.get("question_id", "") - answer = p.get("answer", "") - if not qid or not answer: - return "需要问题ID和回答内容" - - # 先精确匹配,再前缀匹配兜底(兼容截断 ID) - recs = await sor.sqlExe( - "SELECT id, task_id FROM pipeline_agent_questions WHERE id=${qid}$", {"qid": qid}) - if not recs and len(qid) >= 6: - recs = await sor.sqlExe( - "SELECT id, task_id FROM pipeline_agent_questions WHERE id LIKE ${prefix}$ LIMIT 1", - {"prefix": qid + "%"}) - if not recs: - return f"问题不存在: {qid}" - full_qid = getattr(recs[0], "id", qid) - task_id = getattr(recs[0], "task_id", "") or "" - - await sor.sqlExe( - "UPDATE pipeline_agent_questions SET answer=${a}$, answer_source='main_agent', " - "answered_by='main_agent', status='answered' WHERE id=${qid}$", - {"a": answer, "qid": full_qid}) - - # 恢复任务:waiting → submitted,并清 claimed_by。 - # 否则任务永久卡在 waiting(poll 器要求 state='submitted' AND claimed_by IS NULL 才会重新认领)。 - resumed = False - if task_id: - await sor.sqlExe( - "UPDATE pipeline_tasks SET state='submitted', claimed_by=NULL " - "WHERE id=${tid}$ AND state='waiting'", - {"tid": task_id}) - resumed = True - - suffix = ",任务已恢复执行" if resumed else "" - return "OK: 已回答" + suffix - - async def _t_add_repo(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - from appPublic.uniqueID import getID - - url = p.get("repo_url", "").strip() - name = p.get("repo_name", "").strip() - branch = p.get("default_branch", "main") - if not url: - return "需要仓库URL" - - await sor.C("sd_project_repos", { - "id": getID(), "project_id": pid, - "repo_url": url, "repo_name": name or url.split("/")[-1].replace(".git", ""), - "default_branch": branch, - }) - return f"OK: 已添加仓库 {name or url}" - - async def _t_list_repos(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - recs = await sor.sqlExe( - "SELECT repo_name, repo_url, default_branch FROM sd_project_repos WHERE project_id=${pid}$", - {"pid": pid}) - if not recs: - return "暂无关联仓库" - - lines = [] - for r in recs: - lines.append(f"- {getattr(r, 'repo_name', '')}: {getattr(r, 'repo_url', '')}") - return "\n".join(lines) - - async def _t_clone_repo(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - import os - from pipeline_service.agent_loop import ( - _setup_repos, _get_workspace_dir, _git_clone) - - workspace_dir = await _get_workspace_dir(sor, pid) - url = (p.get("repo_url", "") or "").strip() - - if url: - name = url.rstrip("/").split("/")[-1].replace(".git", "") - target = os.path.join(workspace_dir, "repos", name) - r = await _git_clone(url, target, p.get("branch", "main")) - return f"rc={r['rc']} {r['message']}" - - results = await _setup_repos(sor, workspace_dir, pid) - if not results: - return "无关联仓库可克隆" - return "\n".join( - f"- {x['repo']}: rc={x.get('rc', '?')} {x.get('message', '')}" - for x in results) - async def _t_run_command(self, sor, p, pid): cmd = p.get("command", "") if not cmd: @@ -1042,58 +708,6 @@ class AgentExecutor: except Exception as e: return f"ERROR: {str(e)[:300]}" - async def _t_check_progress(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - recs = await sor.sqlExe( - "SELECT role, state, COUNT(*) as c FROM pipeline_tasks " - "WHERE tenant_id=${pid}$ GROUP BY role, state ORDER BY role, state", - {"pid": pid}) - if not recs: - return "暂无进度数据" - - lines = [] - for r in recs: - lines.append( - f"- {getattr(r, 'role', '?')}: {getattr(r, 'state', '?')} x{getattr(r, 'c', 0)}" - ) - return "\n".join(lines) - - async def _t_add_bug(self, sor, p, pid): - if not pid: - return "请先切换到项目" - - from appPublic.uniqueID import getID - - title = p.get("title", "").strip() - desc = p.get("description", "") - severity = p.get("severity", "medium") - if not title: - return "需要Bug标题" - - # 找项目的迭代(iteration_id 必填) - iteration_id = p.get("iteration_id", "") - if not iteration_id: - its = await sor.sqlExe( - "SELECT id FROM sd_iterations WHERE project_id=${pid}$ ORDER BY created_at ASC LIMIT 1", - {"pid": pid}) - if its: - iteration_id = getattr(its[0], "id", "") - if not iteration_id: - return "ERROR: 项目无迭代,无法提交Bug" - - await sor.C("sd_bugs", { - "id": getID(), "title": title, - "description": desc, "severity": severity, - "status": "open", "iteration_id": iteration_id, - }) - return f"OK: 已提交Bug: {title}" - - # ═══════════════════════════════════════════════════════ - # 通用工具集(Hermes CLI 能力子集) - # ═══════════════════════════════════════════════════════ - def _resolve_ws_path(self, path: str) -> str: """解析相对路径为工作空间内绝对路径(越界返回 '')。""" import os diff --git a/pipeline_service/sdlc_ability.py b/pipeline_service/sdlc_ability.py new file mode 100644 index 0000000..181960d --- /dev/null +++ b/pipeline_service/sdlc_ability.py @@ -0,0 +1,535 @@ +""" +pipeline-service: sdlc_ability — SDLC 开发产线能力包 + +以「可插拔能力包」形式定义 SDLC 产线专属能力: + - 工具定义(ToolDefinition):任务/交付件/问答/仓库/进度/Bug + - prompt 片段(追加到 core 通用心智之后) + - 工具 handler(独立函数,签名 async def handler(sor, params, ctx) -> str) + +通过 register_ability 注册到 pipeline-core 的 PipelineAbility 注册表, +AgentExecutor 按 pipeline_id 动态挂载。 + +未来 B 产线 = 新增 b_ability.py + register_ability,零侵入 core / AgentExecutor。 +""" + +import json +import os + +from pipeline_core import ( + ToolDefinition, + PipelineAbility, + register_ability, +) + +# ── 工具定义 ── + +SDL_TOOLS = [ + ToolDefinition( + name="create_task", + description="创建开发任务,分配给角色agent执行", + parameters={ + "title": "任务标题", + "role": "目标角色(requirement/design/develop/test/deploy)", + "description": "任务详细描述", + }, + category="task", + ), + ToolDefinition( + name="list_tasks", + description="查看当前项目的任务列表", + parameters={"role": "按角色筛选(可选)", "state": "按状态筛选(可选)"}, + category="task", + ), + ToolDefinition( + name="task_detail", + description="查看任务详情", + parameters={"task_id": "任务ID"}, + category="task", + ), + ToolDefinition( + name="start_agents", + description="启动角色agent执行已提交的任务", + parameters={}, + category="agent", + ), + ToolDefinition( + name="diagnose_project", + description="诊断项目:汇总卡点/失败/问题/运行中任务", + parameters={}, + category="agent", + ), + ToolDefinition( + name="list_deliverables", + description="查看交付件列表", + parameters={"role": "按角色筛选(可选)"}, + category="task", + ), + ToolDefinition( + name="view_deliverable", + description="查看交付件详细内容", + parameters={"deliverable_id": "交付件ID"}, + category="task", + ), + ToolDefinition( + name="list_questions", + description="查看待回答的问题", + parameters={}, + category="agent", + ), + ToolDefinition( + name="answer_question", + description="回答agent提出的问题", + parameters={"question_id": "问题ID", "answer": "回答内容"}, + category="agent", + ), + ToolDefinition( + name="add_repo", + description="添加项目关联的Git仓库", + parameters={"repo_url": "仓库URL", "repo_name": "仓库名称", "default_branch": "默认分支(默认main)"}, + category="repo", + ), + ToolDefinition( + name="list_repos", + description="查看项目关联的仓库列表", + parameters={}, + category="repo", + ), + ToolDefinition( + name="clone_repo", + description="克隆项目关联的仓库到工作空间repos/下(不传url则clone全部,传则clone指定)", + parameters={"repo_url": "仓库URL(可选)"}, + category="repo", + ), + ToolDefinition( + name="check_progress", + description="查看项目整体进展", + parameters={}, + category="agent", + ), + ToolDefinition( + name="add_bug", + description="提交Bug", + parameters={"title": "Bug标题", "description": "描述", "severity": "严重程度"}, + category="task", + ), +] + + +# ── prompt 片段(追加到 core 通用心智之后) ── + +SDL_PROMPT = """你是「开发产线」的驾驶舱 agent,负责软件项目的全流程管理(需求 → 设计 → 开发 → 测试 → 部署)。 + +## 典型场景(帮助判断何时用哪个工具,非强制) +- 用户提出新的开发需求/功能("实现XX""XX系统要做XX")→ 用 create_task 创建任务(可按 requirement→design→develop 拆分),再 start_agents 启动。 +- 用户问进展/状态 → list_tasks / diagnose_project / check_progress。 +- 用户问待回答问题 → list_questions / answer_question。 +- 用户报告异常/故障 → 先用 diagnose_project 定位根因,再用 task_detail 查详情。 +- 仓库管理 → add_repo / list_repos / clone_repo。 +- 发现卡点(任务卡死、审核超时、失败、僵尸 claimed_by)时,主动定位根因并推动修复。""" + + +# ── handler 独立函数(签名统一 async def handler(sor, params, ctx) -> str) ── + +async def _h_create_task(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + from appPublic.uniqueID import getID + + title = p.get("title", "").strip() + role = p.get("role", "develop") + desc = p.get("description", "") + if not title: + return "需要任务标题" + + tid = getID() + await sor.C("pipeline_tasks", { + "id": tid, "tenant_id": pid, "pipeline_id": "role_task", + "owner_id": ctx.get("user_id") or "user", "title": title, + "state": "submitted", "role": role, + "params": json.dumps({"description": desc}, ensure_ascii=False), + }) + return f"OK: 已创建任务 {title}(角色: {role})" + + +async def _h_list_tasks(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + role = p.get("role", "") + state = p.get("state", "") + + sql = "SELECT id, title, role, state, created_at FROM pipeline_tasks WHERE tenant_id=${pid}$" + params_dict = {"pid": pid} + if role: + sql += " AND role=${role}$" + params_dict["role"] = role + if state: + sql += " AND state=${state}$" + params_dict["state"] = state + sql += " ORDER BY created_at DESC LIMIT 30" + + recs = await sor.sqlExe(sql, params_dict) + if not recs: + return "暂无任务" + + icons = { + "submitted": "⏳", "running": "🔄", "review": "👀", + "approved": "✅", "completed": "✔️", "failed": "❌", + "waiting": "⏸️", + } + lines = ["| 状态 | 任务 | 角色 | ID |", "|------|------|------|----|"] + for r in recs: + st = getattr(r, 'state', '?') + icon = icons.get(st, "❓") + title = (getattr(r, 'title', '') or '')[:60] + rrole = getattr(r, 'role', '') + tid = getattr(r, 'id', '') + lines.append(f"| {icon} {st} | {title} | {rrole} | {tid} |") + return "\n".join(lines) + + +async def _h_task_detail(sor, p, ctx): + tid = p.get("task_id", "") + if not tid: + return "需要任务ID" + recs = await sor.sqlExe( + "SELECT * FROM pipeline_tasks WHERE id=${tid}$", {"tid": tid}) + if not recs and len(tid) >= 6: + recs = await sor.sqlExe( + "SELECT * FROM pipeline_tasks WHERE id LIKE ${prefix}$ LIMIT 1", + {"prefix": tid + "%"}) + if not recs: + return f"任务不存在: {tid}" + r = recs[0] + return json.dumps({ + "id": getattr(r, "id", ""), + "title": getattr(r, "title", ""), + "state": getattr(r, "state", ""), + "role": getattr(r, "role", ""), + "params": getattr(r, "params", "{}"), + }, ensure_ascii=False, indent=2) + + +async def _h_start_agents(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + recs = await sor.sqlExe( + "SELECT id, title, role FROM pipeline_tasks " + "WHERE tenant_id=${pid}$ AND state='submitted' ORDER BY created_at ASC LIMIT 10", + {"pid": pid}) + if not recs: + return "没有待执行的任务" + + results = [] + for r in recs: + tid = getattr(r, "id", "") + role = getattr(r, "role", "develop") + try: + from pipeline_service.agent_loop import role_agent_run + result = await role_agent_run(pid, role) + results.append(f"- {getattr(r, 'title', '')}: {result.get('status', '?')}") + except Exception as e: + results.append(f"- {getattr(r, 'title', '')}: error={str(e)[:100]}") + return "Agent 执行结果:\n" + "\n".join(results) + + +async def _h_diagnose_project(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + + async def _cnt(state): + recs = await sor.sqlExe( + "SELECT COUNT(*) as c FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state=${st}$", + {"pid": pid, "st": state}) + return getattr(recs[0], "c", 0) if recs else 0 + + sc = await _cnt("submitted") + rc = await _cnt("running") + fc = await _cnt("failed") + vc = await _cnt("review") + + qc = 0 + try: + questions = await sor.sqlExe( + "SELECT COUNT(*) as c FROM pipeline_agent_questions WHERE tenant_id=${pid}$ AND status='pending'", + {"pid": pid}) + qc = getattr(questions[0], "c", 0) if questions else 0 + except Exception: + pass + + lines = [ + "项目诊断:", + f"- 待执行: {sc}", + f"- 运行中: {rc}", + f"- 待审核: {vc}", + f"- 失败: {fc}", + f"- 待回答问题: {qc}", + ] + + if fc: + failed_rows = await sor.sqlExe( + "SELECT id, title, role FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='failed' ORDER BY created_at DESC LIMIT 10", + {"pid": pid}) + lines.append("失败任务:") + for fr in (failed_rows or []): + lines.append(f" - [{getattr(fr, 'role', '?')}] {getattr(fr, 'title', '')} (id={getattr(fr, 'id', '')})") + if sc: + submitted_rows = await sor.sqlExe( + "SELECT id, title, role FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='submitted' ORDER BY created_at ASC LIMIT 10", + {"pid": pid}) + lines.append("待执行任务:") + for sr in (submitted_rows or []): + lines.append(f" - [{getattr(sr, 'role', '?')}] {getattr(sr, 'title', '')} (id={getattr(sr, 'id', '')})") + if rc: + running_rows = await sor.sqlExe( + "SELECT id, title, role, claimed_by, " + "TIMESTAMPDIFF(MINUTE, updated_at, NOW()) AS mins " + "FROM pipeline_tasks WHERE tenant_id=${pid}$ AND state='running' " + "ORDER BY updated_at ASC LIMIT 10", + {"pid": pid}) + lines.append("运行中任务:") + for rr in (running_rows or []): + rrole = getattr(rr, 'role', '?') + rtitle = getattr(rr, 'title', '') + rid = getattr(rr, 'id', '') + rcb = getattr(rr, 'claimed_by', '') or '' + try: + mins = int(getattr(rr, 'mins', 0) or 0) + except (TypeError, ValueError): + mins = 0 + flag = f"⚠️心跳超时{mins}分钟(疑似僵尸,将被回收重跑)" if mins >= 20 else f"已运行{mins}分钟" + lines.append(f" - [{rrole}] {rtitle} (id={rid}) | {flag} | claimed_by={rcb[:8]}") + return "\n".join(lines) + + +async def _h_list_deliverables(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + sql = "SELECT id, title, deliverable_type, review_status FROM pipeline_deliverables WHERE project_id=${pid}$ ORDER BY created_at DESC LIMIT 20" + recs = await sor.sqlExe(sql, {"pid": pid}) + if not recs: + return "暂无交付件" + lines = [] + for r in recs: + lines.append( + f"- [{getattr(r, 'review_status', '?')}] {getattr(r, 'title', '')} " + f"({getattr(r, 'deliverable_type', '')}) id={getattr(r, 'id', '')}") + return "\n".join(lines) + + +async def _h_view_deliverable(sor, p, ctx): + did = p.get("deliverable_id", "") + if not did: + return "需要交付件ID" + recs = await sor.sqlExe( + "SELECT * FROM pipeline_deliverables WHERE id=${did}$", {"did": did}) + if not recs and len(did) >= 6: + recs = await sor.sqlExe( + "SELECT * FROM pipeline_deliverables WHERE id LIKE ${prefix}$ LIMIT 1", + {"prefix": did + "%"}) + if not recs: + return f"交付件不存在: {did}" + return getattr(recs[0], "content", "")[:4000] or "(空)" + + +async def _h_list_questions(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + try: + recs = await sor.sqlExe( + "SELECT id, question, answer_source FROM pipeline_agent_questions " + "WHERE tenant_id=${pid}$ AND status='pending' ORDER BY created_at DESC LIMIT 10", + {"pid": pid}) + except Exception: + try: + recs = await sor.sqlExe( + "SELECT id, question, answer_source FROM pipeline_agent_questions " + "WHERE tenant_id=${pid}$ ORDER BY created_at DESC LIMIT 10", + {"pid": pid}) + except Exception: + return "暂无问题数据" + if not recs: + return "没有待回答问题" + lines = [] + for r in recs: + lines.append(f"- [{getattr(r, 'id', '')}] {getattr(r, 'question', '')[:200]}") + return "\n".join(lines) + + +async def _h_answer_question(sor, p, ctx): + qid = p.get("question_id", "") + answer = p.get("answer", "") + if not qid or not answer: + return "需要问题ID和回答内容" + recs = await sor.sqlExe( + "SELECT id, task_id FROM pipeline_agent_questions WHERE id=${qid}$", {"qid": qid}) + if not recs and len(qid) >= 6: + recs = await sor.sqlExe( + "SELECT id, task_id FROM pipeline_agent_questions WHERE id LIKE ${prefix}$ LIMIT 1", + {"prefix": qid + "%"}) + if not recs: + return f"问题不存在: {qid}" + full_qid = getattr(recs[0], "id", qid) + task_id = getattr(recs[0], "task_id", "") or "" + + await sor.sqlExe( + "UPDATE pipeline_agent_questions SET answer=${a}$, answer_source='main_agent', " + "answered_by='main_agent', status='answered' WHERE id=${qid}$", + {"a": answer, "qid": full_qid}) + + resumed = False + if task_id: + await sor.sqlExe( + "UPDATE pipeline_tasks SET state='submitted', claimed_by=NULL " + "WHERE id=${tid}$ AND state='waiting'", + {"tid": task_id}) + resumed = True + suffix = ",任务已恢复执行" if resumed else "" + return "OK: 已回答" + suffix + + +async def _h_add_repo(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + from appPublic.uniqueID import getID + + url = p.get("repo_url", "").strip() + name = p.get("repo_name", "").strip() + branch = p.get("default_branch", "main") + if not url: + return "需要仓库URL" + await sor.C("sd_project_repos", { + "id": getID(), "project_id": pid, + "repo_url": url, "repo_name": name or url.split("/")[-1].replace(".git", ""), + "default_branch": branch, + }) + return f"OK: 已添加仓库 {name or url}" + + +async def _h_list_repos(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + recs = await sor.sqlExe( + "SELECT repo_name, repo_url, default_branch FROM sd_project_repos WHERE project_id=${pid}$", + {"pid": pid}) + if not recs: + return "暂无关联仓库" + lines = [] + for r in recs: + lines.append(f"- {getattr(r, 'repo_name', '')}: {getattr(r, 'repo_url', '')}") + return "\n".join(lines) + + +async def _h_clone_repo(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + from pipeline_service.agent_loop import ( + _setup_repos, _get_workspace_dir, _git_clone) + + workspace_dir = await _get_workspace_dir(sor, pid) + url = (p.get("repo_url", "") or "").strip() + + if url: + name = url.rstrip("/").split("/")[-1].replace(".git", "") + target = os.path.join(workspace_dir, "repos", name) + r = await _git_clone(url, target, p.get("branch", "main")) + return f"rc={r['rc']} {r['message']}" + + results = await _setup_repos(sor, workspace_dir, pid) + if not results: + return "无关联仓库可克隆" + return "\n".join( + f"- {x['repo']}: rc={x.get('rc', '?')} {x.get('message', '')}" + for x in results) + + +async def _h_check_progress(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + recs = await sor.sqlExe( + "SELECT role, state, COUNT(*) as c FROM pipeline_tasks " + "WHERE tenant_id=${pid}$ GROUP BY role, state ORDER BY role, state", + {"pid": pid}) + if not recs: + return "暂无进度数据" + lines = [] + for r in recs: + lines.append( + f"- {getattr(r, 'role', '?')}: {getattr(r, 'state', '?')} x{getattr(r, 'c', 0)}") + return "\n".join(lines) + + +async def _h_add_bug(sor, p, ctx): + pid = ctx.get("project_id", "") + if not pid: + return "请先切换到项目" + from appPublic.uniqueID import getID + + title = p.get("title", "").strip() + desc = p.get("description", "") + severity = p.get("severity", "medium") + if not title: + return "需要Bug标题" + + iteration_id = p.get("iteration_id", "") + if not iteration_id: + its = await sor.sqlExe( + "SELECT id FROM sd_iterations WHERE project_id=${pid}$ ORDER BY created_at ASC LIMIT 1", + {"pid": pid}) + if its: + iteration_id = getattr(its[0], "id", "") + if not iteration_id: + return "ERROR: 项目无迭代,无法提交Bug" + + await sor.C("sd_bugs", { + "id": getID(), "title": title, + "description": desc, "severity": severity, + "status": "open", "iteration_id": iteration_id, + }) + return f"OK: 已提交Bug: {title}" + + +# ── 注册能力包 ── + +SDL_HANDLERS = { + "create_task": _h_create_task, + "list_tasks": _h_list_tasks, + "task_detail": _h_task_detail, + "start_agents": _h_start_agents, + "diagnose_project": _h_diagnose_project, + "list_deliverables": _h_list_deliverables, + "view_deliverable": _h_view_deliverable, + "list_questions": _h_list_questions, + "answer_question": _h_answer_question, + "add_repo": _h_add_repo, + "list_repos": _h_list_repos, + "clone_repo": _h_clone_repo, + "check_progress": _h_check_progress, + "add_bug": _h_add_bug, +} + + +def register_sdlc_ability(): + """注册 SDLC 产线能力包(幂等)。""" + ability = PipelineAbility( + pipeline_id="sdlc_general", + name="通用软件开发产线", + tools=SDL_TOOLS, + system_prompt=SDL_PROMPT, + handlers=SDL_HANDLERS, + ) + register_ability(ability) + return ability + + +# import 即注册(模块加载副作用),pipeline-app 启动时 import pipeline_service 即生效 +register_sdlc_ability()