package controlplane import ( "database/sql" "encoding/json" "errors" "fmt" "net/url" "regexp" "strings" "time" ) // MCPTool is one controller-owned tool definition. The controller stores only // the definition; execution stays on the Agent that the tool is assigned to. type MCPTool struct { ID string `json:"id"` Name string `json:"name"` Description string `json:"description"` Kind string `json:"kind"` Endpoint string `json:"endpoint"` Command string `json:"command"` Args []string `json:"args"` Version string `json:"version"` UpdatedAt string `json:"updatedAt"` AgentCount int `json:"agentCount"` } // MCPToolInput is accepted only from the local Wails UI. type MCPToolInput struct { ID string `json:"id"` Name string `json:"name"` Description string `json:"description"` Kind string `json:"kind"` Endpoint string `json:"endpoint"` Command string `json:"command"` Args []string `json:"args"` Version string `json:"version"` } // AgentMCPStatus is one Agent's application state for one assigned tool. type AgentMCPStatus struct { AgentID string `json:"agentID"` AgentName string `json:"agentName"` ToolID string `json:"toolID"` ToolName string `json:"toolName"` Kind string `json:"kind"` Endpoint string `json:"endpoint"` Command string `json:"command"` Version string `json:"version"` DeployVersion int `json:"deployVersion"` ApplyStatus string `json:"applyStatus"` LastError string `json:"lastError"` UpdatedAt string `json:"updatedAt"` } // AgentDeploymentStatus summarizes controller-owned LLM, Skill and MCP application state. type AgentDeploymentStatus struct { AgentID string `json:"agentID"` LLMProfile string `json:"llmProfile"` LLMStatus string `json:"llmStatus"` SkillTotal int `json:"skillTotal"` SkillApplied int `json:"skillApplied"` SkillFailed int `json:"skillFailed"` MCPStatus string `json:"mcpStatus"` MCPError string `json:"mcpError"` LastUpdatedAt string `json:"lastUpdatedAt"` } const ( toolKindBuiltin = "builtin" toolKindHTTP = "http" toolKindStdio = "stdio" ) // toolArgumentLimits bound a controller-supplied stdio command line before it is // stored, so a malformed definition can never reach an Agent process spawn. const ( maxToolCommandRunes = 500 maxToolArguments = 32 maxToolArgumentRunes = 500 ) // builtinToolIDs are the tool ids the Agent implements natively. Any other id // must be an HTTP MCP endpoint, so an unknown builtin can never be deployed. var builtinToolIDs = map[string]struct{}{"desktop-filesystem": {}} var toolIDPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9._-]{0,63}$`) func seedBuiltinMCPTools(db *sql.DB) error { now := time.Now().UTC().Format(time.RFC3339Nano) // The legacy narrow desktop tool stays available as a builtin tool so an // existing deployment keeps working after the repository migration. _, err := db.Exec( "INSERT OR IGNORE INTO mcp_tools(tool_id, name, description, kind, endpoint, version, updated_at) VALUES ('desktop-filesystem', '桌面文件工具(旧版兼容)', '仅识别在桌面直接创建文件夹的窄工具;新版 Agent 默认由完全访问执行器处理本机操作。', ?, '', '1.0.0', ?)", toolKindBuiltin, now, ) if err != nil { return fmt.Errorf("无法初始化内置工具: %w", err) } return nil } // MCPTools lists the controller-owned tool repository with assignment counts. func (c *ControlPlane) MCPTools() ([]MCPTool, error) { rows, err := c.db.Query(`SELECT t.tool_id, t.name, t.description, t.kind, t.endpoint, t.command, t.args, t.version, t.updated_at, COALESCE((SELECT COUNT(*) FROM agent_mcp_assignments a WHERE a.tool_id = t.tool_id), 0) FROM mcp_tools t ORDER BY t.name COLLATE NOCASE, t.tool_id`) if err != nil { return nil, fmt.Errorf("无法读取工具仓库: %w", err) } defer rows.Close() tools := make([]MCPTool, 0) for rows.Next() { var tool MCPTool var encodedArgs string if err := rows.Scan(&tool.ID, &tool.Name, &tool.Description, &tool.Kind, &tool.Endpoint, &tool.Command, &encodedArgs, &tool.Version, &tool.UpdatedAt, &tool.AgentCount); err != nil { return nil, err } tool.Args = decodeToolArgs(encodedArgs) tools = append(tools, tool) } return tools, rows.Err() } func decodeToolArgs(encoded string) []string { args := make([]string, 0) if encoded == "" || json.Unmarshal([]byte(encoded), &args) != nil { return make([]string, 0) } if args == nil { return make([]string, 0) } return args } func encodeToolArgs(args []string) string { if len(args) == 0 { return "[]" // One canonical empty form keeps definition comparison stable. } encoded, err := json.Marshal(args) if err != nil { return "[]" } return string(encoded) } // SaveMCPTool creates or updates one repository entry. A changed definition // bumps the deployment version of every Agent that already has the tool, so a // rollback or an endpoint change is re-applied instead of being ignored as stale. func (c *ControlPlane) SaveMCPTool(input MCPToolInput) (MCPTool, error) { tool, err := normaliseMCPTool(input) if err != nil { return MCPTool{}, err } var previousKind, previousEndpoint, previousCommand, previousArgs, previousVersion string update := c.db.QueryRow("SELECT kind, endpoint, command, args, version FROM mcp_tools WHERE tool_id = ?", tool.ID). Scan(&previousKind, &previousEndpoint, &previousCommand, &previousArgs, &previousVersion) == nil definitionChanged := !update || previousKind != tool.Kind || previousEndpoint != tool.Endpoint || previousCommand != tool.Command || previousArgs != encodeToolArgs(tool.Args) || previousVersion != tool.Version now := time.Now().UTC().Format(time.RFC3339Nano) if _, err := c.db.Exec( "INSERT INTO mcp_tools(tool_id, name, description, kind, endpoint, command, args, version, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(tool_id) DO UPDATE SET name = excluded.name, description = excluded.description, kind = excluded.kind, endpoint = excluded.endpoint, command = excluded.command, args = excluded.args, version = excluded.version, updated_at = excluded.updated_at", tool.ID, tool.Name, tool.Description, tool.Kind, tool.Endpoint, tool.Command, encodeToolArgs(tool.Args), tool.Version, now, ); err != nil { return MCPTool{}, fmt.Errorf("无法保存工具: %w", err) } tool.UpdatedAt = now c.recordAudit("mcp.tool_saved", "controller", "mcp-tool:"+tool.ID, "completed", "已保存工具定义;未写入任何凭据。") if update && definitionChanged { c.recordAudit("mcp.tool_redeployment_requested", "controller", "mcp-tool:"+tool.ID, "pending", "工具定义已变更,将重新下发到已分配的 Agent。") } if definitionChanged { c.redeployMCPTool(tool.ID) } return tool, nil } func normaliseMCPTool(input MCPToolInput) (MCPTool, error) { tool := MCPTool{ ID: strings.TrimSpace(strings.ToLower(input.ID)), Name: strings.TrimSpace(input.Name), Description: strings.TrimSpace(input.Description), Kind: strings.TrimSpace(strings.ToLower(input.Kind)), Endpoint: strings.TrimRight(strings.TrimSpace(input.Endpoint), "/"), Command: strings.TrimSpace(input.Command), Version: strings.TrimSpace(input.Version), } if !toolIDPattern.MatchString(tool.ID) { return MCPTool{}, errors.New("工具标识只能使用小写字母、数字、点、下划线和连字符") } if tool.Name == "" { return MCPTool{}, errors.New("请填写工具名称") } if tool.Version == "" { tool.Version = "1.0.0" } if len([]rune(tool.Name)) > 80 || len([]rune(tool.Description)) > 500 || len([]rune(tool.Version)) > 40 { return MCPTool{}, errors.New("工具名称、说明或版本过长") } args, err := normaliseToolArgs(input.Args) if err != nil { return MCPTool{}, err } switch tool.Kind { case toolKindBuiltin: if _, ok := builtinToolIDs[tool.ID]; !ok { return MCPTool{}, errors.New("该标识不是 Agent 内置工具,请选择 HTTP MCP 或本机命令类型") } tool.Endpoint, tool.Command, tool.Args = "", "", nil case toolKindHTTP: parsed, parseErr := url.Parse(tool.Endpoint) if parseErr != nil || parsed.Host == "" || (parsed.Scheme != "https" && parsed.Scheme != "http") { return MCPTool{}, errors.New("HTTP MCP 工具必须填写有效的 http:// 或 https:// 地址") } if len(tool.Endpoint) > 500 { return MCPTool{}, errors.New("MCP 地址过长") } tool.Command, tool.Args = "", nil case toolKindStdio: if tool.Command == "" { return MCPTool{}, errors.New("本机命令工具必须填写要执行的命令") } if len([]rune(tool.Command)) > maxToolCommandRunes || strings.ContainsAny(tool.Command, "\x00\r\n") { return MCPTool{}, errors.New("命令内容无效或过长") } tool.Endpoint, tool.Args = "", args default: return MCPTool{}, errors.New("工具类型必须是内置、HTTP MCP 或本机命令") } return tool, nil } func normaliseToolArgs(input []string) ([]string, error) { if len(input) > maxToolArguments { return nil, fmt.Errorf("参数最多 %d 个", maxToolArguments) } args := make([]string, 0, len(input)) for _, argument := range input { if strings.ContainsAny(argument, "\x00\r\n") { return nil, errors.New("参数内容无效") } if len([]rune(argument)) > maxToolArgumentRunes { return nil, errors.New("单个参数过长") } args = append(args, argument) } return args, nil } // DeleteMCPTool removes the definition and every assignment, then clears the // tool on the Agents that had it (a rollback of the deployment). func (c *ControlPlane) DeleteMCPTool(toolID string) error { if !toolIDPattern.MatchString(strings.TrimSpace(toolID)) { return errors.New("工具标识无效") } agentIDs, err := c.assignedAgentIDs(toolID) if err != nil { return err } transaction, err := c.db.Begin() if err != nil { return err } defer transaction.Rollback() if _, err := transaction.Exec("DELETE FROM agent_mcp_assignments WHERE tool_id = ?", toolID); err != nil { return fmt.Errorf("无法移除工具分配: %w", err) } result, err := transaction.Exec("DELETE FROM mcp_tools WHERE tool_id = ?", toolID) if err != nil { return fmt.Errorf("无法删除工具: %w", err) } if changed, _ := result.RowsAffected(); changed != 1 { return errors.New("未找到要删除的工具") } if err := transaction.Commit(); err != nil { return fmt.Errorf("无法删除工具: %w", err) } c.recordAudit("mcp.tool_deleted", "controller", "mcp-tool:"+toolID, "completed", "已删除工具定义并回收其分配。") for _, agentID := range agentIDs { c.recordAudit("mcp.deployment_withdrawn", "controller", "agent:"+agentID+":"+toolID, "pending", "已撤销该工具的下发。") go c.pushMCP(agentID, true) } return nil } func (c *ControlPlane) assignedAgentIDs(toolID string) ([]string, error) { rows, err := c.db.Query("SELECT agent_id FROM agent_mcp_assignments WHERE tool_id = ?", toolID) if err != nil { return nil, fmt.Errorf("无法读取工具分配: %w", err) } defer rows.Close() agentIDs := make([]string, 0) for rows.Next() { var agentID string if err := rows.Scan(&agentID); err != nil { return nil, err } agentIDs = append(agentIDs, agentID) } return agentIDs, rows.Err() } // redeployMCPTool bumps the deployment version of every assignment so a changed // definition is not rejected as a stale deployment. func (c *ControlPlane) redeployMCPTool(toolID string) { now := time.Now().UTC().Format(time.RFC3339Nano) if _, err := c.db.Exec( "UPDATE agent_mcp_assignments SET deployment_version = deployment_version + 1, apply_status = 'pending', last_error = '', updated_at = ? WHERE tool_id = ?", now, toolID, ); err != nil { return } agentIDs, err := c.assignedAgentIDs(toolID) if err != nil { return } for _, agentID := range agentIDs { go c.pushMCP(agentID, true) } } // SetAgentMCPTools atomically replaces one Agent's assigned tools. Existing // assignments keep their deployment version so an unchanged tool stays applied. func (c *ControlPlane) SetAgentMCPTools(agentID string, toolIDs []string) error { if strings.TrimSpace(agentID) == "" { return errors.New("Agent 标识无效") } var agentExists int if err := c.db.QueryRow("SELECT COUNT(*) FROM managed_agents WHERE agent_id = ?", agentID).Scan(&agentExists); err != nil || agentExists != 1 { return errors.New("目标 Agent 不存在") } unique := make(map[string]struct{}, len(toolIDs)) for _, toolID := range toolIDs { toolID = strings.TrimSpace(toolID) if toolID == "" { continue } if _, seen := unique[toolID]; seen { continue } var exists int if err := c.db.QueryRow("SELECT COUNT(*) FROM mcp_tools WHERE tool_id = ?", toolID).Scan(&exists); err != nil || exists != 1 { return errors.New("选择了不存在的工具") } unique[toolID] = struct{}{} } transaction, err := c.db.Begin() if err != nil { return err } defer transaction.Rollback() rows, err := transaction.Query("SELECT tool_id FROM agent_mcp_assignments WHERE agent_id = ?", agentID) if err != nil { return fmt.Errorf("无法读取 Agent 工具分配: %w", err) } removed := 0 for rows.Next() { var toolID string if rows.Scan(&toolID) != nil { continue } if _, kept := unique[toolID]; kept { continue } if _, err := transaction.Exec("DELETE FROM agent_mcp_assignments WHERE agent_id = ? AND tool_id = ?", agentID, toolID); err != nil { rows.Close() return fmt.Errorf("无法更新 Agent 工具分配: %w", err) } removed++ } rows.Close() now := time.Now().UTC().Format(time.RFC3339Nano) for toolID := range unique { if _, err := transaction.Exec( "INSERT INTO agent_mcp_assignments(agent_id, tool_id, deployment_version, apply_status, last_error, updated_at) VALUES (?, ?, 1, 'pending', '', ?) ON CONFLICT(agent_id, tool_id) DO UPDATE SET apply_status = CASE WHEN agent_mcp_assignments.apply_status = 'applied' THEN 'applied' ELSE 'pending' END, last_error = '', updated_at = excluded.updated_at", agentID, toolID, now, ); err != nil { return fmt.Errorf("无法保存 Agent 工具分配: %w", err) } } if err := transaction.Commit(); err != nil { return fmt.Errorf("无法保存 Agent 工具分配: %w", err) } c.recordAudit("mcp.selection_updated", "controller", "agent:"+agentID, "pending", fmt.Sprintf("已更新为 %d 个工具。", len(unique))) go c.pushMCP(agentID, removed > 0) return nil } // AgentMCPStatuses lists every Agent/tool deployment state for the tools page. func (c *ControlPlane) AgentMCPStatuses() ([]AgentMCPStatus, error) { rows, err := c.db.Query(`SELECT a.agent_id, a.name, t.tool_id, t.name, t.kind, t.endpoint, t.command, t.version, x.deployment_version, x.apply_status, x.last_error, x.updated_at FROM agent_mcp_assignments x JOIN managed_agents a ON a.agent_id = x.agent_id JOIN mcp_tools t ON t.tool_id = x.tool_id ORDER BY a.name COLLATE NOCASE, t.name COLLATE NOCASE`) if err != nil { return nil, fmt.Errorf("无法读取工具下发状态: %w", err) } defer rows.Close() statuses := make([]AgentMCPStatus, 0) for rows.Next() { var status AgentMCPStatus if err := rows.Scan(&status.AgentID, &status.AgentName, &status.ToolID, &status.ToolName, &status.Kind, &status.Endpoint, &status.Command, &status.Version, &status.DeployVersion, &status.ApplyStatus, &status.LastError, &status.UpdatedAt); err != nil { return nil, err } status.LastError = sanitizeDisplayText(status.LastError) statuses = append(statuses, status) } return statuses, rows.Err() } // pushMCP sends the Agent's complete tool set. An empty set is only transmitted // when the controller withdrew a tool, so a freshly paired Agent stays silent. func (c *ControlPlane) pushMCP(agentID string, force bool) { client := c.clientFor(agentID) if client == nil { return } rows, err := c.db.Query(`SELECT t.tool_id, t.name, t.description, t.kind, t.endpoint, t.command, t.args, t.version, x.deployment_version FROM agent_mcp_assignments x JOIN mcp_tools t ON t.tool_id = x.tool_id WHERE x.agent_id = ? ORDER BY t.tool_id`, agentID) if err != nil { return } defer rows.Close() tools := make([]map[string]any, 0) versions := make(map[string]int) for rows.Next() { var toolID, name, description, kind, endpoint, command, encodedArgs, version string var deploymentVersion int if rows.Scan(&toolID, &name, &description, &kind, &endpoint, &command, &encodedArgs, &version, &deploymentVersion) != nil { return } tool := map[string]any{ "id": toolID, "name": name, "description": description, "kind": kind, "version": version, "deployment_version": deploymentVersion, "enabled": true, } switch kind { case toolKindHTTP: tool["endpoint"] = endpoint case toolKindStdio: tool["command"] = command tool["args"] = decodeToolArgs(encodedArgs) } if kind == toolKindBuiltin && toolID == "desktop-filesystem" { // Kept so an existing Agent build keeps accepting this deployment. tool["scope"] = "desktop-directories" } tools = append(tools, tool) versions[toolID] = deploymentVersion } if rows.Err() != nil || len(tools) == 0 && !force { return } if err := client.send(map[string]any{"type": "mcp_snapshot", "schema_version": protocolVersion, "tools": tools}); err != nil { return } now := time.Now().UTC().Format(time.RFC3339Nano) for toolID, version := range versions { // Every replay awaits a fresh acknowledgement, so an Agent that dropped a // tool while offline cannot stay silently "applied" on the controller. _, _ = c.db.Exec("UPDATE agent_mcp_assignments SET apply_status='sent', updated_at=? WHERE agent_id=? AND tool_id=? AND deployment_version=?", now, agentID, toolID, version) } if len(tools) == 0 { c.recordAudit("mcp.snapshot_sent", "controller", "agent:"+agentID, "sent", "已下发空的工具集合,Agent 会移除全部工具。") return } c.recordAudit("mcp.snapshot_sent", "controller", "agent:"+agentID, "sent", fmt.Sprintf("已通过认证 WSS 下发 %d 个工具。", len(tools))) } func (c *ControlPlane) applyMCPAck(agentID string, message map[string]any) { toolID, _ := message["tool_id"].(string) version, _ := message["deployment_version"].(float64) status, _ := message["status"].(string) errorText, _ := message["error"].(string) if !toolIDPattern.MatchString(toolID) || version < 1 || (status != "applied" && status != "failed") { return } result, _ := c.db.Exec("UPDATE agent_mcp_assignments SET apply_status=?, last_error=?, updated_at=? WHERE agent_id=? AND tool_id=? AND deployment_version=?", status, limit(sanitizeDisplayText(errorText), 1000), time.Now().UTC().Format(time.RFC3339Nano), agentID, toolID, int(version)) if changed, _ := result.RowsAffected(); changed == 0 { return // Stale acknowledgement from an earlier deployment; preserve newer state. } c.recordAudit("mcp.deployment_ack", "agent:"+agentID, "agent:"+agentID+":"+toolID, status, errorText) } // AgentDeploymentStatuses provides the Agent page with consolidated deployment state. func (c *ControlPlane) AgentDeploymentStatuses() ([]AgentDeploymentStatus, error) { rows, err := c.db.Query(`SELECT a.agent_id, COALESCE(p.name, ''), COALESCE(l.apply_status, '未部署'), COALESCE((SELECT COUNT(*) FROM agent_skill_assignments s WHERE s.agent_id = a.agent_id), 0), COALESCE((SELECT COUNT(*) FROM agent_skill_assignments s WHERE s.agent_id = a.agent_id AND s.apply_status = 'applied'), 0), COALESCE((SELECT COUNT(*) FROM agent_skill_assignments s WHERE s.agent_id = a.agent_id AND s.apply_status = 'failed'), 0), ` + mcpStatusExpression("a.agent_id") + `, ` + mcpErrorExpression("a.agent_id") + `, COALESCE(MAX(COALESCE(l.updated_at, ''), COALESCE((SELECT MAX(updated_at) FROM agent_mcp_assignments m WHERE m.agent_id = a.agent_id), '')), '') FROM managed_agents a LEFT JOIN agent_llm_assignments l ON l.agent_id = a.agent_id LEFT JOIN llm_profiles p ON p.profile_id = l.profile_id GROUP BY a.agent_id, p.name, l.apply_status`) if err != nil { return nil, fmt.Errorf("无法读取 Agent 部署状态: %w", err) } defer rows.Close() statuses := make([]AgentDeploymentStatus, 0) for rows.Next() { var status AgentDeploymentStatus if err := rows.Scan(&status.AgentID, &status.LLMProfile, &status.LLMStatus, &status.SkillTotal, &status.SkillApplied, &status.SkillFailed, &status.MCPStatus, &status.MCPError, &status.LastUpdatedAt); err != nil { return nil, err } status.MCPError = sanitizeDisplayText(status.MCPError) statuses = append(statuses, status) } return statuses, rows.Err() } // mcpStatusExpression aggregates every assigned tool into one worst-case state so // a single failed tool cannot be hidden behind applied ones. func mcpStatusExpression(agentColumn string) string { return `COALESCE(NULLIF((SELECT CASE WHEN SUM(CASE WHEN apply_status = 'failed' THEN 1 ELSE 0 END) > 0 THEN 'failed' WHEN SUM(CASE WHEN apply_status = 'pending' THEN 1 ELSE 0 END) > 0 THEN 'pending' WHEN SUM(CASE WHEN apply_status = 'sent' THEN 1 ELSE 0 END) > 0 THEN 'sent' WHEN COUNT(*) > 0 THEN 'applied' ELSE '' END FROM agent_mcp_assignments m WHERE m.agent_id = ` + agentColumn + `), ''), '未启用')` } func mcpErrorExpression(agentColumn string) string { return `COALESCE((SELECT COALESCE(MAX(CASE WHEN apply_status = 'failed' THEN last_error ELSE '' END), '') FROM agent_mcp_assignments m WHERE m.agent_id = ` + agentColumn + `), '')` }