Skip to content
Open
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,7 @@
* [BUGFIX] Compactor: Fix spurious `bucket operation fail after retries` error logs emitted during partial block cleanup. #7749
* [BUGFIX] Alertmanager: Fix panic in `validateAlertmanagerConfig` when receiver config traversal encounters nil interface values. #7751
* [BUGFIX] Parquet Converter: Fix `auto_forget_delay` having no effect. The ring lifecycler was created without the auto-forget delegate, so unhealthy instances were never automatically removed from the ring. #7752
* [BUGFIX] Compactor: Properly handle error from ReadPartitionedGroupInfo in UpdatePartitionedGroupInfo. #7766

## 1.21.1 2026-06-04

Expand Down
8 changes: 5 additions & 3 deletions pkg/compactor/partitioned_group_info.go
Original file line number Diff line number Diff line change
Expand Up @@ -318,9 +318,11 @@ func ReadPartitionedGroupInfoFile(ctx context.Context, bkt objstore.Instrumented
}

func UpdatePartitionedGroupInfo(ctx context.Context, bkt objstore.InstrumentedBucket, logger log.Logger, partitionedGroupInfo PartitionedGroupInfo) (*PartitionedGroupInfo, error) {
// Ignore error in order to always update partitioned group info. There is no harm to put latest version of
// partitioned group info which is supposed to be the correct grouping based on latest bucket store.
existingPartitionedGroup, _ := ReadPartitionedGroupInfo(ctx, bkt, logger, partitionedGroupInfo.PartitionedGroupID)
existingPartitionedGroup, err := ReadPartitionedGroupInfo(ctx, bkt, logger, partitionedGroupInfo.PartitionedGroupID)
if err != nil && !errors.Is(err, ErrorPartitionedGroupInfoNotFound) {
level.Error(logger).Log("msg", "failed to check existing partitioned group info, skipping creation", "partitioned_group_id", partitionedGroupInfo.PartitionedGroupID, "err", err)
return nil, errors.Wrap(err, "unable to check existing partitioned group info")
}
if existingPartitionedGroup != nil {
level.Warn(logger).Log("msg", "partitioned group info already exists", "partitioned_group_id", partitionedGroupInfo.PartitionedGroupID, "partitioned_group_creation_time", partitionedGroupInfo.CreationTimeString())
return existingPartitionedGroup, nil
Expand Down
114 changes: 114 additions & 0 deletions pkg/compactor/partitioned_group_info_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (

"github.com/go-kit/log"
"github.com/oklog/ulid/v2"
"github.com/pkg/errors"
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"
"github.com/thanos-io/objstore"
Expand Down Expand Up @@ -899,3 +900,116 @@ func TestGetPartitionedGroupStatus(t *testing.T) {
})
}
}

func TestUpdatePartitionedGroupInfo_ErrorHandling(t *testing.T) {
partitionedGroupID := uint32(12345)
partitionedGroupInfo := PartitionedGroupInfo{
PartitionedGroupID: partitionedGroupID,
PartitionCount: 1,
Partitions: []Partition{
{
PartitionID: 0,
Blocks: []ulid.ULID{ulid.MustNew(0, nil)},
},
},
RangeStart: (1 * time.Hour).Milliseconds(),
RangeEnd: (2 * time.Hour).Milliseconds(),
Version: PartitionedGroupInfoVersion1,
}

t.Run("should return error when ReadPartitionedGroupInfo fails with non-NotFound error", func(t *testing.T) {
ctx := context.Background()
testBkt, _ := cortex_testutil.PrepareFilesystemBucket(t)
logger := log.NewNopLogger()

// Wrap the bucket to inject a Get failure on the partitioned group file path.
// This simulates S3 throttling or transient errors.
failingBkt := &cortex_testutil.MockBucketFailure{
Bucket: testBkt,
GetFailures: map[string]error{
GetPartitionedGroupFile(partitionedGroupID): errors.New("simulated S3 throttling error"),
},
}

result, err := UpdatePartitionedGroupInfo(ctx, failingBkt, logger, partitionedGroupInfo)
require.Error(t, err)
require.Nil(t, result)
require.Contains(t, err.Error(), "unable to check existing partitioned group info")
})

t.Run("should create partitioned group info when ReadPartitionedGroupInfo returns NotFound", func(t *testing.T) {
ctx := context.Background()
testBkt, _ := cortex_testutil.PrepareFilesystemBucket(t)
bkt := objstore.WithNoopInstr(testBkt)
logger := log.NewNopLogger()

// No pre-existing file — ReadPartitionedGroupInfo should return ErrorPartitionedGroupInfoNotFound,
// and UpdatePartitionedGroupInfo should proceed to create the file.
result, err := UpdatePartitionedGroupInfo(ctx, bkt, logger, partitionedGroupInfo)
require.NoError(t, err)
require.NotNil(t, result)
require.Equal(t, partitionedGroupID, result.PartitionedGroupID)
require.Greater(t, result.CreationTime, int64(0))

// Verify it was actually written to the bucket.
readResult, err := ReadPartitionedGroupInfo(ctx, bkt, logger, partitionedGroupID)
require.NoError(t, err)
require.Equal(t, result.PartitionedGroupID, readResult.PartitionedGroupID)
require.Equal(t, result.CreationTime, readResult.CreationTime)
})

t.Run("should return existing partitioned group info when it already exists", func(t *testing.T) {
ctx := context.Background()
testBkt, _ := cortex_testutil.PrepareFilesystemBucket(t)
bkt := objstore.WithNoopInstr(testBkt)
logger := log.NewNopLogger()

// First, create a partitioned group info.
firstResult, err := UpdatePartitionedGroupInfo(ctx, bkt, logger, partitionedGroupInfo)
require.NoError(t, err)
require.NotNil(t, firstResult)

// Try to update again with the same partitioned group ID — should return the existing one.
modifiedInfo := partitionedGroupInfo
modifiedInfo.PartitionCount = 5 // Different value to prove we get the original back
secondResult, err := UpdatePartitionedGroupInfo(ctx, bkt, logger, modifiedInfo)
require.NoError(t, err)
require.NotNil(t, secondResult)
// Should return the original, not the modified one.
require.Equal(t, firstResult.PartitionCount, secondResult.PartitionCount)
require.Equal(t, firstResult.CreationTime, secondResult.CreationTime)
})

t.Run("should not overwrite existing file when S3 Get fails", func(t *testing.T) {
ctx := context.Background()
testBkt, _ := cortex_testutil.PrepareFilesystemBucket(t)
bkt := objstore.WithNoopInstr(testBkt)
logger := log.NewNopLogger()

// First, create a partitioned group info successfully.
firstResult, err := UpdatePartitionedGroupInfo(ctx, bkt, logger, partitionedGroupInfo)
require.NoError(t, err)
require.NotNil(t, firstResult)
originalCreationTime := firstResult.CreationTime

// Now wrap the bucket to simulate S3 throttling on the existence check.
failingBkt := &cortex_testutil.MockBucketFailure{
Bucket: testBkt,
GetFailures: map[string]error{
GetPartitionedGroupFile(partitionedGroupID): errors.New("simulated S3 throttling error"),
},
}

// Attempt to update — this should fail, NOT overwrite.
modifiedInfo := partitionedGroupInfo
modifiedInfo.CreationTime = time.Now().Unix() + 9999 // A clearly different creation time
secondResult, err := UpdatePartitionedGroupInfo(ctx, failingBkt, logger, modifiedInfo)
require.Error(t, err)
require.Nil(t, secondResult)

// Verify the original file was NOT overwritten.
readResult, err := ReadPartitionedGroupInfo(ctx, bkt, logger, partitionedGroupID)
require.NoError(t, err)
require.Equal(t, originalCreationTime, readResult.CreationTime)
})
}
Loading