squah-confluent commented on code in PR #23542:
URL: https://github.com/apache/kafka/pull/23542#discussion_r4068275874


##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java:
##########
@@ -1616,13 +1633,234 @@ public void 
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
             context.groupMetadataManager.streamsGroup(groupId).staticMembers()
         );
 
-        // Member rejoins with the same member id and a different instance id.
-        context.streamsGroupHeartbeat(
+        // Member rejoins with the same member id and a different instance id. 
A member id
+        // must never acquire a different instance id.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
             staticJoinHeartbeat(groupId, memberId, instanceId2, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a static member with instance id %s.", memberId, 
instanceId2, instanceId1), e.getMessage());
+
+        assertEquals(
+            Map.of(instanceId1, memberId),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
         );
+    }
+
+    @Test
+    public void testDynamicMemberCannotRejoinWithInstanceId() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId = Uuid.randomUuid().toString();
+        String instanceId = "instance-1";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+        TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+        TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch, 
0, 1, 2, 3);
+
+        // Streams group with one dynamic member.
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+                    .setMemberEpoch(groupEpoch)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(assignedTasks)
+                    .build())
+                .withTargetAssignment(memberId, targetAssignment)
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        assertEquals(Map.of(), 
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+
+        // Member rejoins with the same member id and an instance id. A 
dynamic member must
+        // never become a static member.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a dynamic member.", memberId, instanceId), 
e.getMessage());
+
+        assertEquals(Map.of(), 
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+    }
+
+    @Test
+    public void testExistingMemberCannotTakeReleasedInstanceId() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId1 = Uuid.randomUuid().toString();
+        String memberId2 = Uuid.randomUuid().toString();
+        String instanceId1 = "instance-1";
+        String instanceId2 = "instance-2";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+
+        // Streams group with two static members. Member 2 has left the group 
and released
+        // instance id 2.
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId1, 
instanceId1)
+                    .setMemberEpoch(groupEpoch)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+                    .build())
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId2, 
instanceId2)
+                    .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+                    .build())
+                .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+                .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        assertEquals(
+            Map.of(instanceId1, memberId1, instanceId2, memberId2),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+        );
+
+        // Member 1 rejoins with its own member id and the released instance 
id 2. The released
+        // instance id may only be taken by a new member id.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a static member with instance id %s.", memberId1, 
instanceId2, instanceId1), e.getMessage());
+
+        assertEquals(
+            Map.of(instanceId1, memberId1, instanceId2, memberId2),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+        );
+    }
+
+    @Test
+    public void testDynamicMemberCannotTakeReleasedInstanceId() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId1 = Uuid.randomUuid().toString();
+        String memberId2 = Uuid.randomUuid().toString();
+        String instanceId2 = "instance-2";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+
+        // Streams group with a dynamic member and a static member. The static 
member has left
+        // the group and released instance id 2.
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId1)
+                    .setMemberEpoch(groupEpoch)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+                    .build())
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId2, 
instanceId2)
+                    .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+                    .build())
+                .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+                .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        // Member 1 rejoins with its own member id and the released instance 
id 2. The released
+        // instance id may only be taken by a new member id.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a dynamic member.", memberId1, instanceId2), 
e.getMessage());
+
+        assertEquals(
+            Map.of(instanceId2, memberId2),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+        );
+    }
+
+    @Test
+    public void testStaticMemberJoinsAgainWhileActive() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId = Uuid.randomUuid().toString();
+        String instanceId = "instance-1";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+        TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+        TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch, 
0, 1, 2, 3);
+
+        // Streams group with one active static member.
+        StreamsGroupMember activeMember = 
streamsGroupMemberBuilderWithDefaults(memberId, instanceId)
+            .setMemberEpoch(groupEpoch)
+            .setPreviousMemberEpoch(groupEpoch - 1)
+            .setAssignedTasks(assignedTasks)
+            .build();
+
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(activeMember)
+                .withTargetAssignment(memberId, targetAssignment)
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        // The member joins again with the same member id and instance id 
while it is still
+        // active, e.g. because the response to its join was lost or because 
it was fenced.
+        // Like a dynamic member, it is not fenced by its own instance id and 
gets its current
+        // state back.
+        CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord> 
result = context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+        );
+
+        assertEquals(memberId, result.response().data().memberId());
+        assertEquals(groupEpoch, result.response().data().memberEpoch());
+
+        // No tombstones and no replacement records.
+        assertFalse(result.records().contains(
+            
StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentTombstoneRecord(groupId,
 memberId)
+        ));
+        assertFalse(result.records().contains(
+            
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentTombstoneRecord(groupId,
 memberId)
+        ));
+        assertFalse(result.records().contains(
+            
StreamsCoordinatorRecordHelpers.newStreamsGroupMemberTombstoneRecord(groupId, 
memberId)
+        ));
+
+        // The member keeps its epoch and its assignment.
+        StreamsGroupMember member = 
context.groupMetadataManager.streamsGroup(groupId).getMemberOrThrow(memberId);
+        assertEquals(groupEpoch, member.memberEpoch());
+        assertEquals(assignedTasks, member.assignedTasks());

Review Comment:
   nit: There's no equivalent check in the consumer version?



##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -3121,7 +3121,7 @@ public void testLeavingStaticMemberBumpsGroupEpoch() {
     }
 
     @Test
-    public void 
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
+    public void 
testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId() {

Review Comment:
   Separately, we don't have `CannotTakeActiveInstanceId` tests. Do we consider 
these cases covered by the `CannotRejoinWith(Different)InstanceId` tests?



##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java:
##########
@@ -1804,6 +1804,36 @@ private void 
throwIfInstanceIdIsFenced(StreamsGroupMember member, String groupId
         }
     }
     
+    /**
+     * Validates that the member id received in a join request does not 
already belong to a member
+     * with a different instance id. An instance id may move to a new member 
id when a static member
+     * is replaced but a member id must never acquire a different instance id.
+     *
+     * @param groupId               The group id.
+     * @param receivedMemberId      The member id received in the request.
+     * @param existingInstanceId    The instance id of the existing member 
with the received
+     *                              member id, or null if the existing member 
is a dynamic member.
+     * @param receivedInstanceId    The instance id received in the request.
+     *
+     * @throws InvalidRequestException if the received instance id differs 
from the instance id
+     *                                 of the existing member.
+     */
+    private void throwIfMemberIdHasDifferentInstanceId(
+        String groupId,
+        String receivedMemberId,
+        String existingInstanceId,
+        String receivedInstanceId
+    ) {
+        if (!receivedInstanceId.equals(existingInstanceId)) {

Review Comment:
   Shall we also forbid an existing static member losing an instance id? (not 
for this PR, it's close to 1,000 lines already)



##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java:
##########
@@ -3121,7 +3121,7 @@ public void testLeavingStaticMemberBumpsGroupEpoch() {
     }
 
     @Test
-    public void 
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
+    public void 
testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId() {

Review Comment:
   Our test naming is pretty messy. As a start, what do you think about these 
renames?
   | Current | Proposed |
   |---|---|
   | testDynamicMemberCannotRejoinWithInstanceId | 
testDynamicMemberCannotRejoinWithInstanceId |
   | testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId | 
testStaticMemberCannotRejoinWithDifferentInstanceId |
   | testDynamicMemberCannotTakeReleasedInstanceId | 
testDynamicMemberCannotRejoinWithReleasedInstanceId |
   | testExistingMemberCannotTakeReleasedInstanceId | 
testStaticMemberCannotRejoinWithReleasedInstanceId |
   | testStaticMemberRejoinsWithSameMemberIdAfterLeaving | 
testStaticMemberRejoinsWithSameMemberIdAfterLeaving |
   | testStaticMemberJoinsAgainWhileActive | 
testStaticMemberRejoinsWithSameMemberIdWhileActive |
   
   ---
   
   | Join? | Prev member | Instance id | Same instance id? | Result | Test |
   |---|---|---|---|---|---|
   | Yes | None | New | - | return new member | 
testGroupEpochBumpWhenNewStaticMemberJoins |
   | Yes | Dynamic | New | - | throw INVALID_REQUEST 
(throwIfMemberIdHasDifferentInstanceId) | 
testDynamicMemberCannotRejoinWithInstanceId |
   | Yes | Static | New | - | throw INVALID_REQUEST 
(throwIfMemberIdHasDifferentInstanceId) | 
testStaticMemberCannotRejoinWithSameMemberIdAndDifferentInstanceId |
   | Yes | None | Released | - | replace previous, emit records, return new 
member | testStaticMemberGetsBackAssignmentUponRejoin |
   | Yes | Dynamic | Released | - | throw INVALID_REQUEST 
(throwIfMemberIdHasDifferentInstanceId) | 
testDynamicMemberCannotTakeReleasedInstanceId |
   | Yes | Static | Released | No | throw INVALID_REQUEST 
(throwIfMemberIdHasDifferentInstanceId) | 
testExistingMemberCannotTakeReleasedInstanceId |
   | Yes | Static | Released | Yes | return rebuilt member, epoch reset to 0 | 
testStaticMemberRejoinsWithSameMemberIdAfterLeaving |
   | Yes | None | Active | - | throw UNRELEASED_INSTANCE_ID 
(throwIfInstanceIdIsUnreleased); classic proto → replace & return | 
testShouldThrownUnreleasedInstanceIdExceptionWhenNewMemberJoinsWithInUseInstanceId;
 classic variant: testJoiningConsumerGroupReplacingExistingStaticMember |
   | Yes | Dynamic | Active | - | throw INVALID_REQUEST 
(throwIfMemberIdHasDifferentInstanceId) | — |
   | Yes | Static | Active | No | throw INVALID_REQUEST 
(throwIfMemberIdHasDifferentInstanceId) | — |
   | Yes | Static | Active | Yes | return existing member unchanged | 
testStaticMemberJoinsAgainWhileActive; 
testConsumerGroupHeartbeatFromExistingClassicStaticMember |
   



##########
group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java:
##########
@@ -1616,13 +1633,234 @@ public void 
testStaticMemberRejoinsWithSameMemberIdAndDifferentInstanceId() {
             context.groupMetadataManager.streamsGroup(groupId).staticMembers()
         );
 
-        // Member rejoins with the same member id and a different instance id.
-        context.streamsGroupHeartbeat(
+        // Member rejoins with the same member id and a different instance id. 
A member id
+        // must never acquire a different instance id.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
             staticJoinHeartbeat(groupId, memberId, instanceId2, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a static member with instance id %s.", memberId, 
instanceId2, instanceId1), e.getMessage());
+
+        assertEquals(
+            Map.of(instanceId1, memberId),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
         );
+    }
+
+    @Test
+    public void testDynamicMemberCannotRejoinWithInstanceId() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId = Uuid.randomUuid().toString();
+        String instanceId = "instance-1";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+        TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+        TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch, 
0, 1, 2, 3);
+
+        // Streams group with one dynamic member.
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+                    .setMemberEpoch(groupEpoch)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(assignedTasks)
+                    .build())
+                .withTargetAssignment(memberId, targetAssignment)
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        assertEquals(Map.of(), 
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+
+        // Member rejoins with the same member id and an instance id. A 
dynamic member must
+        // never become a static member.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a dynamic member.", memberId, instanceId), 
e.getMessage());
+
+        assertEquals(Map.of(), 
context.groupMetadataManager.streamsGroup(groupId).staticMembers());
+    }
+
+    @Test
+    public void testExistingMemberCannotTakeReleasedInstanceId() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId1 = Uuid.randomUuid().toString();
+        String memberId2 = Uuid.randomUuid().toString();
+        String instanceId1 = "instance-1";
+        String instanceId2 = "instance-2";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+
+        // Streams group with two static members. Member 2 has left the group 
and released
+        // instance id 2.
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId1, 
instanceId1)
+                    .setMemberEpoch(groupEpoch)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+                    .build())
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId2, 
instanceId2)
+                    .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+                    .build())
+                .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+                .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        assertEquals(
+            Map.of(instanceId1, memberId1, instanceId2, memberId2),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+        );
+
+        // Member 1 rejoins with its own member id and the released instance 
id 2. The released
+        // instance id may only be taken by a new member id.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a static member with instance id %s.", memberId1, 
instanceId2, instanceId1), e.getMessage());
+
+        assertEquals(
+            Map.of(instanceId1, memberId1, instanceId2, memberId2),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+        );
+    }
+
+    @Test
+    public void testDynamicMemberCannotTakeReleasedInstanceId() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId1 = Uuid.randomUuid().toString();
+        String memberId2 = Uuid.randomUuid().toString();
+        String instanceId2 = "instance-2";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+
+        // Streams group with a dynamic member and a static member. The static 
member has left
+        // the group and released instance id 2.
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId1)
+                    .setMemberEpoch(groupEpoch)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 0, 1))
+                    .build())
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId2, 
instanceId2)
+                    .setMemberEpoch(LEAVE_GROUP_STATIC_MEMBER_EPOCH)
+                    .setPreviousMemberEpoch(groupEpoch - 1)
+                    .setAssignedTasks(topic.assignedTasks(groupEpoch, 2, 3))
+                    .build())
+                .withTargetAssignment(memberId1, topic.targetAssignment(0, 1))
+                .withTargetAssignment(memberId2, topic.targetAssignment(2, 3))
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        // Member 1 rejoins with its own member id and the released instance 
id 2. The released
+        // instance id may only be taken by a new member id.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId1, instanceId2, topic)
+        ));
+        assertEquals(String.format("Member %s with instance id %s cannot join 
the group because the member id is " +
+            "already used by a dynamic member.", memberId1, instanceId2), 
e.getMessage());
+
+        assertEquals(
+            Map.of(instanceId2, memberId2),
+            context.groupMetadataManager.streamsGroup(groupId).staticMembers()
+        );
+    }
+
+    @Test
+    public void testStaticMemberJoinsAgainWhileActive() {
+        int groupEpoch = DEFAULT_GROUP_EPOCH;
+
+        String groupId = "fooup";
+        String memberId = Uuid.randomUuid().toString();
+        String instanceId = "instance-1";
+
+        StreamsTopicFixture topic = streamsTopicFixture("subtopology1", "foo", 
4);
+        TasksTuple targetAssignment = topic.targetAssignment(0, 1, 2, 3);
+        TasksTupleWithEpochs assignedTasks = topic.assignedTasks(groupEpoch, 
0, 1, 2, 3);
+
+        // Streams group with one active static member.
+        StreamsGroupMember activeMember = 
streamsGroupMemberBuilderWithDefaults(memberId, instanceId)
+            .setMemberEpoch(groupEpoch)
+            .setPreviousMemberEpoch(groupEpoch - 1)
+            .setAssignedTasks(assignedTasks)
+            .build();
+
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(new 
MockTaskAssignor("sticky")))
+            .withMetadataImage(topic.metadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, groupEpoch)
+                .withMember(activeMember)
+                .withTargetAssignment(memberId, targetAssignment)
+                .withTargetAssignmentEpoch(groupEpoch)
+                
.withTopology(StreamsTopology.fromHeartbeatRequest(topic.topology()))
+                .withValidatedTopologyEpoch(0)
+                .withMetadataHash(topic.metadataHash())
+                .withLastAssignmentConfigs(getDefaultAssignmentConfigs()))
+            
.withConfig(GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG,
 GroupCoordinatorConfig.STREAMS_GROUP_INITIAL_REBALANCE_DELAY_MS_DEFAULT)
+            .build();
+
+        // The member joins again with the same member id and instance id 
while it is still
+        // active, e.g. because the response to its join was lost or because 
it was fenced.
+        // Like a dynamic member, it is not fenced by its own instance id and 
gets its current
+        // state back.
+        CoordinatorResult<StreamsGroupHeartbeatResult, CoordinatorRecord> 
result = context.streamsGroupHeartbeat(
+            staticJoinHeartbeat(groupId, memberId, instanceId, topic)
+        );
+
+        assertEquals(memberId, result.response().data().memberId());
+        assertEquals(groupEpoch, result.response().data().memberEpoch());
+
+        // No tombstones and no replacement records.
+        assertFalse(result.records().contains(
+            
StreamsCoordinatorRecordHelpers.newStreamsGroupCurrentAssignmentTombstoneRecord(groupId,
 memberId)
+        ));
+        assertFalse(result.records().contains(
+            
StreamsCoordinatorRecordHelpers.newStreamsGroupTargetAssignmentTombstoneRecord(groupId,
 memberId)
+        ));
+        assertFalse(result.records().contains(
+            
StreamsCoordinatorRecordHelpers.newStreamsGroupMemberTombstoneRecord(groupId, 
memberId)
+        ));

Review Comment:
   Can we assert that no records are written, like in the consumer version?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to