diff --git a/go/core/v2/controller/collections.go b/go/core/v2/controller/collections.go index cc6888b12..820508b9e 100644 --- a/go/core/v2/controller/collections.go +++ b/go/core/v2/controller/collections.go @@ -4,6 +4,7 @@ import ( atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1" "github.com/agent-substrate/substrate/pkg/proto/ateapipb" kagentv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3" + v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "istio.io/istio/pkg/kube" "istio.io/istio/pkg/kube/controllers" "istio.io/istio/pkg/kube/kclient" @@ -16,18 +17,19 @@ import ( // Collections contains the Kubernetes inputs used to resolve an AgentTemplate // and the template/harness pairs derived from Harness admission selectors. type Collections struct { - AgentTemplates krt.Collection[*kagentv1alpha3.AgentTemplate] - Harnesses krt.Collection[*kagentv1alpha3.Harness] - ModelConfigs krt.Collection[*kagentv1alpha3.ModelConfig] - RemoteMCPServers krt.Collection[*kagentv1alpha3.RemoteMCPServer] - ConfigMaps krt.Collection[*corev1.ConfigMap] - Secrets krt.Collection[*corev1.Secret] - WorkerPools krt.Collection[*atev1alpha1.WorkerPool] - ActorTemplates krt.StaticCollection[ObservedActorTemplate] - Pairs krt.Collection[AgentTemplateHarnessPair] - Reconciliations krt.Collection[PairReconciliation] - ModelConfigReconciliations krt.StatusCollection[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus] - AgentTemplateStatuses krt.StatusCollection[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus] + AgentTemplates krt.Collection[*kagentv1alpha3.AgentTemplate] + Harnesses krt.Collection[*kagentv1alpha3.Harness] + ModelConfigs krt.Collection[*kagentv1alpha3.ModelConfig] + RemoteMCPServers krt.Collection[*kagentv1alpha3.RemoteMCPServer] + ConfigMaps krt.Collection[*corev1.ConfigMap] + Secrets krt.Collection[*corev1.Secret] + WorkerPools krt.Collection[*atev1alpha1.WorkerPool] + ActorTemplates krt.StaticCollection[ObservedActorTemplate] + Pairs krt.Collection[AgentTemplateHarnessPair] + Reconciliations krt.Collection[PairReconciliation] + ModelConfigStatuses krt.StatusCollection[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus] + ResolvedModelConfigs krt.Collection[v2translator.ResolvedModelConfig] + AgentTemplateStatuses krt.StatusCollection[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus] } // ObservedActorTemplate adapts an ate-api resource to KRT's keyed collection. @@ -63,23 +65,28 @@ func NewCollections(client kube.Client, watchNamespaces []string, opts krt.Optio workerPools := typedCollection[*atev1alpha1.WorkerPool](client, watchNamespaces, "WorkerPools", opts) actorTemplates := krt.NewStaticCollection[ObservedActorTemplate](nil, nil, opts.WithName("ActorTemplates")...) pairs := newPairCollection(agentTemplates, harnesses, opts) - modelConfigReconciliations := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) - reconciliations := newPairReconciliations(pairs, agentTemplates, modelConfigs, remoteMCPServers, configMaps, secrets, workerPools, actorTemplates, opts) + modelConfigStatuses, resolvedModelConfigs := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) + compilerCollections := v2translator.Collections{ + AgentTemplates: agentTemplates, ResolvedModelConfigs: resolvedModelConfigs, RemoteMCPServers: remoteMCPServers, + ConfigMaps: configMaps, Secrets: secrets, WorkerPools: workerPools, + } + reconciliations := newPairReconciliations(pairs, compilerCollections, actorTemplates, opts) statuses := newAgentTemplateStatuses(agentTemplates, reconciliations, opts) return Collections{ - AgentTemplates: agentTemplates, - Harnesses: harnesses, - ModelConfigs: modelConfigs, - RemoteMCPServers: remoteMCPServers, - ConfigMaps: configMaps, - Secrets: secrets, - WorkerPools: workerPools, - ActorTemplates: actorTemplates, - Pairs: pairs, - Reconciliations: reconciliations, - ModelConfigReconciliations: modelConfigReconciliations, - AgentTemplateStatuses: statuses, + AgentTemplates: agentTemplates, + Harnesses: harnesses, + ModelConfigs: modelConfigs, + RemoteMCPServers: remoteMCPServers, + ConfigMaps: configMaps, + Secrets: secrets, + WorkerPools: workerPools, + ActorTemplates: actorTemplates, + Pairs: pairs, + Reconciliations: reconciliations, + ModelConfigStatuses: modelConfigStatuses, + ResolvedModelConfigs: resolvedModelConfigs, + AgentTemplateStatuses: statuses, } } diff --git a/go/core/v2/controller/collections_test.go b/go/core/v2/controller/collections_test.go index a154f3462..c294abdb5 100644 --- a/go/core/v2/controller/collections_test.go +++ b/go/core/v2/controller/collections_test.go @@ -7,8 +7,10 @@ import ( atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1" "github.com/agent-substrate/substrate/pkg/proto/ateapipb" kagentv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3" + v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "google.golang.org/protobuf/proto" "istio.io/istio/pkg/kube/krt" + "istio.io/istio/pkg/kube/krt/krttest" corev1 "k8s.io/api/core/v1" apimeta "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -29,7 +31,8 @@ func TestAgentTemplateHarnessPairs(t *testing.T) { harness("team-b", "other-namespace", map[string]string{"runtime": "python"}), {ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "no-admission"}}, }, opts.WithName("Harnesses")...) - templates := krt.NewStaticCollection(nil, []*kagentv1alpha3.AgentTemplate{template}, opts.WithName("AgentTemplates")...) + mock := krttest.NewMock(t, []any{template}) + templates := krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock) pairs := newPairCollection(templates, harnesses, opts) if !pairs.WaitUntilSynced(stop) { @@ -64,22 +67,30 @@ func TestReconciliationCollectionsCompileAndObserveRevision(t *testing.T) { SnapshotPolicy: kagentv1alpha3.HarnessSnapshotPolicy{Location: "snapshots"}, } modelConfigs := krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model"}, Spec: kagentv1alpha3.ModelConfigSpec{Provider: kagentv1alpha3.ModelProviderOpenAI, Model: "gpt-5"}}}, opts.WithName("ModelConfigs")...) + mock := krttest.NewMock(t, []any{ + template, + matchingHarness, + &atev1alpha1.WorkerPool{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "default"}}, + }) collections := Collections{ - AgentTemplates: krt.NewStaticCollection(nil, []*kagentv1alpha3.AgentTemplate{template}, opts.WithName("AgentTemplates")...), - Harnesses: krt.NewStaticCollection(nil, []*kagentv1alpha3.Harness{matchingHarness}, opts.WithName("Harnesses")...), + AgentTemplates: krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock), + Harnesses: krttest.GetMockCollection[*kagentv1alpha3.Harness](mock), ModelConfigs: modelConfigs, - RemoteMCPServers: krt.NewStaticCollection[*kagentv1alpha3.RemoteMCPServer](nil, nil, opts.WithName("RemoteMCPServers")...), - ConfigMaps: krt.NewStaticCollection[*corev1.ConfigMap](nil, nil, opts.WithName("ConfigMaps")...), - Secrets: krt.NewStaticCollection[*corev1.Secret](nil, nil, opts.WithName("Secrets")...), - WorkerPools: krt.NewStaticCollection(nil, []*atev1alpha1.WorkerPool{{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "default"}}}, opts.WithName("WorkerPools")...), + RemoteMCPServers: krttest.GetMockCollection[*kagentv1alpha3.RemoteMCPServer](mock), + ConfigMaps: krttest.GetMockCollection[*corev1.ConfigMap](mock), + Secrets: krttest.GetMockCollection[*corev1.Secret](mock), + WorkerPools: krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock), ActorTemplates: krt.NewStaticCollection[ObservedActorTemplate](nil, nil, opts.WithName("ActorTemplates")...), } - collections.ModelConfigReconciliations = newModelConfigReconciliations(collections.ModelConfigs, collections.ConfigMaps, collections.Secrets, opts) + collections.ModelConfigStatuses, collections.ResolvedModelConfigs = newModelConfigReconciliations(collections.ModelConfigs, collections.ConfigMaps, collections.Secrets, opts) collections.Pairs = newPairCollection(collections.AgentTemplates, collections.Harnesses, opts) collections.Reconciliations = newPairReconciliations( - collections.Pairs, collections.AgentTemplates, collections.ModelConfigs, collections.RemoteMCPServers, - collections.ConfigMaps, collections.Secrets, collections.WorkerPools, collections.ActorTemplates, opts, + collections.Pairs, v2translator.Collections{ + AgentTemplates: collections.AgentTemplates, ResolvedModelConfigs: collections.ResolvedModelConfigs, + RemoteMCPServers: collections.RemoteMCPServers, ConfigMaps: collections.ConfigMaps, + Secrets: collections.Secrets, WorkerPools: collections.WorkerPools, + }, collections.ActorTemplates, opts, ) collections.AgentTemplateStatuses = newAgentTemplateStatuses(collections.AgentTemplates, collections.Reconciliations, opts) @@ -141,16 +152,27 @@ func TestClaudeReconciliationCompilesActorTemplate(t *testing.T) { Provider: kagentv1alpha3.ModelProviderAnthropic, Model: "claude-sonnet-4-5", APIKeySecret: "model-auth", APIKeySecretKey: "api-key", }} secret := &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model-auth", UID: "secret-uid"}, Data: map[string][]byte{"api-key": []byte("secret")}} - templates := krt.NewStaticCollection(nil, []*kagentv1alpha3.AgentTemplate{template}, opts.WithName("AgentTemplates")...) - pairs := newPairCollection(templates, krt.NewStaticCollection(nil, []*kagentv1alpha3.Harness{claudeHarness}, opts.WithName("Harnesses")...), opts) + mock := krttest.NewMock(t, []any{ + template, + claudeHarness, + model, + secret, + &atev1alpha1.WorkerPool{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "default"}}, + }) + templates := krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock) + pairs := newPairCollection(templates, krttest.GetMockCollection[*kagentv1alpha3.Harness](mock), opts) + configMaps := krttest.GetMockCollection[*corev1.ConfigMap](mock) + secrets := krttest.GetMockCollection[*corev1.Secret](mock) + _, resolvedModelConfigs := newModelConfigReconciliations( + krttest.GetMockCollection[*kagentv1alpha3.ModelConfig](mock), configMaps, secrets, opts, + ) reconciliations := newPairReconciliations( - pairs, templates, - krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{model}, opts.WithName("ModelConfigs")...), - krt.NewStaticCollection[*kagentv1alpha3.RemoteMCPServer](nil, nil, opts.WithName("RemoteMCPServers")...), - krt.NewStaticCollection[*corev1.ConfigMap](nil, nil, opts.WithName("ConfigMaps")...), - krt.NewStaticCollection(nil, []*corev1.Secret{secret}, opts.WithName("Secrets")...), - krt.NewStaticCollection(nil, []*atev1alpha1.WorkerPool{{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "default"}}}, opts.WithName("WorkerPools")...), - krt.NewStaticCollection[ObservedActorTemplate](nil, nil, opts.WithName("ActorTemplates")...), opts, + pairs, v2translator.Collections{ + AgentTemplates: templates, ResolvedModelConfigs: resolvedModelConfigs, + RemoteMCPServers: krttest.GetMockCollection[*kagentv1alpha3.RemoteMCPServer](mock), + ConfigMaps: configMaps, Secrets: secrets, + WorkerPools: krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock), + }, krttest.GetMockCollection[ObservedActorTemplate](mock), opts, ) waitFor(t, func() bool { states := reconciliations.List() @@ -187,17 +209,23 @@ func TestReconciliationTracksSharedAgentTemplate(t *testing.T) { harness.Spec.Workload.Image = "example.com/kagent@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" harness.Spec.Substrate = kagentv1alpha3.HarnessSubstratePolicy{WorkerPoolRef: corev1.LocalObjectReference{Name: "default"}, SnapshotPolicy: kagentv1alpha3.HarnessSnapshotPolicy{Location: "snapshots"}} templates := krt.NewStaticCollection(nil, []*kagentv1alpha3.AgentTemplate{root, child}, opts.WithName("AgentTemplates")...) - pairs := newPairCollection(templates, krt.NewStaticCollection(nil, []*kagentv1alpha3.Harness{harness}, opts.WithName("Harnesses")...), opts) - modelConfigs := krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model"}, Spec: kagentv1alpha3.ModelConfigSpec{Provider: kagentv1alpha3.ModelProviderOpenAI, Model: "gpt-5"}}}, opts.WithName("ModelConfigs")...) - configMaps := krt.NewStaticCollection[*corev1.ConfigMap](nil, nil, opts.WithName("ConfigMaps")...) - secrets := krt.NewStaticCollection[*corev1.Secret](nil, nil, opts.WithName("Secrets")...) + mock := krttest.NewMock(t, []any{ + harness, + &kagentv1alpha3.ModelConfig{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model"}, Spec: kagentv1alpha3.ModelConfigSpec{Provider: kagentv1alpha3.ModelProviderOpenAI, Model: "gpt-5"}}, + &atev1alpha1.WorkerPool{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "default"}}, + }) + pairs := newPairCollection(templates, krttest.GetMockCollection[*kagentv1alpha3.Harness](mock), opts) + modelConfigs := krttest.GetMockCollection[*kagentv1alpha3.ModelConfig](mock) + configMaps := krttest.GetMockCollection[*corev1.ConfigMap](mock) + secrets := krttest.GetMockCollection[*corev1.Secret](mock) + _, resolvedModelConfigs := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) reconciliations := newPairReconciliations( - pairs, templates, - modelConfigs, - krt.NewStaticCollection[*kagentv1alpha3.RemoteMCPServer](nil, nil, opts.WithName("RemoteMCPServers")...), - configMaps, secrets, - krt.NewStaticCollection(nil, []*atev1alpha1.WorkerPool{{ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "default"}}}, opts.WithName("WorkerPools")...), - krt.NewStaticCollection[ObservedActorTemplate](nil, nil, opts.WithName("ActorTemplates")...), opts, + pairs, v2translator.Collections{ + AgentTemplates: templates, ResolvedModelConfigs: resolvedModelConfigs, + RemoteMCPServers: krttest.GetMockCollection[*kagentv1alpha3.RemoteMCPServer](mock), + ConfigMaps: configMaps, Secrets: secrets, + WorkerPools: krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock), + }, krttest.GetMockCollection[ObservedActorTemplate](mock), opts, ) var initial string waitFor(t, func() bool { diff --git a/go/core/v2/controller/modelconfig.go b/go/core/v2/controller/modelconfig.go index 6775fd72a..f2311d880 100644 --- a/go/core/v2/controller/modelconfig.go +++ b/go/core/v2/controller/modelconfig.go @@ -1,7 +1,6 @@ package controller import ( - "context" "crypto/sha256" "encoding/hex" "slices" @@ -11,41 +10,17 @@ import ( v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" - apiequality "k8s.io/apimachinery/pkg/api/equality" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) -type ModelConfigReconciliation struct { - ModelConfigName krt.Named - // If Translation is nil, Failure is non-nil and describes why the ModelConfig could not be translated. - // note that failure may be not nil even if Translation is non-nil. - Translation *v2translator.ResolvedModelConfig -} - -func (r ModelConfigReconciliation) Equals(other ModelConfigReconciliation) bool { - if r.ModelConfigName != other.ModelConfigName { - return false - } - if !apiequality.Semantic.DeepEqual(r.Translation, other.Translation) { - return false - } - return true -} - -func (r ModelConfigReconciliation) ResourceName() string { - return r.ModelConfigName.ResourceName() -} - func newModelConfigReconciliations( modelConfigs krt.Collection[*kagentv1alpha3.ModelConfig], configMaps krt.Collection[*corev1.ConfigMap], secrets krt.Collection[*corev1.Secret], opts krt.OptionsBuilder, -) krt.StatusCollection[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus] { - statuses, _ := krt.NewStatusCollection(modelConfigs, func(ctx krt.HandlerContext, modelConfig *kagentv1alpha3.ModelConfig) (*kagentv1alpha3.ModelConfigStatus, *ModelConfigReconciliation) { - state := &ModelConfigReconciliation{ModelConfigName: krt.Named{Namespace: modelConfig.Namespace, Name: modelConfig.Name}} - reader := collectionReader{ctx: ctx, configMaps: configMaps, secrets: secrets, modelConfigs: modelConfigs} - translation, err := v2translator.ResolveModelConfig(context.Background(), reader, modelConfig) +) (krt.StatusCollection[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus], krt.Collection[v2translator.ResolvedModelConfig]) { + statuses, reconciliations := krt.NewStatusCollection(modelConfigs, func(ctx krt.HandlerContext, modelConfig *kagentv1alpha3.ModelConfig) (*kagentv1alpha3.ModelConfigStatus, *v2translator.ResolvedModelConfig) { + translation, err := v2translator.ResolveModelConfig(ctx, v2translator.Collections{ConfigMaps: configMaps, Secrets: secrets}, modelConfig) if err != nil { return nil, nil } @@ -60,7 +35,6 @@ func newModelConfigReconciliations( failure := translation.ReferenceFailures[0] resolvedRefsFailure = &ReconciliationFailure{Condition: kagentv1alpha3.ModelConfigConditionTypeResolvedRefs, Reason: failure.Reason, Message: failure.Message} } - state.Translation = translation seenSecrets := map[string]struct{}{} for _, reference := range translation.References { switch reference.Kind { @@ -122,9 +96,9 @@ func newModelConfigReconciliations( ObservedGeneration: modelConfig.Generation, SecretHash: hashModelConfigValues(values), Conditions: conditions, - }, state + }, translation }, opts.WithName("ModelConfigReconciliations")...) - return statuses + return statuses, reconciliations } type hashValue struct { diff --git a/go/core/v2/controller/modelconfig_test.go b/go/core/v2/controller/modelconfig_test.go index 07c2ef405..ba0c4eaba 100644 --- a/go/core/v2/controller/modelconfig_test.go +++ b/go/core/v2/controller/modelconfig_test.go @@ -7,26 +7,21 @@ import ( kagentv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "istio.io/istio/pkg/kube/krt" + "istio.io/istio/pkg/kube/krt/krttest" corev1 "k8s.io/api/core/v1" apimeta "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) -func TestModelConfigReconciliationEquals(t *testing.T) { - left := ModelConfigReconciliation{ - ModelConfigName: krt.Named{Namespace: "team-a", Name: "model"}, - Translation: &v2translator.ResolvedModelConfig{Config: &kagentv1alpha3.ModelConfig{Spec: kagentv1alpha3.ModelConfigSpec{Model: "gpt-5"}}}, - } - right := ModelConfigReconciliation{ - ModelConfigName: krt.Named{Namespace: "team-a", Name: "model"}, - Translation: &v2translator.ResolvedModelConfig{Config: &kagentv1alpha3.ModelConfig{Spec: kagentv1alpha3.ModelConfigSpec{Model: "gpt-5"}}}, - } +func TestResolvedModelConfigEquals(t *testing.T) { + left := v2translator.ResolvedModelConfig{Config: &kagentv1alpha3.ModelConfig{Spec: kagentv1alpha3.ModelConfigSpec{Model: "gpt-5"}}} + right := v2translator.ResolvedModelConfig{Config: &kagentv1alpha3.ModelConfig{Spec: kagentv1alpha3.ModelConfigSpec{Model: "gpt-5"}}} if !krt.Equal(left, right) { - t.Fatal("equal reconciliations were not considered equal") + t.Fatal("equal resolutions were not considered equal") } - left.Translation = &v2translator.ResolvedModelConfig{Config: &kagentv1alpha3.ModelConfig{Spec: kagentv1alpha3.ModelConfigSpec{Model: "gpt-4"}}} + left.Config.Spec.Model = "gpt-4" if krt.Equal(left, right) { - t.Fatal("different translations were considered equal") + t.Fatal("different resolutions were considered equal") } } @@ -35,16 +30,18 @@ func TestModelConfigReconciliationTracksSecret(t *testing.T) { t.Cleanup(func() { close(stop) }) opts := krt.NewOptionsBuilder(stop, "test", nil) - modelConfigs := krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{{ + modelConfig := &kagentv1alpha3.ModelConfig{ ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model"}, Spec: kagentv1alpha3.ModelConfigSpec{Model: "gpt-5", Provider: kagentv1alpha3.ModelProviderOpenAI, APIKeySecret: "credentials"}, - }}, opts.WithName("ModelConfigs")...) + } + mock := krttest.NewMock(t, []any{modelConfig}) + modelConfigs := krttest.GetMockCollection[*kagentv1alpha3.ModelConfig](mock) secrets := krt.NewStaticCollection(nil, []*corev1.Secret{{ ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "credentials"}, Data: map[string][]byte{"key": []byte("before")}, }}, opts.WithName("Secrets")...) - configMaps := krt.NewStaticCollection[*corev1.ConfigMap](nil, nil, opts.WithName("ConfigMaps")...) - reconciliations := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) + configMaps := krttest.GetMockCollection[*corev1.ConfigMap](mock) + reconciliations, _ := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) waitFor(t, func() bool { return len(reconciliations.List()) == 1 }) initial := reconciliations.List()[0].Status @@ -70,7 +67,7 @@ func TestModelConfigReconciliationMissingAPIKeySecretKey(t *testing.T) { t.Cleanup(func() { close(stop) }) opts := krt.NewOptionsBuilder(stop, "test", nil) - modelConfigs := krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{{ + modelConfig := &kagentv1alpha3.ModelConfig{ ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model"}, Spec: kagentv1alpha3.ModelConfigSpec{ Model: "gpt-5", @@ -78,13 +75,16 @@ func TestModelConfigReconciliationMissingAPIKeySecretKey(t *testing.T) { APIKeySecret: "credentials", APIKeySecretKey: "NON_EXISTENT_KEY", }, - }}, opts.WithName("ModelConfigs")...) - secrets := krt.NewStaticCollection(nil, []*corev1.Secret{{ + } + secret := &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "credentials"}, Data: map[string][]byte{"EXISTING_KEY": []byte("secret-value")}, - }}, opts.WithName("Secrets")...) - configMaps := krt.NewStaticCollection[*corev1.ConfigMap](nil, nil, opts.WithName("ConfigMaps")...) - reconciliations := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) + } + mock := krttest.NewMock(t, []any{modelConfig, secret}) + modelConfigs := krttest.GetMockCollection[*kagentv1alpha3.ModelConfig](mock) + secrets := krttest.GetMockCollection[*corev1.Secret](mock) + configMaps := krttest.GetMockCollection[*corev1.ConfigMap](mock) + reconciliations, _ := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) waitFor(t, func() bool { return len(reconciliations.List()) == 1 }) status := reconciliations.List()[0].Status @@ -137,20 +137,21 @@ func TestModelConfigReconciliationValidatesEffectiveProviderReferences(t *testin stop := make(chan struct{}) t.Cleanup(func() { close(stop) }) opts := krt.NewOptionsBuilder(stop, "test", nil) - modelConfigs := krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{{ + modelConfig := &kagentv1alpha3.ModelConfig{ ObjectMeta: metav1.ObjectMeta{Namespace: "team-a", Name: "model"}, Spec: test.spec, - }}, opts.WithName("ModelConfigs")...) - secrets := make([]*corev1.Secret, 0, 1) + } + inputs := []any{modelConfig} if test.secret != nil { - secrets = append(secrets, test.secret) + inputs = append(inputs, test.secret) } - configMaps := make([]*corev1.ConfigMap, 0, 1) if test.configMap != nil { - configMaps = append(configMaps, test.configMap) + inputs = append(inputs, test.configMap) } - secretCollection := krt.NewStaticCollection(nil, secrets, opts.WithName("Secrets")...) - configMapCollection := krt.NewStaticCollection(nil, configMaps, opts.WithName("ConfigMaps")...) - reconciliations := newModelConfigReconciliations(modelConfigs, configMapCollection, secretCollection, opts) + mock := krttest.NewMock(t, inputs) + modelConfigs := krttest.GetMockCollection[*kagentv1alpha3.ModelConfig](mock) + secretCollection := krttest.GetMockCollection[*corev1.Secret](mock) + configMapCollection := krttest.GetMockCollection[*corev1.ConfigMap](mock) + reconciliations, _ := newModelConfigReconciliations(modelConfigs, configMapCollection, secretCollection, opts) waitFor(t, func() bool { return len(reconciliations.List()) == 1 }) status := reconciliations.List()[0].Status diff --git a/go/core/v2/controller/reader.go b/go/core/v2/controller/reader.go deleted file mode 100644 index 1890c7eb2..000000000 --- a/go/core/v2/controller/reader.go +++ /dev/null @@ -1,89 +0,0 @@ -package controller - -import ( - "context" - "fmt" - - atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1" - kagentv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3" - v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" - "istio.io/istio/pkg/kube/controllers" - "istio.io/istio/pkg/kube/krt" - corev1 "k8s.io/api/core/v1" - apierrors "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/runtime" - "k8s.io/apimachinery/pkg/runtime/schema" - "k8s.io/apimachinery/pkg/types" -) - -// collectionReader lets the compiler use ordinary typed reads while KRT tracks -// every object fetched by a transformation as a recomputation dependency. -type collectionReader struct { - ctx krt.HandlerContext - agentTemplates krt.Collection[*kagentv1alpha3.AgentTemplate] - modelConfigs krt.Collection[*kagentv1alpha3.ModelConfig] - remoteMCPServers krt.Collection[*kagentv1alpha3.RemoteMCPServer] - configMaps krt.Collection[*corev1.ConfigMap] - secrets krt.Collection[*corev1.Secret] - workerPools krt.Collection[*atev1alpha1.WorkerPool] -} - -func (r collectionReader) GetResolvedModelConfig(ctx context.Context, key types.NamespacedName) (*v2translator.ResolvedModelConfig, error) { - modelConfig := krt.FetchOne(r.ctx, r.modelConfigs, krt.FilterObjectName(key)) - if modelConfig == nil { - return nil, fmt.Errorf("ModelConfig %s does not exist", key) - } - return v2translator.ResolveModelConfig(ctx, r, *modelConfig) -} - -func (r collectionReader) Get(_ context.Context, key types.NamespacedName, object runtime.Object) error { - switch target := object.(type) { - case *kagentv1alpha3.AgentTemplate: - source, err := r.fetchObject(r.agentTemplates, key, kagentv1alpha3.GroupVersion.WithResource("agenttemplates").GroupResource()) - if err == nil { - *target = *source.DeepCopy() - } - return err - case *kagentv1alpha3.ModelConfig: - source, err := r.fetchObject(r.modelConfigs, key, kagentv1alpha3.GroupVersion.WithResource("modelconfigs").GroupResource()) - if err == nil { - *target = *source.DeepCopy() - } - return err - case *kagentv1alpha3.RemoteMCPServer: - source, err := r.fetchObject(r.remoteMCPServers, key, kagentv1alpha3.GroupVersion.WithResource("remotemcpservers").GroupResource()) - if err == nil { - *target = *source.DeepCopy() - } - return err - case *corev1.ConfigMap: - source, err := r.fetchObject(r.configMaps, key, schema.GroupResource{Resource: "configmaps"}) - if err == nil { - *target = *source.DeepCopy() - } - return err - case *corev1.Secret: - source, err := r.fetchObject(r.secrets, key, schema.GroupResource{Resource: "secrets"}) - if err == nil { - *target = *source.DeepCopy() - } - return err - case *atev1alpha1.WorkerPool: - source, err := r.fetchObject(r.workerPools, key, atev1alpha1.GroupVersion.WithResource("workerpools").GroupResource()) - if err == nil { - *target = *source.DeepCopy() - } - return err - default: - return fmt.Errorf("unsupported KRT read type %T", object) - } -} - -func (r collectionReader) fetchObject[T controllers.ComparableObject](collection krt.Collection[T], key types.NamespacedName, resource schema.GroupResource) (T, error) { - object := krt.FetchOne(r.ctx, collection, krt.FilterObjectName(key)) - if object == nil { - var zero T - return zero, apierrors.NewNotFound(resource, key.Name) - } - return *object, nil -} diff --git a/go/core/v2/controller/reconciler.go b/go/core/v2/controller/reconciler.go index 598821e37..d85a46b3a 100644 --- a/go/core/v2/controller/reconciler.go +++ b/go/core/v2/controller/reconciler.go @@ -7,7 +7,6 @@ import ( "strings" "time" - atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1" "github.com/agent-substrate/substrate/pkg/proto/ateapipb" kagentclient "github.com/kagent-dev/kagent/go/api/clientset/versioned/typed/api/v1alpha3" dbpkg "github.com/kagent-dev/kagent/go/api/database" @@ -21,7 +20,6 @@ import ( "google.golang.org/grpc/status" "istio.io/istio/pkg/kube/controllers" "istio.io/istio/pkg/kube/krt" - corev1 "k8s.io/api/core/v1" apiequality "k8s.io/apimachinery/pkg/api/equality" apimeta "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -53,25 +51,16 @@ type ReconciliationFailure struct { func newPairReconciliations( pairs krt.Collection[AgentTemplateHarnessPair], - agentTemplates krt.Collection[*kagentv1alpha3.AgentTemplate], - modelConfigs krt.Collection[*kagentv1alpha3.ModelConfig], - remoteMCPServers krt.Collection[*kagentv1alpha3.RemoteMCPServer], - configMaps krt.Collection[*corev1.ConfigMap], - secrets krt.Collection[*corev1.Secret], - workerPools krt.Collection[*atev1alpha1.WorkerPool], + collections v2translator.Collections, actorTemplates krt.Collection[ObservedActorTemplate], opts krt.OptionsBuilder, ) krt.Collection[PairReconciliation] { return krt.NewCollection(pairs, func(ctx krt.HandlerContext, pair AgentTemplateHarnessPair) *PairReconciliation { state := &PairReconciliation{Pair: pair} - reader := collectionReader{ - ctx: ctx, agentTemplates: agentTemplates, modelConfigs: modelConfigs, remoteMCPServers: remoteMCPServers, - configMaps: configMaps, secrets: secrets, workerPools: workerPools, - } - revision, err := v2translator.NewCompiler(reader, map[v2translator.HarnessType]v2translator.HarnessCompiler{ - v2translator.HarnessTypeKagent: kagenttranslator.NewCompiler(reader), - v2translator.HarnessTypeClaude: claudetranslator.NewCompiler(reader), - v2translator.HarnessTypeBYO: byotranslator.NewCompiler(reader), + revision, err := v2translator.NewCompiler(ctx, collections, map[v2translator.HarnessType]v2translator.HarnessCompiler{ + v2translator.HarnessTypeKagent: kagenttranslator.NewCompiler(ctx, collections), + v2translator.HarnessTypeClaude: claudetranslator.NewCompiler(ctx, collections), + v2translator.HarnessTypeBYO: byotranslator.NewCompiler(ctx, collections), }).CompileAgentTemplate(context.Background(), pair.Harness, pair.AgentTemplate) if err != nil { condition, reason := kagentv1alpha3.AgentTemplateConditionResolvedRefs, "ReferenceResolutionFailed" @@ -89,10 +78,9 @@ func newPairReconciliations( return state } - workerPool := &atev1alpha1.WorkerPool{} workerKey := types.NamespacedName{Namespace: revision.Namespace, Name: revision.WorkerPoolName} - if err := reader.Get(context.Background(), workerKey, workerPool); err != nil { - state.Failure = &ReconciliationFailure{Condition: kagentv1alpha3.AgentTemplateConditionResolvedRefs, Reason: "WorkerPoolNotFound", Message: err.Error()} + if krt.FetchOne(ctx, collections.WorkerPools, krt.FilterObjectName(workerKey)) == nil { + state.Failure = &ReconciliationFailure{Condition: kagentv1alpha3.AgentTemplateConditionResolvedRefs, Reason: "WorkerPoolNotFound", Message: fmt.Sprintf("WorkerPool %q not found", workerKey.String())} return state } state.DesiredActorTemplate, err = substrate.ActorTemplateForRevision(revision, state.RevisionID) @@ -197,7 +185,7 @@ func newReconciler( } r.agentTemplateStatuses.Add(status.ResourceName()) }) - r.modelConfigStatusHandler = collections.ModelConfigReconciliations.Register(func(event krt.Event[krt.ObjectWithStatus[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus]]) { + r.modelConfigStatusHandler = collections.ModelConfigStatuses.Register(func(event krt.Event[krt.ObjectWithStatus[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus]]) { status := event.Latest() if apiequality.Semantic.DeepEqual(modelConfigStatusWithTransitionTimes(status.Status, status.Obj.Status), status.Obj.Status) { return @@ -346,7 +334,7 @@ func (r *Reconciler) reconcileAgentTemplateStatus(ctx context.Context, key strin } func (r *Reconciler) reconcileModelConfigStatus(ctx context.Context, key string) error { - desired := r.collections.ModelConfigReconciliations.GetKey(key) + desired := r.collections.ModelConfigStatuses.GetKey(key) modelConfig := r.collections.ModelConfigs.GetKey(key) if desired == nil || modelConfig == nil { return nil diff --git a/go/core/v2/controller/reconciler_test.go b/go/core/v2/controller/reconciler_test.go index bd283b93f..741838bab 100644 --- a/go/core/v2/controller/reconciler_test.go +++ b/go/core/v2/controller/reconciler_test.go @@ -14,6 +14,7 @@ import ( "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "istio.io/istio/pkg/kube/krt" + "istio.io/istio/pkg/kube/krt/krttest" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) @@ -38,13 +39,17 @@ func TestReconcilerPersistsPairInOrder(t *testing.T) { status := kagentv1alpha3.AgentTemplateStatus{ObservedGeneration: 1, Harnesses: []kagentv1alpha3.AgentTemplateHarnessStatus{{ Harness: "kagent", Conditions: []metav1.Condition{{Type: kagentv1alpha3.AgentTemplateConditionReady, Status: metav1.ConditionFalse}}, }}} - statuses := krt.NewStaticCollection(nil, []krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]{{Obj: template, Status: status}}, opts.WithName("Statuses")...) + mock := krttest.NewMock(t, []any{ + template, + krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]{Obj: template, Status: status}, + }) + statuses := krttest.GetMockCollection[krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]](mock) store := &fakeRuntimeRevisionStore{} templates := &fakeActorTemplates{} statusClient := kagentfake.NewSimpleClientset(template.DeepCopy()).ApiV1alpha3() reconciler := &Reconciler{ collections: Collections{ - AgentTemplates: krt.NewStaticCollection(nil, []*kagentv1alpha3.AgentTemplate{template}, opts.WithName("AgentTemplates")...), + AgentTemplates: krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock), ActorTemplates: krt.NewStaticCollection[ObservedActorTemplate](nil, nil, opts.WithName("ActorTemplates")...), Reconciliations: reconciliations, AgentTemplateStatuses: statuses, }, @@ -167,19 +172,21 @@ func TestReconcilerUpdatesModelConfigStatusOnSecretHashChange(t *testing.T) { Data: map[string][]byte{"key": []byte("initial-secret")}, } - modelConfigs := krt.NewStaticCollection(nil, []*kagentv1alpha3.ModelConfig{modelConfig}, opts.WithName("ModelConfigs")...) + mock := krttest.NewMock(t, []any{modelConfig}) + modelConfigs := krttest.GetMockCollection[*kagentv1alpha3.ModelConfig](mock) secrets := krt.NewStaticCollection(nil, []*corev1.Secret{secret}, opts.WithName("Secrets")...) - configMaps := krt.NewStaticCollection[*corev1.ConfigMap](nil, nil, opts.WithName("ConfigMaps")...) - modelConfigReconciliations := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) + configMaps := krttest.GetMockCollection[*corev1.ConfigMap](mock) + modelConfigStatuses, resolvedModelConfigs := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts) collections := Collections{ - ModelConfigs: modelConfigs, - Secrets: secrets, - ConfigMaps: configMaps, - ModelConfigReconciliations: modelConfigReconciliations, - AgentTemplates: krt.NewStaticCollection[*kagentv1alpha3.AgentTemplate](nil, nil, opts.WithName("AgentTemplates")...), - Reconciliations: krt.NewStaticCollection[PairReconciliation](nil, nil, opts.WithName("Reconciliations")...), - AgentTemplateStatuses: krt.NewStaticCollection[krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]](nil, nil, opts.WithName("AgentTemplateStatuses")...), + ModelConfigs: modelConfigs, + Secrets: secrets, + ConfigMaps: configMaps, + ModelConfigStatuses: modelConfigStatuses, + ResolvedModelConfigs: resolvedModelConfigs, + AgentTemplates: krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock), + Reconciliations: krttest.GetMockCollection[PairReconciliation](mock), + AgentTemplateStatuses: krttest.GetMockCollection[krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]](mock), } statusClient := kagentfake.NewSimpleClientset(modelConfig.DeepCopy()).ApiV1alpha3() diff --git a/go/core/v2/translator/adkconfig/builder.go b/go/core/v2/translator/adkconfig/builder.go index 6500d6238..6be027d7f 100644 --- a/go/core/v2/translator/adkconfig/builder.go +++ b/go/core/v2/translator/adkconfig/builder.go @@ -12,6 +12,7 @@ import ( "github.com/kagent-dev/kagent/go/api/adk" "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/types" ) @@ -29,10 +30,15 @@ type provenanceEntry struct { } // Builder assembles resolved inputs into an ADK agent configuration. -type Builder struct{ kube v2translator.Reader } +type Builder struct { + ctx krt.HandlerContext + collections v2translator.Collections +} // NewBuilder constructs an ADK configuration builder. -func NewBuilder(kube v2translator.Reader) *Builder { return &Builder{kube: kube} } +func NewBuilder(ctx krt.HandlerContext, collections v2translator.Collections) *Builder { + return &Builder{ctx: ctx, collections: collections} +} type Result struct { Config *adk.AgentConfig @@ -52,9 +58,9 @@ type ModelResult struct { // BuildModel translates a standalone ModelConfig without building an agent. func (c *Builder) BuildModel(ctx context.Context, namespace, name string) (*ModelResult, error) { - resolved, err := c.kube.GetResolvedModelConfig(ctx, types.NamespacedName{Namespace: namespace, Name: name}) - if err != nil { - return nil, err + resolved := krt.FetchOne(c.ctx, c.collections.ResolvedModelConfigs, krt.FilterObjectName(types.NamespacedName{Namespace: namespace, Name: name})) + if resolved == nil { + return nil, fmt.Errorf("model config %q not found", name) } if failures := resolved.SemanticFailures; len(failures) > 0 { return nil, v2translator.NewValidationError("ModelConfig %q: %s", name, failures[0].Message) @@ -170,10 +176,11 @@ func (c *Builder) ResolveEnvironment(ctx context.Context, namespace string, envi return nil, fmt.Errorf("environment variable %q uses an unsupported value source", variable.Name) } ref := variable.ValueFrom.SecretKeyRef - secret := &corev1.Secret{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: namespace, Name: ref.Name}, secret); err != nil { - return nil, err + fetched := krt.FetchOne(c.ctx, c.collections.Secrets, krt.FilterObjectName(types.NamespacedName{Namespace: namespace, Name: ref.Name})) + if fetched == nil { + return nil, fmt.Errorf("secret %q not found", ref.Name) } + secret := *fetched value, ok := secret.Data[ref.Key] if !ok { return nil, fmt.Errorf("secret %q does not contain key %q", ref.Name, ref.Key) @@ -204,10 +211,11 @@ func (c *Builder) BuildProvenance(ctx context.Context, harness *v1alpha3.Harness entries = append(entries, objectProvenance(v1alpha3.GroupVersion.String(), "ModelConfig", model.Name, model.UID, model.Generation, model.Spec)) } for name := range configMaps { - configMap := &corev1.ConfigMap{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: harness.Namespace, Name: name}, configMap); err != nil { - return nil, err + fetched := krt.FetchOne(c.ctx, c.collections.ConfigMaps, krt.FilterObjectName(types.NamespacedName{Namespace: harness.Namespace, Name: name})) + if fetched == nil { + return nil, fmt.Errorf("config map %q not found", name) } + configMap := *fetched entries = append(entries, objectProvenance("v1", "ConfigMap", name, configMap.UID, configMap.Generation, configMap.Data)) } for _, template := range templates { @@ -217,10 +225,11 @@ func (c *Builder) BuildProvenance(ctx context.Context, harness *v1alpha3.Harness } switch binding.MCP.Server.Kind { case "RemoteMCPServer": - server := &v1alpha3.RemoteMCPServer{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: template.Namespace, Name: binding.MCP.Server.Name}, server); err != nil { - return nil, err + fetched := krt.FetchOne(c.ctx, c.collections.RemoteMCPServers, krt.FilterObjectName(types.NamespacedName{Namespace: template.Namespace, Name: binding.MCP.Server.Name})) + if fetched == nil { + return nil, fmt.Errorf("remote MCP server %q not found", binding.MCP.Server.Name) } + server := *fetched entries = append(entries, objectProvenance(v1alpha3.GroupVersion.String(), "RemoteMCPServer", server.Name, server.UID, server.Generation, server.Spec)) } } @@ -238,10 +247,11 @@ func (c *Builder) BuildProvenance(ctx context.Context, harness *v1alpha3.Harness continue } seenSecrets[identity] = struct{}{} - secret := &corev1.Secret{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: harness.Namespace, Name: ref.Name}, secret); err != nil { - return nil, err + fetched := krt.FetchOne(c.ctx, c.collections.Secrets, krt.FilterObjectName(types.NamespacedName{Namespace: harness.Namespace, Name: ref.Name})) + if fetched == nil { + return nil, fmt.Errorf("secret %q not found", ref.Name) } + secret := *fetched value, ok := secret.Data[ref.Key] if !ok { return nil, fmt.Errorf("secret %q does not contain key %q", ref.Name, ref.Key) @@ -295,13 +305,14 @@ func (c *Builder) resolveValueRef(ctx context.Context, namespace string, ref v1a if ref.ValueFrom.Type != v1alpha3.ConfigMapValueSource { return "", "", fmt.Errorf("unsupported value source type %q", ref.ValueFrom.Type) } - configMap := &corev1.ConfigMap{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: namespace, Name: ref.ValueFrom.Name}, configMap); err != nil { - return "", "", err + fetched := krt.FetchOne(c.ctx, c.collections.ConfigMaps, krt.FilterObjectName(types.NamespacedName{Namespace: namespace, Name: ref.ValueFrom.Name})) + if fetched == nil { + return "", "", fmt.Errorf("config map %q not found", ref.ValueFrom.Name) } + configMap := *fetched value, found := configMap.Data[ref.ValueFrom.Key] if !found { - return "", "", fmt.Errorf("ConfigMap %q does not contain key %q", ref.ValueFrom.Name, ref.ValueFrom.Key) + return "", "", fmt.Errorf("config map %q does not contain key %q", ref.ValueFrom.Name, ref.ValueFrom.Key) } return ref.Name, value, nil } diff --git a/go/core/v2/translator/adkconfig/model.go b/go/core/v2/translator/adkconfig/model.go index 9f7f20f44..30c4f6b0f 100644 --- a/go/core/v2/translator/adkconfig/model.go +++ b/go/core/v2/translator/adkconfig/model.go @@ -15,6 +15,7 @@ import ( "github.com/kagent-dev/kagent/go/core/internal/utils" "github.com/kagent-dev/kagent/go/core/pkg/env" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/types" ) @@ -208,7 +209,7 @@ func addTokenExchangeConfiguration(openai *adk.OpenAI, mdd *modelDeploymentData, // resolveFoundryEndpoint returns the Foundry endpoint, preferring the inline // value and otherwise resolving it from the referenced ConfigMap (endpointFrom), // which lets Azure Service Operator own the account endpoint. -func (c *Builder) resolveFoundryEndpoint(ctx context.Context, namespace string, cfg *v1alpha3.FoundryConfig) (string, error) { +func (c *Builder) resolveFoundryEndpoint(_ context.Context, namespace string, cfg *v1alpha3.FoundryConfig) (string, error) { if cfg.Endpoint != "" { return cfg.Endpoint, nil } @@ -216,10 +217,11 @@ func (c *Builder) resolveFoundryEndpoint(ctx context.Context, namespace string, return "", nil } ref := cfg.EndpointFrom - cm := &corev1.ConfigMap{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: namespace, Name: ref.Name}, cm); err != nil { - return "", fmt.Errorf("failed to get Foundry endpoint config map %s: %w", ref.Name, err) + fetched := krt.FetchOne(c.ctx, c.collections.ConfigMaps, krt.FilterObjectName(types.NamespacedName{Namespace: namespace, Name: ref.Name})) + if fetched == nil { + return "", fmt.Errorf("failed to get Foundry endpoint config map %s: not found", ref.Name) } + cm := *fetched value, ok := cm.Data[ref.Key] if !ok { if ref.Optional != nil && *ref.Optional { diff --git a/go/core/v2/translator/adkconfig/model_test.go b/go/core/v2/translator/adkconfig/model_test.go index d508d1abc..0c490bb7a 100644 --- a/go/core/v2/translator/adkconfig/model_test.go +++ b/go/core/v2/translator/adkconfig/model_test.go @@ -8,6 +8,9 @@ import ( "github.com/kagent-dev/kagent/go/core/pkg/env" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "github.com/stretchr/testify/require" + "istio.io/istio/pkg/kube/krt" + "istio.io/istio/pkg/kube/krt/krttest" + corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" ) @@ -56,3 +59,23 @@ func TestRenderBedrockCredentialsFromReferences(t *testing.T) { }) } } + +func TestResolveFoundryEndpointFromConfigMap(t *testing.T) { + configMap := &corev1.ConfigMap{ + ObjectMeta: metav1.ObjectMeta{Name: "account", Namespace: "test"}, + Data: map[string]string{"endpoint": "https://example.services.ai.azure.com"}, + } + mock := krttest.NewMock(t, []any{configMap}) + compiler := NewBuilder(krt.TestingDummyContext{}, v2translator.Collections{ + ConfigMaps: krttest.GetMockCollection[*corev1.ConfigMap](mock), + }) + + endpoint, err := compiler.resolveFoundryEndpoint(context.Background(), "test", &v1alpha3.FoundryConfig{ + EndpointFrom: &corev1.ConfigMapKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{Name: "account"}, + Key: "endpoint", + }, + }) + require.NoError(t, err) + require.Equal(t, configMap.Data["endpoint"], endpoint) +} diff --git a/go/core/v2/translator/byo/compiler.go b/go/core/v2/translator/byo/compiler.go index 81c07e075..d0fd84593 100644 --- a/go/core/v2/translator/byo/compiler.go +++ b/go/core/v2/translator/byo/compiler.go @@ -11,6 +11,7 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "github.com/kagent-dev/kagent/go/core/v2/translator/adkconfig" + "istio.io/istio/pkg/kube/krt" ) // Compiler translates resolved inputs into a BYO A2A runtime revision. @@ -18,8 +19,8 @@ type Compiler struct{ config *adkconfig.Builder } var _ v2translator.HarnessCompiler = (*Compiler)(nil) -func NewCompiler(kube v2translator.Reader) *Compiler { - return &Compiler{config: adkconfig.NewBuilder(kube)} +func NewCompiler(ctx krt.HandlerContext, collections v2translator.Collections) *Compiler { + return &Compiler{config: adkconfig.NewBuilder(ctx, collections)} } func (c *Compiler) Compile(ctx context.Context, input *v2translator.HarnessInput) (*v2translator.Revision, error) { diff --git a/go/core/v2/translator/byo/compiler_test.go b/go/core/v2/translator/byo/compiler_test.go index 3c1f0eff6..83947384b 100644 --- a/go/core/v2/translator/byo/compiler_test.go +++ b/go/core/v2/translator/byo/compiler_test.go @@ -9,20 +9,11 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "github.com/stretchr/testify/require" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime" - "k8s.io/apimachinery/pkg/types" ) -type reader struct{} - -func (reader) Get(context.Context, types.NamespacedName, runtime.Object) error { return nil } - -func (reader) GetResolvedModelConfig(context.Context, types.NamespacedName) (*v2translator.ResolvedModelConfig, error) { - return nil, nil -} - func TestCompileOpaqueImage(t *testing.T) { harness := &v1alpha3.Harness{ObjectMeta: metav1.ObjectMeta{Name: "byo", Namespace: "test"}, Spec: v1alpha3.HarnessSpec{ BYO: &v1alpha3.BYOHarness{}, @@ -36,7 +27,7 @@ func TestCompileOpaqueImage(t *testing.T) { Description: "custom A2A agent", SystemPrompt: "be helpful", }} - revision, err := NewCompiler(reader{}).Compile(context.Background(), &v2translator.HarnessInput{ + revision, err := NewCompiler(krt.TestingDummyContext{}, v2translator.Collections{}).Compile(context.Background(), &v2translator.HarnessInput{ Harness: harness, Root: &v2translator.AgentInput{Template: template, Instruction: template.Spec.SystemPrompt}, }) require.NoError(t, err) diff --git a/go/core/v2/translator/claude/compiler.go b/go/core/v2/translator/claude/compiler.go index 7e0070dfa..23b2167b0 100644 --- a/go/core/v2/translator/claude/compiler.go +++ b/go/core/v2/translator/claude/compiler.go @@ -16,6 +16,7 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" claudeconfig "github.com/kagent-dev/kagent/go/harness/claude/config" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/types" ) @@ -43,9 +44,14 @@ var ownedEnvironment = map[string]struct{}{ sandboxEnv: {}, claudeconfig.GoogleCredentialsJSONEnvName: {}, } -type Compiler struct{ kube v2translator.Reader } +type Compiler struct { + ctx krt.HandlerContext + collections v2translator.Collections +} -func NewCompiler(kube v2translator.Reader) *Compiler { return &Compiler{kube: kube} } +func NewCompiler(ctx krt.HandlerContext, collections v2translator.Collections) *Compiler { + return &Compiler{ctx: ctx, collections: collections} +} func (c *Compiler) Compile(ctx context.Context, input *v2translator.HarnessInput) (*v2translator.Revision, error) { if input == nil || input.Harness == nil || input.Root == nil || input.Root.Template == nil || input.Root.ResolvedModelConfig == nil || input.Root.ResolvedModelConfig.Config == nil { @@ -326,11 +332,11 @@ func (c *Compiler) requireGoogleCredentials(ctx context.Context, model *v1alpha3 } func (c *Compiler) secret(ctx context.Context, namespace, name string) (*corev1.Secret, error) { - secret := &corev1.Secret{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: namespace, Name: name}, secret); err != nil { - return nil, fmt.Errorf("read Claude credential Secret %q: %w", name, err) + secret := krt.FetchOne(c.ctx, c.collections.Secrets, krt.FilterObjectName(types.NamespacedName{Namespace: namespace, Name: name})) + if secret == nil { + return nil, fmt.Errorf("read Claude credential Secret %q: not found", name) } - return secret, nil + return *secret, nil } func secretEnvironment(environmentName, secretName, key string) corev1.EnvVar { @@ -398,11 +404,11 @@ func (c *Compiler) buildProvenance(ctx context.Context, input *v2translator.Harn } addAgent(input.Root) for name := range configMaps { - configMap := &corev1.ConfigMap{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: harness.Namespace, Name: name}, configMap); err != nil { - return nil, err + configMap := krt.FetchOne(c.ctx, c.collections.ConfigMaps, krt.FilterObjectName(types.NamespacedName{Namespace: harness.Namespace, Name: name})) + if configMap == nil { + return nil, fmt.Errorf("ConfigMap %q not found", name) } - entries = append(entries, objectProvenance("v1", "ConfigMap", name, configMap.UID, configMap.Generation, configMap.Data)) + entries = append(entries, objectProvenance("v1", "ConfigMap", name, (*configMap).UID, (*configMap).Generation, (*configMap).Data)) } seen := map[string]struct{}{} for _, variable := range environment { diff --git a/go/core/v2/translator/claude/compiler_test.go b/go/core/v2/translator/claude/compiler_test.go index a3b7cb6f6..015524317 100644 --- a/go/core/v2/translator/claude/compiler_test.go +++ b/go/core/v2/translator/claude/compiler_test.go @@ -13,13 +13,10 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" claudeconfig "github.com/kagent-dev/kagent/go/harness/claude/config" + "istio.io/istio/pkg/kube/krt" + "istio.io/istio/pkg/kube/krt/krttest" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime" - "k8s.io/apimachinery/pkg/types" - schemev1 "k8s.io/client-go/kubernetes/scheme" - "sigs.k8s.io/controller-runtime/pkg/client" - "sigs.k8s.io/controller-runtime/pkg/client/fake" ) const credentialValue = "credential-must-not-be-serialized" @@ -82,7 +79,7 @@ func TestCompileSupportedProviders(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { input, reader := testInput(t, tt.model, tt.secretData) - revision, err := NewCompiler(reader).Compile(context.Background(), input) + revision, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatal(err) } @@ -118,7 +115,7 @@ func TestCompileSupportedProviders(t *testing.T) { t.Fatalf("provenance omits credential Secret: %s", revision.Provenance) } - again, err := NewCompiler(reader).Compile(context.Background(), input) + again, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil || !reflect.DeepEqual(revision, again) { t.Fatalf("compilation is not deterministic: %v", err) } @@ -143,7 +140,7 @@ func TestCompileRejectsUnsupportedConfiguration(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { input, reader := testInput(t, tt.model, map[string][]byte{"api-key": []byte("secret"), awsAccessKeyEnv: []byte("access"), awsSecretKeyEnv: []byte("secret"), "credentials.json": []byte(`{"type":"service_account"}`)}) - _, err := NewCompiler(reader).Compile(context.Background(), input) + _, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) var validation *v2translator.ValidationError if !errors.As(err, &validation) { t.Fatalf("Compile() error = %v, want validation error", err) @@ -160,7 +157,7 @@ func TestCompileRejectsProviderOwnedHarnessEnvironment(t *testing.T) { input, reader := testInput(t, model, map[string][]byte{"api-key": []byte("secret")}) value := "http://mock.example.com" input.Harness.Spec.Env = []v1alpha3.HarnessEnvVar{{Name: anthropicBaseURLEnv, Value: &value}} - _, err := NewCompiler(reader).Compile(context.Background(), input) + _, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) var validation *v2translator.ValidationError if !errors.As(err, &validation) { t.Fatalf("Compile() error = %v, want validation error", err) @@ -183,7 +180,7 @@ func TestCompileRootSkillsAndPluginSelections(t *testing.T) { Skills: []string{"deploy"}, }} - revision, err := NewCompiler(reader).Compile(context.Background(), input) + revision, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatal(err) } @@ -201,7 +198,7 @@ func TestCompileRootSkillsAndPluginSelections(t *testing.T) { } input.Root.Template.Spec.Plugins[0].Skills = []string{"review"} - if _, err := NewCompiler(reader).Compile(context.Background(), input); err == nil || !strings.Contains(err.Error(), "duplicate skill name") { + if _, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input); err == nil || !strings.Contains(err.Error(), "duplicate skill name") { t.Fatalf("duplicate skill Compile() error = %v", err) } } @@ -235,7 +232,7 @@ func TestCompileDirectWholeServerMCP(t *testing.T) { }}} input.Root.MCPTools = []v2translator.ResolvedMCPTool{{Binding: *input.Root.Template.Spec.Tools[0].MCP.DeepCopy(), Server: server}} - revision, err := NewCompiler(reader).Compile(context.Background(), input) + revision, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatal(err) } @@ -286,7 +283,7 @@ func TestCompileWholeServerMCPSelectionWarnings(t *testing.T) { } binding := v1alpha3.MCPToolBinding{Server: corev1.TypedLocalObjectReference{Kind: "RemoteMCPServer", Name: server.Name}} input.Root.MCPTools = []v2translator.ResolvedMCPTool{{Binding: binding, Server: server}} - revision, err := NewCompiler(reader).Compile(context.Background(), input) + revision, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatalf("omitted selection Compile() error = %v", err) } @@ -295,7 +292,7 @@ func TestCompileWholeServerMCPSelectionWarnings(t *testing.T) { } input.Root.MCPTools[0].Binding.Tools = []string{"one"} - revision, err = NewCompiler(reader).Compile(context.Background(), input) + revision, err = NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatalf("partial selection Compile() error = %v", err) } @@ -304,7 +301,7 @@ func TestCompileWholeServerMCPSelectionWarnings(t *testing.T) { } server.Status.ObservedGeneration = 0 - revision, err = NewCompiler(reader).Compile(context.Background(), input) + revision, err = NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatalf("stale discovery Compile() error = %v", err) } @@ -345,7 +342,7 @@ func TestCompileLocalSharedAgent(t *testing.T) { Name: "specialist", Description: "Handles specialist requests", Agent: child, }} - revision, err := NewCompiler(reader).Compile(context.Background(), input) + revision, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err != nil { t.Fatal(err) } @@ -407,7 +404,7 @@ func TestCompileRejectsUnsupportedLocalAgentConfiguration(t *testing.T) { } tt.mutate(&binding) input.Root.Shared = []v2translator.AgentInputBinding{binding} - _, err := NewCompiler(reader).Compile(context.Background(), input) + _, err := NewCompiler(krt.TestingDummyContext{}, reader).Compile(context.Background(), input) if err == nil || !strings.Contains(err.Error(), tt.want) { t.Fatalf("Compile() error = %v, want containing %q", err, tt.want) } @@ -415,11 +412,8 @@ func TestCompileRejectsUnsupportedLocalAgentConfiguration(t *testing.T) { } } -func testInput(t *testing.T, modelSpec v1alpha3.ModelConfigSpec, secretData map[string][]byte) (*v2translator.HarnessInput, v2translator.Reader) { +func testInput(t *testing.T, modelSpec v1alpha3.ModelConfigSpec, secretData map[string][]byte) (*v2translator.HarnessInput, v2translator.Collections) { t.Helper() - if err := v1alpha3.AddToScheme(schemev1.Scheme); err != nil { - t.Fatal(err) - } harness := &v1alpha3.Harness{ObjectMeta: metav1.ObjectMeta{Name: "claude", Namespace: "test", UID: "harness-uid"}, Spec: v1alpha3.HarnessSpec{ Claude: &v1alpha3.ClaudeHarness{}, Workload: v1alpha3.HarnessWorkload{Image: "example.com/claude@sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}, Substrate: v1alpha3.HarnessSubstratePolicy{WorkerPoolRef: corev1.LocalObjectReference{Name: "default"}, SnapshotPolicy: v1alpha3.HarnessSnapshotPolicy{Location: "snapshots"}}, @@ -429,21 +423,10 @@ func testInput(t *testing.T, modelSpec v1alpha3.ModelConfigSpec, secretData map[ }} model := &v1alpha3.ModelConfig{ObjectMeta: metav1.ObjectMeta{Name: "model", Namespace: "test", UID: "model-uid"}, Spec: modelSpec} secret := &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "model-auth", Namespace: "test", UID: "secret-uid"}, Data: secretData} - kube := fake.NewClientBuilder().WithScheme(schemev1.Scheme).WithObjects(secret).Build() - reader := testReader{kube} - return &v2translator.HarnessInput{Harness: harness, Root: &v2translator.AgentInput{Template: template, ResolvedModelConfig: &v2translator.ResolvedModelConfig{Config: model}, Instruction: "help carefully"}}, reader -} - -type testReader struct{ client.Client } - -func (r testReader) Get(ctx context.Context, key types.NamespacedName, object runtime.Object) error { - return r.Client.Get(ctx, key, object.(client.Object)) -} - -func (r testReader) GetResolvedModelConfig(ctx context.Context, key types.NamespacedName) (*v2translator.ResolvedModelConfig, error) { - model := &v1alpha3.ModelConfig{} - if err := r.Get(ctx, key, model); err != nil { - return nil, err + mock := krttest.NewMock(t, []any{secret}) + collections := v2translator.Collections{ + Secrets: krttest.GetMockCollection[*corev1.Secret](mock), + ConfigMaps: krttest.GetMockCollection[*corev1.ConfigMap](mock), } - return v2translator.ResolveModelConfig(ctx, r, model) + return &v2translator.HarnessInput{Harness: harness, Root: &v2translator.AgentInput{Template: template, ResolvedModelConfig: &v2translator.ResolvedModelConfig{Config: model}, Instruction: "help carefully"}}, collections } diff --git a/go/core/v2/translator/claude/mcp.go b/go/core/v2/translator/claude/mcp.go index 6f7fffabd..d47920393 100644 --- a/go/core/v2/translator/claude/mcp.go +++ b/go/core/v2/translator/claude/mcp.go @@ -12,6 +12,7 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" claudeconfig "github.com/kagent-dev/kagent/go/harness/claude/config" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/types" ) @@ -157,11 +158,12 @@ func (c *Compiler) compileMCPHeaders(ctx context.Context, namespace string, refs case ref.ValueFrom == nil: headers[ref.Name] = ref.Value case ref.ValueFrom.Type == v1alpha3.ConfigMapValueSource: - configMap := &corev1.ConfigMap{} key := types.NamespacedName{Namespace: namespace, Name: ref.ValueFrom.Name} - if err := c.kube.Get(ctx, key, configMap); err != nil { - return nil, nil, err + fetched := krt.FetchOne(c.ctx, c.collections.ConfigMaps, krt.FilterObjectName(key)) + if fetched == nil { + return nil, nil, fmt.Errorf("ConfigMap %q not found", ref.ValueFrom.Name) } + configMap := *fetched value, exists := configMap.Data[ref.ValueFrom.Key] if !exists { return nil, nil, fmt.Errorf("ConfigMap %q does not contain key %q", configMap.Name, ref.ValueFrom.Key) diff --git a/go/core/v2/translator/collections.go b/go/core/v2/translator/collections.go new file mode 100644 index 000000000..40960537f --- /dev/null +++ b/go/core/v2/translator/collections.go @@ -0,0 +1,19 @@ +package translator + +import ( + atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1" + "github.com/kagent-dev/kagent/go/api/v1alpha3" + "istio.io/istio/pkg/kube/krt" + corev1 "k8s.io/api/core/v1" +) + +// Collections contains every typed input used while compiling a revision. +// Production supplies informer-backed collections; tests supply KRT mocks. +type Collections struct { + AgentTemplates krt.Collection[*v1alpha3.AgentTemplate] + ResolvedModelConfigs krt.Collection[ResolvedModelConfig] + RemoteMCPServers krt.Collection[*v1alpha3.RemoteMCPServer] + ConfigMaps krt.Collection[*corev1.ConfigMap] + Secrets krt.Collection[*corev1.Secret] + WorkerPools krt.Collection[*atev1alpha1.WorkerPool] +} diff --git a/go/core/v2/translator/compiler.go b/go/core/v2/translator/compiler.go index 90d9c3880..e535b8ed2 100644 --- a/go/core/v2/translator/compiler.go +++ b/go/core/v2/translator/compiler.go @@ -6,6 +6,7 @@ import ( "maps" "github.com/kagent-dev/kagent/go/api/v1alpha3" + "istio.io/istio/pkg/kube/krt" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/types" @@ -15,7 +16,8 @@ import ( // revision. It owns the v2 translation boundary rather than delegating to an // earlier API translator. type Compiler struct { - kube Reader + ctx krt.HandlerContext + collections Collections harnessCompilers map[HarnessType]HarnessCompiler } @@ -83,8 +85,8 @@ type AgentInputBinding struct { } // NewCompiler constructs the v2 runtime compiler. -func NewCompiler(kube Reader, harnessCompilers map[HarnessType]HarnessCompiler) *Compiler { - return &Compiler{kube: kube, harnessCompilers: maps.Clone(harnessCompilers)} +func NewCompiler(ctx krt.HandlerContext, collections Collections, harnessCompilers map[HarnessType]HarnessCompiler) *Compiler { + return &Compiler{ctx: ctx, collections: collections, harnessCompilers: maps.Clone(harnessCompilers)} } // CompileAgentTemplate resolves an API v2 attachment into an immutable runtime @@ -156,12 +158,12 @@ func (c *Compiler) resolveTree(ctx context.Context, harness *v1alpha3.Harness, r return nil, NewValidationError("duplicate Shared AgentTemplate binding name %q", binding.Name) } names[binding.Name] = struct{}{} - childTemplate := &v1alpha3.AgentTemplate{} key := types.NamespacedName{Namespace: template.Namespace, Name: binding.TemplateRef.Name} - if err := c.kube.Get(ctx, key, childTemplate); err != nil { - return nil, fmt.Errorf("resolve AgentTemplate %q: %w", binding.TemplateRef.Name, err) + childTemplate := krt.FetchOne(c.ctx, c.collections.AgentTemplates, krt.FilterObjectName(key)) + if childTemplate == nil { + return nil, fmt.Errorf("resolve AgentTemplate %q: not found", binding.TemplateRef.Name) } - agent, err := resolve(childTemplate, true) + agent, err := resolve(*childTemplate, true) if err != nil { return nil, err } @@ -201,9 +203,9 @@ func (c *Compiler) buildInputs(ctx context.Context, tree *ResolvedTree) (*Harnes } input := &AgentInput{Template: template, Instruction: instruction} if template.Spec.ModelConfig != nil { - input.ResolvedModelConfig, err = c.kube.GetResolvedModelConfig(ctx, types.NamespacedName{Namespace: template.Namespace, Name: template.Spec.ModelConfig.Name}) - if err != nil { - return nil, fmt.Errorf("resolve ModelConfig %q: %w", template.Spec.ModelConfig.Name, err) + input.ResolvedModelConfig = krt.FetchOne(c.ctx, c.collections.ResolvedModelConfigs, krt.FilterObjectName(types.NamespacedName{Namespace: template.Namespace, Name: template.Spec.ModelConfig.Name})) + if input.ResolvedModelConfig == nil { + return nil, fmt.Errorf("resolve ModelConfig %q: not found", template.Spec.ModelConfig.Name) } if failures := input.ResolvedModelConfig.SemanticFailures; len(failures) > 0 { return nil, NewValidationError("ModelConfig %q: %s", template.Spec.ModelConfig.Name, failures[0].Message) @@ -223,12 +225,12 @@ func (c *Compiler) buildInputs(ctx context.Context, tree *ResolvedTree) (*Harnes if tool.MCP.Server.Kind != "RemoteMCPServer" { return nil, NewValidationError("unsupported MCP server kind %q", tool.MCP.Server.Kind) } - server := &v1alpha3.RemoteMCPServer{} key := types.NamespacedName{Namespace: template.Namespace, Name: tool.MCP.Server.Name} - if err := c.kube.Get(ctx, key, server); err != nil { - return nil, fmt.Errorf("resolve %s %q: %w", tool.MCP.Server.Kind, tool.MCP.Server.Name, err) + server := krt.FetchOne(c.ctx, c.collections.RemoteMCPServers, krt.FilterObjectName(key)) + if server == nil { + return nil, fmt.Errorf("resolve %s %q: not found", tool.MCP.Server.Kind, tool.MCP.Server.Name) } - input.MCPTools = append(input.MCPTools, ResolvedMCPTool{Binding: *tool.MCP.DeepCopy(), Server: server}) + input.MCPTools = append(input.MCPTools, ResolvedMCPTool{Binding: *tool.MCP.DeepCopy(), Server: *server}) toolNames = append(toolNames, tool.MCP.Tools...) } if template.Spec.PromptTemplate != nil { @@ -236,7 +238,7 @@ func (c *Compiler) buildInputs(ctx context.Context, tree *ResolvedTree) (*Harnes for _, source := range template.Spec.PromptTemplate.DataSources { refs = append(refs, promptSourceRef{Name: source.Name, Alias: source.Alias}) } - lookup, err := resolvePromptSourceRefs(ctx, c.kube, template.Namespace, refs) + lookup, err := resolvePromptSourceRefs(c.ctx, c.collections.ConfigMaps, template.Namespace, refs) if err != nil { return nil, fmt.Errorf("resolve prompt sources: %w", err) } diff --git a/go/core/v2/translator/compiler_test.go b/go/core/v2/translator/compiler_test.go index 41e65c276..29342e0e9 100644 --- a/go/core/v2/translator/compiler_test.go +++ b/go/core/v2/translator/compiler_test.go @@ -12,13 +12,11 @@ import ( v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" kagenttranslator "github.com/kagent-dev/kagent/go/core/v2/translator/kagent" "github.com/stretchr/testify/require" + "istio.io/istio/pkg/kube/krt" + "istio.io/istio/pkg/kube/krt/krttest" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" - schemev1 "k8s.io/client-go/kubernetes/scheme" - "sigs.k8s.io/controller-runtime/pkg/client" - "sigs.k8s.io/controller-runtime/pkg/client/fake" ) func modelConfig() *v1alpha3.ModelConfig { @@ -106,32 +104,37 @@ func remoteMCPServer(name, url string) *v1alpha3.RemoteMCPServer { }} } -func compiler(t *testing.T, objects ...client.Object) *v2translator.Compiler { +func compiler(t *testing.T, objects ...any) *v2translator.Compiler { t.Helper() - require.NoError(t, v1alpha3.AddToScheme(schemev1.Scheme)) - kube := fake.NewClientBuilder().WithScheme(schemev1.Scheme).WithObjects(objects...).Build() - reader := testReader{kube} - return v2translator.NewCompiler(reader, map[v2translator.HarnessType]v2translator.HarnessCompiler{ - v2translator.HarnessTypeKagent: kagenttranslator.NewCompiler(reader), + collections := mockCollections(t, objects...) + ctx := krt.TestingDummyContext{} + return v2translator.NewCompiler(ctx, collections, map[v2translator.HarnessType]v2translator.HarnessCompiler{ + v2translator.HarnessTypeKagent: kagenttranslator.NewCompiler(ctx, collections), }) } -type testReader struct{ client.Client } - -func (r testReader) Get(ctx context.Context, key types.NamespacedName, object runtime.Object) error { - return r.Client.Get(ctx, key, object.(client.Object)) -} - -func (r testReader) GetResolvedModelConfig(ctx context.Context, key types.NamespacedName) (*v2translator.ResolvedModelConfig, error) { - model := &v1alpha3.ModelConfig{} - if err := r.Get(ctx, key, model); err != nil { - return nil, err +func mockCollections(t *testing.T, objects ...any) v2translator.Collections { + t.Helper() + mock := krttest.NewMock(t, objects) + collections := v2translator.Collections{ + AgentTemplates: krttest.GetMockCollection[*v1alpha3.AgentTemplate](mock), + RemoteMCPServers: krttest.GetMockCollection[*v1alpha3.RemoteMCPServer](mock), + ConfigMaps: krttest.GetMockCollection[*corev1.ConfigMap](mock), + Secrets: krttest.GetMockCollection[*corev1.Secret](mock), } - return v2translator.ResolveModelConfig(ctx, r, model) + models := krttest.GetMockCollection[*v1alpha3.ModelConfig](mock) + resolved := make([]any, 0, len(models.List())) + for _, model := range models.List() { + value, err := v2translator.ResolveModelConfig(krt.TestingDummyContext{}, collections, model) + require.NoError(t, err) + resolved = append(resolved, *value) + } + resolvedMock := krttest.NewMock(t, resolved) + collections.ResolvedModelConfigs = krttest.GetMockCollection[v2translator.ResolvedModelConfig](resolvedMock) + return collections } func TestResolveModelConfigRecordsFoundryEndpointReference(t *testing.T) { - require.NoError(t, v1alpha3.AddToScheme(schemev1.Scheme)) model := &v1alpha3.ModelConfig{ ObjectMeta: metav1.ObjectMeta{Name: "foundry", Namespace: "test"}, Spec: v1alpha3.ModelConfigSpec{ @@ -145,10 +148,8 @@ func TestResolveModelConfigRecordsFoundryEndpointReference(t *testing.T) { ObjectMeta: metav1.ObjectMeta{Name: "account", Namespace: "test"}, Data: map[string]string{"endpoint": "https://example.services.ai.azure.com"}, } - reader := testReader{fake.NewClientBuilder().WithScheme(schemev1.Scheme).WithObjects(model, configMap).Build()} - - resolved, err := reader.GetResolvedModelConfig(context.Background(), types.NamespacedName{Namespace: "test", Name: "foundry"}) - require.NoError(t, err) + collections := mockCollections(t, model, configMap) + resolved := collections.ResolvedModelConfigs.List()[0] require.Equal(t, model.Spec, resolved.Config.Spec) require.Equal(t, []v2translator.ModelConfigReference{{ NamespacedName: types.NamespacedName{Namespace: "test", Name: "account"}, Kind: "ConfigMap", Key: "endpoint", @@ -163,8 +164,7 @@ func (c *testHarnessCompiler) Compile(_ context.Context, input *v2translator.Har } func TestCompilerAcceptsExternalHarnessCompiler(t *testing.T) { - require.NoError(t, v1alpha3.AddToScheme(schemev1.Scheme)) - kube := fake.NewClientBuilder().WithScheme(schemev1.Scheme).WithObjects(modelConfig()).Build() + collections := mockCollections(t, modelConfig()) adapter := &testHarnessCompiler{} harness := &v1alpha3.Harness{ ObjectMeta: metav1.ObjectMeta{Name: "codex", Namespace: "test"}, @@ -172,7 +172,7 @@ func TestCompilerAcceptsExternalHarnessCompiler(t *testing.T) { } template := &v1alpha3.AgentTemplate{ObjectMeta: metav1.ObjectMeta{Name: "assistant", Namespace: "test"}, Spec: v1alpha3.AgentTemplateSpec{ModelConfig: &corev1.LocalObjectReference{Name: "default-model"}}} - revision, err := v2translator.NewCompiler(testReader{kube}, map[v2translator.HarnessType]v2translator.HarnessCompiler{ + revision, err := v2translator.NewCompiler(krt.TestingDummyContext{}, collections, map[v2translator.HarnessType]v2translator.HarnessCompiler{ v2translator.HarnessTypeCodex: adapter, }).CompileAgentTemplate(context.Background(), harness, template) require.NoError(t, err) @@ -182,11 +182,10 @@ func TestCompilerAcceptsExternalHarnessCompiler(t *testing.T) { } func TestCompilerRejectsUnusableModelConfigBeforeHarnessCompiler(t *testing.T) { - require.NoError(t, v1alpha3.AddToScheme(schemev1.Scheme)) model := modelConfig() model.Spec.APIKeySecret = "missing" model.Spec.APIKeySecretKey = "key" - kube := fake.NewClientBuilder().WithScheme(schemev1.Scheme).WithObjects(model).Build() + collections := mockCollections(t, model) adapter := &testHarnessCompiler{} harness := &v1alpha3.Harness{ ObjectMeta: metav1.ObjectMeta{Name: "codex", Namespace: "test"}, @@ -194,7 +193,7 @@ func TestCompilerRejectsUnusableModelConfigBeforeHarnessCompiler(t *testing.T) { } template := &v1alpha3.AgentTemplate{ObjectMeta: metav1.ObjectMeta{Name: "assistant", Namespace: "test"}, Spec: v1alpha3.AgentTemplateSpec{ModelConfig: &corev1.LocalObjectReference{Name: model.Name}}} - _, err := v2translator.NewCompiler(testReader{kube}, map[v2translator.HarnessType]v2translator.HarnessCompiler{ + _, err := v2translator.NewCompiler(krt.TestingDummyContext{}, collections, map[v2translator.HarnessType]v2translator.HarnessCompiler{ v2translator.HarnessTypeCodex: adapter, }).CompileAgentTemplate(context.Background(), harness, template) require.ErrorContains(t, err, `resolve ModelConfig "default-model": secret missing not found`) @@ -202,15 +201,13 @@ func TestCompilerRejectsUnusableModelConfigBeforeHarnessCompiler(t *testing.T) { } func TestCompilerPermitsBYOWithoutModelConfig(t *testing.T) { - require.NoError(t, v1alpha3.AddToScheme(schemev1.Scheme)) - kube := fake.NewClientBuilder().WithScheme(schemev1.Scheme).Build() adapter := &testHarnessCompiler{} harness := &v1alpha3.Harness{ObjectMeta: metav1.ObjectMeta{Name: "byo", Namespace: "test"}, Spec: v1alpha3.HarnessSpec{ BYO: &v1alpha3.BYOHarness{}, AllowedAgentTemplates: &v1alpha3.HarnessAgentTemplateAdmission{Selector: metav1.LabelSelector{}}, }} template := &v1alpha3.AgentTemplate{ObjectMeta: metav1.ObjectMeta{Name: "assistant", Namespace: "test"}} - _, err := v2translator.NewCompiler(testReader{kube}, map[v2translator.HarnessType]v2translator.HarnessCompiler{ + _, err := v2translator.NewCompiler(krt.TestingDummyContext{}, mockCollections(t), map[v2translator.HarnessType]v2translator.HarnessCompiler{ v2translator.HarnessTypeBYO: adapter, }).CompileAgentTemplate(context.Background(), harness, template) require.NoError(t, err) diff --git a/go/core/v2/translator/kagent/agentcard_test.go b/go/core/v2/translator/kagent/agentcard_test.go index f5e19e4cc..a5967bb33 100644 --- a/go/core/v2/translator/kagent/agentcard_test.go +++ b/go/core/v2/translator/kagent/agentcard_test.go @@ -7,11 +7,12 @@ import ( "github.com/kagent-dev/kagent/go/api/v1alpha3" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" + "istio.io/istio/pkg/kube/krt" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) func TestCompilerRequiresModelConfig(t *testing.T) { - _, err := NewCompiler(nil).Compile(context.Background(), &v2translator.HarnessInput{Root: &v2translator.AgentInput{ + _, err := NewCompiler(krt.TestingDummyContext{}, v2translator.Collections{}).Compile(context.Background(), &v2translator.HarnessInput{Root: &v2translator.AgentInput{ Template: &v1alpha3.AgentTemplate{}, }}) if err == nil || !strings.Contains(err.Error(), "kagent ModelConfig is required") { diff --git a/go/core/v2/translator/kagent/compiler.go b/go/core/v2/translator/kagent/compiler.go index 9775f8580..563963466 100644 --- a/go/core/v2/translator/kagent/compiler.go +++ b/go/core/v2/translator/kagent/compiler.go @@ -14,6 +14,7 @@ import ( "github.com/kagent-dev/kagent/go/core/pkg/env" v2translator "github.com/kagent-dev/kagent/go/core/v2/translator" "github.com/kagent-dev/kagent/go/core/v2/translator/adkconfig" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" ) @@ -26,8 +27,8 @@ type Compiler struct { var _ v2translator.HarnessCompiler = (*Compiler)(nil) -func NewCompiler(kube v2translator.Reader) *Compiler { - return &Compiler{config: adkconfig.NewBuilder(kube)} +func NewCompiler(ctx krt.HandlerContext, collections v2translator.Collections) *Compiler { + return &Compiler{config: adkconfig.NewBuilder(ctx, collections)} } func (c *Compiler) Compile(ctx context.Context, input *v2translator.HarnessInput) (*v2translator.Revision, error) { diff --git a/go/core/v2/translator/modelconfig.go b/go/core/v2/translator/modelconfig.go index 30ab8c2e9..7dd7a8b9b 100644 --- a/go/core/v2/translator/modelconfig.go +++ b/go/core/v2/translator/modelconfig.go @@ -1,12 +1,13 @@ package translator import ( - "context" "fmt" "github.com/kagent-dev/kagent/go/api/v1alpha3" "github.com/kagent-dev/kagent/go/core/pkg/env" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" + apiequality "k8s.io/apimachinery/pkg/api/equality" "k8s.io/apimachinery/pkg/types" ) @@ -34,9 +35,41 @@ type ResolvedModelConfig struct { ReferenceFailures []ModelConfigFailure } +func (r ResolvedModelConfig) ResourceName() string { + if r.Config == nil { + return "" + } + return r.Config.Namespace + "/" + r.Config.Name +} + +func (r ResolvedModelConfig) Equals(other ResolvedModelConfig) bool { + return apiequality.Semantic.DeepEqual(r, other) +} + +// Usable reports whether every intrinsic configuration requirement and every +// referenced Kubernetes input was resolved. +func (r *ResolvedModelConfig) Usable() bool { + return r != nil && r.Config != nil && len(r.SemanticFailures) == 0 && len(r.ReferenceFailures) == 0 +} + +// Failure returns the first diagnostic suitable for reporting at a compilation +// boundary. Reconciliation should inspect both failure lists to report status. +func (r *ResolvedModelConfig) Failure() *ModelConfigFailure { + if r == nil || r.Config == nil { + return &ModelConfigFailure{Reason: "ModelConfigMissing", Message: "model config is required"} + } + if len(r.SemanticFailures) > 0 { + return &r.SemanticFailures[0] + } + if len(r.ReferenceFailures) > 0 { + return &r.ReferenceFailures[0] + } + return nil +} + // ResolveModelConfig validates and resolves ModelConfig data shared by every // harness. It does not expose secret values or produce runtime-specific inputs. -func ResolveModelConfig(ctx context.Context, kube Reader, config *v1alpha3.ModelConfig) (*ResolvedModelConfig, error) { +func ResolveModelConfig(ctx krt.HandlerContext, collections Collections, config *v1alpha3.ModelConfig) (*ResolvedModelConfig, error) { if config == nil { return &ResolvedModelConfig{SemanticFailures: []ModelConfigFailure{{Reason: "ModelConfigMissing", Message: "model config is required"}}}, nil } @@ -49,11 +82,12 @@ func ResolveModelConfig(ctx context.Context, kube Reader, config *v1alpha3.Model } requireSecret := func(name, notFoundReason, keyNotFoundReason string, keys ...string) *corev1.Secret { key := types.NamespacedName{Namespace: config.Namespace, Name: name} - secret := &corev1.Secret{} - if err := kube.Get(ctx, key, secret); err != nil { + fetched := krt.FetchOne(ctx, collections.Secrets, krt.FilterObjectName(key)) + if fetched == nil { addReferenceFailure(notFoundReason, fmt.Sprintf("secret %s not found", name)) return nil } + secret := *fetched for _, secretKey := range keys { if _, ok := secret.Data[secretKey]; !ok { addReferenceFailure(keyNotFoundReason, fmt.Sprintf("secret %s does not contain key %q", name, secretKey)) @@ -118,12 +152,13 @@ func ResolveModelConfig(ctx context.Context, kube Reader, config *v1alpha3.Model break } if !config.Spec.APIKeyPassthrough && config.Spec.APIKeySecret != "" { - secret := &corev1.Secret{} key := types.NamespacedName{Namespace: config.Namespace, Name: config.Spec.APIKeySecret} - if err := kube.Get(ctx, key, secret); err != nil { + fetched := krt.FetchOne(ctx, collections.Secrets, krt.FilterObjectName(key)) + if fetched == nil { addReferenceFailure("APIKeySecretNotFound", fmt.Sprintf("secret %s not found", config.Spec.APIKeySecret)) break } + secret := *fetched if _, ok := secret.Data[env.AWSBearerTokenBedrock.Name()]; ok { resolved.References = append(resolved.References, ModelConfigReference{NamespacedName: key, Kind: "Secret", Key: env.AWSBearerTokenBedrock.Name()}) } else { @@ -149,10 +184,11 @@ func ResolveModelConfig(ctx context.Context, kube Reader, config *v1alpha3.Model ref := config.Spec.Foundry.EndpointFrom key := types.NamespacedName{Namespace: config.Namespace, Name: ref.Name} resolved.References = append(resolved.References, ModelConfigReference{NamespacedName: key, Kind: "ConfigMap", Key: ref.Key}) - configMap := &corev1.ConfigMap{} - if err := kube.Get(ctx, key, configMap); err != nil { + fetched := krt.FetchOne(ctx, collections.ConfigMaps, krt.FilterObjectName(key)) + if fetched == nil { addReferenceFailure("EndpointConfigMapNotFound", fmt.Sprintf("config map %s not found", ref.Name)) } else { + configMap := *fetched _, ok := configMap.Data[ref.Key] if !ok { addReferenceFailure("EndpointConfigMapKeyNotFound", fmt.Sprintf("config map %s does not contain key %q", ref.Name, ref.Key)) diff --git a/go/core/v2/translator/reader.go b/go/core/v2/translator/reader.go deleted file mode 100644 index 034413e81..000000000 --- a/go/core/v2/translator/reader.go +++ /dev/null @@ -1,15 +0,0 @@ -package translator - -import ( - "context" - - "k8s.io/apimachinery/pkg/runtime" - "k8s.io/apimachinery/pkg/types" -) - -// Reader is the only Kubernetes capability used while compiling a revision. -// KRT supplies a dependency-tracking implementation; tests may use any reader. -type Reader interface { - Get(context.Context, types.NamespacedName, runtime.Object) error - GetResolvedModelConfig(context.Context, types.NamespacedName) (*ResolvedModelConfig, error) -} diff --git a/go/core/v2/translator/template.go b/go/core/v2/translator/template.go index 49fc1c8a6..3877990ba 100644 --- a/go/core/v2/translator/template.go +++ b/go/core/v2/translator/template.go @@ -8,6 +8,7 @@ import ( "text/template" "github.com/kagent-dev/kagent/go/api/v1alpha3" + "istio.io/istio/pkg/kube/krt" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/types" ) @@ -26,14 +27,14 @@ type PromptTemplateContext struct { ToolNames []string } -func (c *Compiler) resolveAgentTemplatePrompt(ctx context.Context, agentTemplate *v1alpha3.AgentTemplate) (string, error) { +func (c *Compiler) resolveAgentTemplatePrompt(_ context.Context, agentTemplate *v1alpha3.AgentTemplate) (string, error) { if agentTemplate.Spec.SystemPromptFrom != nil { ref := agentTemplate.Spec.SystemPromptFrom - configMap := &corev1.ConfigMap{} - if err := c.kube.Get(ctx, types.NamespacedName{Namespace: agentTemplate.Namespace, Name: ref.Name}, configMap); err != nil { - return "", fmt.Errorf("resolve systemPromptFrom: %w", err) + configMap := krt.FetchOne(c.ctx, c.collections.ConfigMaps, krt.FilterObjectName(types.NamespacedName{Namespace: agentTemplate.Namespace, Name: ref.Name})) + if configMap == nil { + return "", fmt.Errorf("resolve systemPromptFrom: ConfigMap %q not found", ref.Name) } - value, found := configMap.Data[ref.Key] + value, found := (*configMap).Data[ref.Key] if !found { return "", fmt.Errorf("resolve systemPromptFrom: ConfigMap %q does not contain key %q", ref.Name, ref.Key) } @@ -51,18 +52,18 @@ type promptSourceRef struct { // resolvePromptSourceRefs flattens ConfigMap keys into "source/key" identifiers // and rejects collisions before template execution. -func resolvePromptSourceRefs(ctx context.Context, kube Reader, namespace string, sources []promptSourceRef) (map[string]string, error) { +func resolvePromptSourceRefs(ctx krt.HandlerContext, configMaps krt.Collection[*corev1.ConfigMap], namespace string, sources []promptSourceRef) (map[string]string, error) { lookup := make(map[string]string) for _, source := range sources { identifier := source.Name if source.Alias != "" { identifier = source.Alias } - configMap := &corev1.ConfigMap{} - if err := kube.Get(ctx, types.NamespacedName{Namespace: namespace, Name: source.Name}, configMap); err != nil { - return nil, fmt.Errorf("resolve prompt source %q: %w", source.Name, err) + configMap := krt.FetchOne(ctx, configMaps, krt.FilterObjectName(types.NamespacedName{Namespace: namespace, Name: source.Name})) + if configMap == nil { + return nil, fmt.Errorf("resolve prompt source %q: ConfigMap not found", source.Name) } - for key, value := range configMap.Data { + for key, value := range (*configMap).Data { lookupKey := identifier + "/" + key if _, exists := lookup[lookupKey]; exists { return nil, fmt.Errorf("duplicate prompt template identifier %q", lookupKey)