Skip to content

Commit 312dfe6

Browse files
committed
add remotemcpserver back, fixes to make claude mcp work e2e
Signed-off-by: Jet Chiang <pokyuen.jetchiang-ext@solo.io>
1 parent 8c8e1a0 commit 312dfe6

6 files changed

Lines changed: 540 additions & 15 deletions

File tree

go/core/cmd/controller-v2/main.go

Lines changed: 54 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -28,11 +28,17 @@ import (
2828
"syscall"
2929
"time"
3030

31+
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
3132
kagentv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3"
33+
remotemcpcontroller "github.com/kagent-dev/kagent/go/core/internal/controller/remotemcpserver"
3234
"github.com/kagent-dev/kagent/go/core/internal/database"
3335
"github.com/kagent-dev/kagent/go/core/internal/grpcserver"
3436
authimpl "github.com/kagent-dev/kagent/go/core/internal/httpserver/auth"
3537
"github.com/kagent-dev/kagent/go/core/internal/service/kubecrud"
38+
modelservice "github.com/kagent-dev/kagent/go/core/internal/service/model"
39+
prompttemplateservice "github.com/kagent-dev/kagent/go/core/internal/service/prompttemplate"
40+
systemservice "github.com/kagent-dev/kagent/go/core/internal/service/system"
41+
toolservice "github.com/kagent-dev/kagent/go/core/internal/service/tool"
3642
"github.com/kagent-dev/kagent/go/core/pkg/auth"
3743
"github.com/kagent-dev/kagent/go/core/pkg/migrations"
3844
"github.com/kagent-dev/kagent/go/core/v2/a2agateway"
@@ -41,13 +47,17 @@ import (
4147
v2controller "github.com/kagent-dev/kagent/go/core/v2/controller"
4248
v2mcp "github.com/kagent-dev/kagent/go/core/v2/mcp"
4349
"github.com/kagent-dev/kagent/go/core/v2/substrate"
50+
kmcp "github.com/kagent-dev/kmcp/api/v1alpha1"
4451
"go.uber.org/zap/zapcore"
4552
"golang.org/x/sync/errgroup"
53+
corev1 "k8s.io/api/core/v1"
4654
k8sruntime "k8s.io/apimachinery/pkg/runtime"
4755
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
4856
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
4957
"k8s.io/client-go/tools/clientcmd"
5058
ctrl "sigs.k8s.io/controller-runtime"
59+
"sigs.k8s.io/controller-runtime/pkg/cache"
60+
"sigs.k8s.io/controller-runtime/pkg/client"
5161
"sigs.k8s.io/controller-runtime/pkg/log/zap"
5262
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
5363
)
@@ -91,8 +101,21 @@ func main() {
91101
managerScheme := k8sruntime.NewScheme()
92102
utilruntime.Must(clientgoscheme.AddToScheme(managerScheme))
93103
utilruntime.Must(kagentv1alpha3.AddToScheme(managerScheme))
104+
utilruntime.Must(atev1alpha1.AddToScheme(managerScheme))
105+
utilruntime.Must(kmcp.AddToScheme(managerScheme))
106+
watchNamespaces := namespaces(os.Getenv("WATCH_NAMESPACES"))
107+
managerClientOptions := client.Options{}
108+
managerCacheOptions := cache.Options{DefaultNamespaces: namespaceCache(watchNamespaces)}
109+
if len(watchNamespaces) > 0 {
110+
// A namespaced Role cannot list cluster-scoped Namespace objects. Read them
111+
// directly so SystemService can fall back to the configured names on a
112+
// Forbidden response without a failing Namespace informer blocking startup.
113+
managerClientOptions.Cache = &client.CacheOptions{DisableFor: []client.Object{&corev1.Namespace{}}}
114+
}
94115
manager, err := ctrl.NewManager(kubeConfig, ctrl.Options{
95116
Scheme: managerScheme,
117+
Cache: managerCacheOptions,
118+
Client: managerClientOptions,
96119
Metrics: metricsserver.Options{BindAddress: "0"},
97120
LeaderElection: envBool("LEADER_ELECT"),
98121
LeaderElectionID: "0e9f6799.kagent.dev",
@@ -101,7 +124,7 @@ func main() {
101124
if err != nil {
102125
log.Fatalf("create controller manager: %v", err)
103126
}
104-
runtime, err := v2controller.NewRuntime(kubeConfig, namespaces(os.Getenv("WATCH_NAMESPACES")), ctx.Done())
127+
runtime, err := v2controller.NewRuntime(kubeConfig, watchNamespaces, ctx.Done())
105128
if err != nil {
106129
log.Fatal(err)
107130
}
@@ -112,6 +135,11 @@ func main() {
112135
if err := manager.Add(reconciler); err != nil {
113136
log.Fatalf("add reconciler to controller manager: %v", err)
114137
}
138+
mcpClient := toolservice.NewRuntimeMCPClient(manager.GetClient())
139+
remoteMCPDiscovery := remotemcpcontroller.New(manager.GetClient(), mcpClient, store)
140+
if err := remoteMCPDiscovery.SetupWithManager(manager); err != nil {
141+
log.Fatalf("set up RemoteMCPServer discovery: %v", err)
142+
}
115143

116144
actors, err := substrate.Dial(ctx, substrate.Config{
117145
AteAPIEndpoint: env("SUBSTRATE_ATE_API_ENDPOINT", "dns:///api.ate-system.svc:443"),
@@ -126,6 +154,11 @@ func main() {
126154

127155
authenticator := &authimpl.UnsecureAuthenticator{}
128156
authorizer := &authimpl.NoopAuthorizer{}
157+
resourceNamespace := env("KAGENT_NAMESPACE", "kagent")
158+
models := modelservice.NewService(manager.GetClient(), authorizer, resourceNamespace)
159+
tools := toolservice.NewService(manager.GetClient(), store, authorizer, resourceNamespace, mcpClient)
160+
prompts := prompttemplateservice.NewService(manager.GetClient(), authorizer)
161+
system := systemservice.NewService(systemservice.WithInventory(manager.GetClient(), watchNamespaces, authorizer, actors))
129162
instanceWorkflow := agentinstance.NewActorWorkflow(store, actors)
130163
instances := agentinstance.NewService(store, authorizer, instanceWorkflow)
131164
checkpoints := checkpoint.NewService(store, authorizer, actors, instanceWorkflow)
@@ -143,11 +176,15 @@ func main() {
143176
log.Fatal(err)
144177
}
145178
server, err := grpcserver.New(grpcserver.Config{
146-
BindAddress: env("GRPC_BIND_ADDRESS", ":8084"),
147-
Reflection: envBool("GRPC_REFLECTION"),
148-
Authenticator: authenticator,
149-
ShareStore: store,
150-
AgentInstanceService: instances,
179+
BindAddress: env("GRPC_BIND_ADDRESS", ":8084"),
180+
Reflection: envBool("GRPC_REFLECTION"),
181+
Authenticator: authenticator,
182+
ShareStore: store,
183+
ModelService: models,
184+
ToolService: tools,
185+
PromptTemplateService: prompts,
186+
SystemService: system,
187+
AgentInstanceService: instances,
151188
// Both halves of the pair CreateAgentInstance names. Without these two
152189
// the only way to author a Harness or an AgentTemplate is kubectl.
153190
AgentTemplateService: kubecrud.NewService(manager.GetClient(), authorizer, &kagentv1alpha3.AgentTemplate{}, &kagentv1alpha3.AgentTemplateList{}, "AgentTemplate"),
@@ -215,3 +252,14 @@ func namespaces(value string) []string {
215252
}
216253
return result
217254
}
255+
256+
func namespaceCache(names []string) map[string]cache.Config {
257+
if len(names) == 0 {
258+
return nil
259+
}
260+
result := make(map[string]cache.Config, len(names))
261+
for _, name := range names {
262+
result[name] = cache.Config{}
263+
}
264+
return result
265+
}

go/core/cmd/controller-v2/main_test.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,3 +11,19 @@ func TestNamespaces(t *testing.T) {
1111
t.Fatalf("namespaces() = %q, want %q", got, want)
1212
}
1313
}
14+
15+
func TestNamespaceCache(t *testing.T) {
16+
if got := namespaceCache(nil); got != nil {
17+
t.Fatalf("namespaceCache(nil) = %#v, want nil", got)
18+
}
19+
got := namespaceCache([]string{"team-a", "team-b"})
20+
if len(got) != 2 {
21+
t.Fatalf("namespaceCache() = %#v", got)
22+
}
23+
if _, ok := got["team-a"]; !ok {
24+
t.Fatal("namespaceCache() missing team-a")
25+
}
26+
if _, ok := got["team-b"]; !ok {
27+
t.Fatal("namespaceCache() missing team-b")
28+
}
29+
}
Lines changed: 240 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
1+
/*
2+
Copyright 2026.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
// Package remotemcpserver reconciles the discovered tool catalog published in
18+
// RemoteMCPServer status.
19+
package remotemcpserver
20+
21+
import (
22+
"context"
23+
"errors"
24+
"fmt"
25+
"reflect"
26+
"slices"
27+
"strings"
28+
"time"
29+
30+
dbmodel "github.com/kagent-dev/kagent/go/api/database"
31+
"github.com/kagent-dev/kagent/go/api/v1alpha3"
32+
toolservice "github.com/kagent-dev/kagent/go/core/internal/service/tool"
33+
corev1 "k8s.io/api/core/v1"
34+
apierrors "k8s.io/apimachinery/pkg/api/errors"
35+
apiMeta "k8s.io/apimachinery/pkg/api/meta"
36+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
37+
"k8s.io/apimachinery/pkg/types"
38+
ctrl "sigs.k8s.io/controller-runtime"
39+
"sigs.k8s.io/controller-runtime/pkg/builder"
40+
"sigs.k8s.io/controller-runtime/pkg/client"
41+
"sigs.k8s.io/controller-runtime/pkg/handler"
42+
"sigs.k8s.io/controller-runtime/pkg/predicate"
43+
"sigs.k8s.io/controller-runtime/pkg/reconcile"
44+
)
45+
46+
const (
47+
conditionAccepted = "Accepted"
48+
remoteGroupKind = "RemoteMCPServer.kagent.dev"
49+
refreshInterval = 5 * time.Minute
50+
)
51+
52+
// ToolDiscoverer returns the tools currently advertised by one MCP server.
53+
type ToolDiscoverer interface {
54+
ListTools(context.Context, toolservice.MCPServerRef) ([]toolservice.MCPAppTool, error)
55+
}
56+
57+
// CatalogStore keeps the gRPC ToolService catalog aligned with Kubernetes status.
58+
// RemoteMCPServer status remains the source used by harness compilers, while the
59+
// database projection serves list RPCs without making those RPCs perform discovery.
60+
type CatalogStore interface {
61+
StoreToolServer(context.Context, *dbmodel.ToolServer) (*dbmodel.ToolServer, error)
62+
RefreshToolsForServer(context.Context, string, string, ...*v1alpha3.MCPTool) error
63+
DeleteToolsForServer(context.Context, string, string) error
64+
DeleteToolServer(context.Context, string, string) error
65+
}
66+
67+
// Reconciler publishes RemoteMCPServer discovery results to its status.
68+
type Reconciler struct {
69+
client client.Client
70+
discoverer ToolDiscoverer
71+
catalog CatalogStore
72+
}
73+
74+
func New(client client.Client, discoverer ToolDiscoverer, catalog CatalogStore) *Reconciler {
75+
return &Reconciler{client: client, discoverer: discoverer, catalog: catalog}
76+
}
77+
78+
func (r *Reconciler) SetupWithManager(manager ctrl.Manager) error {
79+
return ctrl.NewControllerManagedBy(manager).
80+
For(&v1alpha3.RemoteMCPServer{}, builder.WithPredicates(predicate.GenerationChangedPredicate{})).
81+
Watches(&corev1.Secret{}, handler.EnqueueRequestsFromMapFunc(r.requestsForDependency)).
82+
Watches(&corev1.ConfigMap{}, handler.EnqueueRequestsFromMapFunc(r.requestsForDependency)).
83+
Complete(r)
84+
}
85+
86+
func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) (reconcile.Result, error) {
87+
server := &v1alpha3.RemoteMCPServer{}
88+
if err := r.client.Get(ctx, request.NamespacedName, server); err != nil {
89+
if !apierrors.IsNotFound(err) {
90+
return reconcile.Result{}, err
91+
}
92+
return reconcile.Result{}, r.deleteCatalog(ctx, request.String())
93+
}
94+
95+
tools, err := r.discoverer.ListTools(ctx, toolservice.MCPServerRef{
96+
Ref: request.NamespacedName, GroupKind: remoteGroupKind,
97+
})
98+
if err != nil {
99+
statusErr := r.updateStatus(ctx, server, nil, metav1.ConditionFalse, "DiscoveryFailed", err.Error())
100+
catalogErr := r.updateCatalog(ctx, server, nil, false)
101+
return reconcile.Result{}, errors.Join(
102+
fmt.Errorf("discover RemoteMCPServer tools: %w", err),
103+
wrapError("update RemoteMCPServer discovery failure", statusErr),
104+
wrapError("clear RemoteMCPServer tool catalog", catalogErr),
105+
)
106+
}
107+
108+
discovered, err := normalizeTools(tools)
109+
if err != nil {
110+
statusErr := r.updateStatus(ctx, server, nil, metav1.ConditionFalse, "InvalidDiscovery", err.Error())
111+
catalogErr := r.updateCatalog(ctx, server, nil, false)
112+
return reconcile.Result{}, errors.Join(
113+
err,
114+
wrapError("update invalid RemoteMCPServer discovery", statusErr),
115+
wrapError("clear invalid RemoteMCPServer tool catalog", catalogErr),
116+
)
117+
}
118+
message := fmt.Sprintf("Discovered %d MCP tools", len(discovered))
119+
if err := r.updateStatus(ctx, server, discovered, metav1.ConditionTrue, "DiscoverySucceeded", message); err != nil {
120+
return reconcile.Result{}, fmt.Errorf("update RemoteMCPServer discovery status: %w", err)
121+
}
122+
if err := r.updateCatalog(ctx, server, discovered, true); err != nil {
123+
return reconcile.Result{}, fmt.Errorf("update RemoteMCPServer tool catalog: %w", err)
124+
}
125+
return reconcile.Result{RequeueAfter: refreshInterval}, nil
126+
}
127+
128+
func (r *Reconciler) updateCatalog(ctx context.Context, server *v1alpha3.RemoteMCPServer, tools []*v1alpha3.MCPTool, connected bool) error {
129+
name := client.ObjectKeyFromObject(server).String()
130+
var lastConnected *time.Time
131+
if connected {
132+
now := time.Now().UTC()
133+
lastConnected = &now
134+
}
135+
if _, err := r.catalog.StoreToolServer(ctx, &dbmodel.ToolServer{
136+
Name: name, GroupKind: remoteGroupKind, Description: server.Spec.Description, LastConnected: lastConnected,
137+
}); err != nil {
138+
return fmt.Errorf("store server: %w", err)
139+
}
140+
if err := r.catalog.RefreshToolsForServer(ctx, name, remoteGroupKind, tools...); err != nil {
141+
return fmt.Errorf("refresh tools: %w", err)
142+
}
143+
return nil
144+
}
145+
146+
func (r *Reconciler) deleteCatalog(ctx context.Context, name string) error {
147+
return errors.Join(
148+
wrapError("delete tools", r.catalog.DeleteToolsForServer(ctx, name, remoteGroupKind)),
149+
wrapError("delete server", r.catalog.DeleteToolServer(ctx, name, remoteGroupKind)),
150+
)
151+
}
152+
153+
func wrapError(action string, err error) error {
154+
if err == nil {
155+
return nil
156+
}
157+
return fmt.Errorf("%s: %w", action, err)
158+
}
159+
160+
func normalizeTools(tools []toolservice.MCPAppTool) ([]*v1alpha3.MCPTool, error) {
161+
discovered := make([]*v1alpha3.MCPTool, 0, len(tools))
162+
seen := make(map[string]struct{}, len(tools))
163+
for _, tool := range tools {
164+
if strings.TrimSpace(tool.Name) == "" {
165+
return nil, fmt.Errorf("MCP discovery returned a tool with an empty name")
166+
}
167+
if _, exists := seen[tool.Name]; exists {
168+
return nil, fmt.Errorf("MCP discovery returned duplicate tool %q", tool.Name)
169+
}
170+
seen[tool.Name] = struct{}{}
171+
discovered = append(discovered, &v1alpha3.MCPTool{Name: tool.Name, Description: tool.Description})
172+
}
173+
slices.SortFunc(discovered, func(a, b *v1alpha3.MCPTool) int {
174+
return strings.Compare(a.Name, b.Name)
175+
})
176+
return discovered, nil
177+
}
178+
179+
func (r *Reconciler) updateStatus(
180+
ctx context.Context,
181+
server *v1alpha3.RemoteMCPServer,
182+
tools []*v1alpha3.MCPTool,
183+
conditionStatus metav1.ConditionStatus,
184+
reason string,
185+
message string,
186+
) error {
187+
original := server.DeepCopy()
188+
server.Status.ObservedGeneration = server.Generation
189+
server.Status.DiscoveredTools = tools
190+
apiMeta.SetStatusCondition(&server.Status.Conditions, metav1.Condition{
191+
Type: conditionAccepted, Status: conditionStatus, Reason: reason, Message: message,
192+
ObservedGeneration: server.Generation,
193+
})
194+
if reflect.DeepEqual(original.Status, server.Status) {
195+
return nil
196+
}
197+
return r.client.Status().Patch(ctx, server, client.MergeFrom(original))
198+
}
199+
200+
func (r *Reconciler) requestsForDependency(ctx context.Context, object client.Object) []reconcile.Request {
201+
servers := &v1alpha3.RemoteMCPServerList{}
202+
if err := r.client.List(ctx, servers, client.InNamespace(object.GetNamespace())); err != nil {
203+
return nil
204+
}
205+
requests := make([]reconcile.Request, 0)
206+
for i := range servers.Items {
207+
server := &servers.Items[i]
208+
if referencesDependency(server, object) {
209+
requests = append(requests, reconcile.Request{NamespacedName: types.NamespacedName{
210+
Namespace: server.Namespace, Name: server.Name,
211+
}})
212+
}
213+
}
214+
return requests
215+
}
216+
217+
func referencesDependency(server *v1alpha3.RemoteMCPServer, object client.Object) bool {
218+
switch object.(type) {
219+
case *corev1.Secret:
220+
if server.Spec.TLS != nil && server.Spec.TLS.CACertSecretRef == object.GetName() {
221+
return true
222+
}
223+
for i := range server.Spec.HeadersFrom {
224+
from := server.Spec.HeadersFrom[i].ValueFrom
225+
if from != nil && from.Type == v1alpha3.SecretValueSource && from.Name == object.GetName() {
226+
return true
227+
}
228+
}
229+
case *corev1.ConfigMap:
230+
for i := range server.Spec.HeadersFrom {
231+
from := server.Spec.HeadersFrom[i].ValueFrom
232+
if from != nil && from.Type == v1alpha3.ConfigMapValueSource && from.Name == object.GetName() {
233+
return true
234+
}
235+
}
236+
}
237+
return false
238+
}
239+
240+
var _ reconcile.Reconciler = (*Reconciler)(nil)

0 commit comments

Comments
 (0)