You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 

531 lines
21 KiB

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 + `), '')`
}