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
92 changes: 23 additions & 69 deletions internal/controllers/watch/kind.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,19 +5,18 @@ import (
"fmt"
"math/rand"
"path"
"reflect"
"slices"
"time"

apiv1 "github.com/Azure/eno/api/v1"
"github.com/Azure/eno/internal/flowcontrol"
"github.com/Azure/eno/internal/manager"
"github.com/go-logr/logr"
"golang.org/x/time/rate"
"k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/util/retry"
"k8s.io/client-go/util/workqueue"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
Expand All @@ -36,6 +35,7 @@ func SetKindWatchRateLimit(rps int) {

type KindWatchController struct {
client client.Client
buffer *flowcontrol.CompositionInputRevisionWriteBuffer
gvk schema.GroupVersionKind
cancel context.CancelFunc
}
Expand All @@ -45,6 +45,7 @@ func NewKindWatchController(ctx context.Context, parent *WatchController, resour

k := &KindWatchController{
client: parent.mgr.GetClient(),
buffer: parent.buffer,
gvk: schema.GroupVersionKind{
Group: resource.Group,
Version: resource.Version,
Expand Down Expand Up @@ -219,10 +220,7 @@ func (k *KindWatchController) Reconcile(ctx context.Context, req ctrl.Request) (
logger.Error(err, "failed to list compositions for synthesizer", "synthesizerName", synth.Name)
return ctrl.Result{}, fmt.Errorf("listing compositions: %w", err)
}
modified, err := k.updateCompositions(ctx, logger, &synth, meta, list, isDeleted)
if modified || err != nil {
return ctrl.Result{}, err
}
k.updateCompositions(logger, &synth, meta, list, isDeleted)
}

list := &apiv1.CompositionList{}
Expand All @@ -233,61 +231,42 @@ func (k *KindWatchController) Reconcile(ctx context.Context, req ctrl.Request) (
logger.Error(err, "failed to list compositions by binding", "synthesizerName", synth.Name)
return ctrl.Result{}, fmt.Errorf("listing compositions: %w", err)
}
modified, err := k.updateCompositions(ctx, logger, &synth, meta, list, isDeleted)
if modified || err != nil {
return ctrl.Result{}, err
}
k.updateCompositions(logger, &synth, meta, list, isDeleted)
}

logger.Info("finished reconciling watched resource")
return ctrl.Result{}, nil
}

func (k *KindWatchController) updateCompositions(ctx context.Context, logger logr.Logger, synth *apiv1.Synthesizer, meta *metav1.PartialObjectMetadata, list *apiv1.CompositionList, isDeleted bool) (bool, error) {
for _, comp := range list.Items {
// updateCompositions enqueues an input-revision change for every composition bound to the
// watched resource. Writes are coalesced per composition by the write buffer, so multiple
// inputs changing for the same composition in a short window collapse into a single
// Status().Update instead of one write per input.
func (k *KindWatchController) updateCompositions(logger logr.Logger, synth *apiv1.Synthesizer, meta *metav1.PartialObjectMetadata, list *apiv1.CompositionList, isDeleted bool) {
for i := range list.Items {
comp := list.Items[i]
key := findRefKey(&comp, synth, meta)
if key == "" {
continue
}

compKey := client.ObjectKeyFromObject(&comp)
var modified bool

err := retry.RetryOnConflict(retry.DefaultBackoff, func() error {
// Re-fetch composition to get the latest version
if err := k.client.Get(ctx, compKey, &comp); err != nil {
return err
}

if isDeleted {
// Only remove InputRevisions for optional refs
// Required refs should trigger MissingInputs status instead
if isOptionalRef(synth, key) {
modified = removeInputRevision(&comp, key)
}
// For required refs, do nothing - the composition controller will handle MissingInputs
} else {
revs := apiv1.NewInputRevisions(meta, key)
modified = setInputRevisions(&comp, revs)
}

if !modified {
return nil
if isDeleted {
// Only remove InputRevisions for optional refs.
// Required refs should trigger MissingInputs status instead - the
// composition controller handles that case.
if isOptionalRef(synth, key) {
k.buffer.RemoveInputRevisionAsync(compKey, key)
logger.V(1).Info("buffered input revision removal", "compositionName", comp.Name, "compositionNamespace", comp.Namespace, "ref", key)
}

return k.client.Status().Update(ctx, &comp)
})
if err != nil {
return false, fmt.Errorf("updating input revisions: %w", err)
continue
}

if modified {
logger.V(1).Info("noticed input resource change", "compositionName", comp.Name, "compositionNamespace", comp.Namespace, "ref", key)
return true, nil // wait for requeue
}
revs := apiv1.NewInputRevisions(meta, key)
k.buffer.PatchInputRevisionAsync(compKey, revs)
logger.V(1).Info("buffered input resource change", "compositionName", comp.Name, "compositionNamespace", comp.Namespace, "ref", key)
}

return false, nil
}

func findRefKey(comp *apiv1.Composition, synth *apiv1.Synthesizer, meta *metav1.PartialObjectMetadata) string {
Expand Down Expand Up @@ -320,28 +299,3 @@ func isOptionalRef(synth *apiv1.Synthesizer, key string) bool {
}
return false
}

func removeInputRevision(comp *apiv1.Composition, key string) bool {
for i, ir := range comp.Status.InputRevisions {
if ir.Key == key {
comp.Status.InputRevisions = append(comp.Status.InputRevisions[:i], comp.Status.InputRevisions[i+1:]...)
return true
}
}
return false
}

func setInputRevisions(comp *apiv1.Composition, revs *apiv1.InputRevisions) bool {
for i, ir := range comp.Status.InputRevisions {
if ir.Key != revs.Key {
continue
}
if reflect.DeepEqual(ir, *revs) {
return false
}
comp.Status.InputRevisions[i] = *revs
return true
}
comp.Status.InputRevisions = append(comp.Status.InputRevisions, *revs)
return true
}
194 changes: 0 additions & 194 deletions internal/controllers/watch/kind_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,122 +6,9 @@ import (
apiv1 "github.com/Azure/eno/api/v1"
"github.com/stretchr/testify/assert"
"k8s.io/apimachinery/pkg/types"
"k8s.io/utils/ptr"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
)

func TestSetInputRevisions(t *testing.T) {
tests := []struct {
name string
comp *apiv1.Composition
revs *apiv1.InputRevisions
expected bool
finalRevs []apiv1.InputRevisions
}{
{
name: "add new revision when key is not found",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
},
revs: &apiv1.InputRevisions{
Key: "rev2",
Revision: ptr.To(2),
},
expected: true,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
{Key: "rev2", Revision: ptr.To(2)},
},
},
{
name: "update existing revision",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
},
revs: &apiv1.InputRevisions{
Key: "rev1",
Revision: ptr.To(2),
},
expected: true,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(2)},
},
},
{
name: "no update if revision is identical",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
},
revs: &apiv1.InputRevisions{
Key: "rev1",
Revision: ptr.To(1),
},
expected: false,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
{
name: "no update if revision is identical and synth generation is set",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))},
},
},
},
revs: &apiv1.InputRevisions{
Key: "rev1",
Revision: ptr.To(1),
SynthesizerGeneration: ptr.To(int64(3)),
},
expected: false,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))},
},
},
{
name: "update if revision is identical but synth generation is not",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(3))},
},
},
},
revs: &apiv1.InputRevisions{
Key: "rev1",
Revision: ptr.To(1),
SynthesizerGeneration: ptr.To(int64(5)),
},
expected: true,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1), SynthesizerGeneration: ptr.To(int64(5))},
},
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
result := setInputRevisions(tt.comp, tt.revs)
assert.Equal(t, tt.expected, result)
assert.Equal(t, tt.finalRevs, tt.comp.Status.InputRevisions)
})
}
}

func TestBuildRequests(t *testing.T) {
tests := []struct {
name string
Expand Down Expand Up @@ -266,87 +153,6 @@ func TestBuildRequests(t *testing.T) {
}
}

func TestRemoveInputRevision(t *testing.T) {
tests := []struct {
name string
comp *apiv1.Composition
key string
expected bool
finalRevs []apiv1.InputRevisions
}{
{
name: "remove existing revision",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
{Key: "rev2", Revision: ptr.To(2)},
},
},
},
key: "rev1",
expected: true,
finalRevs: []apiv1.InputRevisions{
{Key: "rev2", Revision: ptr.To(2)},
},
},
{
name: "remove last revision",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
},
key: "rev1",
expected: true,
finalRevs: []apiv1.InputRevisions{},
},
{
name: "no removal if key not found",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
},
key: "rev2",
expected: false,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
},
},
{
name: "remove from middle of list",
comp: &apiv1.Composition{
Status: apiv1.CompositionStatus{
InputRevisions: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
{Key: "rev2", Revision: ptr.To(2)},
{Key: "rev3", Revision: ptr.To(3)},
},
},
},
key: "rev2",
expected: true,
finalRevs: []apiv1.InputRevisions{
{Key: "rev1", Revision: ptr.To(1)},
{Key: "rev3", Revision: ptr.To(3)},
},
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
result := removeInputRevision(tt.comp, tt.key)
assert.Equal(t, tt.expected, result)
assert.Equal(t, tt.finalRevs, tt.comp.Status.InputRevisions)
})
}
}

func TestIsOptionalRef(t *testing.T) {
tests := []struct {
name string
Expand Down
Loading
Loading