From 0f626abd97d5557a3bff9814893efa0197e2a680 Mon Sep 17 00:00:00 2001 From: SungJin1212 Date: Sat, 26 Sep 2026 06:41:15 +0900 Subject: [PATCH 1/2] fix: correct misleading error messages referring to compactor ring in parquet converter (#7755) Signed-off-by: SungJin1212 Signed-off-by: Friedrich Gonzalez <1517449+friedrichg@users.noreply.github.com> Co-authored-by: Friedrich Gonzalez <1517449+friedrichg@users.noreply.github.com> (cherry picked from commit 1f6ce6bafcfb81a08da1c807861c4623c249471a) Signed-off-by: Charlie Le --- CHANGELOG.md | 1 + pkg/parquetconverter/converter.go | 12 ++++++------ 2 files changed, 7 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5dc60180890..85d70f42c67 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -119,6 +119,7 @@ * [BUGFIX] Compactor: Properly handle error from ReadPartitionedGroupInfo in UpdatePartitionedGroupInfo. #7766 * [BUGFIX] Alertmanager: Reject the global `mattermost_webhook_url_file` setting in per-tenant configs, consistent with every other global `*_file` setting. #7768 * [BUGFIX] Alertmanager: Tighten per-tenant config validation to reject additional file-based settings. #7767 +* [BUGFIX] Parquet Converter: Fix misleading error messages referring to the compactor ring instead of the parquet converter ring during startup and sharding checks. #7755 * [BUGFIX] Querier: Fix panic (`index out of range [-1]`) in the active request tracker when truncating a `match[]`/`query` value made entirely of invalid UTF-8 continuation bytes. The backwards scan for a rune boundary now stops at index 0 instead of underflowing. #7743 * [BUGFIX] Config: Fix CSV-list flags/YAML fields (e.g. `-compactor.enabled-tenants`) treating an explicitly empty string as a one-element list containing an empty tenant name instead of an empty list. #7714 * [BUGFIX] Tenant Federation: Fix regex tenant federation dropping tenants when `-blocks-storage.users-scanner.cache-ttl` is set. The regex resolver sorted the user list returned by the users scanner in place, corrupting the scanner cache and progressively losing tenants on every sync until the cache expired. #7812 diff --git a/pkg/parquetconverter/converter.go b/pkg/parquetconverter/converter.go index 3f72a78af29..06e46442562 100644 --- a/pkg/parquetconverter/converter.go +++ b/pkg/parquetconverter/converter.go @@ -45,7 +45,7 @@ import ( ) const ( - // ringKey is the key under which we store the compactors ring in the KVStore. + // ringKey is the key under which we store the parquet converters ring in the KVStore. ringKey = "parquet-converter" converterMetaPrefix = "converter-meta-" @@ -184,12 +184,12 @@ func (c *Converter) starting(ctx context.Context) error { } c.ringLifecycler, err = ring.NewLifecyclerWithDelegate(lifecyclerCfg, ring.NewNoopFlushTransferer(), "parquet-converter", ringKey, true, false, c.logger, prometheus.WrapRegistererWithPrefix("cortex_", c.reg), delegate) if err != nil { - return errors.Wrap(err, "unable to initialize converter ring lifecycler") + return errors.Wrap(err, "unable to initialize parquet converter ring lifecycler") } c.ring, err = ring.New(lifecyclerCfg.RingConfig, "parquet-converter", ringKey, c.logger, prometheus.WrapRegistererWithPrefix("cortex_", c.reg)) if err != nil { - return errors.Wrap(err, "unable to initialize compactor ring") + return errors.Wrap(err, "unable to initialize parquet converter ring") } c.ringSubservices, err = services.NewManager(c.ringLifecycler, c.ring) @@ -200,7 +200,7 @@ func (c *Converter) starting(ctx context.Context) error { err = services.StartManagerAndAwaitHealthy(ctx, c.ringSubservices) } if err != nil { - return errors.Wrap(err, "unable to start compactor ring dependencies") + return errors.Wrap(err, "unable to start parquet converter ring dependencies") } ctxWithTimeout, cancel := context.WithTimeout(ctx, time.Minute*3) @@ -633,14 +633,14 @@ func (c *Converter) ownBlock(ring ring.ReadRing, blockId string) (bool, error) { _, _ = hasher.Write([]byte(blockId)) userHash := hasher.Sum32() - // Check whether this compactor instance owns the user. + // Check whether this parquet converter instance owns the user. rs, err := ring.Get(userHash, RingOp, nil, nil, nil) if err != nil { return false, err } if len(rs.Instances) != 1 { - return false, fmt.Errorf("unexpected number of compactors in the shard (expected 1, got %d)", len(rs.Instances)) + return false, fmt.Errorf("unexpected number of parquet converters in the shard (expected 1, got %d)", len(rs.Instances)) } return rs.Instances[0].Addr == c.ringLifecycler.Addr, nil From 03df5a5a8b9a67fd45a5937e87adee8ec58731d8 Mon Sep 17 00:00:00 2001 From: Sandy Chen Date: Fri, 2 Oct 2026 11:15:34 +0900 Subject: [PATCH 2/2] Compactor: fix spurious tenant cleanup failures and TestBlocksCleaner flake (#7861) TestBlocksCleaner still flakes after #7386. That change made the per-tenant cleanup wait for the visit marker HeartBeat goroutine to exit, closing the race between the heartbeat's final writes and the test's t.Cleanup. The heartbeat's first write can still overlap the cleaner's own deletions under the same tenant: the filesystem bucket's Delete removes each emptied parent directory with os.RemoveAll, and a concurrent Upload of /markers/cleaner-visit-marker.json makes it fail with "unlinkat /user-4/markers: directory not empty". deleteUserMarkedForDeletion then returns before deleting the tenant deletion mark, which is the failure seen on CI (user-4's mark still exists at blocks_cleaner_test.go:275). Object stores have no directories, so the race only affects the filesystem bucket, where a failed run is retried on the next cleanup cycle. Switch TestBlocksCleaner to the in-memory bucket, like #7486 did for TestBlocksCleaner_ShouldRemoveBlocksOutsideRetentionPeriod, and assert that no cleaner visit marker is left once the initial cleanup has completed, so the test keeps guarding the #7386 fix. Doing so exposed a production bug: DeleteTenantDeletionMark deletes the mark from the global and then the legacy per-tenant location and returns any error. The per-tenant mark is never there at that point (since #5676 marks are only written to the global location, and the tenant's markers/ prefix was deleted just before), so on object stores whose Delete fails for a missing object (GCS, Azure, Swift, OCI) the final cleanup of every deleted tenant is reported as failed. Ignore not-found errors, like bucketindex.DeleteIndex and DeleteIndexSyncStatus already do, and assert in TestBlocksCleaner that the deleted-tenants run completes without failures. Fixes #7564 Signed-off-by: Sandy Chen (cherry picked from commit cc57488b21e336c5a5a1f7367723344249e4e6e7) Signed-off-by: Charlie Le --- CHANGELOG.md | 1 + pkg/compactor/blocks_cleaner_test.go | 15 +++++- pkg/util/users/tenant_deletion_mark.go | 7 +-- pkg/util/users/tenant_deletion_mark_test.go | 54 +++++++++++++++++++++ 4 files changed, 73 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 85d70f42c67..3149cb997d1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -125,6 +125,7 @@ * [BUGFIX] Tenant Federation: Fix regex tenant federation dropping tenants when `-blocks-storage.users-scanner.cache-ttl` is set. The regex resolver sorted the user list returned by the users scanner in place, corrupting the scanner cache and progressively losing tenants on every sync until the cache expired. #7812 * [BUGFIX] Tenant Federation: Fix regex tenant federation resolving to an empty user list right after startup. #7811 * [BUGFIX] Ingester: Don't count a forced head compaction skipped because blocks shipping is in progress as a failure. Previously such skips incremented `cortex_ingester_tsdb_compactions_failed_total`, producing spurious alerts. #7842 +* [BUGFIX] Compactor: Fix the final cleanup of a tenant marked for deletion being reported as failed on object stores that return an error when deleting a missing object (GCS, Azure, Swift, OCI). #7861 ## 1.21.1 2026-06-04 diff --git a/pkg/compactor/blocks_cleaner_test.go b/pkg/compactor/blocks_cleaner_test.go index a30a1e665e8..db9e88989eb 100644 --- a/pkg/compactor/blocks_cleaner_test.go +++ b/pkg/compactor/blocks_cleaner_test.go @@ -133,7 +133,9 @@ func TestBlockCleaner_KeyPermissionDenied(t *testing.T) { } func testBlocksCleanerWithOptions(t *testing.T, options testBlocksCleanerOptions) { - bucketClient, _ := cortex_testutil.PrepareFilesystemBucket(t) + // Use an in-memory bucket: the filesystem bucket's Delete also removes emptied parent + // directories, which races with the visit marker heartbeat writing under the same tenant. + bucketClient := objstore.WithNoopInstr(objstore.NewInMemBucket()) // If the markers migration is enabled, then we create the fixture blocks without // writing the deletion marks in the global location, because they will be migrated @@ -224,6 +226,14 @@ func testBlocksCleanerWithOptions(t *testing.T, options testBlocksCleanerOptions require.NoError(t, services.StartAndAwaitRunning(ctx, cleaner)) defer services.StopAndAwaitTerminated(ctx, cleaner) //nolint:errcheck + // The cleanup of each tenant waits for the visit marker heartbeat to delete the cleaner + // visit marker before returning, so none is left once the initial cleanup has completed. + for _, userID := range []string{"user-1", "user-2", "user-3", "user-4", "user-5", "user-6"} { + exists, err := bucketClient.Exists(ctx, path.Join(userID, bucketindex.MarkersPathname, CleanerVisitMarkerName)) + require.NoError(t, err) + assert.False(t, exists, userID) + } + for _, tc := range []struct { path string expectedExists bool @@ -278,6 +288,9 @@ func testBlocksCleanerWithOptions(t *testing.T, options testBlocksCleanerOptions assert.Equal(t, float64(1), prom_testutil.ToFloat64(cleaner.runsStarted.WithLabelValues(activeStatus))) assert.Equal(t, float64(1), prom_testutil.ToFloat64(cleaner.runsCompleted.WithLabelValues(activeStatus))) assert.Equal(t, float64(0), prom_testutil.ToFloat64(cleaner.runsFailed.WithLabelValues(activeStatus))) + assert.Equal(t, float64(1), prom_testutil.ToFloat64(cleaner.runsStarted.WithLabelValues(deletedStatus))) + assert.Equal(t, float64(1), prom_testutil.ToFloat64(cleaner.runsCompleted.WithLabelValues(deletedStatus))) + assert.Equal(t, float64(0), prom_testutil.ToFloat64(cleaner.runsFailed.WithLabelValues(deletedStatus))) assert.Equal(t, float64(7), prom_testutil.ToFloat64(cleaner.blocksCleanedTotal)) assert.Equal(t, float64(0), prom_testutil.ToFloat64(cleaner.blocksFailedTotal)) diff --git a/pkg/util/users/tenant_deletion_mark.go b/pkg/util/users/tenant_deletion_mark.go index 622f1e8bdb1..5c7b1e4970e 100644 --- a/pkg/util/users/tenant_deletion_mark.go +++ b/pkg/util/users/tenant_deletion_mark.go @@ -59,12 +59,13 @@ func ReadTenantDeletionMark(ctx context.Context, bkt objstore.InstrumentedBucket return read(ctx, bkt.WithExpectedErrs(bkt.IsObjNotFoundErr), markerFile, logger) } -// Deletes the tenant deletion mark for given user if it exists. +// Deletes the tenant deletion mark for given user from both the global and the local location, +// if it exists. Not-found errors are ignored for both locations. func DeleteTenantDeletionMark(ctx context.Context, bkt objstore.Bucket, userID string) error { - if err := bkt.Delete(ctx, GetGlobalDeletionMarkPath(userID)); err != nil { + if err := bkt.Delete(ctx, GetGlobalDeletionMarkPath(userID)); err != nil && !bkt.IsObjNotFoundErr(err) { return err } - if err := bkt.Delete(ctx, GetLocalDeletionMarkPath(userID)); err != nil { + if err := bkt.Delete(ctx, GetLocalDeletionMarkPath(userID)); err != nil && !bkt.IsObjNotFoundErr(err) { return err } return nil diff --git a/pkg/util/users/tenant_deletion_mark_test.go b/pkg/util/users/tenant_deletion_mark_test.go index 5da53554260..9c5b0f2ab77 100644 --- a/pkg/util/users/tenant_deletion_mark_test.go +++ b/pkg/util/users/tenant_deletion_mark_test.go @@ -7,6 +7,8 @@ import ( "github.com/stretchr/testify/require" "github.com/thanos-io/objstore" + + "github.com/cortexproject/cortex/pkg/util/testutil" ) func TestTenantDeletionMarkExists(t *testing.T) { @@ -68,3 +70,55 @@ func TestTenantDeletionMarkExists(t *testing.T) { }) } } + +func TestDeleteTenantDeletionMark(t *testing.T) { + const username = "user" + + for name, tc := range map[string]struct { + objects []string + deleteFailures []string + expectedErr string + }{ + "only global mark exists": { + objects: []string{GetGlobalDeletionMarkPath(username)}, + }, + "only local mark exists": { + objects: []string{GetLocalDeletionMarkPath(username)}, + }, + "both marks exist": { + objects: []string{GetGlobalDeletionMarkPath(username), GetLocalDeletionMarkPath(username)}, + }, + "no mark exists": { + objects: nil, + }, + "failure deleting global mark": { + objects: []string{GetGlobalDeletionMarkPath(username)}, + deleteFailures: []string{GetGlobalDeletionMarkPath(username)}, + expectedErr: "mocked delete failure", + }, + "failure deleting local mark": { + objects: []string{GetGlobalDeletionMarkPath(username), GetLocalDeletionMarkPath(username)}, + deleteFailures: []string{GetLocalDeletionMarkPath(username)}, + expectedErr: "mocked delete failure", + }, + } { + t.Run(name, func(t *testing.T) { + // Like GCS, Azure, Swift and OCI, the in-memory bucket returns an error when deleting a missing object. + bkt := objstore.NewInMemBucket() + for _, objName := range tc.objects { + require.NoError(t, bkt.Upload(context.Background(), objName, bytes.NewReader([]byte("data")))) + } + + err := DeleteTenantDeletionMark(context.Background(), &testutil.MockBucketFailure{Bucket: bkt, DeleteFailures: tc.deleteFailures}, username) + if tc.expectedErr != "" { + require.ErrorContains(t, err, tc.expectedErr) + return + } + require.NoError(t, err) + + exists, err := TenantDeletionMarkExists(context.Background(), bkt, username) + require.NoError(t, err) + require.False(t, exists) + }) + } +}