unified-storage: dont use polling notifier with sqlite in sqlkv (#116283)

* unified-storage: dont use polling notifier with sqlite in sqlkv
This commit is contained in:
Will Assis
2026-01-14 18:22:39 +00:00
committed by GitHub
parent 189d50d815
commit ba416eab4e
4 changed files with 45 additions and 27 deletions
+23 -7
View File
@@ -19,13 +19,18 @@ const (
defaultBufferSize = 10000 defaultBufferSize = 10000
) )
type notifier struct { type notifier interface {
Watch(context.Context, watchOptions) <-chan Event
}
type pollingNotifier struct {
eventStore *eventStore eventStore *eventStore
log logging.Logger log logging.Logger
} }
type notifierOptions struct { type notifierOptions struct {
log logging.Logger log logging.Logger
useChannelNotifier bool
} }
type watchOptions struct { 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 { if opts.log == nil {
opts.log = &logging.NoOpLogger{} opts.log = &logging.NoOpLogger{}
} }
return &notifier{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 // 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) e, err := n.eventStore.LastEventKey(ctx)
if err != nil { if err != nil {
return 0, err return 0, err
@@ -60,11 +76,11 @@ func (n *notifier) lastEventResourceVersion(ctx context.Context) (int64, error)
return e.ResourceVersion, nil 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) 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 { if opts.MinBackoff <= 0 {
opts.MinBackoff = defaultMinBackoff opts.MinBackoff = defaultMinBackoff
} }
+12 -12
View File
@@ -13,7 +13,7 @@ import (
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
) )
func setupTestNotifier(t *testing.T) (*notifier, *eventStore) { func setupTestNotifier(t *testing.T) (*pollingNotifier, *eventStore) {
db := setupTestBadgerDB(t) db := setupTestBadgerDB(t)
t.Cleanup(func() { t.Cleanup(func() {
err := db.Close() err := db.Close()
@@ -22,10 +22,10 @@ func setupTestNotifier(t *testing.T) (*notifier, *eventStore) {
kv := NewBadgerKV(db) kv := NewBadgerKV(db)
eventStore := newEventStore(kv) eventStore := newEventStore(kv)
notifier := newNotifier(eventStore, notifierOptions{log: &logging.NoOpLogger{}}) 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) dbstore := db.InitTestDB(t)
eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil) eDB, err := dbimpl.ProvideResourceDB(dbstore, setting.NewCfg(), nil)
require.NoError(t, err) require.NoError(t, err)
@@ -33,7 +33,7 @@ func setupTestNotifierSqlKv(t *testing.T) (*notifier, *eventStore) {
require.NoError(t, err) require.NoError(t, err)
eventStore := newEventStore(kv) eventStore := newEventStore(kv)
notifier := newNotifier(eventStore, notifierOptions{log: &logging.NoOpLogger{}}) notifier := newNotifier(eventStore, notifierOptions{log: &logging.NoOpLogger{}})
return notifier, eventStore return notifier.(*pollingNotifier), eventStore
} }
func TestNewNotifier(t *testing.T) { func TestNewNotifier(t *testing.T) {
@@ -49,7 +49,7 @@ func TestDefaultWatchOptions(t *testing.T) {
assert.Equal(t, defaultBufferSize, opts.BufferSize) 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) { t.Run(storeName, func(t *testing.T) {
ctx := context.Background() ctx := context.Background()
notifier, eventStore := newStoreFn(t) notifier, eventStore := newStoreFn(t)
@@ -62,7 +62,7 @@ func TestNotifier_lastEventResourceVersion(t *testing.T) {
runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierLastEventResourceVersion) 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 // Test with no events
rv, err := notifier.lastEventResourceVersion(ctx) rv, err := notifier.lastEventResourceVersion(ctx)
assert.Error(t, err) assert.Error(t, err)
@@ -113,7 +113,7 @@ func TestNotifier_cachekey(t *testing.T) {
runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierCachekey) 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 { tests := []struct {
name string name string
event Event event Event
@@ -167,7 +167,7 @@ func TestNotifier_Watch_NoEvents(t *testing.T) {
runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchNoEvents) 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) ctx, cancel := context.WithTimeout(ctx, 500*time.Millisecond)
defer cancel() defer cancel()
@@ -208,7 +208,7 @@ func TestNotifier_Watch_WithExistingEvents(t *testing.T) {
runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchWithExistingEvents) 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) ctx, cancel := context.WithTimeout(ctx, 2*time.Second)
defer cancel() defer cancel()
@@ -282,7 +282,7 @@ func TestNotifier_Watch_EventDeduplication(t *testing.T) {
runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchEventDeduplication) 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) ctx, cancel := context.WithTimeout(ctx, 2*time.Second)
defer cancel() defer cancel()
@@ -348,7 +348,7 @@ func TestNotifier_Watch_ContextCancellation(t *testing.T) {
runNotifierTestWith(t, "sqlkv", setupTestNotifierSqlKv, testNotifierWatchContextCancellation) 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) ctx, cancel := context.WithCancel(ctx)
// Add an initial event so that lastEventResourceVersion doesn't return ErrNotFound // 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) 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) ctx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel() defer cancel()
rv := time.Now().UnixNano() rv := time.Now().UnixNano()
@@ -61,7 +61,7 @@ type kvStorageBackend struct {
bulkLock *BulkLock bulkLock *BulkLock
dataStore *dataStore dataStore *dataStore
eventStore *eventStore eventStore *eventStore
notifier *notifier notifier notifier
builder DocumentBuilder builder DocumentBuilder
log logging.Logger log logging.Logger
withPruner bool withPruner bool
@@ -91,6 +91,7 @@ type KVBackendOptions struct {
Tracer trace.Tracer // TODO add tracing Tracer trace.Tracer // TODO add tracing
Reg prometheus.Registerer // TODO add metrics Reg prometheus.Registerer // TODO add metrics
UseChannelNotifier bool
// Adding RvManager overrides the RV generated with snowflake in order to keep backwards compatibility with // Adding RvManager overrides the RV generated with snowflake in order to keep backwards compatibility with
// unified/sql // unified/sql
RvManager *rvmanager.ResourceVersionManager RvManager *rvmanager.ResourceVersionManager
@@ -121,7 +122,7 @@ func NewKVStorageBackend(opts KVBackendOptions) (KVBackend, error) {
bulkLock: NewBulkLock(), bulkLock: NewBulkLock(),
dataStore: newDataStore(kv), dataStore: newDataStore(kv),
eventStore: eventStore, eventStore: eventStore,
notifier: newNotifier(eventStore, notifierOptions{}), notifier: newNotifier(eventStore, notifierOptions{useChannelNotifier: opts.UseChannelNotifier}),
snowflake: s, snowflake: s,
builder: StandardDocumentBuilder(), // For now we use the standard document builder. builder: StandardDocumentBuilder(), // For now we use the standard document builder.
log: &logging.NoOpLogger{}, // Make this configurable log: &logging.NoOpLogger{}, // Make this configurable
+7 -6
View File
@@ -99,6 +99,9 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, error) {
return nil, err return nil, err
} }
isHA := isHighAvailabilityEnabled(opts.Cfg.SectionWithEnvOverrides("database"),
opts.Cfg.SectionWithEnvOverrides("resource_api"))
if opts.Cfg.EnableSQLKVBackend { if opts.Cfg.EnableSQLKVBackend {
sqlkv, err := resource.NewSQLKV(eDB) sqlkv, err := resource.NewSQLKV(eDB)
if err != nil { if err != nil {
@@ -106,9 +109,10 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, error) {
} }
kvBackendOpts := resource.KVBackendOptions{ kvBackendOpts := resource.KVBackendOptions{
KvStore: sqlkv, KvStore: sqlkv,
Tracer: opts.Tracer, Tracer: opts.Tracer,
Reg: opts.Reg, Reg: opts.Reg,
UseChannelNotifier: !isHA,
} }
ctx := context.Background() ctx := context.Background()
@@ -140,9 +144,6 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, error) {
serverOptions.Backend = kvBackend serverOptions.Backend = kvBackend
serverOptions.Diagnostics = kvBackend serverOptions.Diagnostics = kvBackend
} else { } else {
isHA := isHighAvailabilityEnabled(opts.Cfg.SectionWithEnvOverrides("database"),
opts.Cfg.SectionWithEnvOverrides("resource_api"))
backend, err := NewBackend(BackendOptions{ backend, err := NewBackend(BackendOptions{
DBProvider: eDB, DBProvider: eDB,
Reg: opts.Reg, Reg: opts.Reg,