From 41e9b9d766d8eb413583d41d1f27a408f08a75d8 Mon Sep 17 00:00:00 2001 From: Daniel Gonzalez Nothnagel Date: Tue, 4 Aug 2026 10:28:59 +0200 Subject: [PATCH 1/6] Add mutex for watch members Signed-off-by: Daniel Gonzalez Nothnagel --- storeutils/host/store.go | 22 +++++++++++++--------- storeutils/host/watch.go | 11 +++++++---- 2 files changed, 20 insertions(+), 13 deletions(-) diff --git a/storeutils/host/store.go b/storeutils/host/store.go index 20169e4..ced7340 100644 --- a/storeutils/host/store.go +++ b/storeutils/host/store.go @@ -346,23 +346,27 @@ func (s *Store[E]) watchHandlers() []*watch[E] { return s.watches.UnsortedList() } -func (s *Store[E]) send(w *watch[E], evt store.WatchEvent[E]) { - select { - case w.events <- evt: - default: - } -} - func (s *Store[E]) enqueue(evt store.WatchEvent[E]) { id := evt.Object.GetID() for _, handler := range s.watchHandlers() { + var toSend *store.WatchEvent[E] + + handler.membersMu.Lock() if handler.matches(evt.Object) { handler.members.Insert(id) - s.send(handler, evt) + toSend = &evt } else if handler.members.Has(id) { // Object transitioned out of this watch's scope. Send deleted event. handler.members.Delete(id) - s.send(handler, store.WatchEvent[E]{Type: store.WatchEventTypeDeleted, Object: evt.Object}) + toSend = &store.WatchEvent[E]{Type: store.WatchEventTypeDeleted, Object: evt.Object} + } + handler.membersMu.Unlock() + + if toSend != nil { + select { + case handler.events <- *toSend: + default: + } } } } diff --git a/storeutils/host/watch.go b/storeutils/host/watch.go index ade8dc4..fba14e3 100644 --- a/storeutils/host/watch.go +++ b/storeutils/host/watch.go @@ -4,16 +4,19 @@ package host import ( + "sync" + "github.com/ironcore-dev/provider-utils/apiutils/api" "github.com/ironcore-dev/provider-utils/storeutils/store" "k8s.io/apimachinery/pkg/util/sets" ) type watch[E api.Object] struct { - store *Store[E] - events chan store.WatchEvent[E] - opts store.ListOptions - members sets.Set[string] + store *Store[E] + events chan store.WatchEvent[E] + opts store.ListOptions + membersMu sync.Mutex + members sets.Set[string] } func (w *watch[E]) matches(obj E) bool { From cb2145516c127c8ed3ebd93babacc7af482cdf4d Mon Sep 17 00:00:00 2001 From: Daniel Gonzalez Nothnagel Date: Tue, 4 Aug 2026 10:29:36 +0200 Subject: [PATCH 2/6] Populate members on watch creation Signed-off-by: Daniel Gonzalez Nothnagel --- storeutils/host/store.go | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/storeutils/host/store.go b/storeutils/host/store.go index ced7340..c0b8b47 100644 --- a/storeutils/host/store.go +++ b/storeutils/host/store.go @@ -270,7 +270,7 @@ func (s *Store[E]) List(ctx context.Context, opts ...store.ListOption) ([]E, err return objs, nil } -func (s *Store[E]) Watch(_ context.Context, opts ...store.ListOption) (store.Watch[E], error) { +func (s *Store[E]) Watch(ctx context.Context, opts ...store.ListOption) (store.Watch[E], error) { listOpts := &store.ListOptions{} for _, opt := range opts { opt.ApplyToList(listOpts) @@ -280,6 +280,20 @@ func (s *Store[E]) Watch(_ context.Context, opts ...store.ListOption) (store.Wat return nil, err } + // List runs before acquiring watchesMu to avoid a deadlock: List→Get acquires + // idMu, while mutations hold idMu before calling enqueue→watchesMu.RLock. + // This means a mutation between List and watches.Insert may be missed in + // members, but we accept that narrow gap. + existing, err := s.List(ctx, opts...) + if err != nil { + return nil, fmt.Errorf("failed to list existing objects for watch: %w", err) + } + + members := sets.New[string]() + for _, obj := range existing { + members.Insert(obj.GetID()) + } + s.watchesMu.Lock() defer s.watchesMu.Unlock() @@ -287,7 +301,7 @@ func (s *Store[E]) Watch(_ context.Context, opts ...store.ListOption) (store.Wat store: s, events: make(chan store.WatchEvent[E], s.watchBufferSize), opts: *listOpts, - members: sets.New[string](), + members: members, } s.watches.Insert(w) From 45224c84fbf2ad4fd05e6b76fa7b2bbaee55e759 Mon Sep 17 00:00:00 2001 From: Daniel Gonzalez Nothnagel Date: Tue, 4 Aug 2026 11:08:11 +0200 Subject: [PATCH 3/6] Handle deleted events before filter check in Watch Deleted events previously went through the filter matcher, which could silently drop them if the object no longer matched the watch's predicate. Signed-off-by: Daniel Gonzalez Nothnagel --- storeutils/host/store.go | 21 ++++++++++++++------- 1 file changed, 14 insertions(+), 7 deletions(-) diff --git a/storeutils/host/store.go b/storeutils/host/store.go index c0b8b47..61f6fb2 100644 --- a/storeutils/host/store.go +++ b/storeutils/host/store.go @@ -363,22 +363,29 @@ func (s *Store[E]) watchHandlers() []*watch[E] { func (s *Store[E]) enqueue(evt store.WatchEvent[E]) { id := evt.Object.GetID() for _, handler := range s.watchHandlers() { - var toSend *store.WatchEvent[E] + var toSend store.WatchEvent[E] handler.membersMu.Lock() - if handler.matches(evt.Object) { + if evt.Type == store.WatchEventTypeDeleted { + // Object was deleted; forward only if we were tracking it. + if handler.members.Has(id) { + handler.members.Delete(id) + toSend = evt + } + } else if handler.matches(evt.Object) { + // Object matches the filter; track it and forward the event. handler.members.Insert(id) - toSend = &evt + toSend = evt } else if handler.members.Has(id) { - // Object transitioned out of this watch's scope. Send deleted event. + // Object no longer matches the filter. Send deleted event. handler.members.Delete(id) - toSend = &store.WatchEvent[E]{Type: store.WatchEventTypeDeleted, Object: evt.Object} + toSend = store.WatchEvent[E]{Type: store.WatchEventTypeDeleted, Object: evt.Object} } handler.membersMu.Unlock() - if toSend != nil { + if toSend.Type != "" { select { - case handler.events <- *toSend: + case handler.events <- toSend: default: } } From bc7038b8ba43d33921d07b6182ca067d6c91dc9a Mon Sep 17 00:00:00 2001 From: Daniel Gonzalez Nothnagel Date: Wed, 5 Aug 2026 10:26:32 +0200 Subject: [PATCH 4/6] Refactor if-else chain to switch statement in enqueue Signed-off-by: Daniel Gonzalez Nothnagel --- storeutils/host/store.go | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/storeutils/host/store.go b/storeutils/host/store.go index 61f6fb2..9903841 100644 --- a/storeutils/host/store.go +++ b/storeutils/host/store.go @@ -366,17 +366,18 @@ func (s *Store[E]) enqueue(evt store.WatchEvent[E]) { var toSend store.WatchEvent[E] handler.membersMu.Lock() - if evt.Type == store.WatchEventTypeDeleted { + switch { + case evt.Type == store.WatchEventTypeDeleted: // Object was deleted; forward only if we were tracking it. if handler.members.Has(id) { handler.members.Delete(id) toSend = evt } - } else if handler.matches(evt.Object) { + case handler.matches(evt.Object): // Object matches the filter; track it and forward the event. handler.members.Insert(id) toSend = evt - } else if handler.members.Has(id) { + case handler.members.Has(id): // Object no longer matches the filter. Send deleted event. handler.members.Delete(id) toSend = store.WatchEvent[E]{Type: store.WatchEventTypeDeleted, Object: evt.Object} From 2ef3f2ccdb00ffd8fd40d7695351486198288656 Mon Sep 17 00:00:00 2001 From: Daniel Gonzalez Nothnagel Date: Wed, 5 Aug 2026 15:49:21 +0200 Subject: [PATCH 5/6] Add watch filter behavior tests for delete and scope transitions Signed-off-by: Daniel Gonzalez Nothnagel --- storeutils/host/store_test.go | 90 +++++++++++++++++++++++++++++++++++ 1 file changed, 90 insertions(+) diff --git a/storeutils/host/store_test.go b/storeutils/host/store_test.go index 8db5e61..069239f 100644 --- a/storeutils/host/store_test.go +++ b/storeutils/host/store_test.go @@ -210,5 +210,95 @@ var _ = Describe("Store", func() { Expect(err).To(HaveOccurred()) Expect(err.Error()).To(ContainSubstring("unindexed field")) }) + + It("should receive a deleted event for an object that existed before the watch was created", func(ctx SpecContext) { + obj, err := dummyStore.Create(ctx, &Dummy{Metadata: api.Metadata{ID: "preexisting-deleted"}}) + Expect(err).NotTo(HaveOccurred()) + + watch, err := dummyStore.Watch(ctx) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(watch.Stop) + + Expect(dummyStore.Delete(ctx, obj.GetID())).To(Succeed()) + + Eventually(watch.Events()).Should(Receive(HaveField("Type", store.WatchEventTypeDeleted))) + }) + + It("should emit a deleted event when a tracked object transitions out of scope", func(ctx SpecContext) { + watch, err := dummyStore.Watch(ctx, store.MatchingLabels{"app": "watched"}) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(watch.Stop) + + obj := createDummy(ctx, dummyStore, "transitions-out", map[string]string{"app": "watched"}, "") + Eventually(watch.Events()).Should(Receive(HaveField("Type", store.WatchEventTypeCreated))) + + obj.Labels["app"] = "other" + _, err = dummyStore.Update(ctx, obj) + Expect(err).NotTo(HaveOccurred()) + + Eventually(watch.Events()).Should(Receive(HaveField("Type", store.WatchEventTypeDeleted))) + }) + + It("should not forward a delete event for an object that was never tracked", func(ctx SpecContext) { + watch, err := dummyStore.Watch(ctx, store.MatchingLabels{"app": "watched"}) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(watch.Stop) + + obj, err := dummyStore.Create(ctx, &Dummy{ + Metadata: api.Metadata{ID: "delete-untracked", Labels: map[string]string{"app": "other"}}, + }) + Expect(err).NotTo(HaveOccurred()) + Consistently(watch.Events()).ShouldNot(Receive()) + + Expect(dummyStore.Delete(ctx, obj.GetID())).To(Succeed()) + Consistently(watch.Events()).ShouldNot(Receive()) + }) + + It("should forward a delete event for a tracked object", func(ctx SpecContext) { + watch, err := dummyStore.Watch(ctx, store.MatchingLabels{"app": "watched"}) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(watch.Stop) + + obj, err := dummyStore.Create(ctx, &Dummy{ + Metadata: api.Metadata{ID: "delete-tracked", Labels: map[string]string{"app": "watched"}}, + }) + Expect(err).NotTo(HaveOccurred()) + Eventually(watch.Events()).Should(Receive(HaveField("Type", store.WatchEventTypeCreated))) + + Expect(dummyStore.Delete(ctx, obj.GetID())).To(Succeed()) + Eventually(watch.Events()).Should(Receive(HaveField("Type", store.WatchEventTypeDeleted))) + }) + + It("should receive events when an object transitions into scope", func(ctx SpecContext) { + watch, err := dummyStore.Watch(ctx, store.MatchingLabels{"app": "watched"}) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(watch.Stop) + + obj := createDummy(ctx, dummyStore, "transitions-in", map[string]string{"app": "other"}, "") + Consistently(watch.Events()).ShouldNot(Receive()) + + obj.Labels["app"] = "watched" + updated, err := dummyStore.Update(ctx, obj) + Expect(err).NotTo(HaveOccurred()) + + Eventually(watch.Events()).Should(Receive(Equal( + store.WatchEvent[*Dummy]{Type: store.WatchEventTypeUpdated, Object: updated}, + ))) + }) + + It("should emit a deleted event when a pre-existing object transitions out of scope", func(ctx SpecContext) { + obj := createDummy(ctx, dummyStore, "preexisting-transitions-out", map[string]string{"app": "watched"}, "") + + watch, err := dummyStore.Watch(ctx, store.MatchingLabels{"app": "watched"}) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(watch.Stop) + + obj.Labels["app"] = "other" + _, err = dummyStore.Update(ctx, obj) + Expect(err).NotTo(HaveOccurred()) + + Eventually(watch.Events()).Should(Receive(HaveField("Type", store.WatchEventTypeDeleted))) + Consistently(watch.Events()).ShouldNot(Receive()) + }) }) }) From 7a20385a37701d1b777bb4cc7f1ebaf6cb8f3277 Mon Sep 17 00:00:00 2001 From: Daniel Gonzalez Nothnagel Date: Thu, 6 Aug 2026 14:54:43 +0200 Subject: [PATCH 6/6] Log dropped watch events when handler channel is full Signed-off-by: Daniel Gonzalez Nothnagel --- storeutils/host/store.go | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/storeutils/host/store.go b/storeutils/host/store.go index 9903841..fad9df9 100644 --- a/storeutils/host/store.go +++ b/storeutils/host/store.go @@ -13,6 +13,7 @@ import ( "sync" "time" + "github.com/go-logr/logr" "github.com/ironcore-dev/provider-utils/apiutils/api" "github.com/ironcore-dev/provider-utils/storeutils/store" utilssync "github.com/ironcore-dev/provider-utils/storeutils/sync" @@ -29,6 +30,7 @@ type Options[E api.Object] struct { CreateStrategy CreateStrategy[E] WatchBufferSize int FieldIndexers map[string]store.IndexerFunc[E] + Log logr.Logger } func (o *Options[E]) Defaults() { @@ -54,6 +56,7 @@ func NewStore[E api.Object](opts Options[E]) (*Store[E], error) { } return &Store[E]{ dir: opts.Dir, + log: opts.Log, idMu: utilssync.NewMutexMap[string](), @@ -69,6 +72,7 @@ func NewStore[E api.Object](opts Options[E]) (*Store[E], error) { type Store[E api.Object] struct { dir string + log logr.Logger idMu *utilssync.MutexMap[string] @@ -388,6 +392,8 @@ func (s *Store[E]) enqueue(evt store.WatchEvent[E]) { select { case handler.events <- toSend: default: + // TODO: switch to `ctx` to backpressure, if channel size is not enough + s.log.Info("Dropping watch event, due to full channel", "event", evt) } } }