diff --git a/pipeline_service/agent_loop_v2.py b/pipeline_service/agent_loop_v2.py index de4886e..50197ec 100644 --- a/pipeline_service/agent_loop_v2.py +++ b/pipeline_service/agent_loop_v2.py @@ -26,6 +26,24 @@ from typing import AsyncGenerator, Dict, List, Optional logger = logging.getLogger("pipeline.agent_executor") +# 项目管理工具别名表:LLM 常编造近义名(停止/终止/删除项目的各种说法), +# 归一到真实工具名。只加「语义等价」的别名,不扩大权限(权限校验在代码层)。 +# 注意:不能映射既有真实工具名(如 list_projects 是 sdlc_ability 的产线项目列表 +# 工具,语义与 list_my_projects 不同),否则会劫持既有功能。 +_PROJECT_TOOL_ALIASES = { + "stop_project": "pause_project", + "terminate_project": "pause_project", + "suspend_project": "pause_project", + "halt_project": "pause_project", + "unpause_project": "resume_project", + "continue_project": "resume_project", + "remove_project": "delete_project", + "destroy_project": "delete_project", + "get_project_info": "project_info", + "show_project": "project_info", + "my_projects": "list_my_projects", +} + from .workspace import WORKSPACE_BASE, GENERAL_SPACE @@ -669,6 +687,9 @@ class AgentExecutor: 2. SDLC 内建工具(兼容 cockpit_chat.dspy 的 TOOLS) 3. ask_user(始终可用) """ + # 别名归一化:LLM 会编造工具名(2026-08-31 项目管理的常见变体), + # 映射到真实工具,避免「未知工具」导致用户指令落空。 + tool_name = _PROJECT_TOOL_ALIASES.get(tool_name, tool_name) # 1. 注册的 handler if self._tool_registry: handler = self._tool_registry.get_handler(tool_name) @@ -761,6 +782,12 @@ class AgentExecutor: handlers = { "switch_project": self._t_switch_project, "create_project": self._t_create_project, + # ── 项目管理(终止/删除类代码层校验人类 owner)── + "project_info": self._t_project_info, + "list_my_projects": self._t_list_my_projects, + "pause_project": self._t_pause_project, + "resume_project": self._t_resume_project, + "delete_project": self._t_delete_project, "run_command": self._t_run_command, # ── 通用工具集(Hermes CLI 能力子集)── "read_file": self._t_read_file, @@ -943,6 +970,154 @@ class AgentExecutor: await self._persist_project(sor, pid_val) return f"OK: 已创建项目 {name}" + # ── 项目管理工具(2026-08-31 新增:会话 agent 获得项目管理能力)── + # 安全模型:终止(暂停)/删除等指令必须来自项目的人类 owner。 + # 校验在代码层做(_resolve_owned_project 内调 check_project_owner), + # 不依赖 prompt——LLM 可以被注入/诱导,但校验函数只认数据库事实。 + # owner 定义:sd_projects.created_by == 当前登录用户(self.user_id)。 + + async def _resolve_owned_project(self, sor, p): + """解析目标项目并校验操作者是否为人类 owner。 + + project_name 缺省时用当前项目(self.project_id);给了名称则在当前产线 + 空间内精确匹配。返回 (project_id, project_name, owner_err): + owner_err 非空 = 无权限或项目不存在,调用方必须原样拒绝。 + """ + from .project_capability import check_project_owner + name = (p.get("project_name") or "").strip() + if not name: + pid = self.project_id + if not pid: + return "", "", "当前会话没有选中项目,请先说明要操作哪个项目(或先切换项目)" + recs = await sor.sqlExe( + "SELECT id, name, created_by FROM sd_projects WHERE id=${pid}$", + {"pid": pid}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return "", "", "当前项目已不存在" + pname = getattr(recs[0], "name", "") or pid + else: + recs = await sor.sqlExe( + "SELECT id, name, created_by FROM sd_projects " + "WHERE name=${n}$ AND pipeline_id=${s}$", + {"n": name, "s": self.space}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return "", "", f"项目不存在:{name}(当前产线空间 {self.space})" + pid = getattr(recs[0], "id", "") + pname = getattr(recs[0], "name", "") or pid + # 代码层 owner 校验:指令必须来自该项目的人类 owner(创建者) + ok, err = await check_project_owner(pid, self.user_id or "", sor) + if not ok: + return "", "", f"权限不足:{err}(项目「{pname}」的 owner 才可执行此操作)" + return pid, pname, "" + + async def _t_project_info(self, sor, p, pid): + tpid, pname, err = await self._resolve_owned_project(sor, p) + if err: + return f"FAIL: {err}" + recs = await sor.sqlExe( + "SELECT id, name, description, status, pipeline_id, org_id, " + "created_by, created_at, updated_at FROM sd_projects WHERE id=${pid}$", + {"pid": tpid}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return "FAIL: 项目不存在" + r = recs[0] + lines = [ + f"项目:{getattr(r, 'name', '')}", + f"状态:{getattr(r, 'status', '')}", + f"描述:{getattr(r, 'description', '') or '(无)'}", + f"产线:{getattr(r, 'pipeline_id', '') or '(未绑定)'}", + f"创建:{getattr(r, 'created_at', '')} | 更新:{getattr(r, 'updated_at', '')}", + ] + # 迭代概况 + iters = await sor.sqlExe( + "SELECT iteration_name, status FROM sd_iterations WHERE project_id=${pid}$ ORDER BY seq_no", + {"pid": tpid}) + await sor.sqlExe("COMMIT", {}) + if iters: + lines.append("迭代:" + ";".join( + f"{getattr(i, 'iteration_name', '')}({getattr(i, 'status', '')})" for i in iters[:10])) + return "\n".join(lines) + + async def _t_list_my_projects(self, sor, p, pid): + # 只列当前用户创建的项目(不暴露他人项目) + uid = self.user_id or "" + if not uid: + return "FAIL: 未登录,无法列出你的项目" + status = (p.get("status") or "").strip() + sql = ("SELECT id, name, status, pipeline_id, created_at FROM sd_projects " + "WHERE created_by=${uid}$") + params = {"uid": uid} + if status: + sql += " AND status=${st}$" + params["st"] = status + sql += " ORDER BY created_at DESC LIMIT 30" + recs = await sor.sqlExe(sql, params) + await sor.sqlExe("COMMIT", {}) + if not recs: + return f"你名下没有项目{('(状态=' + status + ')') if status else ''}" + lines = [] + for r in recs: + lines.append(f"- {getattr(r, 'name', '')} | {getattr(r, 'status', '')} | " + f"产线 {getattr(r, 'pipeline_id', '') or '—'} | {getattr(r, 'created_at', '')}") + return f"你创建的项目({len(lines)} 个):\n" + "\n".join(lines) + + async def _t_pause_project(self, sor, p, pid): + tpid, pname, err = await self._resolve_owned_project(sor, p) + if err: + return f"FAIL: {err}" + from .project_capability import pause_project + ok, msg = await pause_project(tpid, who=self.user_id, agent_id="cockpit") + if not ok: + return f"FAIL: {msg}" + logger.info("cockpit pause_project: %s by user=%s", tpid, self.user_id) + return f"OK: 项目「{pname}」已暂停(终止推进)。需要时可用 resume_project 恢复。" + + async def _t_resume_project(self, sor, p, pid): + tpid, pname, err = await self._resolve_owned_project(sor, p) + if err: + return f"FAIL: {err}" + from .project_capability import resume_project + ok, msg = await resume_project(tpid, who=self.user_id, agent_id="cockpit") + if not ok: + return f"FAIL: {msg}" + logger.info("cockpit resume_project: %s by user=%s", tpid, self.user_id) + return f"OK: 项目「{pname}」已恢复推进。" + + async def _t_delete_project(self, sor, p, pid): + confirm = p.get("confirm") + if str(confirm).lower() not in ("true", "1", "yes"): + return ("FAIL: 删除项目不可撤销。请先向用户复述项目名与后果," + "取得明确确认后再带 confirm=true 调用本工具。") + tpid, pname, err = await self._resolve_owned_project(sor, p) + if err: + return f"FAIL: {err}" + # 删除后清掉本会话的项目上下文(否则留下悬空指针,下次打开报「请先切换项目」) + from .project_capability import delete_project + ok, msg = await delete_project(tpid, who=self.user_id, agent_id="cockpit", confirm=True) + if not ok: + return f"FAIL: {msg}" + if self.project_id == tpid: + self.project_id = "" + # 清会话级+全局「当前项目」指针(_persist_project 对空 pid 直接 return, + # 这里直接清,否则留下悬空指针,下次打开报「请先切换项目」) + try: + await sor.sqlExe( + "UPDATE pipeline_session_settings SET current_project_id='' " + "WHERE user_id=${uid}$ AND session_id=${sid}$", + {"uid": self.user_id or "", "sid": self.session_id or ""}) + await sor.sqlExe( + "UPDATE pipeline_agent_settings SET current_project_id='' " + "WHERE user_id=${uid}$", + {"uid": self.user_id or ""}) + await sor.sqlExe("COMMIT", {}) + except Exception as e: + logger.warning(f"clear project pointer after delete failed: {e}") + logger.info("cockpit delete_project: %s by user=%s", tpid, self.user_id) + return f"OK: 项目「{pname}」已删除(含归档备份)。{msg}" + async def _t_run_command(self, sor, p, pid): cmd = p.get("command", "") if not cmd: @@ -1433,7 +1608,11 @@ class AgentExecutor: 会话级「全部确认」(approve_all)时直接放行;否则 run_command 只对 危险命令确认,普通命令自动执行,其他 requires_confirmation 工具一律确认。 + + 别名先归一再查注册表(2026-08-31 修复):否则 LLM 用别名(如 stop_project) + 调需确认工具时,注册表查不到 → 误判无需确认 → 绕过确认门直接执行。 """ + tool_name = _PROJECT_TOOL_ALIASES.get(tool_name, tool_name) if self._session and getattr(self._session, "approve_all", False): return False if self._tool_registry: diff --git a/pipeline_service/sdlc_ability.py b/pipeline_service/sdlc_ability.py index 110ef66..f26fb16 100644 --- a/pipeline_service/sdlc_ability.py +++ b/pipeline_service/sdlc_ability.py @@ -184,10 +184,10 @@ SDL_TOOLS = [ ToolDefinition(name="complete_project", description="完成项目(active→completed)", parameters={"project_id": "项目ID"}, category="project"), ToolDefinition(name="archive_project", description="归档项目(completed→archived)", parameters={"project_id": "项目ID"}, category="project"), ToolDefinition(name="reopen_project", description="重新打开归档项目(archived→active)", parameters={"project_id": "项目ID"}, category="project"), - ToolDefinition(name="pause_project", description="暂停项目推进(active→paused,PM 不再推进,仅用户明确指令暂停时用)", parameters={"project_id": "项目ID"}, category="project"), - ToolDefinition(name="resume_project", description="恢复项目推进(paused→active)", parameters={"project_id": "项目ID"}, category="project"), + ToolDefinition(name="pause_project", description="暂停项目推进(active→paused,PM 不再推进,仅用户明确指令暂停时用;仅项目人类 owner 可执行)", parameters={"project_id": "项目ID"}, category="project", requires_confirmation=True), + ToolDefinition(name="resume_project", description="恢复项目推进(paused→active;仅项目人类 owner 可执行)", parameters={"project_id": "项目ID"}, category="project", requires_confirmation=True), ToolDefinition(name="backup_project", description="归档备份项目(打包工作目录+数据记录为tgz到机构工作目录_archive/)", parameters={"project_id": "项目ID"}, category="project"), - ToolDefinition(name="delete_project", description="删除项目(需二次确认confirm=true,删除前自动归档备份tgz)", parameters={"project_id": "项目ID", "confirm": "二次确认,必须为true"}, category="project"), + ToolDefinition(name="delete_project", description="删除项目(需二次确认confirm=true,删除前自动归档备份tgz;仅项目人类 owner 可执行)", parameters={"project_id": "项目ID", "confirm": "二次确认,必须为true"}, category="project", requires_confirmation=True), ToolDefinition(name="cleanup_orphans", description="清扫指向已删项目/迭代/任务的孤儿残留记录(需confirm=true,dry_run=true只统计不删)", parameters={"confirm": "二次确认,必须为true", "dry_run": "true=只统计不删除(可选)"}, category="admin"), # ── 迭代 ── ToolDefinition(name="list_iterations", description="查看当前项目迭代列表", parameters={"status": "按状态筛选(可选)"}, category="iteration"), @@ -1021,6 +1021,29 @@ async def _resolve_id(sor, table, fid): return "" +async def _require_project_owner(sor, ctx, pid): + """项目生命周期变更的 owner 门禁(2026-08-31 安全加固)。 + + 会话 agent 的项目终止/删除等指令必须来自该项目的**人类 owner**。 + 校验在代码层做(不靠 prompt):ctx.user_id 是网关层登录态,无法被 LLM 伪造; + 仅当 created_by == 登录用户时放行。 + + 无 user_id(无人值守/agent 直调)一律拒绝——生命周期变更只允许人类 owner + 从会话发起;产线内部如需暂停(如 failed-poller 熔断)直接调 + project_capability,不走本工具层。 + + 返回 (ok, err):ok=True 放行;ok=False 时 err 为拒绝原因,原样返回给调用方。 + """ + from .project_capability import check_project_owner + uid = ctx.get("user_id", "") or "" + if not uid: + return False, "权限不足:项目生命周期变更(暂停/恢复/归档/删除)必须由登录用户发起,且须为项目 owner。当前调用无登录身份,已拒绝。" + ok, err = await check_project_owner(pid, uid, sor) + if not ok: + return False, f"权限不足:{err}。终止/删除等项目指令只接受该项目的创建者(人类 owner)。" + return True, "" + + async def _resolve_task_id(sor, task_id, pid): """任务 ID 前缀兜底(租户隔离):精确匹配失败且入参≥6位时 LIKE 前缀解析完整 ID。""" if not task_id: @@ -1112,8 +1135,11 @@ async def _h_start_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import start_project - ok, msg = await start_project(pid, who="agent.main_agent") + ok, msg = await start_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1167,8 +1193,11 @@ async def _h_complete_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import complete_project - ok, msg = await complete_project(pid, who="agent.main_agent") + ok, msg = await complete_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1177,8 +1206,11 @@ async def _h_archive_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import archive_project - ok, msg = await archive_project(pid, who="agent.main_agent") + ok, msg = await archive_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1187,8 +1219,11 @@ async def _h_reopen_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import reopen_project - ok, msg = await reopen_project(pid, who="agent.main_agent") + ok, msg = await reopen_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1197,8 +1232,11 @@ async def _h_pause_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import pause_project - ok, msg = await pause_project(pid, who="agent.main_agent") + ok, msg = await pause_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1207,8 +1245,11 @@ async def _h_resume_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import resume_project - ok, msg = await resume_project(pid, who="agent.main_agent") + ok, msg = await resume_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1217,8 +1258,11 @@ async def _h_backup_project(sor, p, ctx): if not pid: return "需要项目ID" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import backup_project - ok, msg = await backup_project(pid, who="agent.main_agent") + ok, msg = await backup_project(pid, who=ctx.get("user_id") or "agent.main_agent") return f"OK: {msg}" if ok else f"ERROR: {msg}" @@ -1230,8 +1274,11 @@ async def _h_delete_project(sor, p, ctx): if confirm != "true": return "ERROR: 删除项目需二次确认,请传 confirm=true 后再执行(删除前会自动归档备份到机构工作目录 _archive/)" pid = await _resolve_id(sor, "sd_projects", pid) or pid + ok, err = await _require_project_owner(sor, ctx, pid) + if not ok: + return f"ERROR: {err}" from .project_capability import delete_project - ok, msg = await delete_project(pid, who="agent.main_agent", confirm=True) + ok, msg = await delete_project(pid, who=ctx.get("user_id") or "agent.main_agent", confirm=True) return f"OK: {msg}" if ok else f"ERROR: {msg}"