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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 33 additions & 26 deletions go/core/v2/controller/collections.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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.
Expand Down Expand Up @@ -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,
}
}

Expand Down
86 changes: 57 additions & 29 deletions go/core/v2/controller/collections_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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) {
Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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 {
Expand Down
36 changes: 5 additions & 31 deletions go/core/v2/controller/modelconfig.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
package controller

import (
"context"
"crypto/sha256"
"encoding/hex"
"slices"
Expand All @@ -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
}
Expand All @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading