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
60 changes: 46 additions & 14 deletions storeutils/host/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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() {
Expand All @@ -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](),

Expand All @@ -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]

Expand Down Expand Up @@ -270,7 +274,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)
Expand All @@ -280,14 +284,28 @@ 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)
}
Comment thread
gonzolino marked this conversation as resolved.

members := sets.New[string]()
for _, obj := range existing {
members.Insert(obj.GetID())
}

s.watchesMu.Lock()
defer s.watchesMu.Unlock()

w := &watch[E]{
store: s,
events: make(chan store.WatchEvent[E], s.watchBufferSize),
opts: *listOpts,
members: sets.New[string](),
members: members,
}

s.watches.Insert(w)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Expand Down Expand Up @@ -346,23 +364,37 @@ 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() {
if handler.matches(evt.Object) {
var toSend store.WatchEvent[E]

handler.membersMu.Lock()
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
}
case handler.matches(evt.Object):
// Object matches the filter; track it and forward the event.
handler.members.Insert(id)
s.send(handler, evt)
} else if handler.members.Has(id) {
// Object transitioned out of this watch's scope. Send deleted event.
toSend = evt
case handler.members.Has(id):
// Object no longer matches the filter. 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}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
handler.membersMu.Unlock()

if toSend.Type != "" {
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)
}
}
}
}
90 changes: 90 additions & 0 deletions storeutils/host/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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())
})
})
})
11 changes: 7 additions & 4 deletions storeutils/host/watch.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down