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
34 changes: 25 additions & 9 deletions flagd-proxy/pkg/service/subscriptions/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -220,16 +220,32 @@ func (s *Coordinator) cleanup() {
case <-s.ctx.Done():
return
case <-time.After(5 * time.Second):
s.mu.Lock()
for k, v := range s.multiplexers {
// delete any multiplexers with 0 active subscriptions through cancelling its context
s.logger.Debug(fmt.Sprintf("multiplexer for target %s has %d subscriptions", k, len(v.subs)))
if len(v.subs) == 0 {
s.logger.Debug(fmt.Sprintf("shutting down multiplexer %s", k))
s.multiplexers[k].cancelFunc()
}
s.cleanupOnce()
}
}
}

// cleanupOnce removes multiplexers without active subscriptions. Called on a
// fixed interval by cleanup, and directly by tests.
func (s *Coordinator) cleanupOnce() {
s.mu.Lock()
defer s.mu.Unlock()
for k, v := range s.multiplexers {
// delete any multiplexers with 0 active subscriptions through cancelling its context
s.logger.Debug(fmt.Sprintf("multiplexer for target %s has %d subscriptions", k, len(v.subs)))
if len(v.subs) == 0 {
s.logger.Debug(fmt.Sprintf("shutting down multiplexer %s", k))
if v.cancelFunc != nil {
// watcher is running; cancelling its context triggers the async
// removal of the multiplexer
v.cancelFunc()
} else {
// the watcher goroutine never assigned cancelFunc (the last
// subscription was cancelled before watchResource started), so
// there is no context to cancel; drop the entry directly to
// avoid leaving a dead multiplexer behind
delete(s.multiplexers, k)
}
s.mu.Unlock()
}
}
}
Expand Down
26 changes: 26 additions & 0 deletions flagd-proxy/pkg/service/subscriptions/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,32 @@ func Test_watchResource_SyncHandlerDoesNotExist(_ *testing.T) {
syncStore.watchResource(target)
}

func Test_cleanupOnce_multiplexerWithoutStartedWatcher(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
syncStore := NewManager(ctx, logger.NewLogger(nil, false))

target := "test-target"
// a multiplexer whose watcher never started (cancelFunc still nil): this is
// the state left behind when the only subscription is cancelled before
// watchResource had a chance to assign cancelFunc
syncStore.mu.Lock()
syncStore.multiplexers[target] = &multiplexer{
subs: map[interface{}]storedChannels{},
mu: &sync.RWMutex{},
}
syncStore.mu.Unlock()

// must not panic and must remove the dead entry
syncStore.cleanupOnce()

syncStore.mu.Lock()
defer syncStore.mu.Unlock()
if syncStore.multiplexers[target] != nil {
t.Error("multiplexer without a started watcher has not been removed by cleanup")
}
}

func Test_watchResource_Cleanup(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
Expand Down