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
31 changes: 16 additions & 15 deletions backend/pkg/console/consumer_group_offsets.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,10 +97,13 @@ func (s *Service) getConsumerGroupOffsets(ctx context.Context, adminCl *kadm.Cli

topicsEndOffsets, err := adminCl.ListEndOffsets(ctx, metadata.Topics.Names()...)
if err != nil {
return nil, fmt.Errorf("failed to list end offsets for topics: %w", err)
}
if tErr := topicsEndOffsets.Error(); tErr != nil && !errors.Is(tErr, kerr.UnknownTopicOrPartition) {
return nil, fmt.Errorf("failed to list end offsets for topics: %w", topicsEndOffsets.Error())
var shardErrors *kadm.ShardErrors
if !errors.As(err, &shardErrors) || shardErrors.AllFailed {
return nil, fmt.Errorf("failed to list end offsets for topics: %w", err)
}
s.logger.WarnContext(ctx, "failed to list end offsets from some shards",
slog.Int("failed_shards", len(shardErrors.Errs)),
slog.Any("error", err))
}

// Collect topics that don't exist (franz-go return an UnknownTopicOrPartition for these).
Expand All @@ -111,9 +114,6 @@ func (s *Service) getConsumerGroupOffsets(ctx context.Context, adminCl *kadm.Cli
}
})

// topicPartitions whose high watermark shall be requested
topicPartitions := make(map[string][]int32, len(metadata.Topics))

type partitionInfo struct {
PartitionID int32
Error string
Expand All @@ -125,22 +125,23 @@ func (s *Service) getConsumerGroupOffsets(ctx context.Context, adminCl *kadm.Cli
partitionInfoByIDAndTopic[td.Topic] = make(map[int32]partitionInfo)
for _, partition := range td.Partitions {
var errMsg string
endOffset, hasEndOffset := topicsEndOffsets.Lookup(td.Topic, partition.Partition)
if partition.Err != nil {
errMsg = partition.Err.Error()
} else {
// Not an offline partition nor unauthorized, so let's add it to the list
// of partition high watermarks we want to request
topicPartitions[td.Topic] = append(topicPartitions[td.Topic], partition.Partition)
} else if hasEndOffset && endOffset.Err != nil {
errMsg = endOffset.Err.Error()
} else if !hasEndOffset && err != nil {
errMsg = err.Error()
}

endOffset := int64(-1)
if offset, exists := topicsEndOffsets.Lookup(td.Topic, partition.Partition); exists {
endOffset = offset.Offset
endOffsetValue := int64(-1)
if hasEndOffset && endOffset.Err == nil {
endOffsetValue = endOffset.Offset
}
partitionInfoByIDAndTopic[td.Topic][partition.Partition] = partitionInfo{
PartitionID: partition.Partition,
Error: errMsg,
HighWaterMark: endOffset,
HighWaterMark: endOffsetValue,
}
}
}
Expand Down
112 changes: 112 additions & 0 deletions backend/pkg/console/consumer_group_offsets_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,10 @@ import (
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/twmb/franz-go/pkg/kadm"
"github.com/twmb/franz-go/pkg/kerr"
"github.com/twmb/franz-go/pkg/kfake"
"github.com/twmb/franz-go/pkg/kgo"
"github.com/twmb/franz-go/pkg/kmsg"

"github.com/redpanda-data/console/backend/pkg/config"
"github.com/redpanda-data/console/backend/pkg/testutil"
Expand Down Expand Up @@ -205,3 +207,113 @@ func TestGetConsumerGroupOffsets_NonExistentTopicInMetadata(t *testing.T) {

ass.Equal(expected, result, "Result should match expected (deleted topic skipped)")
}

func TestGetConsumerGroupOffsets_PartialListEndOffsetsFailure(t *testing.T) {
testCases := []struct {
Comment on lines 208 to +212
name string
offsetError *kerr.Error
errorMessage string
}{
{
name: "listener not found",
offsetError: kerr.ListenerNotFound,
errorMessage: "LISTENER_NOT_FOUND",
},
{
name: "leader not available",
offsetError: kerr.LeaderNotAvailable,
errorMessage: "LEADER_NOT_AVAILABLE",
},
}

for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
ctx := context.Background()
req := require.New(t)
ass := assert.New(t)

fakeCluster, err := kfake.NewCluster(kfake.NumBrokers(2))
req.NoError(err)
defer fakeCluster.Close()

client, adminClient := testutil.CreateClients(t, fakeCluster.ListenAddrs())
defer client.Close()

topicName := "test-partial-list-offsets-failure"
_, err = adminClient.CreateTopics(ctx, 2, 1, nil, topicName)
req.NoError(err)
// Keep one partition on each broker so one ListOffsets request can succeed while the other returns partition-level errors.
req.NoError(fakeCluster.MoveTopicPartition(topicName, 0, 0))
req.NoError(fakeCluster.MoveTopicPartition(topicName, 1, 1))

produceResults := client.ProduceSync(ctx,
&kgo.Record{Topic: topicName, Partition: 0, Value: []byte("p0")},
&kgo.Record{Topic: topicName, Partition: 1, Value: []byte("p1")},
)
req.NoError(produceResults.FirstErr())

groupID := "test-partial-list-offsets-group"
offsets := kadm.OffsetsList{
{Topic: topicName, Partition: 0, At: 1, LeaderEpoch: -1},
{Topic: topicName, Partition: 1, At: 1, LeaderEpoch: -1},
}.Offsets()
commitResponses, err := adminClient.CommitOffsets(ctx, groupID, offsets)
req.NoError(err)
req.True(commitResponses.Ok())

// Simulate a broker that is reachable for metadata, but cannot return end offsets for partitions it leads.
fakeCluster.ControlKey(int16(kmsg.ListOffsets), func(kreq kmsg.Request) (kmsg.Response, error, bool) {
fakeCluster.KeepControl()
if fakeCluster.CurrentNode() != 1 {
return nil, nil, false
}

listOffsetsReq := kreq.(*kmsg.ListOffsetsRequest)
resp := listOffsetsReq.ResponseKind().(*kmsg.ListOffsetsResponse)
for _, reqTopic := range listOffsetsReq.Topics {
respTopic := kmsg.NewListOffsetsResponseTopic()
respTopic.Topic = reqTopic.Topic
for _, reqPartition := range reqTopic.Partitions {
respPartition := kmsg.NewListOffsetsResponseTopicPartition()
respPartition.Partition = reqPartition.Partition
respPartition.ErrorCode = tc.offsetError.Code
respTopic.Partitions = append(respTopic.Partitions, respPartition)
}
resp.Topics = append(resp.Topics, respTopic)
}

return resp, nil, true
})

consoleSvc := &Service{
cfg: &config.Config{},
logger: slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelDebug})),
}

result, err := consoleSvc.getConsumerGroupOffsets(ctx, adminClient, []string{groupID})

req.NoError(err)
groupOffsets := result[groupID]
req.Len(groupOffsets, 1)
ass.Equal(topicName, groupOffsets[0].Topic)

partitionsByID := make(map[int32]PartitionOffsets)
for _, partition := range groupOffsets[0].PartitionOffsets {
partitionsByID[partition.PartitionID] = partition
}

partition0, exists := partitionsByID[0]
req.True(exists)
req.NotNil(partition0.GroupOffset)
// Partition 0 is served by the healthy broker, so lag data should still be available even though partition 1 fails.
ass.Empty(partition0.Error)
ass.Equal(int64(1), *partition0.GroupOffset)
ass.NotEqual(int64(-1), partition0.HighWaterMark)

partition1, exists := partitionsByID[1]
req.True(exists)
// The failed partition should be represented as a partition error, not as an error for the whole consumer group response.
ass.Contains(partition1.Error, tc.errorMessage)
})
}
}