diff --git a/pkg/storage/unified/entitybridge/entitybridge.go b/pkg/storage/unified/entitybridge/entitybridge.go index abc5fc16f9c..7d47c001b27 100644 --- a/pkg/storage/unified/entitybridge/entitybridge.go +++ b/pkg/storage/unified/entitybridge/entitybridge.go @@ -28,7 +28,7 @@ func ProvideResourceServer(db db.DB, cfg *setting.Cfg, features featuremgmt.Feat Tracer: tracer, } - useEntitySQL := true + useEntitySQL := false if useEntitySQL { eDB, err := dbimpl.ProvideEntityDB(db, cfg, features, tracer) if err != nil { diff --git a/pkg/storage/unified/resource/cdk_backend.go b/pkg/storage/unified/resource/cdk_backend.go index c37ce5c518c..fc201b13413 100644 --- a/pkg/storage/unified/resource/cdk_backend.go +++ b/pkg/storage/unified/resource/cdk_backend.go @@ -69,8 +69,9 @@ type cdkBackend struct { nextRV NextResourceVersion mutex sync.Mutex - // Typically one... the server wrapper - subscribers []chan *WrittenEvent + // Simple watch stream -- NOTE, this only works for single tenant! + broadcaster Broadcaster[*WrittenEvent] + stream chan<- *WrittenEvent } func (s *cdkBackend) getPath(key *ResourceKey, rv int64) string { @@ -123,24 +124,19 @@ func (s *cdkBackend) WriteEvent(ctx context.Context, event WriteEvent) (rv int64 } // Async notify all subscribers - if s.subscribers != nil { + if s.stream != nil { go func() { write := &WrittenEvent{ - WriteEvent: event, - + WriteEvent: event, Timestamp: time.Now().UnixMilli(), ResourceVersion: rv, } - for _, sub := range s.subscribers { - sub <- write - } + s.stream <- write }() } - return rv, err } -// Read implements ResourceStoreServer. func (s *cdkBackend) Read(ctx context.Context, req *ReadRequest) (*ReadResponse, error) { rv := req.ResourceVersion @@ -191,7 +187,6 @@ func isDeletedMarker(raw []byte) bool { return false } -// List implements AppendingStore. func (s *cdkBackend) PrepareList(ctx context.Context, req *ListRequest) (*ListResponse, error) { resources, err := buildTree(ctx, s, req.Options.Key) if err != nil { @@ -215,36 +210,21 @@ func (s *cdkBackend) PrepareList(ctx context.Context, req *ListRequest) (*ListRe return rsp, nil } -// Watch implements AppendingStore. func (s *cdkBackend) WatchWriteEvents(ctx context.Context) (<-chan *WrittenEvent, error) { - stream := make(chan *WrittenEvent, 10) - { - s.mutex.Lock() - defer s.mutex.Unlock() + s.mutex.Lock() + defer s.mutex.Unlock() - // Add the event stream - s.subscribers = append(s.subscribers, stream) - } - - // Wait for context done - go func() { - // Wait till the context is done - <-ctx.Done() - - // Then remove the subscription - s.mutex.Lock() - defer s.mutex.Unlock() - - // Copy all streams without our listener - subs := []chan *WrittenEvent{} - for _, sub := range s.subscribers { - if sub != stream { - subs = append(subs, sub) - } + if s.broadcaster == nil { + var err error + s.broadcaster, err = NewBroadcaster(context.Background(), func(c chan<- *WrittenEvent) error { + s.stream = c + return nil + }) + if err != nil { + return nil, err } - s.subscribers = subs - }() - return stream, nil + } + return s.broadcaster.Subscribe(ctx) } // group > resource > namespace > name > versions diff --git a/pkg/storage/unified/resource/server_test.go b/pkg/storage/unified/resource/server_test.go index 4589b071dcc..e3bfd646e08 100644 --- a/pkg/storage/unified/resource/server_test.go +++ b/pkg/storage/unified/resource/server_test.go @@ -2,7 +2,6 @@ package resource import ( "context" - "embed" "encoding/json" "fmt" "os" @@ -39,7 +38,6 @@ func TestSimpleServer(t *testing.T) { Metadata: fileblob.MetadataDontWrite, // skip }) require.NoError(t, err) - fmt.Printf("ROOT: %s\n\n", tmp) } store, err := NewCDKBackend(ctx, CDKBackendOptions{ @@ -53,7 +51,30 @@ func TestSimpleServer(t *testing.T) { require.NoError(t, err) t.Run("playlist happy CRUD paths", func(t *testing.T) { - raw := testdata(t, "01_create_playlist.json") + raw := []byte(`{ + "apiVersion": "playlist.grafana.app/v0alpha1", + "kind": "Playlist", + "metadata": { + "name": "fdgsv37qslr0ga", + "namespace": "default", + "annotations": { + "grafana.app/originName": "elsewhere", + "grafana.app/originPath": "path/to/item", + "grafana.app/originTimestamp": "2024-02-02T00:00:00Z" + } + }, + "spec": { + "title": "hello", + "interval": "5m", + "items": [ + { + "type": "dashboard_by_uid", + "value": "vmie2cmWz" + } + ] + } + }`) + key := &ResourceKey{ Group: "playlist.grafana.app", Resource: "rrrr", // can be anything :( @@ -144,13 +165,3 @@ func TestSimpleServer(t *testing.T) { require.Len(t, all.Items, 0) // empty }) } - -//go:embed testdata/* -var testdataFS embed.FS - -func testdata(t *testing.T, filename string) []byte { - t.Helper() - b, err := testdataFS.ReadFile(`testdata/` + filename) - require.NoError(t, err) - return b -} diff --git a/pkg/storage/unified/resource/testdata/01_create_playlist.json b/pkg/storage/unified/resource/testdata/01_create_playlist.json deleted file mode 100644 index 151435fe1c6..00000000000 --- a/pkg/storage/unified/resource/testdata/01_create_playlist.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "apiVersion": "playlist.grafana.app/v0alpha1", - "kind": "Playlist", - "metadata": { - "name": "fdgsv37qslr0ga", - "namespace": "default", - "annotations": { - "grafana.app/originName": "elsewhere", - "grafana.app/originPath": "path/to/item", - "grafana.app/originTimestamp": "2024-02-02T00:00:00Z" - } - }, - "spec": { - "title": "hello", - "interval": "5m", - "items": [ - { - "type": "dashboard_by_uid", - "value": "vmie2cmWz" - } - ] - } -} \ No newline at end of file