From 2f520454aeedd1d221c1bb6ac263e06e642f8898 Mon Sep 17 00:00:00 2001 From: Renato Costa <103441181+renatolabs@users.noreply.github.com> Date: Wed, 14 Jan 2026 14:21:46 -0500 Subject: [PATCH] unified-storage: save `previous_resource_version` as microsecond timestamp in compat mode (#116179) * unified-storage: save `previous_resource_version` as microsecond timestamp in compat mode * debug test * fix test * fix typo * Fix misleading var name. `ms` is typically the abbreviation for millisecond, not microsecond. --- pkg/storage/unified/resource/datastore.go | 27 ++++++++++-- pkg/storage/unified/resource/notifier_test.go | 24 ++++------- pkg/storage/unified/resource/sqlkv.go | 2 - .../unified/resource/storage_backend.go | 2 +- .../unified/sql/rvmanager/rv_manager.go | 16 +++++-- .../unified/sql/rvmanager/rv_manager_test.go | 10 +++++ .../storage_backend_sql_compatibility.go | 42 ++++++++----------- 7 files changed, 73 insertions(+), 50 deletions(-) diff --git a/pkg/storage/unified/resource/datastore.go b/pkg/storage/unified/resource/datastore.go index 931a5b20560..0cc56e51d72 100644 --- a/pkg/storage/unified/resource/datastore.go +++ b/pkg/storage/unified/resource/datastore.go @@ -14,6 +14,7 @@ import ( "github.com/grafana/grafana/pkg/apimachinery/validation" "github.com/grafana/grafana/pkg/storage/unified/sql/db" "github.com/grafana/grafana/pkg/storage/unified/sql/dbutil" + "github.com/grafana/grafana/pkg/storage/unified/sql/rvmanager" "github.com/grafana/grafana/pkg/storage/unified/sql/sqltemplate" gocache "github.com/patrickmn/go-cache" ) @@ -868,10 +869,18 @@ func (d *dataStore) applyBackwardsCompatibleChanges(ctx context.Context, tx db.T if key.Action == DataActionDeleted { generation = 0 } + + // In compatibility mode, the previous RV, when available, is saved as a microsecond + // timestamp, as is done in the SQL backend. + previousRV := event.PreviousRV + if event.PreviousRV > 0 && isSnowflake(event.PreviousRV) { + previousRV = rvmanager.RVFromSnowflake(event.PreviousRV) + } + _, err := dbutil.Exec(ctx, tx, sqlKVUpdateLegacyResourceHistory, sqlKVLegacyUpdateHistoryRequest{ SQLTemplate: sqltemplate.New(kv.dialect), GUID: key.GUID, - PreviousRV: event.PreviousRV, + PreviousRV: previousRV, Generation: generation, }) @@ -900,7 +909,7 @@ func (d *dataStore) applyBackwardsCompatibleChanges(ctx context.Context, tx db.T Name: key.Name, Action: action, Folder: key.Folder, - PreviousRV: event.PreviousRV, + PreviousRV: previousRV, }) if err != nil { @@ -916,7 +925,7 @@ func (d *dataStore) applyBackwardsCompatibleChanges(ctx context.Context, tx db.T Name: key.Name, Action: action, Folder: key.Folder, - PreviousRV: event.PreviousRV, + PreviousRV: previousRV, }) if err != nil { @@ -938,3 +947,15 @@ func (d *dataStore) applyBackwardsCompatibleChanges(ctx context.Context, tx db.T return nil } + +// isSnowflake returns whether the argument passed is a snowflake ID (new) or a microsecond timestamp (old). +// We try to interpret the number as a microsecond timestamp first. If it represents a time in the past, +// it is considered a microsecond timestamp. Snowflake IDs are much larger integers and would lead +// to dates in the future if interpreted as a microsecond timestamp. +func isSnowflake(rv int64) bool { + ts := time.UnixMicro(rv) + oneHourFromNow := time.Now().Add(time.Hour) + isMicroSecRV := ts.Before(oneHourFromNow) + + return !isMicroSecRV +} diff --git a/pkg/storage/unified/resource/notifier_test.go b/pkg/storage/unified/resource/notifier_test.go index 5c217bf0a21..5e1c0dacea8 100644 --- a/pkg/storage/unified/resource/notifier_test.go +++ b/pkg/storage/unified/resource/notifier_test.go @@ -456,33 +456,27 @@ func testNotifierWatchMultipleEvents(t *testing.T, ctx context.Context, notifier }, } + errCh := make(chan error) go func() { for _, event := range testEvents { - err := eventStore.Save(ctx, event) - require.NoError(t, err) + errCh <- eventStore.Save(ctx, event) } }() // Receive events - receivedEvents := make([]Event, 0, len(testEvents)) - for i := 0; i < len(testEvents); i++ { + receivedEvents := make([]string, 0, len(testEvents)) + for len(receivedEvents) != len(testEvents) { select { case event := <-events: - receivedEvents = append(receivedEvents, event) + receivedEvents = append(receivedEvents, event.Name) + case err := <-errCh: + require.NoError(t, err) case <-time.After(1 * time.Second): - t.Fatalf("Timed out waiting for event %d", i+1) + t.Fatalf("Timed out waiting for event %d", len(receivedEvents)+1) } } - // Verify all events were received - assert.Len(t, receivedEvents, len(testEvents)) - // Verify the events match and ordered by resource version - receivedNames := make([]string, len(receivedEvents)) - for i, event := range receivedEvents { - receivedNames[i] = event.Name - } - expectedNames := []string{"test-resource-1", "test-resource-2", "test-resource-3"} - assert.ElementsMatch(t, expectedNames, receivedNames) + assert.ElementsMatch(t, expectedNames, receivedEvents) } diff --git a/pkg/storage/unified/resource/sqlkv.go b/pkg/storage/unified/resource/sqlkv.go index 9651c65f2f6..a0d47b69c35 100644 --- a/pkg/storage/unified/resource/sqlkv.go +++ b/pkg/storage/unified/resource/sqlkv.go @@ -473,8 +473,6 @@ func (k *sqlKV) Delete(ctx context.Context, section string, key string) error { return ErrNotFound } - // TODO reflect change to resource table - return nil } diff --git a/pkg/storage/unified/resource/storage_backend.go b/pkg/storage/unified/resource/storage_backend.go index f285daf0cb8..56272662c24 100644 --- a/pkg/storage/unified/resource/storage_backend.go +++ b/pkg/storage/unified/resource/storage_backend.go @@ -347,7 +347,7 @@ func (k *kvStorageBackend) WriteEvent(ctx context.Context, event WriteEvent) (in return 0, fmt.Errorf("failed to write data: %w", err) } - rv = rvmanager.SnowflakeFromRv(rv) + rv = rvmanager.SnowflakeFromRV(rv) dataKey.ResourceVersion = rv } else { err := k.dataStore.Save(ctx, dataKey, bytes.NewReader(event.Value)) diff --git a/pkg/storage/unified/sql/rvmanager/rv_manager.go b/pkg/storage/unified/sql/rvmanager/rv_manager.go index b10685f22ad..42e1f9fb74c 100644 --- a/pkg/storage/unified/sql/rvmanager/rv_manager.go +++ b/pkg/storage/unified/sql/rvmanager/rv_manager.go @@ -307,7 +307,7 @@ func (m *ResourceVersionManager) execBatch(ctx context.Context, group, resource // Allocate the RVs for i, guid := range guids { guidToRV[guid] = rv - guidToSnowflakeRV[guid] = SnowflakeFromRv(rv) + guidToSnowflakeRV[guid] = SnowflakeFromRV(rv) rvs[i] = rv rv++ } @@ -364,12 +364,20 @@ func (m *ResourceVersionManager) execBatch(ctx context.Context, group, resource } } -// takes a unix microsecond rv and transforms into a snowflake format. The timestamp is converted from microsecond to +// takes a unix microsecond RV and transforms into a snowflake format. The timestamp is converted from microsecond to // millisecond (the integer division) and the remainder is saved in the stepbits section. machine id is always 0 -func SnowflakeFromRv(rv int64) int64 { +func SnowflakeFromRV(rv int64) int64 { return (((rv / 1000) - snowflake.Epoch) << (snowflake.NodeBits + snowflake.StepBits)) + (rv % 1000) } +// It is generally not possible to convert from a snowflakeID to a microsecond RV due to the loss in precision +// (snowflake ID stores timestamp in milliseconds). However, this implementation stores the microsecond fraction +// in the step bits (see SnowflakeFromRV), allowing us to compute the microsecond timestamp. +func RVFromSnowflake(snowflakeID int64) int64 { + microSecFraction := snowflakeID & ((1 << snowflake.StepBits) - 1) + return ((snowflakeID>>(snowflake.NodeBits+snowflake.StepBits))+snowflake.Epoch)*1000 + microSecFraction +} + // helper utility to compare two RVs. The first RV must be in snowflake format. Will convert rv2 to snowflake and retry // if comparison fails func IsRvEqual(rv1, rv2 int64) bool { @@ -377,7 +385,7 @@ func IsRvEqual(rv1, rv2 int64) bool { return true } - return rv1 == SnowflakeFromRv(rv2) + return rv1 == SnowflakeFromRV(rv2) } // Lock locks the resource version for the given key diff --git a/pkg/storage/unified/sql/rvmanager/rv_manager_test.go b/pkg/storage/unified/sql/rvmanager/rv_manager_test.go index 9a2e105aa2f..9da98d967d0 100644 --- a/pkg/storage/unified/sql/rvmanager/rv_manager_test.go +++ b/pkg/storage/unified/sql/rvmanager/rv_manager_test.go @@ -63,3 +63,13 @@ func TestResourceVersionManager(t *testing.T) { require.Equal(t, rv, int64(200)) }) } + +func TestSnowflakeFromRVRoundtrips(t *testing.T) { + // 2026-01-12 19:33:58.806211 +0000 UTC + offset := int64(1768246438806211) // in microseconds + + for n := range int64(100) { + ts := offset + n + require.Equal(t, ts, RVFromSnowflake(SnowflakeFromRV(ts))) + } +} diff --git a/pkg/storage/unified/testing/storage_backend_sql_compatibility.go b/pkg/storage/unified/testing/storage_backend_sql_compatibility.go index 5e1423fa1ec..d498ebba2f4 100644 --- a/pkg/storage/unified/testing/storage_backend_sql_compatibility.go +++ b/pkg/storage/unified/testing/storage_backend_sql_compatibility.go @@ -200,7 +200,7 @@ func verifyKeyPath(t *testing.T, db sqldb.DB, ctx context.Context, key *resource var keyPathRV int64 if isSqlBackend { // Convert microsecond RV to snowflake for key_path construction - keyPathRV = rvmanager.SnowflakeFromRv(resourceVersion) + keyPathRV = rvmanager.SnowflakeFromRV(resourceVersion) } else { // KV backend already provides snowflake RV keyPathRV = resourceVersion @@ -434,9 +434,6 @@ func verifyResourceHistoryTable(t *testing.T, db sqldb.DB, namespace string, res rows, err := db.QueryContext(ctx, query, namespace) require.NoError(t, err) - defer func() { - _ = rows.Close() - }() var records []ResourceHistoryRecord for rows.Next() { @@ -460,33 +457,34 @@ func verifyResourceHistoryTable(t *testing.T, db sqldb.DB, namespace string, res for resourceIdx, res := range resources { // Check create record (action=1, generation=1) createRecord := records[recordIndex] - verifyResourceHistoryRecord(t, createRecord, res, resourceIdx, 1, 0, 1, resourceVersions[resourceIdx][0]) + verifyResourceHistoryRecord(t, createRecord, namespace, res, resourceIdx, 1, 0, 1, resourceVersions[resourceIdx][0]) recordIndex++ } for resourceIdx, res := range resources { // Check update record (action=2, generation=2) updateRecord := records[recordIndex] - verifyResourceHistoryRecord(t, updateRecord, res, resourceIdx, 2, resourceVersions[resourceIdx][0], 2, resourceVersions[resourceIdx][1]) + verifyResourceHistoryRecord(t, updateRecord, namespace, res, resourceIdx, 2, resourceVersions[resourceIdx][0], 2, resourceVersions[resourceIdx][1]) recordIndex++ } for resourceIdx, res := range resources[:2] { // Check delete record (action=3, generation=0) - only first 2 resources were deleted deleteRecord := records[recordIndex] - verifyResourceHistoryRecord(t, deleteRecord, res, resourceIdx, 3, resourceVersions[resourceIdx][1], 0, resourceVersions[resourceIdx][2]) + verifyResourceHistoryRecord(t, deleteRecord, namespace, res, resourceIdx, 3, resourceVersions[resourceIdx][1], 0, resourceVersions[resourceIdx][2]) recordIndex++ } } // verifyResourceHistoryRecord validates a single resource_history record -func verifyResourceHistoryRecord(t *testing.T, record ResourceHistoryRecord, expectedRes struct{ name, folder string }, resourceIdx, expectedAction int, expectedPrevRV int64, expectedGeneration int, expectedRV int64) { +func verifyResourceHistoryRecord(t *testing.T, record ResourceHistoryRecord, namespace string, expectedRes struct{ name, folder string }, resourceIdx, expectedAction int, expectedPrevRV int64, expectedGeneration int, expectedRV int64) { // Validate GUID (should be non-empty) require.NotEmpty(t, record.GUID, "GUID should not be empty") // Validate group/resource/namespace/name require.Equal(t, "playlist.grafana.app", record.Group) require.Equal(t, "playlists", record.Resource) + require.Equal(t, namespace, record.Namespace) require.Equal(t, expectedRes.name, record.Name) // Validate value contains expected JSON - server modifies/formats the JSON differently for different operations @@ -513,8 +511,12 @@ func verifyResourceHistoryRecord(t *testing.T, record ResourceHistoryRecord, exp // For KV backend operations, expectedPrevRV is now in snowflake format (returned by KV backend) // but resource_history table stores microsecond RV, so we need to use IsRvEqual for comparison if strings.Contains(record.Namespace, "-kv") { - require.True(t, rvmanager.IsRvEqual(expectedPrevRV, record.PreviousResourceVersion), - "Previous resource version should match (KV backend snowflake format)") + if expectedPrevRV == 0 { + require.Zero(t, record.PreviousResourceVersion) + } else { + require.Equal(t, expectedPrevRV, rvmanager.SnowflakeFromRV(record.PreviousResourceVersion), + "Previous resource version should match (KV backend snowflake format)") + } } else { require.Equal(t, expectedPrevRV, record.PreviousResourceVersion) } @@ -546,9 +548,6 @@ func verifyResourceTable(t *testing.T, db sqldb.DB, namespace string, resources rows, err := db.QueryContext(ctx, query, namespace) require.NoError(t, err) - defer func() { - _ = rows.Close() - }() var records []ResourceRecord for rows.Next() { @@ -612,9 +611,6 @@ func verifyResourceVersionTable(t *testing.T, db sqldb.DB, namespace string, res // Check that we have exactly one entry for playlist.grafana.app/playlists rows, err := db.QueryContext(ctx, query, "playlist.grafana.app", "playlists") require.NoError(t, err) - defer func() { - _ = rows.Close() - }() var records []ResourceVersionRecord for rows.Next() { @@ -649,7 +645,7 @@ func verifyResourceVersionTable(t *testing.T, db sqldb.DB, namespace string, res isKvBackend := strings.Contains(namespace, "-kv") recordResourceVersion := record.ResourceVersion if isKvBackend { - recordResourceVersion = rvmanager.SnowflakeFromRv(record.ResourceVersion) + recordResourceVersion = rvmanager.SnowflakeFromRV(record.ResourceVersion) } require.Less(t, recordResourceVersion, int64(9223372036854775807), "resource_version should be reasonable") @@ -841,24 +837,20 @@ func runMixedConcurrentOperations(t *testing.T, sqlServer, kvServer resource.Res } // SQL backend operations - wg.Add(1) - go func() { - defer wg.Done() + wg.Go(func() { <-startBarrier // Wait for signal to start if err := runBackendOperationsWithCounts(ctx, sqlServer, namespace+"-sql", "sql", opCounts); err != nil { errors <- fmt.Errorf("SQL backend operations failed: %w", err) } - }() + }) // KV backend operations - wg.Add(1) - go func() { - defer wg.Done() + wg.Go(func() { <-startBarrier // Wait for signal to start if err := runBackendOperationsWithCounts(ctx, kvServer, namespace+"-kv", "kv", opCounts); err != nil { errors <- fmt.Errorf("KV backend operations failed: %w", err) } - }() + }) // Start both goroutines simultaneously close(startBarrier)