From e1716668f7450cdbf44c31bb3c62f2568614cd1c Mon Sep 17 00:00:00 2001 From: constanca-m Date: Fri, 14 Aug 2026 12:03:01 +0200 Subject: [PATCH] Add ListEndOffsets to kafka.Manager. --- kafka/manager.go | 8 ++++++++ kafka/manager_test.go | 46 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+) diff --git a/kafka/manager.go b/kafka/manager.go index 6bdce133..8298dec4 100644 --- a/kafka/manager.go +++ b/kafka/manager.go @@ -389,3 +389,11 @@ func (m *Manager) ListTopics(ctx context.Context, prefix string) ([]string, erro slices.Sort(topics) return topics, errors.Join(errs...) } + +// ListEndOffsets returns the high watermark for each partition of the given +// topics. Topic names are broker names, the same as ListTopics returns. +// Namespace is not prepended. If no topics are specified, all topics are +// listed. +func (m *Manager) ListEndOffsets(ctx context.Context, topics ...string) (kadm.ListedOffsets, error) { + return m.adminClient.ListEndOffsets(ctx, topics...) +} diff --git a/kafka/manager_test.go b/kafka/manager_test.go index 106f47c8..0ecb10d4 100644 --- a/kafka/manager_test.go +++ b/kafka/manager_test.go @@ -23,6 +23,7 @@ import ( "sort" "strings" "testing" + "time" "github.com/google/go-cmp/cmp" "github.com/stretchr/testify/assert" @@ -30,6 +31,7 @@ import ( "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" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" @@ -807,6 +809,50 @@ func TestListTopics(t *testing.T) { assert.Equal(t, []string{"name_space-mytopic", "name_space-topic1", "name_space-topic3"}, topics) } +func TestListEndOffsets(t *testing.T) { + _, commonConfig := newFakeCluster(t) + m, err := NewManager(ManagerConfig{CommonConfig: commonConfig}) + require.NoError(t, err) + t.Cleanup(func() { m.Close() }) + + client, err := kgo.NewClient(kgo.SeedBrokers(commonConfig.Brokers...)) + require.NoError(t, err) + t.Cleanup(client.Close) + + admin := kadm.NewClient(client) + _, err = admin.CreateTopics(t.Context(), 2, 1, nil, "logs", "metrics") + require.NoError(t, err) + + ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + defer cancel() + for range 3 { + require.NoError(t, client.ProduceSync(ctx, &kgo.Record{ + Topic: "logs", + Value: []byte("x"), + }).FirstErr()) + } + + offsets, err := m.ListEndOffsets(ctx, "logs", "metrics") + require.NoError(t, err) + require.NoError(t, offsets.Error()) + + var logsEnd int64 + var logsPartitions int + offsets.Each(func(o kadm.ListedOffset) { + if o.Topic != "logs" { + return + } + logsPartitions++ + logsEnd += o.Offset + }) + require.Equal(t, 2, logsPartitions) + require.Equal(t, int64(3), logsEnd) + + metrics0, ok := offsets.Lookup("metrics", 0) + require.True(t, ok) + require.Equal(t, int64(0), metrics0.Offset) +} + func TestUnknownTopicOrPartition(t *testing.T) { testCases := []struct { desc string