Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
15 changes: 14 additions & 1 deletion pkg/compactor/blocks_cleaner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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))

Expand Down
12 changes: 6 additions & 6 deletions pkg/parquetconverter/converter.go
Original file line number Diff line number Diff line change
Expand Up @@ -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-"
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand Down
7 changes: 4 additions & 3 deletions pkg/util/users/tenant_deletion_mark.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
54 changes: 54 additions & 0 deletions pkg/util/users/tenant_deletion_mark_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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)
})
}
}