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()

Reply via email to