diff --git a/pkg/storage/unified/resource/notifier.go b/pkg/storage/unified/resource/notifier.go index 3d3b2024d7e..12cefa5add5 100644 --- a/pkg/storage/unified/resource/notifier.go +++ b/pkg/storage/unified/resource/notifier.go @@ -19,13 +19,18 @@ const ( defaultBufferSize = 10000 ) -type notifier struct { +type notifier interface { + Watch(context.Context, watchOptions) <-chan Event +} + +type pollingNotifier struct { eventStore *eventStore log logging.Logger } type notifierOptions struct { - log logging.Logger + log logging.Logger + useChannelNotifier bool } type watchOptions struct { @@ -44,15 +49,26 @@ func defaultWatchOptions() watchOptions { } } -func newNotifier(eventStore *eventStore, opts notifierOptions) *notifier { +func newNotifier(eventStore *eventStore, opts notifierOptions) notifier { if opts.log == nil { opts.log = &logging.NoOpLogger{} } - return ¬ifier{eventStore: eventStore, log: opts.log} + + if opts.useChannelNotifier { + return &channelNotifier{} + } + + return &pollingNotifier{eventStore: eventStore, log: opts.log} +} + +type channelNotifier struct{} + +func (cn *channelNotifier) Watch(ctx context.Context, opts watchOptions) <-chan Event { + return nil } // Return the last resource version from the event store -func (n *notifier) lastEventResourceVersion(ctx context.Context) (int64, error) { +func (n *pollingNotifier) lastEventResourceVersion(ctx context.Context) (int64, error) { e, err := n.eventStore.LastEventKey(ctx) if err != nil { return 0, err @@ -60,11 +76,11 @@ func (n *notifier) lastEventResourceVersion(ctx context.Context) (int64, error) return e.ResourceVersion, nil } -func (n *notifier) cacheKey(evt Event) string { +func (n *pollingNotifier) cacheKey(evt Event) string { return fmt.Sprintf("%s~%s~%s~%s~%d", evt.Namespace, evt.Group, evt.Resource, evt.Name, evt.ResourceVersion) } -func (n *notifier) Watch(ctx context.Context, opts watchOptions) <-chan Event { +func (n *pollingNotifier) Watch(ctx context.Context, opts watchOptions) <-chan Event { if opts.MinBackoff <= 0 { opts.MinBackoff = defaultMinBackoff } diff --git a/pkg/storage/unified/resource/notifier_test.go b/pkg/storage/unified/resource/notifier_test.go index f78629ebeb7..5c217bf0a21 100644 --- a/pkg/storage/unified/resource/notifier_test.go +++ b/pkg/storage/unified/resource/notifier_test.go @@ -13,7 +13,7 @@ import ( "github.com/stretchr/testify/require" ) -func setupTestNotifier(t *testing.T) (*notifier, *eventStore) { +func setupTestNotifier(t *testing.T) (*pollingNotifier, *eventStore) { db := setupTestBadgerDB(t) t.Cleanup(func() { err := db.Close() @@ -22,10 +22,10 @@ func setupTestNotifier(t *testing.T) (*notifier, *eventStore) { kv := NewBadgerKV(db) eventStore := newEventStore(kv) notifier := newNotifier(eventStore, notifierOptions{log: &logging.NoOpLogger{}}) - return notifier, eventStore + return notifier.(*pollingNotifier), eventStore } -func setupTestNotifierSqlKv(t *testing.T) (*notifier, *eventStore) { +func setupTestNotifierSqlKv(t *testing.T) (*pollingNotifier, *eventStore) { dbstore := db.InitTestDB(t) eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil) require.NoError(t, err) @@ -33,7 +33,7 @@ func setupTestNotifierSqlKv(t *testing.T) (*notifier, *eventStore) { require.NoError(t, err) eventStore := newEventStore(kv) notifier := newNotifier(eventStore, notifierOptions{log: &logging.NoOpLogger{}}) - return notifier, eventStore + return notifier.(*pollingNotifier), eventStore } func TestNewNotifier(t *testing.T) { @@ -49,7 +49,7 @@ func TestDefaultWatchOptions(t *testing.T) { assert.Equal(t, defaultBufferSize, opts.BufferSize) } -func runNotifierTestWith(t *testing.T, storeName string, newStoreFn func(*testing.T) (*notifier, *eventStore), testFn func(*testing.T, context.Context, *notifier, *eventStore)) { +func runNotifierTestWith(t *testing.T, storeName string, newStoreFn func(*testing.T) (*pollingNotifier, *eventStore), testFn func(*testing.T, context.Context, *pollingNotifier, *eventStore)) { t.Run(storeName, func(t *testing.T) { ctx := context.Background() notifier, eventStore := newStoreFn(t) @@ -62,7 +62,7 @@ func TestNotifier_lastEventResourceVersion(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierLastEventResourceVersion) } -func testNotifierLastEventResourceVersion(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierLastEventResourceVersion(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { // Test with no events rv, err := notifier.lastEventResourceVersion(ctx) assert.Error(t, err) @@ -113,7 +113,7 @@ func TestNotifier_cachekey(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierCachekey) } -func testNotifierCachekey(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierCachekey(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { tests := []struct { name string event Event @@ -167,7 +167,7 @@ func TestNotifier_Watch_NoEvents(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchNoEvents) } -func testNotifierWatchNoEvents(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierWatchNoEvents(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { ctx, cancel := context.WithTimeout(ctx, 500*time.Millisecond) defer cancel() @@ -208,7 +208,7 @@ func TestNotifier_Watch_WithExistingEvents(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchWithExistingEvents) } -func testNotifierWatchWithExistingEvents(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierWatchWithExistingEvents(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { ctx, cancel := context.WithTimeout(ctx, 2*time.Second) defer cancel() @@ -282,7 +282,7 @@ func TestNotifier_Watch_EventDeduplication(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchEventDeduplication) } -func testNotifierWatchEventDeduplication(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierWatchEventDeduplication(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { ctx, cancel := context.WithTimeout(ctx, 2*time.Second) defer cancel() @@ -348,7 +348,7 @@ func TestNotifier_Watch_ContextCancellation(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchContextCancellation) } -func testNotifierWatchContextCancellation(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierWatchContextCancellation(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { ctx, cancel := context.WithCancel(ctx) // Add an initial event so that lastEventResourceVersion doesn't return ErrNotFound @@ -394,7 +394,7 @@ func TestNotifier_Watch_MultipleEvents(t *testing.T) { runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchMultipleEvents) } -func testNotifierWatchMultipleEvents(t *testing.T, ctx context.Context, notifier *notifier, eventStore *eventStore) { +func testNotifierWatchMultipleEvents(t *testing.T, ctx context.Context, notifier *pollingNotifier, eventStore *eventStore) { ctx, cancel := context.WithTimeout(ctx, 3*time.Second) defer cancel() rv := time.Now().UnixNano() diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index 4db6da89d9a..f285daf0cb8 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -61,7 +61,7 @@ type kvStorageBackend struct { bulkLock *BulkLock dataStore *dataStore eventStore *eventStore - notifier *notifier + notifier notifier builder DocumentBuilder log logging.Logger withPruner bool @@ -91,6 +91,7 @@ type KVBackendOptions struct { Tracer trace.Tracer // TODO add tracing Reg prometheus.Registerer // TODO add metrics + UseChannelNotifier bool // Adding RvManager overrides the RV generated with snowflake in order to keep backwards compatibility with // unified/sql RvManager *rvmanager.ResourceVersionManager @@ -121,7 +122,7 @@ func NewKVStorageBackend(opts KVBackendOptions) (KVBackend, error) { bulkLock: NewBulkLock(), dataStore: newDataStore(kv), eventStore: eventStore, - notifier: newNotifier(eventStore, notifierOptions{}), + notifier: newNotifier(eventStore, notifierOptions{useChannelNotifier: opts.UseChannelNotifier}), snowflake: s, builder: StandardDocumentBuilder(), // For now we use the standard document builder. log: &logging.NoOpLogger{}, // Make this configurable diff --git a/pkg/storage/unified/sql/server.go b/pkg/storage/unified/sql/server.go index f4a1ee3ce77..b96d559ba55 100644 --- a/pkg/storage/unified/sql/server.go +++ b/pkg/storage/unified/sql/server.go @@ -99,6 +99,9 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, error) { return nil, err } + isHA := isHighAvailabilityEnabled(opts.Cfg.SectionWithEnvOverrides("database"), + opts.Cfg.SectionWithEnvOverrides("resource_api")) + if opts.Cfg.EnableSQLKVBackend { sqlkv, err := resource.NewSQLKV(eDB) if err != nil { @@ -106,9 +109,10 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, error) { } kvBackendOpts := resource.KVBackendOptions{ - KvStore: sqlkv, - Tracer: opts.Tracer, - Reg: opts.Reg, + KvStore: sqlkv, + Tracer: opts.Tracer, + Reg: opts.Reg, + UseChannelNotifier: !isHA, } ctx := context.Background() @@ -140,9 +144,6 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, error) { serverOptions.Backend = kvBackend serverOptions.Diagnostics = kvBackend } else { - isHA := isHighAvailabilityEnabled(opts.Cfg.SectionWithEnvOverrides("database"), - opts.Cfg.SectionWithEnvOverrides("resource_api")) - backend, err := NewBackend(BackendOptions{ DBProvider: eDB, Reg: opts.Reg,