This is an automated email from the ASF dual-hosted git repository.
hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new 11dd94c62 test(go): fail when a manual group commit on a split primary
drops membership (#4293)
11dd94c62 is described below
commit 11dd94c6283c8070cb295bb3064dc33397577bb2
Author: Mark Aron Szulyovszky <[email protected]>
AuthorDate: Mon Sep 28 11:46:00 2026 +0200
test(go): fail when a manual group commit on a split primary drops
membership (#4293)
---
.../tests/cluster/partition_primary_routing.rs | 29 +++++++++----
foreign/go/client/tcp/tcp_offset_management.go | 9 ++++
foreign/go/tests/e2e_test.go | 49 ++++++++++++++++++++++
3 files changed, 79 insertions(+), 8 deletions(-)
diff --git a/core/integration/tests/cluster/partition_primary_routing.rs
b/core/integration/tests/cluster/partition_primary_routing.rs
index 66f3b5eb0..c0ee05ae3 100644
--- a/core/integration/tests/cluster/partition_primary_routing.rs
+++ b/core/integration/tests/cluster/partition_primary_routing.rs
@@ -250,16 +250,29 @@ async fn
given_split_primaries_when_http_auto_commits_on_a_backup_should_replica
async fn
given_split_primaries_when_go_group_auto_commits_should_preserve_membership(
harness: &mut TestHarness,
) {
+ run_go_split_primary_test(
+ harness,
+ "^TestE2E_SplitPrimaryPollsPreserveCoordinatorMembership$",
+ )
+ .await;
+}
+
+#[iggy_harness(cluster_nodes = 3, server(metadata.journal_slots = "256"))]
+#[ignore = "requires Go; run this test explicitly with --ignored"]
+async fn
given_split_primaries_when_go_group_commits_manually_should_preserve_membership(
+ harness: &mut TestHarness,
+) {
+ run_go_split_primary_test(
+ harness,
+ "^TestE2E_SplitPrimaryManualCommitPreservesMembership$",
+ )
+ .await;
+}
+
+async fn run_go_split_primary_test(harness: &mut TestHarness, test: &str) {
let (_, metadata_primary, _) = seed_split_primaries(harness).await;
let output = tokio::process::Command::new("go")
- .args([
- "test",
- "./tests",
- "-run",
- "^TestE2E_SplitPrimaryPollsPreserveCoordinatorMembership$",
- "-count=1",
- "-v",
- ])
+ .args(["test", "./tests", "-run", test, "-count=1", "-v"])
.current_dir(std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../foreign/go"))
.env(
"IGGY_TCP_ADDRESS",
diff --git a/foreign/go/client/tcp/tcp_offset_management.go
b/foreign/go/client/tcp/tcp_offset_management.go
index c2551ae8d..8479ee966 100644
--- a/foreign/go/client/tcp/tcp_offset_management.go
+++ b/foreign/go/client/tcp/tcp_offset_management.go
@@ -40,6 +40,15 @@ func (c *IggyTcpClient) GetConsumerOffset(ctx
context.Context, consumer iggcon.C
}
func (c *IggyTcpClient) StoreConsumerOffset(ctx context.Context, consumer
iggcon.Consumer, streamId iggcon.Identifier, topicId iggcon.Identifier, offset
uint64, partitionId *uint32) error {
+ // TODO(#4292): a group commit for a partition whose primary is not the
+ // coordinator goes out on the coordinator session, is refused as not
+ // admitted, and sendFrame walks the roster to the primary. That
reconnect
+ // registers a new client identity, which is not a member of the group,
so
+ // the replayed commit fails with ConsumerGroupPartitionNotOwned and the
+ // membership is gone. Route clustered group commits (and deletes) to
the
+ // partition primary through the attached consumer session, as
pollPrimary
+ // does for auto-commit polls and the Rust SDK's
PollRouter::write_offset
+ // does for offset writes.
_, err := c.do(ctx, &command.StoreConsumerOffsetRequest{
StreamId: streamId,
TopicId: topicId,
diff --git a/foreign/go/tests/e2e_test.go b/foreign/go/tests/e2e_test.go
index ca52cfebe..ff5169f70 100644
--- a/foreign/go/tests/e2e_test.go
+++ b/foreign/go/tests/e2e_test.go
@@ -269,6 +269,55 @@ func
TestE2E_SplitPrimaryPollsPreserveCoordinatorMembership(t *testing.T) {
primaryAddress,
int(details.PartitionsCount)*messagesPerPartition, afterClient.ID,
afterClient.ConsumerGroupsCount)
}
+// A manual group commit on the same split-primary fixture. The Rust SDK routes
+// it to the partition primary over the consumer-session data connection and
+// keeps its coordinator membership; the Go client must do the same rather than
+// move its session to the primary, which registers a new client identity that
+// is not a member and gets the commit refused.
+func TestE2E_SplitPrimaryManualCommitPreservesMembership(t *testing.T) {
+ streamName := os.Getenv("IGGY_POLL_ROUTING_STREAM")
+ if streamName == "" {
+ t.Skip("set IGGY_POLL_ROUTING_STREAM and
IGGY_POLL_ROUTING_TOPIC to a split-primary topic")
+ }
+ stream, err := iggcon.NewIdentifier(streamName)
+ require.NoError(t, err)
+ topic, err := iggcon.NewIdentifier(os.Getenv("IGGY_POLL_ROUTING_TOPIC"))
+ require.NoError(t, err)
+ connected := connect(t)
+ ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
+ defer cancel()
+ group, err := connected.CreateConsumerGroup(ctx, stream, topic,
fmt.Sprintf("go-commit-%d", time.Now().UnixNano()))
+ require.NoError(t, err)
+ groupID, err := iggcon.NewIdentifier(group.Id)
+ require.NoError(t, err)
+ t.Cleanup(func() { _ =
connected.DeleteConsumerGroup(context.Background(), stream, topic, groupID) })
+ require.NoError(t, connected.JoinConsumerGroup(ctx, stream, topic,
groupID))
+ consumer := iggcon.NewGroupConsumer(groupID)
+ coordinator := connected.GetConnectionInfo().ServerAddress
+ before, err := connected.SendBinaryRequest(ctx,
uint32(command.GetMeCode), nil)
+ require.NoError(t, err)
+ beforeClient := binaryserialization.DeserializeClient(before)
+ require.Equal(t, uint32(1), beforeClient.ConsumerGroupsCount)
+
+ partition := uint32(0)
+ for offset := range uint64(2) {
+ err := connected.StoreConsumerOffset(ctx, consumer, stream,
topic, offset, &partition)
+ require.NoError(t, err, "manual group commit of offset %d must
reach the partition primary as a member (session %s -> %s)",
+ offset, coordinator,
connected.GetConnectionInfo().ServerAddress)
+ assert.Equal(t, coordinator,
connected.GetConnectionInfo().ServerAddress,
+ "the commit must not move the coordinator session")
+ }
+ stored, err := connected.GetConsumerOffset(ctx, consumer, stream,
topic, &partition)
+ require.NoError(t, err)
+ require.NotNil(t, stored)
+ assert.Equal(t, uint64(1), stored.StoredOffset)
+ after, err := connected.SendBinaryRequest(ctx,
uint32(command.GetMeCode), nil)
+ require.NoError(t, err)
+ afterClient := binaryserialization.DeserializeClient(after)
+ assert.Equal(t, beforeClient.ID, afterClient.ID, "the commit registered
a new client identity")
+ assert.Equal(t, beforeClient.ConsumerGroupsCount,
afterClient.ConsumerGroupsCount, "the commit dropped the group membership")
+}
+
func TestE2E_RawRequestsDoNotGapMetadataRequestIDs(t *testing.T) {
connected := connect(t)
ctx := context.Background()