Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 18 additions & 3 deletions go/core/internal/a2a/agent_client_registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -39,13 +40,27 @@ func (r *AgentClientRegistry) Register(namespace, name string, c *a2aclient.Clie
}

// 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
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
}
r.mu.RLock()
c, ok := r.clients[key]
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)
}
}
72 changes: 72 additions & 0 deletions go/core/internal/a2a/agent_client_registry_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
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"
"k8s.io/apimachinery/pkg/runtime/schema"
)

func TestAgentClientRegistrySendMessageRoutesByGroupKind(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()
sandboxGroupKind := schema.GroupKind{Group: "kagent.dev", Kind: "SandboxAgent"}.String()
registry.set(routeKey(true, "default", "sandbox-agent"), client)

_, err = registry.SendMessage(
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(),
schema.GroupKind{Group: "kagent.dev", Kind: "Agent"}.String(),
"default",
"sandbox-agent",
&a2atype.SendMessageRequest{Message: a2atype.NewMessage(a2atype.MessageRoleUser, a2atype.NewTextPart("hello"))},
)
require.Error(t, err)
}
101 changes: 78 additions & 23 deletions go/core/internal/mcp/mcp_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,13 @@ 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"
"k8s.io/apimachinery/pkg/runtime/schema"
"sigs.k8s.io/controller-runtime/pkg/client"
ctrllog "sigs.k8s.io/controller-runtime/pkg/log"
)
Expand All @@ -37,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"`
}
Expand Down Expand Up @@ -129,33 +135,72 @@ 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(),
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 {
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
Expand Down Expand Up @@ -185,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)
Expand Down Expand Up @@ -245,16 +295,20 @@ 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 != "" {
message.ContextID = input.ContextID
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.SendMessage(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)},
Expand Down Expand Up @@ -292,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
Expand Down
63 changes: 60 additions & 3 deletions go/core/internal/mcp/mcp_handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,13 @@ 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"
"k8s.io/apimachinery/pkg/runtime/schema"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)

Expand Down Expand Up @@ -82,6 +85,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", 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)
}

// a2aBackend is a fake A2A server that records whether it was called.
type a2aBackend struct {
server *httptest.Server
Expand Down Expand Up @@ -126,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{
{
Expand All @@ -136,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
Expand Down