diff --git a/CHANGELOG.md b/CHANGELOG.md index abd8bcc792..8044de4e4b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/pkg/compactor/partitioned_group_info.go b/pkg/compactor/partitioned_group_info.go index e3d90a4e30..91154e8e7a 100644 --- a/pkg/compactor/partitioned_group_info.go +++ b/pkg/compactor/partitioned_group_info.go @@ -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 diff --git a/pkg/compactor/partitioned_group_info_test.go b/pkg/compactor/partitioned_group_info_test.go index 6db8a4e017..da6af033b6 100644 --- a/pkg/compactor/partitioned_group_info_test.go +++ b/pkg/compactor/partitioned_group_info_test.go @@ -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" @@ -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) + }) +}