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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions go/core/internal/controller/collections.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ type Collections struct {
ModelConfigStatuses krt.StatusCollection[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus]
ResolvedModelConfigs krt.Collection[v2translator.ResolvedModelConfig]
AgentTemplateStatuses krt.StatusCollection[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]
HarnessStatuses krt.StatusCollection[*kagentv1alpha3.Harness, kagentv1alpha3.HarnessStatus]
}

// PairRuntimeObservation records a pair's preparation. The revision prevents a
Expand Down Expand Up @@ -76,6 +77,7 @@ func NewCollections(client kube.Client, watchNamespaces []string, opts krt.Optio
}
reconciliations := newPairReconciliations(pairs, compilerCollections, pairRuntimeObservations, opts)
statuses := newAgentTemplateStatuses(agentTemplates, reconciliations, opts)
harnessStatuses := newHarnessStatuses(harnesses, workerPools, opts)

return Collections{
AgentTemplates: agentTemplates,
Expand All @@ -91,6 +93,7 @@ func NewCollections(client kube.Client, watchNamespaces []string, opts krt.Optio
ModelConfigStatuses: modelConfigStatuses,
ResolvedModelConfigs: resolvedModelConfigs,
AgentTemplateStatuses: statuses,
HarnessStatuses: harnessStatuses,
}
}

Expand Down
48 changes: 47 additions & 1 deletion go/core/internal/controller/reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -138,9 +138,11 @@ type Reconciler struct {
pairs controllers.Queue
agentTemplateStatuses controllers.Queue
modelConfigStatuses controllers.Queue
harnessStatuses controllers.Queue
pairHandler krt.HandlerRegistration
agentTemplateStatusHandler krt.HandlerRegistration
modelConfigStatusHandler krt.HandlerRegistration
harnessStatusHandler krt.HandlerRegistration
}

// NewReconciler creates the Kubernetes and database write boundary. Run starts
Expand Down Expand Up @@ -174,6 +176,9 @@ func newReconciler(
r.modelConfigStatuses = newReconciliationQueue("v2-model-config-status", func(item any) error {
return r.reconcileModelConfigStatus(context.Background(), item.(string))
})
r.harnessStatuses = newReconciliationQueue("v2-harness-status", func(item any) error {
return r.reconcileHarnessStatus(context.Background(), item.(string))
})

r.pairHandler = collections.Reconciliations.Register(func(event krt.Event[PairReconciliation]) {
r.pairs.Add(krt.GetKey(event.Latest()))
Expand All @@ -192,6 +197,13 @@ func newReconciler(
}
r.modelConfigStatuses.Add(status.ResourceName())
})
r.harnessStatusHandler = collections.HarnessStatuses.Register(func(event krt.Event[krt.ObjectWithStatus[*kagentv1alpha3.Harness, kagentv1alpha3.HarnessStatus]]) {
status := event.Latest()
if apiequality.Semantic.DeepEqual(harnessStatusWithTransitionTimes(status.Status, status.Obj.Status), status.Obj.Status) {
return
}
r.harnessStatuses.Add(status.ResourceName())
})
return r
}

Expand All @@ -207,15 +219,17 @@ func newReconciliationQueue(name string, reconcile func(any) error) controllers.
// Run waits for the graph boundary to observe initial state, then processes
// pair and status writes until stop closes.
func (r *Reconciler) Run(stop <-chan struct{}) {
if !r.pairHandler.WaitUntilSynced(stop) || !r.agentTemplateStatusHandler.WaitUntilSynced(stop) || !r.modelConfigStatusHandler.WaitUntilSynced(stop) {
if !r.pairHandler.WaitUntilSynced(stop) || !r.agentTemplateStatusHandler.WaitUntilSynced(stop) || !r.modelConfigStatusHandler.WaitUntilSynced(stop) || !r.harnessStatusHandler.WaitUntilSynced(stop) {
r.pairs.ShutDownEarly()
r.agentTemplateStatuses.ShutDownEarly()
r.modelConfigStatuses.ShutDownEarly()
r.harnessStatuses.ShutDownEarly()
return
}
go r.pollPendingTemplates(stop)
go r.agentTemplateStatuses.Run(stop)
go r.modelConfigStatuses.Run(stop)
go r.harnessStatuses.Run(stop)
r.pairs.Run(stop)
}

Expand Down Expand Up @@ -370,6 +384,38 @@ func (r *Reconciler) reconcileModelConfigStatus(ctx context.Context, key string)
return nil
}

func (r *Reconciler) reconcileHarnessStatus(ctx context.Context, key string) error {
desired := r.collections.HarnessStatuses.GetKey(key)
harness := r.collections.Harnesses.GetKey(key)
if desired == nil || harness == nil {
return nil
}
updated := (*harness).DeepCopy()
updated.Status = harnessStatusWithTransitionTimes(desired.Status, updated.Status)
if apiequality.Semantic.DeepEqual(updated.Status, (*harness).Status) {
return nil
}
if _, err := r.status.Harnesses(updated.Namespace).UpdateStatus(ctx, updated, metav1.UpdateOptions{}); err != nil {
return fmt.Errorf("update Harness %s status: %w", key, err)
}
return nil
}

func harnessStatusWithTransitionTimes(desired, current kagentv1alpha3.HarnessStatus) kagentv1alpha3.HarnessStatus {
desired.Conditions = append([]metav1.Condition(nil), desired.Conditions...)
for conditionIndex := range desired.Conditions {
condition := &desired.Conditions[conditionIndex]
if previous := apimeta.FindStatusCondition(current.Conditions, condition.Type); previous != nil &&
previous.Status == condition.Status && previous.Reason == condition.Reason &&
previous.Message == condition.Message && previous.ObservedGeneration == condition.ObservedGeneration {
condition.LastTransitionTime = previous.LastTransitionTime
continue
}
condition.LastTransitionTime = metav1.Now()
}
return desired
}

func statusWithTransitionTimes(desired, current kagentv1alpha3.AgentTemplateStatus) kagentv1alpha3.AgentTemplateStatus {
desired.Harnesses = append([]kagentv1alpha3.AgentTemplateHarnessStatus(nil), desired.Harnesses...)
for harnessIndex := range desired.Harnesses {
Expand Down
80 changes: 80 additions & 0 deletions go/core/internal/controller/reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"time"

a2apb "github.com/a2aproject/a2a-go/v2/a2apb/v1"
atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
"google.golang.org/protobuf/encoding/protowire"
"google.golang.org/protobuf/proto"

Expand All @@ -26,6 +27,7 @@ import (
"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 Down Expand Up @@ -314,6 +316,78 @@ func (s *fakeRuntimeRevisionStore) RetirePairIdentities(_ context.Context, names
return nil
}

func TestReconcilerUpdatesHarnessStatus(t *testing.T) {
stop := make(chan struct{})
t.Cleanup(func() { close(stop) })
opts := krt.NewOptionsBuilder(stop, "test-harness-status-reconciler", nil)

harness := &kagentv1alpha3.Harness{
ObjectMeta: metav1.ObjectMeta{
Namespace: "team-a",
Name: "kagent",
Generation: 1,
},
Spec: kagentv1alpha3.HarnessSpec{
Substrate: kagentv1alpha3.HarnessSubstratePolicy{
WorkerPoolRef: corev1.LocalObjectReference{Name: "default"},
},
},
}
workerPool := &atev1alpha1.WorkerPool{
ObjectMeta: metav1.ObjectMeta{
Namespace: "team-a",
Name: "default",
},
}

mock := krttest.NewMock(t, []any{harness, workerPool})
harnesses := krttest.GetMockCollection[*kagentv1alpha3.Harness](mock)
workerPools := krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock)
harnessStatuses := newHarnessStatuses(harnesses, workerPools, opts)

statusClient := kagentfake.NewSimpleClientset(harness.DeepCopy()).ApiV1alpha3()

collections := Collections{
Harnesses: harnesses,
WorkerPools: workerPools,
HarnessStatuses: harnessStatuses,
AgentTemplates: krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock),
Reconciliations: krttest.GetMockCollection[PairReconciliation](mock),
AgentTemplateStatuses: krttest.GetMockCollection[krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]](mock),
ModelConfigStatuses: krttest.GetMockCollection[krt.ObjectWithStatus[*kagentv1alpha3.ModelConfig, kagentv1alpha3.ModelConfigStatus]](mock),
}

reconciler := newReconciler(
collections,
&fakeActorTemplates{},
&fakeRuntimeRevisionStore{},
statusClient,
)

go reconciler.Run(stop)

var updated *kagentv1alpha3.Harness
var err error
require.Eventually(t, func() bool {
updated, err = statusClient.Harnesses(harness.Namespace).Get(
context.Background(),
harness.Name,
metav1.GetOptions{},
)
if err != nil {
return false
}

condition := apimeta.FindStatusCondition(
updated.Status.Conditions,
kagentv1alpha3.HarnessConditionTypeReady,
)
return condition != nil &&
condition.Status == metav1.ConditionTrue &&
!condition.LastTransitionTime.IsZero()
}, 3*time.Second, 10*time.Millisecond)
}

func TestReconcilerUpdatesModelConfigStatusOnSecretHashChange(t *testing.T) {
stop := make(chan struct{})
t.Cleanup(func() { close(stop) })
Expand All @@ -337,6 +411,9 @@ func TestReconcilerUpdatesModelConfigStatusOnSecretHashChange(t *testing.T) {
secrets := krt.NewStaticCollection(nil, []*corev1.Secret{secret}, opts.WithName("Secrets")...)
configMaps := krttest.GetMockCollection[*corev1.ConfigMap](mock)
modelConfigStatuses, resolvedModelConfigs := newModelConfigReconciliations(modelConfigs, configMaps, secrets, opts)
harnesses := krttest.GetMockCollection[*kagentv1alpha3.Harness](mock)
workerPools := krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock)
harnessStatuses := newHarnessStatuses(harnesses, workerPools, opts)

collections := Collections{
ModelConfigs: modelConfigs,
Expand All @@ -347,6 +424,9 @@ func TestReconcilerUpdatesModelConfigStatusOnSecretHashChange(t *testing.T) {
AgentTemplates: krttest.GetMockCollection[*kagentv1alpha3.AgentTemplate](mock),
Reconciliations: krttest.GetMockCollection[PairReconciliation](mock),
AgentTemplateStatuses: krttest.GetMockCollection[krt.ObjectWithStatus[*kagentv1alpha3.AgentTemplate, kagentv1alpha3.AgentTemplateStatus]](mock),
Harnesses: harnesses,
WorkerPools: workerPools,
HarnessStatuses: harnessStatuses,
}

statusClient := kagentfake.NewSimpleClientset(modelConfig.DeepCopy()).ApiV1alpha3()
Expand Down
41 changes: 41 additions & 0 deletions go/core/internal/controller/status.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,11 @@ import (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"slices"
"strings"

atev1alpha1 "github.com/agent-substrate/substrate/pkg/api/v1alpha1"
kagentv1alpha3 "github.com/kagent-dev/kagent/go/api/v1alpha3"
"istio.io/istio/pkg/kube/krt"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
Expand Down Expand Up @@ -35,6 +37,45 @@ func newAgentTemplateStatuses(templates krt.Collection[*kagentv1alpha3.AgentTemp
return statuses
}

func newHarnessStatuses(
harnesses krt.Collection[*kagentv1alpha3.Harness],
workerPools krt.Collection[*atev1alpha1.WorkerPool],
opts krt.OptionsBuilder,
) krt.StatusCollection[*kagentv1alpha3.Harness, kagentv1alpha3.HarnessStatus] {
statuses, _ := krt.NewStatusManyCollection(harnesses, func(ctx krt.HandlerContext, harness *kagentv1alpha3.Harness) (*kagentv1alpha3.HarnessStatus, []*atev1alpha1.WorkerPool) {
workerPoolKey := types.NamespacedName{
Namespace: harness.Namespace,
Name: harness.Spec.Substrate.WorkerPoolRef.Name,
}
workerPool := krt.FetchOne(ctx, workerPools, krt.FilterObjectName(workerPoolKey))

status := &kagentv1alpha3.HarnessStatus{
ObservedGeneration: harness.Generation,
}

if workerPool == nil || *workerPool == nil {
status.Conditions = []metav1.Condition{{
Type: kagentv1alpha3.HarnessConditionTypeReady,
Status: metav1.ConditionFalse,
Reason: "WorkerPoolNotFound",
Message: fmt.Sprintf("WorkerPool %q not found", workerPoolKey.String()),
ObservedGeneration: harness.Generation,
}}
return status, nil
}

status.Conditions = []metav1.Condition{{
Type: kagentv1alpha3.HarnessConditionTypeReady,
Status: metav1.ConditionTrue,
Reason: "WorkerPoolResolved",
Message: fmt.Sprintf("WorkerPool %q exists", workerPoolKey.String()),
ObservedGeneration: harness.Generation,
}}
return status, nil
}, opts.WithName("HarnessStatuses")...)
return statuses
}

func statusForPair(state PairReconciliation, generation int64, latestSuccessful string) kagentv1alpha3.AgentTemplateHarnessStatus {
desired := state.RevisionID.String()
if state.RevisionID.IsZero() {
Expand Down
86 changes: 86 additions & 0 deletions go/core/internal/controller/status_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,14 @@ package controller
import (
"testing"

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/internal/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 TestStatusForPairPublishesCompilationWarnings(t *testing.T) {
Expand All @@ -24,3 +30,83 @@ func TestStatusForPairPublishesCompilationWarnings(t *testing.T) {
t.Fatal("status aliases mutable compilation warnings")
}
}

func TestHarnessStatusReadyWhenWorkerPoolExists(t *testing.T) {
stop := make(chan struct{})
t.Cleanup(func() { close(stop) })
opts := krt.NewOptionsBuilder(stop, "test-harness-status", nil)

harness := &kagentv1alpha3.Harness{
ObjectMeta: metav1.ObjectMeta{
Namespace: "team-a",
Name: "kagent",
},
Spec: kagentv1alpha3.HarnessSpec{
Substrate: kagentv1alpha3.HarnessSubstratePolicy{
WorkerPoolRef: corev1.LocalObjectReference{Name: "default"},
},
},
}
workerPool := &atev1alpha1.WorkerPool{
ObjectMeta: metav1.ObjectMeta{
Namespace: "team-a",
Name: "default",
},
}

mock := krttest.NewMock(t, []any{harness, workerPool})
harnesses := krttest.GetMockCollection[*kagentv1alpha3.Harness](mock)
workerPools := krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock)

statuses := newHarnessStatuses(harnesses, workerPools, opts)

waitFor(t, func() bool {
status := statuses.GetKey("team-a/kagent")
if status == nil {
return false
}
condition := apimeta.FindStatusCondition(
status.Status.Conditions,
kagentv1alpha3.HarnessConditionTypeReady,
)
return condition != nil && condition.Status == metav1.ConditionTrue
})
}

func TestHarnessStatusNotReadyWhenWorkerPoolMissing(t *testing.T) {
stop := make(chan struct{})
t.Cleanup(func() { close(stop) })
opts := krt.NewOptionsBuilder(stop, "test-harness-status-missing", nil)

harness := &kagentv1alpha3.Harness{
ObjectMeta: metav1.ObjectMeta{
Namespace: "team-a",
Name: "kagent",
},
Spec: kagentv1alpha3.HarnessSpec{
Substrate: kagentv1alpha3.HarnessSubstratePolicy{
WorkerPoolRef: corev1.LocalObjectReference{Name: "missing"},
},
},
}

mock := krttest.NewMock(t, []any{harness})
harnesses := krttest.GetMockCollection[*kagentv1alpha3.Harness](mock)
workerPools := krttest.GetMockCollection[*atev1alpha1.WorkerPool](mock)

statuses := newHarnessStatuses(harnesses, workerPools, opts)

waitFor(t, func() bool {
status := statuses.GetKey("team-a/kagent")
if status == nil {
return false
}
condition := apimeta.FindStatusCondition(
status.Status.Conditions,
kagentv1alpha3.HarnessConditionTypeReady,
)
return condition != nil &&
condition.Status == metav1.ConditionFalse &&
condition.Reason == "WorkerPoolNotFound"
})
}
Loading