diff --git a/CHANGELOG.md b/CHANGELOG.md index 5dc60180890..3149cb997d1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -119,11 +119,13 @@ * [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 * [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/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 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) + }) + } +}