From d49340494d21037fd898d41ddc4e6b7d05e7e31c Mon Sep 17 00:00:00 2001 From: lubingtan Date: Fri, 10 Jul 2026 22:08:39 +0800 Subject: [PATCH 1/3] fix(mcp): include sandbox agents in agent listing Signed-off-by: lubingtan --- go/core/internal/a2a/agent_client_registry.go | 3 + .../a2a/agent_client_registry_test.go | 60 +++++++++++++++++ go/core/internal/mcp/mcp_handler.go | 67 ++++++++++++++----- go/core/internal/mcp/mcp_handler_test.go | 50 ++++++++++++++ 4 files changed, 162 insertions(+), 18 deletions(-) create mode 100644 go/core/internal/a2a/agent_client_registry_test.go diff --git a/go/core/internal/a2a/agent_client_registry.go b/go/core/internal/a2a/agent_client_registry.go index e4b799e36..ef424e3fd 100644 --- a/go/core/internal/a2a/agent_client_registry.go +++ b/go/core/internal/a2a/agent_client_registry.go @@ -43,6 +43,9 @@ func (r *AgentClientRegistry) SendMessage(ctx context.Context, namespace, name s key := namespace + "/" + name r.mu.RLock() c, ok := r.clients[key] + if !ok { + c, ok = r.clients[routeKey(true, namespace, name)] + } r.mu.RUnlock() if !ok { return nil, fmt.Errorf("agent %s/%s not found or not ready", namespace, name) diff --git a/go/core/internal/a2a/agent_client_registry_test.go b/go/core/internal/a2a/agent_client_registry_test.go new file mode 100644 index 000000000..6bb9fc864 --- /dev/null +++ b/go/core/internal/a2a/agent_client_registry_test.go @@ -0,0 +1,60 @@ +package a2a + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + + a2atype "github.com/a2aproject/a2a-go/v2/a2a" + a2aclient "github.com/a2aproject/a2a-go/v2/a2aclient" + "github.com/stretchr/testify/require" +) + +func TestAgentClientRegistrySendMessageFallsBackToSandboxRoute(t *testing.T) { + var called atomic.Bool + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + called.Store(true) + var rpcReq map[string]any + require.NoError(t, json.NewDecoder(r.Body).Decode(&rpcReq)) + resp := map[string]any{ + "jsonrpc": "2.0", + "id": rpcReq["id"], + "result": map[string]any{ + "message": map[string]any{ + "messageId": "test-msg", + "role": "ROLE_AGENT", + "parts": []any{map[string]any{"text": "hello from sandbox"}}, + }, + }, + } + w.Header().Set("Content-Type", "application/json") + require.NoError(t, json.NewEncoder(w).Encode(resp)) + })) + t.Cleanup(server.Close) + + client, err := a2aclient.NewFromEndpoints( + context.Background(), + []*a2atype.AgentInterface{{ + URL: server.URL, + ProtocolBinding: a2atype.TransportProtocolJSONRPC, + ProtocolVersion: a2atype.Version, + }}, + a2aclient.WithJSONRPCTransport(&http.Client{}), + ) + require.NoError(t, err) + + registry := NewAgentClientRegistry() + registry.set(routeKey(true, "default", "sandbox-agent"), client) + + _, err = registry.SendMessage( + context.Background(), + "default", + "sandbox-agent", + &a2atype.SendMessageRequest{Message: a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart("hello"))}, + ) + require.NoError(t, err) + require.True(t, called.Load()) +} diff --git a/go/core/internal/mcp/mcp_handler.go b/go/core/internal/mcp/mcp_handler.go index 940482803..4263662d4 100644 --- a/go/core/internal/mcp/mcp_handler.go +++ b/go/core/internal/mcp/mcp_handler.go @@ -11,10 +11,12 @@ import ( "github.com/google/jsonschema-go/jsonschema" "github.com/kagent-dev/kagent/go/api/v1alpha2" "github.com/kagent-dev/kagent/go/core/internal/a2a" + "github.com/kagent-dev/kagent/go/core/internal/controller/reconciler" "github.com/kagent-dev/kagent/go/core/internal/version" "github.com/kagent-dev/kagent/go/core/pkg/auth" "github.com/kagent-dev/kagent/go/core/pkg/env" mcpsdk "github.com/modelcontextprotocol/go-sdk/mcp" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "sigs.k8s.io/controller-runtime/pkg/client" ctrllog "sigs.k8s.io/controller-runtime/pkg/log" ) @@ -129,33 +131,62 @@ func NewMCPHandler(kubeClient client.Client, agentClients *a2a.AgentClientRegist return handler, nil } -// listReadyAgents returns agents that are accepted and deployment-ready. +// listReadyAgents returns agents that are accepted and workload-ready. func (h *MCPHandler) listReadyAgents(ctx context.Context) ([]AgentSummary, error) { agentList := &v1alpha2.AgentList{} if err := h.kubeClient.List(ctx, agentList); err != nil { return nil, err } - agents := make([]AgentSummary, 0, len(agentList.Items)) - for _, agent := range agentList.Items { - deploymentReady := false - accepted := false - for _, condition := range agent.Status.Conditions { - if condition.Type == "Ready" && condition.Reason == "DeploymentReady" && condition.Status == "True" { - deploymentReady = true - } - if condition.Type == "Accepted" && condition.Status == "True" { - accepted = true + sandboxAgentList := &v1alpha2.SandboxAgentList{} + if err := h.kubeClient.List(ctx, sandboxAgentList); err != nil { + return nil, err + } + + agents := make([]AgentSummary, 0, len(agentList.Items)+len(sandboxAgentList.Items)) + for i := range agentList.Items { + agents = appendReadyAgentSummary(agents, &agentList.Items[i]) + } + for i := range sandboxAgentList.Items { + agents = appendReadyAgentSummary(agents, &sandboxAgentList.Items[i]) + } + return agents, nil +} + +func appendReadyAgentSummary(agents []AgentSummary, agent v1alpha2.AgentObject) []AgentSummary { + if !isReadyForMCP(agent) { + return agents + } + spec := agent.GetAgentSpec() + description := "" + if spec != nil { + description = spec.Description + } + return append(agents, AgentSummary{ + Ref: agent.GetNamespace() + "/" + agent.GetName(), + Description: description, + }) +} + +func isReadyForMCP(agent v1alpha2.AgentObject) bool { + status := agent.GetAgentStatus() + if status == nil { + return false + } + + workloadReady := false + accepted := false + for _, condition := range status.Conditions { + if condition.Type == v1alpha2.AgentConditionTypeReady && condition.Status == metav1.ConditionTrue { + switch condition.Reason { + case reconciler.AgentReadyReasonDeploymentReady, reconciler.AgentReadyReasonWorkloadReady: + workloadReady = true } } - if !accepted || !deploymentReady { - continue + if condition.Type == v1alpha2.AgentConditionTypeAccepted && condition.Status == metav1.ConditionTrue { + accepted = true } - agents = append(agents, AgentSummary{ - Ref: agent.Namespace + "/" + agent.Name, - Description: agent.Spec.Description, - }) } - return agents, nil + return accepted && workloadReady } // handleListAgents handles the list_agents MCP tool diff --git a/go/core/internal/mcp/mcp_handler_test.go b/go/core/internal/mcp/mcp_handler_test.go index d274461f3..78fede209 100644 --- a/go/core/internal/mcp/mcp_handler_test.go +++ b/go/core/internal/mcp/mcp_handler_test.go @@ -13,9 +13,11 @@ import ( a2aclient "github.com/a2aproject/a2a-go/v2/a2aclient" "github.com/kagent-dev/kagent/go/api/v1alpha2" "github.com/kagent-dev/kagent/go/core/internal/a2a" + "github.com/kagent-dev/kagent/go/core/internal/controller/reconciler" mcpsdk "github.com/modelcontextprotocol/go-sdk/mcp" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "sigs.k8s.io/controller-runtime/pkg/client/fake" ) @@ -82,6 +84,54 @@ func TestListAgentsInputSchemaHasProperties(t *testing.T) { require.Equal(t, false, schema["additionalProperties"], "additionalProperties must remain false") } +func TestListReadyAgentsIncludesSandboxAgents(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, v1alpha2.AddToScheme(scheme)) + + regularAgent := &v1alpha2.Agent{ + ObjectMeta: metav1.ObjectMeta{Name: "regular-agent", Namespace: "default"}, + Spec: v1alpha2.AgentSpec{ + Description: "regular", + }, + Status: v1alpha2.AgentStatus{ + Conditions: []metav1.Condition{ + {Type: v1alpha2.AgentConditionTypeAccepted, Status: metav1.ConditionTrue}, + {Type: v1alpha2.AgentConditionTypeReady, Status: metav1.ConditionTrue, Reason: reconciler.AgentReadyReasonDeploymentReady}, + }, + }, + } + sandboxAgent := &v1alpha2.SandboxAgent{ + ObjectMeta: metav1.ObjectMeta{Name: "sandbox-agent", Namespace: "default"}, + Spec: v1alpha2.SandboxAgentSpec{ + AgentSpec: v1alpha2.AgentSpec{ + Description: "sandbox", + }, + }, + Status: v1alpha2.AgentStatus{ + Conditions: []metav1.Condition{ + {Type: v1alpha2.AgentConditionTypeAccepted, Status: metav1.ConditionTrue}, + {Type: v1alpha2.AgentConditionTypeReady, Status: metav1.ConditionTrue, Reason: reconciler.AgentReadyReasonWorkloadReady}, + }, + }, + } + kubeClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithStatusSubresource(&v1alpha2.Agent{}, &v1alpha2.SandboxAgent{}). + WithObjects(regularAgent, sandboxAgent). + Build() + + h, err := NewMCPHandler(kubeClient, nil, nil) + require.NoError(t, err) + + agents, err := h.listReadyAgents(context.Background()) + require.NoError(t, err) + + assert.ElementsMatch(t, []AgentSummary{ + {Ref: "default/regular-agent", Description: "regular"}, + {Ref: "default/sandbox-agent", Description: "sandbox"}, + }, agents) +} + // a2aBackend is a fake A2A server that records whether it was called. type a2aBackend struct { server *httptest.Server From 74df43128e446a1195924500c528c9c2dca57e0e Mon Sep 17 00:00:00 2001 From: lubingtan Date: Fri, 10 Jul 2026 22:08:39 +0800 Subject: [PATCH 2/3] fix(mcp): route agent invocations by group kind Signed-off-by: lubingtan --- go/core/internal/a2a/agent_client_registry.go | 37 ++++++++++-- .../a2a/agent_client_registry_test.go | 17 +++++- go/core/internal/mcp/mcp_handler.go | 34 +++++++++-- go/core/internal/mcp/mcp_handler_test.go | 57 +++++++++++++++++-- 4 files changed, 127 insertions(+), 18 deletions(-) diff --git a/go/core/internal/a2a/agent_client_registry.go b/go/core/internal/a2a/agent_client_registry.go index ef424e3fd..effbbc132 100644 --- a/go/core/internal/a2a/agent_client_registry.go +++ b/go/core/internal/a2a/agent_client_registry.go @@ -7,6 +7,7 @@ import ( a2atype "github.com/a2aproject/a2a-go/v2/a2a" a2aclient "github.com/a2aproject/a2a-go/v2/a2aclient" + "k8s.io/apimachinery/pkg/runtime/schema" ) // AgentClientRegistry maps agent route keys to their A2A clients. @@ -38,17 +39,43 @@ func (r *AgentClientRegistry) Register(namespace, name string, c *a2aclient.Clie r.set(namespace+"/"+name, c) } +// RegisterForGroupKind adds or replaces the A2A client for the given agent GroupKind. +func (r *AgentClientRegistry) RegisterForGroupKind(groupKind, namespace, name string, c *a2aclient.Client) error { + key, err := routeKeyForGroupKind(groupKind, namespace, name) + if err != nil { + return err + } + r.set(key, c) + return nil +} + // SendMessage invokes an agent directly via its cached A2A client. func (r *AgentClientRegistry) SendMessage(ctx context.Context, namespace, name string, req *a2atype.SendMessageRequest) (a2atype.SendMessageResult, error) { - key := namespace + "/" + name + return r.SendMessageForGroupKind(ctx, schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String(), namespace, name, req) +} + +// SendMessageForGroupKind invokes an agent directly via its cached A2A client. +func (r *AgentClientRegistry) SendMessageForGroupKind(ctx context.Context, groupKind, namespace, name string, req *a2atype.SendMessageRequest) (a2atype.SendMessageResult, error) { + key, err := routeKeyForGroupKind(groupKind, namespace, name) + if err != nil { + return nil, err + } r.mu.RLock() c, ok := r.clients[key] - if !ok { - c, ok = r.clients[routeKey(true, namespace, name)] - } r.mu.RUnlock() if !ok { - return nil, fmt.Errorf("agent %s/%s not found or not ready", namespace, name) + return nil, fmt.Errorf("agent %s %s/%s not found or not ready", groupKind, namespace, name) } return c.SendMessage(ctx, req) } + +func routeKeyForGroupKind(groupKind, namespace, name string) (string, error) { + switch groupKind { + case "", schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String(), "Agent": + return routeKey(false, namespace, name), nil + case schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String(), "SandboxAgent": + return routeKey(true, namespace, name), nil + default: + return "", fmt.Errorf("unsupported agent groupKind %q", groupKind) + } +} diff --git a/go/core/internal/a2a/agent_client_registry_test.go b/go/core/internal/a2a/agent_client_registry_test.go index 6bb9fc864..181dde791 100644 --- a/go/core/internal/a2a/agent_client_registry_test.go +++ b/go/core/internal/a2a/agent_client_registry_test.go @@ -11,9 +11,10 @@ import ( a2atype "github.com/a2aproject/a2a-go/v2/a2a" a2aclient "github.com/a2aproject/a2a-go/v2/a2aclient" "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/runtime/schema" ) -func TestAgentClientRegistrySendMessageFallsBackToSandboxRoute(t *testing.T) { +func TestAgentClientRegistrySendMessageRoutesByGroupKind(t *testing.T) { var called atomic.Bool server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { called.Store(true) @@ -47,14 +48,24 @@ func TestAgentClientRegistrySendMessageFallsBackToSandboxRoute(t *testing.T) { require.NoError(t, err) registry := NewAgentClientRegistry() - registry.set(routeKey(true, "default", "sandbox-agent"), client) + sandboxGroupKind := schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String() + require.NoError(t, registry.RegisterForGroupKind(sandboxGroupKind, "default", "sandbox-agent", client)) - _, err = registry.SendMessage( + _, err = registry.SendMessageForGroupKind( context.Background(), + sandboxGroupKind, "default", "sandbox-agent", &a2atype.SendMessageRequest{Message: a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart("hello"))}, ) require.NoError(t, err) require.True(t, called.Load()) + + _, err = registry.SendMessage( + context.Background(), + "default", + "sandbox-agent", + &a2atype.SendMessageRequest{Message: a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart("hello"))}, + ) + require.Error(t, err) } diff --git a/go/core/internal/mcp/mcp_handler.go b/go/core/internal/mcp/mcp_handler.go index 4263662d4..014c7152d 100644 --- a/go/core/internal/mcp/mcp_handler.go +++ b/go/core/internal/mcp/mcp_handler.go @@ -17,6 +17,7 @@ import ( "github.com/kagent-dev/kagent/go/core/pkg/env" mcpsdk "github.com/modelcontextprotocol/go-sdk/mcp" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" "sigs.k8s.io/controller-runtime/pkg/client" ctrllog "sigs.k8s.io/controller-runtime/pkg/log" ) @@ -39,17 +40,20 @@ type ListAgentsOutput struct { type AgentSummary struct { Ref string `json:"ref"` + GroupKind string `json:"groupKind"` Description string `json:"description,omitempty"` } type InvokeAgentInput struct { Agent string `json:"agent" jsonschema:"Agent reference in format namespace/name. To find a list of available sources, use the 'agents' resource."` + GroupKind string `json:"groupKind,omitempty" jsonschema:"Optional Kubernetes GroupKind from list_agents, such as Agent.kagent.dev or SandboxAgent.kagent.dev. Defaults to Agent.kagent.dev."` Task string `json:"task" jsonschema:"Task to run"` ContextID string `json:"context_id,omitempty" jsonschema:"Optional A2A context ID to continue a conversation"` } type InvokeAgentOutput struct { Agent string `json:"agent"` + GroupKind string `json:"groupKind,omitempty"` Text string `json:"text"` ContextID string `json:"context_id,omitempty"` } @@ -163,10 +167,20 @@ func appendReadyAgentSummary(agents []AgentSummary, agent v1alpha2.AgentObject) } return append(agents, AgentSummary{ Ref: agent.GetNamespace() + "/" + agent.GetName(), + GroupKind: groupKindForAgent(agent), Description: description, }) } +func groupKindForAgent(agent v1alpha2.AgentObject) string { + switch agent.(type) { + case *v1alpha2.SandboxAgent: + return schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String() + default: + return schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String() + } +} + func isReadyForMCP(agent v1alpha2.AgentObject) bool { status := agent.GetAgentStatus() if status == nil { @@ -216,6 +230,11 @@ func (h *MCPHandler) handleListAgents(ctx context.Context, req *mcpsdk.CallToolR fallbackText.WriteByte('\n') } fallbackText.WriteString(agent.Ref) + if agent.GroupKind != "" { + fallbackText.WriteString(" (") + fallbackText.WriteString(agent.GroupKind) + fallbackText.WriteString(")") + } if agent.Description != "" { fallbackText.WriteString(" - ") fallbackText.WriteString(agent.Description) @@ -276,6 +295,10 @@ func (h *MCPHandler) handleInvokeAgent(ctx context.Context, req *mcpsdk.CallTool } agentNS, agentName := parts[0], parts[1] agentRef := agentNS + "/" + agentName + groupKind := input.GroupKind + if groupKind == "" { + groupKind = schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String() + } message := a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart(input.Task)) if input.ContextID != "" { @@ -283,9 +306,9 @@ func (h *MCPHandler) handleInvokeAgent(ctx context.Context, req *mcpsdk.CallTool log.V(1).Info("Using context_id from client request", "context_id", input.ContextID) } - result, err := h.agentClients.SendMessage(ctx, agentNS, agentName, &a2atype.SendMessageRequest{Message: message}) + result, err := h.agentClients.SendMessageForGroupKind(ctx, groupKind, agentNS, agentName, &a2atype.SendMessageRequest{Message: message}) if err != nil { - log.Error(err, "Failed to send A2A message", "agent", agentRef) + log.Error(err, "Failed to send A2A message", "agent", agentRef, "groupKind", groupKind) return &mcpsdk.CallToolResult{ Content: []mcpsdk.Content{ &mcpsdk.TextContent{Text: fmt.Sprintf("Failed to send A2A message: %v", err)}, @@ -323,12 +346,13 @@ func (h *MCPHandler) handleInvokeAgent(ctx context.Context, req *mcpsdk.CallTool responseText = string(raw) } - log.Info("Invoked agent", "agent", agentRef, "hasContextID", newContextID != "") + log.Info("Invoked agent", "agent", agentRef, "groupKind", groupKind, "hasContextID", newContextID != "") // Return context_id in response so client can store it for stateless operation output := InvokeAgentOutput{ - Agent: agentRef, - Text: responseText, + Agent: agentRef, + GroupKind: groupKind, + Text: responseText, } if newContextID != "" { output.ContextID = newContextID diff --git a/go/core/internal/mcp/mcp_handler_test.go b/go/core/internal/mcp/mcp_handler_test.go index 78fede209..e48bb82c1 100644 --- a/go/core/internal/mcp/mcp_handler_test.go +++ b/go/core/internal/mcp/mcp_handler_test.go @@ -19,6 +19,7 @@ import ( "github.com/stretchr/testify/require" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" "sigs.k8s.io/controller-runtime/pkg/client/fake" ) @@ -127,8 +128,8 @@ func TestListReadyAgentsIncludesSandboxAgents(t *testing.T) { require.NoError(t, err) assert.ElementsMatch(t, []AgentSummary{ - {Ref: "default/regular-agent", Description: "regular"}, - {Ref: "default/sandbox-agent", Description: "sandbox"}, + {Ref: "default/regular-agent", GroupKind: schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String(), Description: "regular"}, + {Ref: "default/sandbox-agent", GroupKind: schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String(), Description: "sandbox"}, }, agents) } @@ -176,6 +177,14 @@ func newA2ABackend(t *testing.T) *a2aBackend { // newTestRegistry builds an AgentClientRegistry with a single agent pre-registered. func newTestRegistry(t *testing.T, namespace, name, backendURL string) *a2a.AgentClientRegistry { + t.Helper() + c := newTestA2AClient(t, namespace, name, backendURL) + registry := a2a.NewAgentClientRegistry() + registry.Register(namespace, name, c) + return registry +} + +func newTestA2AClient(t *testing.T, namespace, name, backendURL string) *a2aclient.Client { t.Helper() interfaces := []*a2atype.AgentInterface{ { @@ -186,9 +195,7 @@ func newTestRegistry(t *testing.T, namespace, name, backendURL string) *a2a.Agen } c, err := a2aclient.NewFromEndpoints(context.Background(), interfaces, a2aclient.WithJSONRPCTransport(&http.Client{})) require.NoError(t, err) - registry := a2a.NewAgentClientRegistry() - registry.Register(namespace, name, c) - return registry + return c } // TestInvokeAgent_InvalidAgentRef verifies that invoke_agent returns a tool @@ -284,3 +291,43 @@ func TestInvokeAgent_RoutesViaRegistry(t *testing.T) { assert.False(t, result.IsError, "unexpected tool error: %v", result.Content) assert.True(t, backend.wasCalled(), "A2A backend should have received the forwarded request") } + +func TestInvokeAgent_RoutesSandboxAgentByGroupKind(t *testing.T) { + regularBackend := newA2ABackend(t) + sandboxBackend := newA2ABackend(t) + + registry := newTestRegistry(t, "default", "shared-name", regularBackend.server.URL) + sandboxClient := newTestA2AClient(t, "default", "shared-name", sandboxBackend.server.URL) + sandboxGroupKind := schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String() + require.NoError(t, registry.RegisterForGroupKind(sandboxGroupKind, "default", "shared-name", sandboxClient)) + + mcpHandler, err := NewMCPHandler(nil, registry, nil) + require.NoError(t, err) + + mcpServer := httptest.NewServer(mcpHandler) + t.Cleanup(mcpServer.Close) + + transport := &mcpsdk.StreamableClientTransport{ + Endpoint: mcpServer.URL, + DisableStandaloneSSE: true, + } + + ctx := context.Background() + cs, err := mcpsdk.NewClient(&mcpsdk.Implementation{Name: "test", Version: "1.0"}, nil). + Connect(ctx, transport, nil) + require.NoError(t, err) + t.Cleanup(func() { cs.Close() }) + + result, err := cs.CallTool(ctx, &mcpsdk.CallToolParams{ + Name: "invoke_agent", + Arguments: map[string]any{ + "agent": "default/shared-name", + "groupKind": sandboxGroupKind, + "task": "say hello", + }, + }) + require.NoError(t, err) + assert.False(t, result.IsError, "unexpected tool error: %v", result.Content) + assert.False(t, regularBackend.wasCalled(), "regular Agent backend should not receive SandboxAgent invocation") + assert.True(t, sandboxBackend.wasCalled(), "SandboxAgent backend should receive the invocation") +} From a6e6be8d3cbcbbf771ba9205cf4140fca2771e4f Mon Sep 17 00:00:00 2001 From: lubingtan Date: Wed, 15 Jul 2026 13:20:13 +0800 Subject: [PATCH 3/3] fix(mcp): disambiguate agent invocation by group kind Signed-off-by: lubingtan --- go/core/internal/a2a/agent_client_registry.go | 17 +------- .../a2a/agent_client_registry_test.go | 5 ++- go/core/internal/mcp/mcp_handler.go | 2 +- go/core/internal/mcp/mcp_handler_test.go | 40 ------------------- 4 files changed, 5 insertions(+), 59 deletions(-) diff --git a/go/core/internal/a2a/agent_client_registry.go b/go/core/internal/a2a/agent_client_registry.go index effbbc132..f074c15d2 100644 --- a/go/core/internal/a2a/agent_client_registry.go +++ b/go/core/internal/a2a/agent_client_registry.go @@ -39,23 +39,8 @@ func (r *AgentClientRegistry) Register(namespace, name string, c *a2aclient.Clie r.set(namespace+"/"+name, c) } -// RegisterForGroupKind adds or replaces the A2A client for the given agent GroupKind. -func (r *AgentClientRegistry) RegisterForGroupKind(groupKind, namespace, name string, c *a2aclient.Client) error { - key, err := routeKeyForGroupKind(groupKind, namespace, name) - if err != nil { - return err - } - r.set(key, c) - return nil -} - // SendMessage invokes an agent directly via its cached A2A client. -func (r *AgentClientRegistry) SendMessage(ctx context.Context, namespace, name string, req *a2atype.SendMessageRequest) (a2atype.SendMessageResult, error) { - return r.SendMessageForGroupKind(ctx, schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String(), namespace, name, req) -} - -// SendMessageForGroupKind invokes an agent directly via its cached A2A client. -func (r *AgentClientRegistry) SendMessageForGroupKind(ctx context.Context, groupKind, namespace, name string, req *a2atype.SendMessageRequest) (a2atype.SendMessageResult, error) { +func (r *AgentClientRegistry) SendMessage(ctx context.Context, groupKind, namespace, name string, req *a2atype.SendMessageRequest) (a2atype.SendMessageResult, error) { key, err := routeKeyForGroupKind(groupKind, namespace, name) if err != nil { return nil, err diff --git a/go/core/internal/a2a/agent_client_registry_test.go b/go/core/internal/a2a/agent_client_registry_test.go index 181dde791..c576d3b3b 100644 --- a/go/core/internal/a2a/agent_client_registry_test.go +++ b/go/core/internal/a2a/agent_client_registry_test.go @@ -49,9 +49,9 @@ func TestAgentClientRegistrySendMessageRoutesByGroupKind(t *testing.T) { registry := NewAgentClientRegistry() sandboxGroupKind := schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String() - require.NoError(t, registry.RegisterForGroupKind(sandboxGroupKind, "default", "sandbox-agent", client)) + registry.set(routeKey(true, "default", "sandbox-agent"), client) - _, err = registry.SendMessageForGroupKind( + _, err = registry.SendMessage( context.Background(), sandboxGroupKind, "default", @@ -63,6 +63,7 @@ func TestAgentClientRegistrySendMessageRoutesByGroupKind(t *testing.T) { _, err = registry.SendMessage( context.Background(), + schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String(), "default", "sandbox-agent", &a2atype.SendMessageRequest{Message: a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart("hello"))}, diff --git a/go/core/internal/mcp/mcp_handler.go b/go/core/internal/mcp/mcp_handler.go index 014c7152d..7c2b9e87b 100644 --- a/go/core/internal/mcp/mcp_handler.go +++ b/go/core/internal/mcp/mcp_handler.go @@ -306,7 +306,7 @@ func (h *MCPHandler) handleInvokeAgent(ctx context.Context, req *mcpsdk.CallTool log.V(1).Info("Using context_id from client request", "context_id", input.ContextID) } - result, err := h.agentClients.SendMessageForGroupKind(ctx, groupKind, agentNS, agentName, &a2atype.SendMessageRequest{Message: message}) + result, err := h.agentClients.SendMessage(ctx, groupKind, agentNS, agentName, &a2atype.SendMessageRequest{Message: message}) if err != nil { log.Error(err, "Failed to send A2A message", "agent", agentRef, "groupKind", groupKind) return &mcpsdk.CallToolResult{ diff --git a/go/core/internal/mcp/mcp_handler_test.go b/go/core/internal/mcp/mcp_handler_test.go index e48bb82c1..2f6b20462 100644 --- a/go/core/internal/mcp/mcp_handler_test.go +++ b/go/core/internal/mcp/mcp_handler_test.go @@ -291,43 +291,3 @@ func TestInvokeAgent_RoutesViaRegistry(t *testing.T) { assert.False(t, result.IsError, "unexpected tool error: %v", result.Content) assert.True(t, backend.wasCalled(), "A2A backend should have received the forwarded request") } - -func TestInvokeAgent_RoutesSandboxAgentByGroupKind(t *testing.T) { - regularBackend := newA2ABackend(t) - sandboxBackend := newA2ABackend(t) - - registry := newTestRegistry(t, "default", "shared-name", regularBackend.server.URL) - sandboxClient := newTestA2AClient(t, "default", "shared-name", sandboxBackend.server.URL) - sandboxGroupKind := schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String() - require.NoError(t, registry.RegisterForGroupKind(sandboxGroupKind, "default", "shared-name", sandboxClient)) - - mcpHandler, err := NewMCPHandler(nil, registry, nil) - require.NoError(t, err) - - mcpServer := httptest.NewServer(mcpHandler) - t.Cleanup(mcpServer.Close) - - transport := &mcpsdk.StreamableClientTransport{ - Endpoint: mcpServer.URL, - DisableStandaloneSSE: true, - } - - ctx := context.Background() - cs, err := mcpsdk.NewClient(&mcpsdk.Implementation{Name: "test", Version: "1.0"}, nil). - Connect(ctx, transport, nil) - require.NoError(t, err) - t.Cleanup(func() { cs.Close() }) - - result, err := cs.CallTool(ctx, &mcpsdk.CallToolParams{ - Name: "invoke_agent", - Arguments: map[string]any{ - "agent": "default/shared-name", - "groupKind": sandboxGroupKind, - "task": "say hello", - }, - }) - require.NoError(t, err) - assert.False(t, result.IsError, "unexpected tool error: %v", result.Content) - assert.False(t, regularBackend.wasCalled(), "regular Agent backend should not receive SandboxAgent invocation") - assert.True(t, sandboxBackend.wasCalled(), "SandboxAgent backend should receive the invocation") -}