From bf1d0809542302b7b2905a946b948e9eb3e84dd0 Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Thu, 20 Jul 2023 00:12:44 -0400 Subject: [PATCH 1/2] Implement Heartbeat in new Coordinator --- .../group/GroupCoordinatorService.java | 27 +- .../group/GroupMetadataManager.java | 76 + .../group/ReplicatedGroupCoordinator.java | 20 + .../group/GroupCoordinatorServiceTest.java | 101 +- .../group/GroupMetadataManagerTest.java | 1288 ++++++++++++++++- 5 files changed, 1507 insertions(+), 5 deletions(-) diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java index 6783eff79f250..4bf1d382d35ab 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.config.TopicConfig; +import org.apache.kafka.common.errors.CoordinatorLoadInProgressException; import org.apache.kafka.common.errors.InvalidFetchSizeException; import org.apache.kafka.common.errors.KafkaStorageException; import org.apache.kafka.common.errors.NotEnoughReplicasException; @@ -369,9 +370,29 @@ public CompletableFuture heartbeat( return FutureUtils.failedFuture(Errors.COORDINATOR_NOT_AVAILABLE.exception()); } - return FutureUtils.failedFuture(Errors.UNSUPPORTED_VERSION.exception( - "This API is not implemented yet." - )); + if (!isGroupIdNotEmpty(request.groupId())) { + return CompletableFuture.completedFuture(new HeartbeatResponseData() + .setErrorCode(Errors.INVALID_GROUP_ID.code())); + } + + return runtime.scheduleReadOperation("generic-group-heartbeat", + topicPartitionFor(request.groupId()), + (coordinator, offset) -> coordinator.genericGroupHeartbeat(context, request) + ).exceptionally(exception -> { + if (!(exception instanceof KafkaException)) { + log.error("Heartbeat request {} hit an unexpected exception: {}", + request, exception.getMessage()); + } + + if (exception instanceof CoordinatorLoadInProgressException) { + // The group is still loading, so blindly respond + return new HeartbeatResponseData() + .setErrorCode(Errors.NONE.code()); + } + + return new HeartbeatResponseData() + .setErrorCode(Errors.forException(exception).code()); + }); } /** diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java index 4ed7a4d2bf646..d70f2f7b1dbbe 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java @@ -20,9 +20,11 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.ApiException; +import org.apache.kafka.common.errors.CoordinatorNotAvailableException; import org.apache.kafka.common.errors.FencedMemberEpochException; import org.apache.kafka.common.errors.GroupIdNotFoundException; import org.apache.kafka.common.errors.GroupMaxSizeReachedException; +import org.apache.kafka.common.errors.IllegalGenerationException; import org.apache.kafka.common.errors.InvalidRequestException; import org.apache.kafka.common.errors.NotCoordinatorException; import org.apache.kafka.common.errors.UnknownMemberIdException; @@ -30,6 +32,8 @@ import org.apache.kafka.common.errors.UnsupportedAssignorException; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData; +import org.apache.kafka.common.message.HeartbeatRequestData; +import org.apache.kafka.common.message.HeartbeatResponseData; import org.apache.kafka.common.message.JoinGroupRequestData.JoinGroupRequestProtocol; import org.apache.kafka.common.message.JoinGroupRequestData.JoinGroupRequestProtocolCollection; import org.apache.kafka.common.message.JoinGroupRequestData; @@ -90,6 +94,7 @@ import java.util.stream.Collectors; import static org.apache.kafka.common.protocol.Errors.COORDINATOR_NOT_AVAILABLE; +import static org.apache.kafka.common.protocol.Errors.ILLEGAL_GENERATION; import static org.apache.kafka.common.protocol.Errors.NOT_COORDINATOR; import static org.apache.kafka.common.protocol.Errors.UNKNOWN_SERVER_ERROR; import static org.apache.kafka.common.requests.JoinGroupRequest.UNKNOWN_MEMBER_ID; @@ -2835,6 +2840,77 @@ private void removePendingSyncMember( } } + /** + * Handle a generic group HeartbeatRequest. + * + * @param context The request context. + * @param request The actual Heartbeat request. + * + * @return The Heartbeat response. + */ + public HeartbeatResponseData genericGroupHeartbeat( + RequestContext context, + HeartbeatRequestData request + ) { + GenericGroup group = getOrMaybeCreateGenericGroup(request.groupId(), false); + + validateGenericGroupHeartbeat(group, request.memberId(), request.groupInstanceId(), request.generationId()); + + switch (group.currentState()) { + case EMPTY: + return new HeartbeatResponseData().setErrorCode(Errors.UNKNOWN_MEMBER_ID.code()); + + case PREPARING_REBALANCE: + rescheduleGenericGroupMemberHeartbeat(group, group.member(request.memberId())); + return new HeartbeatResponseData().setErrorCode(Errors.REBALANCE_IN_PROGRESS.code()); + + case COMPLETING_REBALANCE: + // Consumers may start sending heartbeat after join-group response, in which case + // we should treat them as normal hb request and reset the timer + + case STABLE: + rescheduleGenericGroupMemberHeartbeat(group, group.member(request.memberId())); + return new HeartbeatResponseData(); + + default: + throw new IllegalStateException("Reached unexpected state " + + group.currentState() + " for group " + group.groupId()); + } + } + + /** + * Validates a generic group heartbeat request. + * + * @param group The group. + * @param memberId The member id. + * @param groupInstanceId The group instance id. + * @param generationId The generation id. + * + * @throws CoordinatorNotAvailableException If group is Dead. + * @throws IllegalGenerationException If the generation id in the request and the generation id of the + * group does not match. + */ + private void validateGenericGroupHeartbeat( + GenericGroup group, + String memberId, + String groupInstanceId, + int generationId + ) throws CoordinatorNotAvailableException, IllegalGenerationException { + if (group.isInState(DEAD)) { + throw COORDINATOR_NOT_AVAILABLE.exception(); + } else { + group.validateMember( + memberId, + groupInstanceId, + "heartbeat" + ); + + if (generationId != group.generationId()) { + throw ILLEGAL_GENERATION.exception(); + } + } + } + /** * Checks whether the given protocol type or name in the request is inconsistent with the group's. * diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java index 26a318f42348f..392af291bc5db 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/ReplicatedGroupCoordinator.java @@ -19,6 +19,8 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData; +import org.apache.kafka.common.message.HeartbeatRequestData; +import org.apache.kafka.common.message.HeartbeatResponseData; import org.apache.kafka.common.message.JoinGroupRequestData; import org.apache.kafka.common.message.JoinGroupResponseData; import org.apache.kafka.common.message.OffsetCommitRequestData; @@ -249,6 +251,24 @@ public CoordinatorResult genericGroupSync( ); } + /** + * Handles a generic group HeartbeatRequest. + * + * @param context The request context. + * @param request The actual Heartbeat request. + * + * @return The HeartbeatResponse. + */ + public HeartbeatResponseData genericGroupHeartbeat( + RequestContext context, + HeartbeatRequestData request + ) { + return groupMetadataManager.genericGroupHeartbeat( + context, + request + ); + } + /** * Handles a OffsetCommit request. * diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java index 2ceebeceff714..235f713da615e 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java @@ -18,17 +18,21 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.config.TopicConfig; +import org.apache.kafka.common.errors.CoordinatorLoadInProgressException; import org.apache.kafka.common.errors.CoordinatorNotAvailableException; import org.apache.kafka.common.errors.InvalidFetchSizeException; import org.apache.kafka.common.errors.InvalidRequestException; import org.apache.kafka.common.errors.KafkaStorageException; import org.apache.kafka.common.errors.NotEnoughReplicasException; import org.apache.kafka.common.errors.NotLeaderOrFollowerException; +import org.apache.kafka.common.errors.RebalanceInProgressException; import org.apache.kafka.common.errors.RecordBatchTooLargeException; import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData; +import org.apache.kafka.common.message.HeartbeatRequestData; +import org.apache.kafka.common.message.HeartbeatResponseData; import org.apache.kafka.common.message.JoinGroupRequestData; import org.apache.kafka.common.message.JoinGroupResponseData; import org.apache.kafka.common.message.SyncGroupRequestData; @@ -495,4 +499,99 @@ public void testSyncGroupInvalidGroupId() throws Exception { assertEquals(expectedResponse, response.get()); } -} + + @Test + public void testHeartbeat() throws Exception { + CoordinatorRuntime runtime = mockRuntime(); + GroupCoordinatorService service = new GroupCoordinatorService( + new LogContext(), + createConfig(), + runtime + ); + + HeartbeatRequestData request = new HeartbeatRequestData() + .setGroupId("foo"); + + service.startup(() -> 1); + + when(runtime.scheduleReadOperation( + ArgumentMatchers.eq("generic-group-heartbeat"), + ArgumentMatchers.eq(new TopicPartition("__consumer_offsets", 0)), + ArgumentMatchers.any() + )).thenReturn(CompletableFuture.completedFuture( + new HeartbeatResponseData() + )); + + CompletableFuture future = service.heartbeat( + requestContext(ApiKeys.HEARTBEAT), + request + ); + + assertTrue(future.isDone()); + assertEquals(new HeartbeatResponseData(), future.get()); + } + + @Test + public void testHeartbeatCoordinatorNotAvailableException() throws Exception { + CoordinatorRuntime runtime = mockRuntime(); + GroupCoordinatorService service = new GroupCoordinatorService( + new LogContext(), + createConfig(), + runtime + ); + + HeartbeatRequestData request = new HeartbeatRequestData() + .setGroupId("foo"); + + service.startup(() -> 1); + + when(runtime.scheduleReadOperation( + ArgumentMatchers.eq("generic-group-heartbeat"), + ArgumentMatchers.eq(new TopicPartition("__consumer_offsets", 0)), + ArgumentMatchers.any() + )).thenReturn(FutureUtils.failedFuture( + new CoordinatorLoadInProgressException(null) + )); + + CompletableFuture future = service.heartbeat( + requestContext(ApiKeys.HEARTBEAT), + request + ); + + assertTrue(future.isDone()); + assertEquals(new HeartbeatResponseData(), future.get()); + } + + @Test + public void testHeartbeatCoordinatorException() throws Exception { + CoordinatorRuntime runtime = mockRuntime(); + GroupCoordinatorService service = new GroupCoordinatorService( + new LogContext(), + createConfig(), + runtime + ); + + HeartbeatRequestData request = new HeartbeatRequestData() + .setGroupId("foo"); + + service.startup(() -> 1); + + when(runtime.scheduleReadOperation( + ArgumentMatchers.eq("generic-group-heartbeat"), + ArgumentMatchers.eq(new TopicPartition("__consumer_offsets", 0)), + ArgumentMatchers.any() + )).thenReturn(FutureUtils.failedFuture( + new RebalanceInProgressException() + )); + + CompletableFuture future = service.heartbeat( + requestContext(ApiKeys.HEARTBEAT), + request + ); + + assertTrue(future.isDone()); + assertEquals( + new HeartbeatResponseData().setErrorCode(Errors.REBALANCE_IN_PROGRESS.code()), + future.get() + ); + }} diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java index 792e4ec777a21..d80f749144c9b 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java @@ -20,9 +20,12 @@ import org.apache.kafka.clients.consumer.internals.ConsumerProtocol; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.errors.CoordinatorNotAvailableException; +import org.apache.kafka.common.errors.FencedInstanceIdException; import org.apache.kafka.common.errors.FencedMemberEpochException; import org.apache.kafka.common.errors.GroupIdNotFoundException; import org.apache.kafka.common.errors.GroupMaxSizeReachedException; +import org.apache.kafka.common.errors.IllegalGenerationException; import org.apache.kafka.common.errors.InvalidRequestException; import org.apache.kafka.common.errors.UnknownMemberIdException; import org.apache.kafka.common.errors.UnknownServerException; @@ -30,6 +33,8 @@ import org.apache.kafka.common.errors.UnsupportedAssignorException; import org.apache.kafka.common.message.ConsumerGroupHeartbeatRequestData; import org.apache.kafka.common.message.ConsumerGroupHeartbeatResponseData; +import org.apache.kafka.common.message.HeartbeatRequestData; +import org.apache.kafka.common.message.HeartbeatResponseData; import org.apache.kafka.common.message.JoinGroupRequestData; import org.apache.kafka.common.message.JoinGroupResponseData; import org.apache.kafka.common.message.JoinGroupResponseData.JoinGroupResponseMember; @@ -120,6 +125,7 @@ import static org.apache.kafka.coordinator.group.GroupMetadataManager.consumerGroupSessionTimeoutKey; import static org.apache.kafka.coordinator.group.GroupMetadataManager.EMPTY_RESULT; import static org.apache.kafka.coordinator.group.GroupMetadataManager.genericGroupHeartbeatKey; +import static org.apache.kafka.coordinator.group.GroupMetadataManager.genericGroupSyncKey; import static org.apache.kafka.coordinator.group.generic.GenericGroupState.COMPLETING_REBALANCE; import static org.apache.kafka.coordinator.group.generic.GenericGroupState.DEAD; import static org.apache.kafka.coordinator.group.generic.GenericGroupState.EMPTY; @@ -945,6 +951,123 @@ public void verifySessionExpiration(GenericGroup group, int timeoutMs) { assertEquals(0, group.size()); } + public HeartbeatResponseData sendGenericGroupHeartbeat( + HeartbeatRequestData request + ) { + RequestContext context = new RequestContext( + new RequestHeader( + ApiKeys.HEARTBEAT, + ApiKeys.HEARTBEAT.latestVersion(), + "client", + 0 + ), + "1", + InetAddress.getLoopbackAddress(), + KafkaPrincipal.ANONYMOUS, + ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT), + SecurityProtocol.PLAINTEXT, + ClientInformation.EMPTY, + false + ); + + return groupMetadataManager.genericGroupHeartbeat( + context, + request + ); + } + + public void verifyHeartbeat( + String groupId, + JoinGroupResponseData joinResponse, + Errors expectedError + ) { + HeartbeatRequestData request = new HeartbeatRequestData() + .setGroupId(groupId) + .setMemberId(joinResponse.memberId()) + .setGenerationId(joinResponse.generationId()); + + if (expectedError == Errors.UNKNOWN_MEMBER_ID) { + assertThrows(UnknownMemberIdException.class, () -> sendGenericGroupHeartbeat(request)); + } else { + HeartbeatResponseData response = sendGenericGroupHeartbeat(request); + assertEquals(expectedError.code(), response.errorCode()); + } + } + + public List joinWithNMembers( + GenericGroup group, + int numMembers, + int rebalanceTimeoutMs, + int sessionTimeoutMs + ) { + boolean requireKnownMemberId = true; + + // First join requests + JoinGroupRequestData request = new JoinGroupRequestBuilder() + .withGroupId(group.groupId()) + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(rebalanceTimeoutMs) + .withSessionTimeoutMs(sessionTimeoutMs) + .build(); + + List> joinFutures = IntStream.range(0, numMembers) + .mapToObj(__ -> new CompletableFuture()) + .collect(Collectors.toList()); + + IntStream.range(0, numMembers).forEach(i -> { + CoordinatorResult result = sendGenericGroupJoin( + request, joinFutures.get(i), requireKnownMemberId); + + assertTrue(result.records().isEmpty()); + }); + + List memberIds = joinFutures.stream().map(future -> { + assertTrue(future.isDone()); + try { + return future.get().memberId(); + } catch (Exception e) { + fail("Unexpected exception: " + e.getMessage()); + } + return null; + }).collect(Collectors.toList()); + + // Second join requests + List> secondJoinFutures = IntStream.range(0, numMembers) + .mapToObj(__ -> new CompletableFuture()) + .collect(Collectors.toList()); + + IntStream.range(0, numMembers).forEach(i -> { + CoordinatorResult result = sendGenericGroupJoin( + request.setMemberId(memberIds.get(i)), secondJoinFutures.get(i), requireKnownMemberId); + + assertTrue(result.records().isEmpty()); + }); + secondJoinFutures.forEach(future -> assertFalse(future.isDone())); + + // Advance clock by initial rebalance delay. + assertNoOrEmptyResult(sleep(genericGroupInitialRebalanceDelayMs)); + secondJoinFutures.forEach(future -> assertFalse(future.isDone())); + // Advance clock by rebalance timeout to complete join phase. + assertNoOrEmptyResult(sleep(rebalanceTimeoutMs)); + + List joinResponses = secondJoinFutures.stream().map(future -> { + assertTrue(future.isDone()); + try { + assertEquals(Errors.NONE.code(), future.get().errorCode()); + return future.get(); + } catch (Exception e) { + fail("Unexpected exception: " + e.getMessage()); + } + return null; + }).collect(Collectors.toList()); + + assertEquals(numMembers, group.size()); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + return joinResponses; + } + private ApiMessage messageOrNull(ApiMessageAndVersion apiMessageAndVersion) { if (apiMessageAndVersion == null) { return null; @@ -4784,7 +4907,7 @@ public void testHeartbeatExpirationShouldRemoveMember() throws Exception { timeouts.forEach(timeout -> { assertEquals(genericGroupHeartbeatKey("group-id", memberId), timeout.key); assertEquals(Collections.singletonList( - RecordHelpers.newGroupMetadataRecord(group, Collections.emptyMap(), MetadataVersion.latest())), + newGroupMetadataRecord(group, MetadataVersion.latest())), timeout.result.records()); }); @@ -8323,5 +8446,1168 @@ public SyncResult( this.appendFuture = coordinatorResult.appendFuture(); } } + + @Test + public void testStaticMemberHeartbeatLeaderWithInvalidMemberId() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + context.createGenericGroup("group-id"); + + RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( + "group-id", + "leader-instance-id", + "follower-instance-id" + ); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withGroupInstanceId("leader-instance-id") + .withMemberId(rebalanceResult.leaderId) + .withGenerationId(rebalanceResult.generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertTrue(syncResult.records.isEmpty()); + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(rebalanceResult.leaderId) + .setGenerationId(rebalanceResult.generationId); + + HeartbeatResponseData validHeartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.NONE.code(), validHeartbeatResponse.errorCode()); + + assertThrows(FencedInstanceIdException.class, () -> context.sendGenericGroupHeartbeat( + heartbeatRequest + .setGroupInstanceId("leader-instance-id") + .setMemberId("invalid-member-id") + )); + } + + @Test + public void testHeartbeatUnknownGroup() { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId("member-id") + .setGenerationId(-1); + + assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + } + + @Test + public void testHeartbeatDeadGroup() { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + group.transitionTo(DEAD); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId("member-id") + .setGenerationId(-1); + + assertThrows(CoordinatorNotAvailableException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + } + + @Test + public void testHeartbeatEmptyGroup() { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestProtocolCollection protocols = new JoinGroupRequestProtocolCollection(); + protocols.add(new JoinGroupRequestProtocol() + .setName("range") + .setMetadata(new byte[]{0})); + + group.add(new GenericGroupMember( + "member-id", + Optional.empty(), + "client-id", + "client-host", + 10000, + 5000, + "consumer", + protocols + )); + + group.transitionTo(PREPARING_REBALANCE); + group.transitionTo(EMPTY); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId("member-id") + .setGenerationId(0); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.UNKNOWN_MEMBER_ID.code(), heartbeatResponse.errorCode()); + } + + @Test + public void testHeartbeatUnknownMemberExistingGroup() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId("unknown-member-id") + .setGenerationId(generationId); + + assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + } + + @Test + public void testHeartbeatDuringPreparingRebalance() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build(); + + CompletableFuture joinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin(joinRequest, joinFuture, true); + + assertTrue(result.records().isEmpty()); + assertTrue(joinFuture.isDone()); + assertEquals(Errors.MEMBER_ID_REQUIRED.code(), joinFuture.get().errorCode()); + + String memberId = joinFuture.get().memberId(); + + joinFuture = new CompletableFuture<>(); + context.sendGenericGroupJoin(joinRequest.setMemberId(memberId), joinFuture); + + assertTrue(group.isInState(PREPARING_REBALANCE)); + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(memberId) + .setGenerationId(0); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); + } + + @Test + public void testHeartbeatDuringCompletingRebalance() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(new HeartbeatResponseData(), heartbeatResponse); + } + + @Test + public void testHeartbeatIllegalGeneration() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId + 1); + + assertThrows(IllegalGenerationException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + } + + @Test + public void testValidHeartbeat() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); + } + + @Test + public void testGenericGroupMemberSessionTimeout() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // Advance clock by session timeout to kick member out. + context.verifySessionExpiration(group, 5000); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + } + + @Test + public void testGenericGroupMemberMaintainsSession() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // Advance clock by 1/2 of session timeout. + assertNoOrEmptyResult(context.sleep(2500)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); + + + assertNoOrEmptyResult(context.sleep(2500)); + + heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); + } + + @Test + public void testGenericGroupMemberSessionTimeoutDuringRebalance() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // Add a new member. This should trigger a rebalance. The new member has the + // 'genericGroupNewMemberJoinTimeoutMs` session timeout, so it has a longer expiration than the existing member. + CompletableFuture joinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin( + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID), + joinFuture + ); + + assertTrue(result.records().isEmpty()); + assertFalse(joinFuture.isDone()); + assertTrue(group.isInState(PREPARING_REBALANCE)); + + // Advance clock by 1/2 of session timeout. + assertNoOrEmptyResult(context.sleep(2500)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); + + // Advance clock by first member's session timeout. + assertNoOrEmptyResult(context.sleep(5000)); + + assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + + // Advance clock by remaining rebalance timeout to complete join phase. + assertNoOrEmptyResult(context.sleep(2500)); + + assertTrue(joinFuture.isDone()); + assertEquals(Errors.NONE.code(), joinFuture.get().errorCode()); + assertEquals(1, group.size()); + assertEquals(2, group.generationId()); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + } + + @Test + public void testRebalanceCompletesBeforeMemberJoins() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + // Create a group with a single member + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withGroupInstanceId("leader-instance-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAndCompleteJoin(joinRequest, true, true); + + String firstMemberId = leaderJoinResponse.memberId(); + int firstGenerationId = leaderJoinResponse.generationId(); + + assertEquals(1, firstGenerationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(firstMemberId) + .withGenerationId(firstGenerationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // Add a new dynamic member. This should trigger a rebalance. The new member has the + // 'genericGroupNewMemberJoinTimeoutMs` session timeout, so it has a longer expiration than the existing member. + CompletableFuture joinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin( + joinRequest.setMemberId(UNKNOWN_MEMBER_ID) + .setGroupInstanceId(null) + .setSessionTimeoutMs(2500), + joinFuture); + + assertTrue(result.records().isEmpty()); + assertFalse(joinFuture.isDone()); + assertTrue(group.isInState(PREPARING_REBALANCE)); + + // Send a couple heartbeats to keep the first member alive while the rebalance finishes. + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(firstMemberId) + .setGenerationId(firstGenerationId); + + for (int i = 0; i < 2; i++) { + assertNoOrEmptyResult(context.sleep(2500)); + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); + } + + // Advance clock by remaining rebalance timeout to complete join phase. + // The second member will become the leader. However, as the first member is a static member + // it will not be kicked out. + assertNoOrEmptyResult(context.sleep(8000)); + + assertTrue(joinFuture.isDone()); + assertEquals(Errors.NONE.code(), joinFuture.get().errorCode()); + assertEquals(2, group.size()); + assertEquals(2, group.generationId()); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + String otherMemberId = joinFuture.get().memberId(); + + syncResult = context.sendGenericGroupSync( + syncRequest.setGroupInstanceId(null) + .setMemberId(otherMemberId) + .setGenerationId(2)); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // The unjoined static member should be remained in the group before session timeout. + assertThrows(IllegalGenerationException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + + // Now session timeout the unjoined member. Still keeping the new member. + List expectedErrors = Arrays.asList(Errors.NONE, Errors.NONE, Errors.REBALANCE_IN_PROGRESS); + for (Errors expectedError : expectedErrors) { + assertNoOrEmptyResult(context.sleep(2000)); + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat( + heartbeatRequest.setMemberId(otherMemberId) + .setGenerationId(2)); + + assertEquals(expectedError.code(), heartbeatResponse.errorCode()); + } + assertEquals(1, group.size()); + assertTrue(group.isInState(PREPARING_REBALANCE)); + + CompletableFuture otherMemberRejoinFuture = new CompletableFuture<>(); + result = context.sendGenericGroupJoin( + joinRequest.setMemberId(otherMemberId) + .setGroupInstanceId(null) + .setSessionTimeoutMs(2500), + otherMemberRejoinFuture); + + assertTrue(result.records().isEmpty()); + assertTrue(otherMemberRejoinFuture.isDone()); + assertEquals(Errors.NONE.code(), otherMemberRejoinFuture.get().errorCode()); + assertEquals(3, otherMemberRejoinFuture.get().generationId()); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncResult otherMemberResyncResult = context.sendGenericGroupSync( + syncRequest.setGroupInstanceId(null) + .setMemberId(otherMemberId) + .setGenerationId(3)); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + otherMemberResyncResult.records + ); + // Simulate a successful write to the log. + otherMemberResyncResult.appendFuture.complete(null); + + assertTrue(otherMemberResyncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), otherMemberResyncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // The joined member should get heart beat response with no error. Let the new member keep + // heartbeating for a while to verify that no new rebalance is triggered unexpectedly. + for (int i = 0; i < 20; i++) { + assertNoOrEmptyResult(context.sleep(2000)); + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat( + heartbeatRequest.setMemberId(otherMemberId) + .setGenerationId(3)); + + assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); + } + } + + @Test + public void testSyncGroupEmptyAssignment() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderId) + .withGenerationId(generationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertEquals(0, syncResult.syncFuture.get().assignment().length); + assertTrue(group.isInState(STABLE)); + + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); + } + + @Test + public void testSecondMemberPartiallyJoinAndTimeout() throws Exception { + // Test if the following scenario completes a rebalance correctly: A new member starts a JoinGroup request with + // an UNKNOWN_MEMBER_ID, attempting to join a stable group. But never initiates the second JoinGroup request with + // the provided member ID and times out. The test checks if original member remains the sole member in this group, + // which should remain stable throughout this test. + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + // Create a group with a single member + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withGroupInstanceId("leader-instance-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAndCompleteJoin(joinRequest, true, true); + + String firstMemberId = leaderJoinResponse.memberId(); + int firstGenerationId = leaderJoinResponse.generationId(); + + assertEquals(1, firstGenerationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(firstMemberId) + .withGenerationId(firstGenerationId) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // Add a new dynamic pending member. + CompletableFuture joinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin( + joinRequest.setMemberId(UNKNOWN_MEMBER_ID) + .setGroupInstanceId(null) + .setSessionTimeoutMs(5000), + joinFuture, + true, + true + ); + + assertTrue(result.records().isEmpty()); + assertTrue(joinFuture.isDone()); + assertEquals(Errors.MEMBER_ID_REQUIRED.code(), joinFuture.get().errorCode()); + assertEquals(1, group.numPendingJoinMembers()); + assertTrue(group.isInState(STABLE)); + + // Heartbeat from the leader to maintain session while timing out pending member. + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(firstMemberId) + .setGenerationId(firstGenerationId); + + for (int i = 0; i < 2; i++) { + assertNoOrEmptyResult(context.sleep(2500)); + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); + } + + // At this point the second member should have been removed from pending list (session timeout), + // and the group should be in Stable state with only the first member in it. + assertEquals(1, group.size()); + assertTrue(group.hasMemberId(firstMemberId)); + assertEquals(1, group.generationId()); + assertTrue(group.isInState(STABLE)); + } + + @Test + public void testRebalanceTimesOutWhenSyncRequestIsNotReceived() throws Exception { + // This test case ensure that the pending sync expiration does kick out all members + // if they don't send sync requests before the rebalance timeout. The + // group is in the Empty state in this case. + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + int rebalanceTimeoutMs = 5000; + int sessionTimeoutMs = 5000; + List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + + // Advance clock by 1/2 rebalance timeout. + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); + + // Heartbeats to ensure that heartbeating does not interfere with the + // delayed sync operation. + joinResponses.forEach(response -> context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); + + // Advance clock by 1/2 rebalance timeout to expire the pending sync. Members should be removed. + // The group becomes empty, generating an empty group metadata record. + List> timeouts = context.sleep(rebalanceTimeoutMs / 2); + assertEquals(1, timeouts.size()); + ExpiredTimeout timeout = timeouts.get(0); + assertEquals(genericGroupSyncKey("group-id"), timeout.key); + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + timeout.result.records() + ); + + // Simulate a successful write to the log. + timeout.result.appendFuture().complete(null); + + // Heartbeats fail because none of the members have sent the sync request + joinResponses.forEach(response -> context.verifyHeartbeat(group.groupId(), response, Errors.UNKNOWN_MEMBER_ID)); + assertTrue(group.isInState(EMPTY)); + } + + @Test + public void testRebalanceTimesOutWhenSyncRequestIsNotReceivedFromFollowers() throws Exception { + // This test case ensure that the pending sync expiration does kick out the followers + // if they don't send a sync request before the rebalance timeout. The + // group is in the PreparingRebalance state in this case. + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + int rebalanceTimeoutMs = 5000; + int sessionTimeoutMs = 5000; + List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + + // Advance clock by 1/2 rebalance timeout. + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); + + // Heartbeats to ensure that heartbeating does not interfere with the + // delayed sync operation. + joinResponses.forEach(response -> context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); + + // Leader sends a sync group request. + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withGenerationId(1) + .withMemberId(joinResponses.get(0).memberId()) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + // Leader should be able to heartbeat + context.verifyHeartbeat(group.groupId(), joinResponses.get(0), Errors.NONE); + + // Advance clock by 1/2 rebalance timeout to expire the pending sync. Followers should be removed. + List> timeouts = context.sleep(rebalanceTimeoutMs / 2); + assertEquals(1, timeouts.size()); + ExpiredTimeout timeout = timeouts.get(0); + assertEquals(genericGroupSyncKey("group-id"), timeout.key); + assertTrue(timeout.result.records().isEmpty()); + + // Leader should be able to heartbeat + joinResponses.subList(0, 1).forEach(response -> + context.verifyHeartbeat(group.groupId(), response, Errors.REBALANCE_IN_PROGRESS)); + + // Heartbeats fail because none of the followers have sent the sync request + joinResponses.subList(1, 3).forEach(response -> + context.verifyHeartbeat(group.groupId(), response, Errors.UNKNOWN_MEMBER_ID)); + + assertTrue(group.isInState(PREPARING_REBALANCE)); + } + + @Test + public void testRebalanceTimesOutWhenSyncRequestIsNotReceivedFromLeaders() throws Exception { + // This test case ensure that the pending sync expiration does kick out the leader + // if it does not send a sync request before the rebalance timeout. The + // group is in the PreparingRebalance state in this case. + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + int rebalanceTimeoutMs = 5000; + int sessionTimeoutMs = 5000; + List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + + // Advance clock by 1/2 rebalance timeout. + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); + + // Heartbeats to ensure that heartbeating does not interfere with the + // delayed sync operation. + joinResponses.forEach(response -> context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); + + // Followers send sync group requests. + List> followerSyncFutures = joinResponses.subList(1, 3).stream() + .map(response -> { + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withGenerationId(1) + .withMemberId(response.memberId()) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + assertTrue(syncResult.records.isEmpty()); + assertFalse(syncResult.syncFuture.isDone()); + return syncResult.syncFuture; + }).collect(Collectors.toList()); + + // Advance clock by 1/2 rebalance timeout to expire the pending sync. Leader should be kicked out. + List> timeouts = context.sleep(rebalanceTimeoutMs / 2); + assertEquals(1, timeouts.size()); + ExpiredTimeout timeout = timeouts.get(0); + assertEquals(genericGroupSyncKey("group-id"), timeout.key); + assertTrue(timeout.result.records().isEmpty()); + + // Follower sync responses should fail. + followerSyncFutures.forEach(future -> { + assertTrue(future.isDone()); + try { + assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), future.get().errorCode()); + } catch (Exception e) { + fail("Unexpected exception: " + e.getMessage()); + } + }); + + // Leader heartbeat should fail. + joinResponses.subList(0, 1).forEach(response -> + context.verifyHeartbeat(group.groupId(), response, Errors.UNKNOWN_MEMBER_ID)); + + // Follower heartbeats should succeed. + joinResponses.subList(1, 3).forEach(response -> + context.verifyHeartbeat(group.groupId(), response, Errors.REBALANCE_IN_PROGRESS)); + + assertTrue(group.isInState(PREPARING_REBALANCE)); + } + + @Test + public void testRebalanceDoesNotTimeOutWhenAllSyncAreReceived() throws Exception { + // This test case ensure that the pending sync expiration does not kick any + // members out when they have all sent their sync requests. Group should be in Stable state. + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + int rebalanceTimeoutMs = 5000; + int sessionTimeoutMs = 5000; + List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + String leaderId = joinResponses.get(0).memberId(); + + // Advance clock by 1/2 rebalance timeout. + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); + + // Heartbeats to ensure that heartbeating does not interfere with the + // delayed sync operation. + joinResponses.forEach(response -> context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); + + // All members send sync group requests. + List> syncFutures = joinResponses.stream().map(response -> { + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withGenerationId(1) + .withMemberId(response.memberId()) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + + if (response.memberId().equals(leaderId)) { + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + } else { + assertTrue(syncResult.records.isEmpty()); + } + assertTrue(syncResult.syncFuture.isDone()); + return syncResult.syncFuture; + }).collect(Collectors.toList()); + + for (CompletableFuture syncFuture : syncFutures) { + assertEquals(Errors.NONE.code(), syncFuture.get().errorCode()); + } + + // Advance clock by 1/2 rebalance timeout. Pending sync should already have been cancelled. + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); + + // All member heartbeats should succeed. + joinResponses.forEach(response -> + context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); + + // Advance clock a bit more + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); + + // All member heartbeats should succeed. + joinResponses.forEach(response -> + context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); + + assertTrue(group.isInState(STABLE)); + } + + @Test + public void testHeartbeatDuringRebalanceCausesRebalanceInProgress() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); + + // First start up a group (with a slightly larger timeout to give us time to heartbeat when the rebalance starts) + JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(); + + JoinGroupResponseData leaderJoinResponse = + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + + String leaderId = leaderJoinResponse.memberId(); + int generationId = leaderJoinResponse.generationId(); + + assertEquals(1, generationId); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + // Then join with a new consumer to trigger a rebalance + CompletableFuture joinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin( + joinRequest.setMemberId(UNKNOWN_MEMBER_ID), joinFuture); + + assertTrue(result.records().isEmpty()); + assertFalse(joinFuture.isDone()); + + // We should be in the middle of a rebalance, so the heartbeat should return rebalance in progress. + HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderId) + .setGenerationId(generationId); + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); + } + + private List joinWithNMembers( + GroupMetadataManagerTestContext context, + GenericGroup group, + int numMembers, + int rebalanceTimeoutMs, + int sessionTimeoutMs + ) { + boolean requireKnownMemberId = true; + + // First join requests + JoinGroupRequestData request = new JoinGroupRequestBuilder() + .withGroupId(group.groupId()) + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(rebalanceTimeoutMs) + .withSessionTimeoutMs(sessionTimeoutMs) + .build(); + + List> joinFutures = IntStream.range(0, numMembers) + .mapToObj(__ -> new CompletableFuture()) + .collect(Collectors.toList()); + + IntStream.range(0, numMembers).forEach(i -> { + CoordinatorResult result = context.sendGenericGroupJoin( + request, joinFutures.get(i), requireKnownMemberId); + + assertTrue(result.records().isEmpty()); + }); + + List memberIds = joinFutures.stream().map(future -> { + assertTrue(future.isDone()); + try { + return future.get().memberId(); + } catch (Exception e) { + fail("Unexpected exception: " + e.getMessage()); + } + return null; + }).collect(Collectors.toList()); + + // Second join requests + List> secondJoinFutures = IntStream.range(0, numMembers) + .mapToObj(__ -> new CompletableFuture()) + .collect(Collectors.toList()); + + IntStream.range(0, numMembers).forEach(i -> { + CoordinatorResult result = context.sendGenericGroupJoin( + request.setMemberId(memberIds.get(i)), secondJoinFutures.get(i), requireKnownMemberId); + + assertTrue(result.records().isEmpty()); + }); + secondJoinFutures.forEach(future -> assertFalse(future.isDone())); + + // Advance clock by initial rebalance delay. + assertNoOrEmptyResult(context.sleep(context.genericGroupInitialRebalanceDelayMs)); + secondJoinFutures.forEach(future -> assertFalse(future.isDone())); + // Advance clock by rebalance timeout to complete join phase. + assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs)); + + List joinResponses = secondJoinFutures.stream().map(future -> { + assertTrue(future.isDone()); + try { + assertEquals(Errors.NONE.code(), future.get().errorCode()); + return future.get(); + } catch (Exception e) { + fail("Unexpected exception: " + e.getMessage()); + } + return null; + }).collect(Collectors.toList()); + + assertEquals(numMembers, group.size()); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + return joinResponses; + } + + private void verifyHeartbeat( + GroupMetadataManagerTestContext context, + String groupId, + JoinGroupResponseData joinResponse, + Errors expectedError + ) { + HeartbeatRequestData request = new HeartbeatRequestData() + .setGroupId(groupId) + .setMemberId(joinResponse.memberId()) + .setGenerationId(joinResponse.generationId()); + + if (expectedError != Errors.NONE) { + assertThrows(expectedError.exception().getClass(), () -> context.sendGenericGroupHeartbeat(request)); + } else { + HeartbeatResponseData response = context.sendGenericGroupHeartbeat(request); + assertEquals(Errors.NONE.code(), response.errorCode()); + } + } } From 5c49dcde025da9fced5ff75e1f10daf00514585c Mon Sep 17 00:00:00 2001 From: Jeff Kim Date: Thu, 27 Jul 2023 21:06:30 -0400 Subject: [PATCH 2/2] address comments --- checkstyle/suppressions.xml | 2 +- .../group/GroupCoordinatorService.java | 4 +- .../group/GroupMetadataManager.java | 6 +- .../group/GroupCoordinatorServiceTest.java | 3 +- .../group/GroupMetadataManagerTest.java | 1070 +++++------------ 5 files changed, 319 insertions(+), 766 deletions(-) diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml index 0c40bcf446b4f..6982920e155d0 100644 --- a/checkstyle/suppressions.xml +++ b/checkstyle/suppressions.xml @@ -329,7 +329,7 @@ + files="(RecordHelpersTest|GroupMetadataManagerTest|GroupCoordinatorServiceTest).java"/> diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java index 4bf1d382d35ab..33450766b4258 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java @@ -375,9 +375,11 @@ public CompletableFuture heartbeat( .setErrorCode(Errors.INVALID_GROUP_ID.code())); } + // Using a read operation is okay here as we ignore the last committed offset in the snapshot registry. + // This means we will read whatever is in the latest snapshot, which is how the old coordinator behaves. return runtime.scheduleReadOperation("generic-group-heartbeat", topicPartitionFor(request.groupId()), - (coordinator, offset) -> coordinator.genericGroupHeartbeat(context, request) + (coordinator, __) -> coordinator.genericGroupHeartbeat(context, request) ).exceptionally(exception -> { if (!(exception instanceof KafkaException)) { log.error("Heartbeat request {} hit an unexpected exception: {}", diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java index d70f2f7b1dbbe..1b490c55e1006 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java @@ -2865,10 +2865,10 @@ public HeartbeatResponseData genericGroupHeartbeat( return new HeartbeatResponseData().setErrorCode(Errors.REBALANCE_IN_PROGRESS.code()); case COMPLETING_REBALANCE: - // Consumers may start sending heartbeat after join-group response, in which case - // we should treat them as normal hb request and reset the timer - case STABLE: + // Consumers may start sending heartbeats after join-group response, while the group + // is in CompletingRebalance state. In this case, we should treat them as + // normal heartbeat requests and reset the timer rescheduleGenericGroupMemberHeartbeat(group, group.member(request.memberId())); return new HeartbeatResponseData(); diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java index 235f713da615e..36dc455f9b429 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java @@ -594,4 +594,5 @@ public void testHeartbeatCoordinatorException() throws Exception { new HeartbeatResponseData().setErrorCode(Errors.REBALANCE_IN_PROGRESS.code()), future.get() ); - }} + } +} diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java index d80f749144c9b..8f6ae07b37f46 100644 --- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java +++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java @@ -573,6 +573,43 @@ public CoordinatorResult sendGenericGroupJoin( ); } + public JoinGroupResponseData joinGenericGroupAsDynamicMemberAndCompleteRebalance( + String groupId + ) throws Exception { + GenericGroup group = createGenericGroup(groupId); + + JoinGroupResponseData leaderJoinResponse = + joinGenericGroupAsDynamicMemberAndCompleteJoin(new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build()); + + assertEquals(1, leaderJoinResponse.generationId()); + assertTrue(group.isInState(COMPLETING_REBALANCE)); + + SyncResult syncResult = sendGenericGroupSync(new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderJoinResponse.memberId()) + .withGenerationId(leaderJoinResponse.generationId()) + .build()); + + assertEquals( + Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), + syncResult.records + ); + // Simulate a successful write to the log. + syncResult.appendFuture.complete(null); + + assertTrue(syncResult.syncFuture.isDone()); + assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + assertTrue(group.isInState(STABLE)); + + return leaderJoinResponse; + } + public JoinGroupResponseData joinGenericGroupAsDynamicMemberAndCompleteJoin( JoinGroupRequestData request ) throws ExecutionException, InterruptedException { @@ -582,11 +619,12 @@ public JoinGroupResponseData joinGenericGroupAsDynamicMemberAndCompleteJoin( if (request.memberId().equals(UNKNOWN_MEMBER_ID)) { // Since member id is required, we need another round to get the successful join group result. CompletableFuture firstJoinFuture = new CompletableFuture<>(); - sendGenericGroupJoin( + CoordinatorResult result = sendGenericGroupJoin( request, firstJoinFuture, requireKnownMemberId ); + assertTrue(result.records().isEmpty()); assertTrue(firstJoinFuture.isDone()); assertEquals(Errors.MEMBER_ID_REQUIRED.code(), firstJoinFuture.get().errorCode()); newMemberId = firstJoinFuture.get().memberId(); @@ -603,12 +641,13 @@ public JoinGroupResponseData joinGenericGroupAsDynamicMemberAndCompleteJoin( .setRebalanceTimeoutMs(request.rebalanceTimeoutMs()) .setReason(request.reason()); - sendGenericGroupJoin( + CoordinatorResult result = sendGenericGroupJoin( secondRequest, secondJoinFuture, requireKnownMemberId ); + assertTrue(result.records().isEmpty()); List> timeouts = sleep(genericGroupInitialRebalanceDelayMs); assertEquals(1, timeouts.size()); assertTrue(secondJoinFuture.isDone()); @@ -727,7 +766,8 @@ public RebalanceResult staticMembersJoinAndRebalance( CompletableFuture followerJoinResponseFuture = new CompletableFuture<>(); result = sendGenericGroupJoin( - joinRequest.setGroupInstanceId(followerInstanceId), + joinRequest + .setGroupInstanceId(followerInstanceId), followerJoinResponseFuture); assertTrue(result.records().isEmpty()); @@ -821,15 +861,12 @@ public JoinGroupResponseData setupGroupWithPendingMember(GenericGroup group) thr JoinGroupResponseData leaderJoinResponse = joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - List assignment = new ArrayList<>(); - assignment.add(new SyncGroupRequestAssignment().setMemberId(leaderId)); + assignment.add(new SyncGroupRequestAssignment().setMemberId(leaderJoinResponse.memberId())); SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) + .withMemberId(leaderJoinResponse.memberId()) + .withGenerationId(leaderJoinResponse.generationId()) .withAssignment(assignment) .build(); @@ -849,7 +886,8 @@ public JoinGroupResponseData setupGroupWithPendingMember(GenericGroup group) thr // Start the join for the second member CompletableFuture followerJoinFuture = new CompletableFuture<>(); CoordinatorResult result = sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID), followerJoinFuture ); @@ -857,7 +895,10 @@ public JoinGroupResponseData setupGroupWithPendingMember(GenericGroup group) thr assertFalse(followerJoinFuture.isDone()); CompletableFuture leaderJoinFuture = new CompletableFuture<>(); - result = sendGenericGroupJoin(joinRequest.setMemberId(leaderId), leaderJoinFuture); + result = sendGenericGroupJoin( + joinRequest + .setMemberId(leaderJoinResponse.memberId()), + leaderJoinFuture); assertTrue(result.records().isEmpty()); assertTrue(group.isInState(COMPLETING_REBALANCE)); @@ -866,8 +907,8 @@ public JoinGroupResponseData setupGroupWithPendingMember(GenericGroup group) thr assertEquals(Errors.NONE.code(), leaderJoinFuture.get().errorCode()); assertEquals(Errors.NONE.code(), followerJoinFuture.get().errorCode()); assertEquals(leaderJoinFuture.get().generationId(), followerJoinFuture.get().generationId()); - assertEquals(leaderId, leaderJoinFuture.get().leader()); - assertEquals(leaderId, followerJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), leaderJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), followerJoinFuture.get().leader()); int nextGenerationId = leaderJoinFuture.get().generationId(); String followerId = followerJoinFuture.get().memberId(); @@ -888,7 +929,10 @@ public JoinGroupResponseData setupGroupWithPendingMember(GenericGroup group) thr // Re-join an existing member, to transition the group to PreparingRebalance state. leaderJoinFuture = new CompletableFuture<>(); - result = sendGenericGroupJoin(joinRequest.setMemberId(leaderId), leaderJoinFuture); + result = sendGenericGroupJoin( + joinRequest + .setMemberId(leaderJoinResponse.memberId()), + leaderJoinFuture); assertTrue(result.records().isEmpty()); assertFalse(leaderJoinFuture.isDone()); @@ -897,7 +941,9 @@ public JoinGroupResponseData setupGroupWithPendingMember(GenericGroup group) thr // Create a pending member in the group CompletableFuture pendingMemberJoinFuture = new CompletableFuture<>(); result = sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID).setSessionTimeoutMs(2500), + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID) + .setSessionTimeoutMs(2500), pendingMemberJoinFuture, true ); @@ -995,16 +1041,17 @@ public void verifyHeartbeat( } public List joinWithNMembers( - GenericGroup group, + String groupId, int numMembers, int rebalanceTimeoutMs, int sessionTimeoutMs ) { + GenericGroup group = createGenericGroup(groupId); boolean requireKnownMemberId = true; // First join requests JoinGroupRequestData request = new JoinGroupRequestBuilder() - .withGroupId(group.groupId()) + .withGroupId(groupId) .withMemberId(UNKNOWN_MEMBER_ID) .withDefaultProtocolTypeAndProtocols() .withRebalanceTimeoutMs(rebalanceTimeoutMs) @@ -5960,15 +6007,15 @@ public void testNewMemberTimeoutCompletion() throws Exception { .build(); GenericGroup group = context.createGenericGroup("group-id"); - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() + CompletableFuture joinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin(new JoinGroupRequestBuilder() .withGroupId("group-id") .withMemberId(UNKNOWN_MEMBER_ID) .withDefaultProtocolTypeAndProtocols() .withSessionTimeoutMs(context.genericGroupNewMemberJoinTimeoutMs + 5000) - .build(); + .build(), + joinFuture); - CompletableFuture joinFuture = new CompletableFuture<>(); - CoordinatorResult result = context.sendGenericGroupJoin(joinRequest, joinFuture); assertTrue(result.records().isEmpty()); assertFalse(joinFuture.isDone()); @@ -5982,13 +6029,11 @@ public void testNewMemberTimeoutCompletion() throws Exception { assertEquals(0, group.allMembers().stream().filter(GenericGroupMember::isNew).count()); - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withMemberId(joinFuture.get().memberId()) .withGenerationId(joinFuture.get().generationId()) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); // Simulate a successful write to the log. syncResult.appendFuture.complete(null); @@ -6051,13 +6096,11 @@ public void testNewMemberFailureAfterJoinGroupCompletion() throws Exception { assertEquals(memberId, joinResponse.leader()); assertEquals(1, joinResponse.generationId()); - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withMemberId(memberId) .withGenerationId(1) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); // Simulate a successful write to the log. syncResult.appendFuture.complete(null); @@ -6070,12 +6113,14 @@ public void testNewMemberFailureAfterJoinGroupCompletion() throws Exception { CompletableFuture otherJoinFuture = new CompletableFuture<>(); CoordinatorResult otherJoinResult = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID), otherJoinFuture); CompletableFuture joinFuture = new CompletableFuture<>(); CoordinatorResult joinResult = context.sendGenericGroupJoin( - joinRequest.setMemberId(memberId), + joinRequest + .setMemberId(memberId), joinFuture); assertTrue(otherJoinResult.records().isEmpty()); @@ -6092,13 +6137,13 @@ public void testNewMemberFailureAfterJoinGroupCompletion() throws Exception { public void testStaticMemberFenceDuplicateRejoinedFollower() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A third member joins. Trigger a rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -6152,13 +6197,13 @@ public void testStaticMemberFenceDuplicateRejoinedFollower() throws Exception { public void testStaticMemberFenceDuplicateSyncingFollowerAfterMemberIdChanged() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Known leader rejoins will trigger rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -6232,14 +6277,12 @@ public void testStaticMemberFenceDuplicateSyncingFollowerAfterMemberIdChanged() // Duplicate follower joins group with unknown member id will trigger member.id replacement, // and will also trigger a rebalance under CompletingRebalance state; the old follower sync callback // will return fenced exception while broker replaces the member identity with the duplicate follower joins. - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult oldFollowerSyncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withGroupInstanceId("follower-instance-id") .withGenerationId(oldFollowerJoinFuture.get().generationId()) .withMemberId(oldFollowerJoinFuture.get().memberId()) - .build(); - - SyncResult oldFollowerSyncResult = context.sendGenericGroupSync(syncRequest); + .build()); assertTrue(result.records().isEmpty()); assertFalse(oldFollowerSyncResult.syncFuture.isDone()); @@ -6273,13 +6316,13 @@ public void testStaticMemberFenceDuplicateSyncingFollowerAfterMemberIdChanged() public void testStaticMemberFenceDuplicateRejoiningFollowerAfterMemberIdChanged() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Known leader rejoins will trigger rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -6401,14 +6444,12 @@ public void testStaticMemberRejoinWithKnownMemberId() throws Exception { assertTrue(rejoinResponseFuture.isDone()); assertEquals(Errors.NONE.code(), rejoinResponseFuture.get().errorCode()); - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withMemberId(memberId) .withGenerationId(joinResponse.generationId()) .withGroupInstanceId("group-instance-id") - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); // Successful write to the log. syncResult.appendFuture.complete(null); @@ -6425,13 +6466,13 @@ public void testStaticMemberRejoinWithLeaderIdAndUnknownMemberId( ) throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A static leader rejoin with unknown id will not trigger rebalance, and no assignment will be returned. // As the group was in Stable state and the member id was updated, this will generate records. @@ -6483,7 +6524,8 @@ public void testStaticMemberRejoinWithLeaderIdAndUnknownMemberId( CompletableFuture oldLeaderJoinFuture = new CompletableFuture<>(); result = context.sendGenericGroupJoin( - joinRequest.setMemberId(rebalanceResult.leaderId), + joinRequest + .setMemberId(rebalanceResult.leaderId), oldLeaderJoinFuture, true, supportSkippingAssignment); @@ -6519,13 +6561,13 @@ public void testStaticMemberRejoinWithLeaderIdAndUnknownMemberId( public void testStaticMemberRejoinWithLeaderIdAndKnownMemberId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Known static leader rejoin will trigger rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -6563,14 +6605,13 @@ public void testStaticMemberRejoinWithLeaderIdAndKnownMemberId() throws Exceptio public void testStaticMemberRejoinWithLeaderIdAndUnexpectedDeadGroup() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); - + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); group.transitionTo(DEAD); JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -6592,13 +6633,13 @@ public void testStaticMemberRejoinWithLeaderIdAndUnexpectedDeadGroup() throws Ex public void testStaticMemberRejoinWithLeaderIdAndUnexpectedEmptyGroup() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); group.transitionTo(PREPARING_REBALANCE); group.transitionTo(EMPTY); @@ -6624,7 +6665,6 @@ public void testStaticMemberRejoinWithFollowerIdAndChangeOfProtocol() throws Exc int sessionTimeoutMs = 15000; GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", @@ -6633,6 +6673,7 @@ public void testStaticMemberRejoinWithFollowerIdAndChangeOfProtocol() throws Exc rebalanceTimeoutMs, sessionTimeoutMs ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A static follower rejoin with changed protocol will trigger rebalance. JoinGroupRequestProtocolCollection protocols = toProtocols("roundrobin"); @@ -6678,14 +6719,12 @@ public void testStaticMemberRejoinWithFollowerIdAndChangeOfProtocol() throws Exc } @Test - public void testStaticMemberRejoinWithUnknownMemberIdAndChangeOfProtocolWithSelectedProtocolChanged() - throws Exception { + public void testStaticMemberRejoinWithUnknownMemberIdAndChangeOfProtocolWithSelectedProtocolChanged() throws Exception { int rebalanceTimeoutMs = 10000; int sessionTimeoutMs = 15000; GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", @@ -6694,6 +6733,7 @@ public void testStaticMemberRejoinWithUnknownMemberIdAndChangeOfProtocolWithSele rebalanceTimeoutMs, sessionTimeoutMs ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); assertNotEquals("roundrobin", group.selectProtocol()); @@ -6746,13 +6786,13 @@ public void testStaticMemberRejoinWithUnknownMemberIdAndChangeOfProtocolWhileSel GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); JoinGroupRequestProtocolCollection protocols = toProtocols(group.selectProtocol()); @@ -6853,16 +6893,15 @@ public void testStaticMemberRejoinWithUnknownMemberIdAndChangeOfProtocolWhileSel @Test public void testStaticMemberRejoinWithUnknownMemberIdAndChangeOfProtocolWhileSelectProtocolUnchanged() throws Exception { - GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A static follower rejoin with protocol changing to leader protocol subset won't trigger rebalance if updated // group's selectProtocol remain unchanged. @@ -6948,13 +6987,13 @@ public void testStaticMemberRejoinWithKnownLeaderIdToTriggerRebalanceAndFollower GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A static leader rejoin with known member id will trigger rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -7076,13 +7115,13 @@ public void testStaticMemberRejoinWithKnownLeaderIdToTriggerRebalanceAndFollower public void testStaticMemberRejoinAsFollowerWithUnknownMemberId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A static follower rejoin with no protocol change will not trigger rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -7130,14 +7169,12 @@ public void testStaticMemberRejoinAsFollowerWithUnknownMemberId() throws Excepti ); assertNotEquals(rebalanceResult.followerId, followerJoinFuture.get().memberId()); - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withGroupInstanceId("follower-instance-id") .withGenerationId(rebalanceResult.generationId) .withMemberId(followerJoinFuture.get().memberId()) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); assertTrue(syncResult.records.isEmpty()); assertTrue(syncResult.syncFuture.isDone()); @@ -7149,13 +7186,13 @@ public void testStaticMemberRejoinAsFollowerWithUnknownMemberId() throws Excepti public void testStaticMemberRejoinAsFollowerWithKnownMemberIdAndNoProtocolChange() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // A static follower rejoin with no protocol change will not trigger rebalance. JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -7202,7 +7239,6 @@ public void testStaticMemberRejoinAsFollowerWithKnownMemberIdAndNoProtocolChange public void testStaticMemberRejoinAsFollowerWithMismatchedInstanceId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", @@ -7233,7 +7269,6 @@ public void testStaticMemberRejoinAsFollowerWithMismatchedInstanceId() throws Ex public void testStaticMemberRejoinAsLeaderWithMismatchedInstanceId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", @@ -7264,7 +7299,6 @@ public void testStaticMemberRejoinAsLeaderWithMismatchedInstanceId() throws Exce public void testStaticMemberSyncAsLeaderWithInvalidMemberId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); context.staticMembersJoinAndRebalance( "group-id", @@ -7289,13 +7323,13 @@ public void testStaticMemberSyncAsLeaderWithInvalidMemberId() throws Exception { public void testGetDifferentStaticMemberIdAfterEachRejoin() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); String lastMemberId = rebalanceResult.leaderId; for (int i = 0; i < 5; i++) { @@ -7332,7 +7366,6 @@ public void testGetDifferentStaticMemberIdAfterEachRejoin() throws Exception { public void testStaticMemberJoinWithUnknownInstanceIdAndKnownMemberId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", @@ -7363,13 +7396,13 @@ public void testStaticMemberJoinWithUnknownInstanceIdAndKnownMemberId() throws E public void testStaticMemberReJoinWithIllegalStateAsUnknownMember() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); context.staticMembersJoinAndRebalance( "group-id", "leader-instance-id", "follower-instance-id" ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); group.transitionTo(PREPARING_REBALANCE); group.transitionTo(EMPTY); @@ -7397,7 +7430,6 @@ public void testStaticMemberReJoinWithIllegalStateAsUnknownMember() throws Excep public void testStaticMemberFollowerFailToRejoinBeforeRebalanceTimeout() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); // Increase session timeout so that the follower won't be evicted when rebalance timeout is reached. RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( @@ -7407,6 +7439,7 @@ public void testStaticMemberFollowerFailToRejoinBeforeRebalanceTimeout() throws 10000, 15000 ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); String newMemberInstanceId = "new-member-instance-id"; String leaderId = rebalanceResult.leaderId; @@ -7487,7 +7520,6 @@ public void testStaticMemberFollowerFailToRejoinBeforeRebalanceTimeout() throws public void testStaticMemberLeaderFailToRejoinBeforeRebalanceTimeout() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); // Increase session timeout so that the leader won't be evicted when rebalance timeout is reached. RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( @@ -7497,6 +7529,7 @@ public void testStaticMemberLeaderFailToRejoinBeforeRebalanceTimeout() throws Ex 10000, 15000 ); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); String newMemberInstanceId = "new-member-instance-id"; JoinGroupRequestData request = new JoinGroupRequestBuilder() @@ -7648,7 +7681,8 @@ private void testSyncGroupProtocolTypeAndNameWith( // JoinGroup(follower) with the Protocol Type of the group CompletableFuture followerJoinFuture = new CompletableFuture<>(); - result = context.sendGenericGroupJoin(joinRequest.setGroupInstanceId("follower-instance-id"), followerJoinFuture); + result = context.sendGenericGroupJoin(joinRequest + .setGroupInstanceId("follower-instance-id"), followerJoinFuture); assertTrue(result.records().isEmpty()); assertFalse(followerJoinFuture.isDone()); @@ -7700,13 +7734,11 @@ public void testSyncGroupFromUnknownGroup() throws Exception { .build(); // SyncGroup with the provided Protocol Type and Name - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withMemberId("member-id") .withGenerationId(1) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); assertTrue(syncResult.records.isEmpty()); assertTrue(syncResult.syncFuture.isDone()); @@ -7719,14 +7751,16 @@ public void testSyncGroupFromUnknownMember() throws Exception { .build(); context.createGenericGroup("group-id"); - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withGroupInstanceId("leader-instance-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); - - JoinGroupResponseData joinResponse = context.joinGenericGroupAndCompleteJoin(joinRequest, true, true); + JoinGroupResponseData joinResponse = context.joinGenericGroupAndCompleteJoin( + new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withGroupInstanceId("leader-instance-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build(), + true, + true + ); String memberId = joinResponse.memberId(); int generationId = joinResponse.generationId(); @@ -7749,6 +7783,7 @@ public void testSyncGroupFromUnknownMember() throws Exception { assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); assertEquals(assignment.get(0).assignment(), syncResult.syncFuture.get().assignment()); + // Sync with unknown member. syncResult = context.sendGenericGroupSync(syncRequest.setMemberId("unknown-member-id")); assertTrue(syncResult.records.isEmpty()); @@ -7776,13 +7811,11 @@ public void testSyncGroupFromIllegalGeneration() throws Exception { int generationId = joinResponse.generationId(); // Send the sync group with an invalid generation - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withMemberId(memberId) .withGenerationId(generationId + 1) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); assertTrue(syncResult.records.isEmpty()); assertTrue(syncResult.syncFuture.isDone()); @@ -7797,7 +7830,7 @@ public void testJoinGroupFromUnchangedFollowerDoesNotRebalance() throws Exceptio GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() .withGroupId("group-id") @@ -7805,29 +7838,9 @@ public void testJoinGroupFromUnchangedFollowerDoesNotRebalance() throws Exceptio .withDefaultProtocolTypeAndProtocols() .build(); - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - // Simulate a successful write to log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - CompletableFuture followerJoinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), + joinRequest, followerJoinFuture ); @@ -7836,7 +7849,8 @@ public void testJoinGroupFromUnchangedFollowerDoesNotRebalance() throws Exceptio CompletableFuture leaderJoinFuture = new CompletableFuture<>(); result = context.sendGenericGroupJoin( - joinRequest.setMemberId(leaderId), + joinRequest + .setMemberId(leaderJoinResponse.memberId()), leaderJoinFuture ); @@ -7846,8 +7860,8 @@ public void testJoinGroupFromUnchangedFollowerDoesNotRebalance() throws Exceptio assertEquals(Errors.NONE.code(), leaderJoinFuture.get().errorCode()); assertEquals(Errors.NONE.code(), followerJoinFuture.get().errorCode()); assertEquals(leaderJoinFuture.get().generationId(), followerJoinFuture.get().generationId()); - assertEquals(leaderId, leaderJoinFuture.get().leader()); - assertEquals(leaderId, followerJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), leaderJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), followerJoinFuture.get().leader()); int nextGenerationId = leaderJoinFuture.get().generationId(); String followerId = followerJoinFuture.get().memberId(); @@ -7873,7 +7887,8 @@ public void testLeaderFailureInSyncGroup() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() .withGroupId("group-id") @@ -7883,29 +7898,9 @@ public void testLeaderFailureInSyncGroup() throws Exception { .withSessionTimeoutMs(5000) .build(); - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - // Simulate a successful write to log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - CompletableFuture followerJoinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), + joinRequest, followerJoinFuture ); @@ -7914,7 +7909,8 @@ public void testLeaderFailureInSyncGroup() throws Exception { CompletableFuture leaderJoinFuture = new CompletableFuture<>(); result = context.sendGenericGroupJoin( - joinRequest.setMemberId(leaderId), + joinRequest + .setMemberId(leaderJoinResponse.memberId()), leaderJoinFuture ); @@ -7924,8 +7920,8 @@ public void testLeaderFailureInSyncGroup() throws Exception { assertEquals(Errors.NONE.code(), leaderJoinFuture.get().errorCode()); assertEquals(Errors.NONE.code(), followerJoinFuture.get().errorCode()); assertEquals(leaderJoinFuture.get().generationId(), followerJoinFuture.get().generationId()); - assertEquals(leaderId, leaderJoinFuture.get().leader()); - assertEquals(leaderId, followerJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), leaderJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), followerJoinFuture.get().leader()); assertTrue(group.isInState(COMPLETING_REBALANCE)); int nextGenerationId = leaderJoinFuture.get().generationId(); @@ -7933,8 +7929,11 @@ public void testLeaderFailureInSyncGroup() throws Exception { // With no leader SyncGroup, the follower's sync request should fail with an error indicating // that it should rejoin - SyncResult followerSyncResult = context.sendGenericGroupSync(syncRequest.setMemberId(followerId) - .setGenerationId(nextGenerationId)); + SyncResult followerSyncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(followerId) + .withGenerationId(nextGenerationId) + .build()); assertTrue(followerSyncResult.records.isEmpty()); assertFalse(followerSyncResult.syncFuture.isDone()); @@ -7961,7 +7960,8 @@ public void testSyncGroupFollowerAfterLeader() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() .withGroupId("group-id") @@ -7971,29 +7971,10 @@ public void testSyncGroupFollowerAfterLeader() throws Exception { .withSessionTimeoutMs(5000) .build(); - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - // Simulate a successful write to log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - CompletableFuture followerJoinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID), followerJoinFuture ); @@ -8002,7 +7983,8 @@ public void testSyncGroupFollowerAfterLeader() throws Exception { CompletableFuture leaderJoinFuture = new CompletableFuture<>(); result = context.sendGenericGroupJoin( - joinRequest.setMemberId(leaderId), + joinRequest + .setMemberId(leaderJoinResponse.memberId()), leaderJoinFuture ); @@ -8012,8 +7994,8 @@ public void testSyncGroupFollowerAfterLeader() throws Exception { assertEquals(Errors.NONE.code(), leaderJoinFuture.get().errorCode()); assertEquals(Errors.NONE.code(), followerJoinFuture.get().errorCode()); assertEquals(leaderJoinFuture.get().generationId(), followerJoinFuture.get().generationId()); - assertEquals(leaderId, leaderJoinFuture.get().leader()); - assertEquals(leaderId, followerJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), leaderJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), followerJoinFuture.get().leader()); assertTrue(group.isInState(COMPLETING_REBALANCE)); int nextGenerationId = leaderJoinFuture.get().generationId(); @@ -8024,7 +8006,7 @@ public void testSyncGroupFollowerAfterLeader() throws Exception { // Sync group with leader to get new assignment. List assignment = new ArrayList<>(); assignment.add(new SyncGroupRequestAssignment() - .setMemberId(leaderId) + .setMemberId(leaderJoinResponse.memberId()) .setAssignment(leaderAssignment) ); assignment.add(new SyncGroupRequestAssignment() @@ -8032,8 +8014,15 @@ public void testSyncGroupFollowerAfterLeader() throws Exception { .setAssignment(followerAssignment) ); - syncResult = context.sendGenericGroupSync(syncRequest.setGenerationId(nextGenerationId) - .setAssignments(assignment)); + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderJoinResponse.memberId()) + .withGenerationId(leaderJoinResponse.generationId()) + .withAssignment(assignment) + .build(); + + SyncResult syncResult = context.sendGenericGroupSync(syncRequest + .setGenerationId(nextGenerationId)); // Simulate a successful write to log. This will update the group's assignment with the new assignment. syncResult.appendFuture.complete(null); @@ -8061,7 +8050,8 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() .withGroupId("group-id") @@ -8071,29 +8061,9 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { .withSessionTimeoutMs(5000) .build(); - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - // Simulate a successful write to log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - CompletableFuture followerJoinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), + joinRequest, followerJoinFuture ); @@ -8102,7 +8072,8 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { CompletableFuture leaderJoinFuture = new CompletableFuture<>(); result = context.sendGenericGroupJoin( - joinRequest.setMemberId(leaderId), + joinRequest + .setMemberId(leaderJoinResponse.memberId()), leaderJoinFuture ); @@ -8112,8 +8083,8 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { assertEquals(Errors.NONE.code(), leaderJoinFuture.get().errorCode()); assertEquals(Errors.NONE.code(), followerJoinFuture.get().errorCode()); assertEquals(leaderJoinFuture.get().generationId(), followerJoinFuture.get().generationId()); - assertEquals(leaderId, leaderJoinFuture.get().leader()); - assertEquals(leaderId, followerJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), leaderJoinFuture.get().leader()); + assertEquals(leaderJoinResponse.memberId(), followerJoinFuture.get().leader()); assertTrue(group.isInState(COMPLETING_REBALANCE)); int nextGenerationId = leaderJoinFuture.get().generationId(); @@ -8122,6 +8093,12 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { byte[] followerAssignment = new byte[]{1}; // Sync group with follower to get new assignment. + SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderJoinResponse.memberId()) + .withGenerationId(leaderJoinResponse.generationId()) + .build(); + SyncResult followerSyncResult = context.sendGenericGroupSync(syncRequest.setMemberId(followerId) .setGenerationId(nextGenerationId)); @@ -8131,7 +8108,7 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { // Sync group with leader to get new assignment. List assignment = new ArrayList<>(); assignment.add(new SyncGroupRequestAssignment() - .setMemberId(leaderId) + .setMemberId(leaderJoinResponse.memberId()) .setAssignment(leaderAssignment) ); assignment.add(new SyncGroupRequestAssignment() @@ -8139,7 +8116,7 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { .setAssignment(followerAssignment) ); - syncResult = context.sendGenericGroupSync(syncRequest.setMemberId(leaderId) + SyncResult syncResult = context.sendGenericGroupSync(syncRequest.setMemberId(leaderJoinResponse.memberId()) .setGenerationId(nextGenerationId) .setAssignments(assignment)); @@ -8170,53 +8147,31 @@ public void testSyncGroupLeaderAfterFollower() throws Exception { public void testJoinGroupFromUnchangedLeaderShouldRebalance() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); + // Join group from the leader should force the group to rebalance, which allows the + // leader to push new assignment when local metadata changes + CompletableFuture leaderRejoinFuture = new CompletableFuture<>(); + CoordinatorResult result = context.sendGenericGroupJoin( + new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderJoinResponse.memberId()) + .withDefaultProtocolTypeAndProtocols() + .build(), + leaderRejoinFuture + ); - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); + assertTrue(result.records().isEmpty()); + assertTrue(leaderRejoinFuture.isDone()); + assertEquals(Errors.NONE.code(), leaderRejoinFuture.get().errorCode()); + assertEquals(leaderJoinResponse.generationId() + 1, leaderRejoinFuture.get().generationId()); + } - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - // Simulate a successful write to log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - - // Join group from the leader should force the group to rebalance, which allows the - // leader to push new assignment when local metadata changes - CompletableFuture leaderJoinFuture = new CompletableFuture<>(); - CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(leaderId), - leaderJoinFuture - ); - - assertTrue(result.records().isEmpty()); - assertTrue(leaderJoinFuture.isDone()); - assertEquals(Errors.NONE.code(), leaderJoinFuture.get().errorCode()); - assertEquals(generationId + 1, leaderJoinFuture.get().generationId()); - } - - @Test - public void testJoinGroupCompletionWhenPendingMemberJoins() throws Exception { - GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() - .build(); - GenericGroup group = context.createGenericGroup("group-id"); + @Test + public void testJoinGroupCompletionWhenPendingMemberJoins() throws Exception { + GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() + .build(); + GenericGroup group = context.createGenericGroup("group-id"); // Set up a group in with a pending member. The test checks if the pending member joining // completes the rebalancing operation @@ -8263,38 +8218,16 @@ public void testJoinGroupCompletionWhenPendingMemberTimesOut() throws Exception public void testGenerationIdIncrementsOnRebalance() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - // Simulate a successful write to log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); CompletableFuture joinFuture = new CompletableFuture<>(); - CoordinatorResult result = context.sendGenericGroupJoin(joinRequest.setMemberId(leaderId), joinFuture); + CoordinatorResult result = context.sendGenericGroupJoin( + new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(leaderJoinResponse.memberId()) + .withDefaultProtocolTypeAndProtocols() + .build(), + joinFuture); assertTrue(result.records().isEmpty()); assertTrue(joinFuture.isDone()); @@ -8451,7 +8384,6 @@ public SyncResult( public void testStaticMemberHeartbeatLeaderWithInvalidMemberId() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - context.createGenericGroup("group-id"); RebalanceResult rebalanceResult = context.staticMembersJoinAndRebalance( "group-id", @@ -8459,14 +8391,12 @@ public void testStaticMemberHeartbeatLeaderWithInvalidMemberId() throws Exceptio "follower-instance-id" ); - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withGroupInstanceId("leader-instance-id") .withMemberId(rebalanceResult.leaderId) .withGenerationId(rebalanceResult.generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); assertTrue(syncResult.records.isEmpty()); assertTrue(syncResult.syncFuture.isDone()); @@ -8483,8 +8413,7 @@ public void testStaticMemberHeartbeatLeaderWithInvalidMemberId() throws Exceptio assertThrows(FencedInstanceIdException.class, () -> context.sendGenericGroupHeartbeat( heartbeatRequest .setGroupInstanceId("leader-instance-id") - .setMemberId("invalid-member-id") - )); + .setMemberId("invalid-member-id"))); } @Test @@ -8522,11 +8451,6 @@ public void testHeartbeatEmptyGroup() { .build(); GenericGroup group = context.createGenericGroup("group-id"); - JoinGroupRequestProtocolCollection protocols = new JoinGroupRequestProtocolCollection(); - protocols.add(new JoinGroupRequestProtocol() - .setName("range") - .setMetadata(new byte[]{0})); - group.add(new GenericGroupMember( "member-id", Optional.empty(), @@ -8535,12 +8459,9 @@ public void testHeartbeatEmptyGroup() { 10000, 5000, "consumer", - protocols + toProtocols("range") )); - group.transitionTo(PREPARING_REBALANCE); - group.transitionTo(EMPTY); - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() .setGroupId("group-id") .setMemberId("member-id") @@ -8554,47 +8475,12 @@ public void testHeartbeatEmptyGroup() { public void testHeartbeatUnknownMemberExistingGroup() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertTrue(group.isInState(STABLE)); - - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(new HeartbeatRequestData() .setGroupId("group-id") .setMemberId("unknown-member-id") - .setGenerationId(generationId); - - assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + .setGenerationId(leaderJoinResponse.generationId()))); } @Test @@ -8619,15 +8505,16 @@ public void testHeartbeatDuringPreparingRebalance() throws Exception { String memberId = joinFuture.get().memberId(); joinFuture = new CompletableFuture<>(); - context.sendGenericGroupJoin(joinRequest.setMemberId(memberId), joinFuture); + context.sendGenericGroupJoin(joinRequest + .setMemberId(memberId), joinFuture); assertTrue(group.isInState(PREPARING_REBALANCE)); - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(new HeartbeatRequestData() .setGroupId("group-id") .setMemberId(memberId) - .setGenerationId(0); + .setGenerationId(0)); - HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); } @@ -8637,27 +8524,21 @@ public void testHeartbeatDuringCompletingRebalance() throws Exception { .build(); GenericGroup group = context.createGenericGroup("group-id"); - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); + context.joinGenericGroupAsDynamicMemberAndCompleteJoin(new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .build()); - assertEquals(1, generationId); + assertEquals(1, leaderJoinResponse.generationId()); assertTrue(group.isInState(COMPLETING_REBALANCE)); - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId())); - HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(new HeartbeatResponseData(), heartbeatResponse); } @@ -8665,96 +8546,27 @@ public void testHeartbeatDuringCompletingRebalance() throws Exception { public void testHeartbeatIllegalGeneration() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertTrue(group.isInState(STABLE)); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() - .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId + 1); - - assertThrows(IllegalGenerationException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + assertThrows(IllegalGenerationException.class, () -> { + context.sendGenericGroupHeartbeat(new HeartbeatRequestData() + .setGroupId("group-id") + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId() + 1)); + }); } @Test public void testValidHeartbeat() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertTrue(group.isInState(STABLE)); - - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId())); - HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); } @@ -8762,109 +8574,35 @@ public void testValidHeartbeat() throws Exception { public void testGenericGroupMemberSessionTimeout() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .withRebalanceTimeoutMs(10000) - .withSessionTimeoutMs(5000) - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertTrue(group.isInState(STABLE)); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Advance clock by session timeout to kick member out. context.verifySessionExpiration(group, 5000); - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); - - assertThrows(UnknownMemberIdException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId()))); } @Test - public void testGenericGroupMemberMaintainsSession() throws Exception { + public void testGenericGroupMemberHeartbeatMaintainSession() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .withRebalanceTimeoutMs(10000) - .withSessionTimeoutMs(5000) - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertTrue(group.isInState(STABLE)); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); // Advance clock by 1/2 of session timeout. assertNoOrEmptyResult(context.sleep(2500)); HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId()); HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); - assertNoOrEmptyResult(context.sleep(2500)); heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); @@ -8875,55 +8613,25 @@ public void testGenericGroupMemberMaintainsSession() throws Exception { public void testGenericGroupMemberSessionTimeoutDuringRebalance() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .withRebalanceTimeoutMs(10000) - .withSessionTimeoutMs(5000) - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertTrue(group.isInState(STABLE)); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Add a new member. This should trigger a rebalance. The new member has the // 'genericGroupNewMemberJoinTimeoutMs` session timeout, so it has a longer expiration than the existing member. - CompletableFuture joinFuture = new CompletableFuture<>(); + CompletableFuture otherJoinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest - .setMemberId(UNKNOWN_MEMBER_ID), - joinFuture + new JoinGroupRequestBuilder() + .withGroupId("group-id") + .withMemberId(UNKNOWN_MEMBER_ID) + .withDefaultProtocolTypeAndProtocols() + .withRebalanceTimeoutMs(10000) + .withSessionTimeoutMs(5000) + .build(), + otherJoinFuture ); assertTrue(result.records().isEmpty()); - assertFalse(joinFuture.isDone()); + assertFalse(otherJoinFuture.isDone()); assertTrue(group.isInState(PREPARING_REBALANCE)); // Advance clock by 1/2 of session timeout. @@ -8931,8 +8639,8 @@ public void testGenericGroupMemberSessionTimeoutDuringRebalance() throws Excepti HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId()); HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); @@ -8945,8 +8653,8 @@ public void testGenericGroupMemberSessionTimeoutDuringRebalance() throws Excepti // Advance clock by remaining rebalance timeout to complete join phase. assertNoOrEmptyResult(context.sleep(2500)); - assertTrue(joinFuture.isDone()); - assertEquals(Errors.NONE.code(), joinFuture.get().errorCode()); + assertTrue(otherJoinFuture.isDone()); + assertEquals(Errors.NONE.code(), otherJoinFuture.get().errorCode()); assertEquals(1, group.size()); assertEquals(2, group.generationId()); assertTrue(group.isInState(COMPLETING_REBALANCE)); @@ -8985,10 +8693,6 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); // Simulate a successful write to the log. syncResult.appendFuture.complete(null); @@ -8998,26 +8702,27 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { // Add a new dynamic member. This should trigger a rebalance. The new member has the // 'genericGroupNewMemberJoinTimeoutMs` session timeout, so it has a longer expiration than the existing member. - CompletableFuture joinFuture = new CompletableFuture<>(); + CompletableFuture secondMemberJoinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID) + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID) .setGroupInstanceId(null) .setSessionTimeoutMs(2500), - joinFuture); + secondMemberJoinFuture); assertTrue(result.records().isEmpty()); - assertFalse(joinFuture.isDone()); + assertFalse(secondMemberJoinFuture.isDone()); assertTrue(group.isInState(PREPARING_REBALANCE)); // Send a couple heartbeats to keep the first member alive while the rebalance finishes. - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + HeartbeatRequestData firstMemberHeartbeatRequest = new HeartbeatRequestData() .setGroupId("group-id") .setMemberId(firstMemberId) .setGenerationId(firstGenerationId); for (int i = 0; i < 2; i++) { assertNoOrEmptyResult(context.sleep(2500)); - HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(firstMemberHeartbeatRequest); assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); } @@ -9026,23 +8731,20 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { // it will not be kicked out. assertNoOrEmptyResult(context.sleep(8000)); - assertTrue(joinFuture.isDone()); - assertEquals(Errors.NONE.code(), joinFuture.get().errorCode()); + assertTrue(secondMemberJoinFuture.isDone()); + assertEquals(Errors.NONE.code(), secondMemberJoinFuture.get().errorCode()); assertEquals(2, group.size()); assertEquals(2, group.generationId()); assertTrue(group.isInState(COMPLETING_REBALANCE)); - String otherMemberId = joinFuture.get().memberId(); + String otherMemberId = secondMemberJoinFuture.get().memberId(); syncResult = context.sendGenericGroupSync( - syncRequest.setGroupInstanceId(null) + syncRequest + .setGroupInstanceId(null) .setMemberId(otherMemberId) .setGenerationId(2)); - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); // Simulate a successful write to the log. syncResult.appendFuture.complete(null); @@ -9050,15 +8752,16 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); assertTrue(group.isInState(STABLE)); - // The unjoined static member should be remained in the group before session timeout. - assertThrows(IllegalGenerationException.class, () -> context.sendGenericGroupHeartbeat(heartbeatRequest)); + // The unjoined (first) static member should be remained in the group before session timeout. + assertThrows(IllegalGenerationException.class, () -> context.sendGenericGroupHeartbeat(firstMemberHeartbeatRequest)); - // Now session timeout the unjoined member. Still keeping the new member. + // Now session timeout the unjoined (first) member. Still keeping the new member. List expectedErrors = Arrays.asList(Errors.NONE, Errors.NONE, Errors.REBALANCE_IN_PROGRESS); for (Errors expectedError : expectedErrors) { assertNoOrEmptyResult(context.sleep(2000)); HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat( - heartbeatRequest.setMemberId(otherMemberId) + firstMemberHeartbeatRequest + .setMemberId(otherMemberId) .setGenerationId(2)); assertEquals(expectedError.code(), heartbeatResponse.errorCode()); @@ -9068,7 +8771,8 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { CompletableFuture otherMemberRejoinFuture = new CompletableFuture<>(); result = context.sendGenericGroupJoin( - joinRequest.setMemberId(otherMemberId) + joinRequest + .setMemberId(otherMemberId) .setGroupInstanceId(null) .setSessionTimeoutMs(2500), otherMemberRejoinFuture); @@ -9080,14 +8784,11 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { assertTrue(group.isInState(COMPLETING_REBALANCE)); SyncResult otherMemberResyncResult = context.sendGenericGroupSync( - syncRequest.setGroupInstanceId(null) + syncRequest + .setGroupInstanceId(null) .setMemberId(otherMemberId) .setGenerationId(3)); - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - otherMemberResyncResult.records - ); // Simulate a successful write to the log. otherMemberResyncResult.appendFuture.complete(null); @@ -9100,7 +8801,8 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { for (int i = 0; i < 20; i++) { assertNoOrEmptyResult(context.sleep(2000)); HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat( - heartbeatRequest.setMemberId(otherMemberId) + firstMemberHeartbeatRequest + .setMemberId(otherMemberId) .setGenerationId(3)); assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); @@ -9111,51 +8813,13 @@ public void testRebalanceCompletesBeforeMemberJoins() throws Exception { public void testSyncGroupEmptyAssignment() throws Exception { GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - - JoinGroupRequestData joinRequest = new JoinGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .withRebalanceTimeoutMs(10000) - .withSessionTimeoutMs(5000) - .build(); - - JoinGroupResponseData leaderJoinResponse = - context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() - .withGroupId("group-id") - .withMemberId(leaderId) - .withGenerationId(generationId) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); - // Simulate a successful write to the log. - syncResult.appendFuture.complete(null); - - assertTrue(syncResult.syncFuture.isDone()); - assertEquals(Errors.NONE.code(), syncResult.syncFuture.get().errorCode()); - assertEquals(0, syncResult.syncFuture.get().assignment().length); - assertTrue(group.isInState(STABLE)); + JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteRebalance("group-id"); - HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() + HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId())); - HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(Errors.NONE.code(), heartbeatResponse.errorCode()); } @@ -9188,18 +8852,12 @@ public void testSecondMemberPartiallyJoinAndTimeout() throws Exception { assertEquals(1, firstGenerationId); assertTrue(group.isInState(COMPLETING_REBALANCE)); - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withMemberId(firstMemberId) .withGenerationId(firstGenerationId) - .build(); + .build()); - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); // Simulate a successful write to the log. syncResult.appendFuture.complete(null); @@ -9210,7 +8868,8 @@ public void testSecondMemberPartiallyJoinAndTimeout() throws Exception { // Add a new dynamic pending member. CompletableFuture joinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID) + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID) .setGroupInstanceId(null) .setSessionTimeoutMs(5000), joinFuture, @@ -9251,11 +8910,10 @@ public void testRebalanceTimesOutWhenSyncRequestIsNotReceived() throws Exception // group is in the Empty state in this case. GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - int rebalanceTimeoutMs = 5000; int sessionTimeoutMs = 5000; - List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + List joinResponses = context.joinWithNMembers("group-id", 3, rebalanceTimeoutMs, sessionTimeoutMs); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Advance clock by 1/2 rebalance timeout. assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); @@ -9290,11 +8948,10 @@ public void testRebalanceTimesOutWhenSyncRequestIsNotReceivedFromFollowers() thr // group is in the PreparingRebalance state in this case. GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - int rebalanceTimeoutMs = 5000; int sessionTimeoutMs = 5000; - List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + List joinResponses = context.joinWithNMembers("group-id", 3, rebalanceTimeoutMs, sessionTimeoutMs); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Advance clock by 1/2 rebalance timeout. assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); @@ -9304,18 +8961,12 @@ public void testRebalanceTimesOutWhenSyncRequestIsNotReceivedFromFollowers() thr joinResponses.forEach(response -> context.verifyHeartbeat(group.groupId(), response, Errors.NONE)); // Leader sends a sync group request. - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withGenerationId(1) .withMemberId(joinResponses.get(0).memberId()) - .build(); + .build()); - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); - - assertEquals( - Collections.singletonList(newGroupMetadataRecord(group, MetadataVersion.latest())), - syncResult.records - ); // Simulate a successful write to the log. syncResult.appendFuture.complete(null); @@ -9351,11 +9002,10 @@ public void testRebalanceTimesOutWhenSyncRequestIsNotReceivedFromLeaders() throw // group is in the PreparingRebalance state in this case. GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - int rebalanceTimeoutMs = 5000; int sessionTimeoutMs = 5000; - List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + List joinResponses = context.joinWithNMembers("group-id", 3, rebalanceTimeoutMs, sessionTimeoutMs); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); // Advance clock by 1/2 rebalance timeout. assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs / 2)); @@ -9367,13 +9017,11 @@ public void testRebalanceTimesOutWhenSyncRequestIsNotReceivedFromLeaders() throw // Followers send sync group requests. List> followerSyncFutures = joinResponses.subList(1, 3).stream() .map(response -> { - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withGenerationId(1) .withMemberId(response.memberId()) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); assertTrue(syncResult.records.isEmpty()); assertFalse(syncResult.syncFuture.isDone()); @@ -9414,11 +9062,10 @@ public void testRebalanceDoesNotTimeOutWhenAllSyncAreReceived() throws Exception // members out when they have all sent their sync requests. Group should be in Stable state. GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder() .build(); - GenericGroup group = context.createGenericGroup("group-id"); - int rebalanceTimeoutMs = 5000; int sessionTimeoutMs = 5000; - List joinResponses = context.joinWithNMembers(group, 3, rebalanceTimeoutMs, sessionTimeoutMs); + List joinResponses = context.joinWithNMembers("group-id", 3, rebalanceTimeoutMs, sessionTimeoutMs); + GenericGroup group = context.groupMetadataManager.getOrMaybeCreateGenericGroup("group-id", false); String leaderId = joinResponses.get(0).memberId(); // Advance clock by 1/2 rebalance timeout. @@ -9430,13 +9077,11 @@ public void testRebalanceDoesNotTimeOutWhenAllSyncAreReceived() throws Exception // All members send sync group requests. List> syncFutures = joinResponses.stream().map(response -> { - SyncGroupRequestData syncRequest = new SyncGroupRequestBuilder() + SyncResult syncResult = context.sendGenericGroupSync(new SyncGroupRequestBuilder() .withGroupId("group-id") .withGenerationId(1) .withMemberId(response.memberId()) - .build(); - - SyncResult syncResult = context.sendGenericGroupSync(syncRequest); + .build()); if (response.memberId().equals(leaderId)) { assertEquals( @@ -9492,16 +9137,16 @@ public void testHeartbeatDuringRebalanceCausesRebalanceInProgress() throws Excep JoinGroupResponseData leaderJoinResponse = context.joinGenericGroupAsDynamicMemberAndCompleteJoin(joinRequest); - String leaderId = leaderJoinResponse.memberId(); - int generationId = leaderJoinResponse.generationId(); - - assertEquals(1, generationId); + assertEquals(1, leaderJoinResponse.generationId()); assertTrue(group.isInState(COMPLETING_REBALANCE)); // Then join with a new consumer to trigger a rebalance CompletableFuture joinFuture = new CompletableFuture<>(); CoordinatorResult result = context.sendGenericGroupJoin( - joinRequest.setMemberId(UNKNOWN_MEMBER_ID), joinFuture); + joinRequest + .setMemberId(UNKNOWN_MEMBER_ID), + joinFuture + ); assertTrue(result.records().isEmpty()); assertFalse(joinFuture.isDone()); @@ -9509,105 +9154,10 @@ public void testHeartbeatDuringRebalanceCausesRebalanceInProgress() throws Excep // We should be in the middle of a rebalance, so the heartbeat should return rebalance in progress. HeartbeatRequestData heartbeatRequest = new HeartbeatRequestData() .setGroupId("group-id") - .setMemberId(leaderId) - .setGenerationId(generationId); + .setMemberId(leaderJoinResponse.memberId()) + .setGenerationId(leaderJoinResponse.generationId()); HeartbeatResponseData heartbeatResponse = context.sendGenericGroupHeartbeat(heartbeatRequest); assertEquals(Errors.REBALANCE_IN_PROGRESS.code(), heartbeatResponse.errorCode()); } - - private List joinWithNMembers( - GroupMetadataManagerTestContext context, - GenericGroup group, - int numMembers, - int rebalanceTimeoutMs, - int sessionTimeoutMs - ) { - boolean requireKnownMemberId = true; - - // First join requests - JoinGroupRequestData request = new JoinGroupRequestBuilder() - .withGroupId(group.groupId()) - .withMemberId(UNKNOWN_MEMBER_ID) - .withDefaultProtocolTypeAndProtocols() - .withRebalanceTimeoutMs(rebalanceTimeoutMs) - .withSessionTimeoutMs(sessionTimeoutMs) - .build(); - - List> joinFutures = IntStream.range(0, numMembers) - .mapToObj(__ -> new CompletableFuture()) - .collect(Collectors.toList()); - - IntStream.range(0, numMembers).forEach(i -> { - CoordinatorResult result = context.sendGenericGroupJoin( - request, joinFutures.get(i), requireKnownMemberId); - - assertTrue(result.records().isEmpty()); - }); - - List memberIds = joinFutures.stream().map(future -> { - assertTrue(future.isDone()); - try { - return future.get().memberId(); - } catch (Exception e) { - fail("Unexpected exception: " + e.getMessage()); - } - return null; - }).collect(Collectors.toList()); - - // Second join requests - List> secondJoinFutures = IntStream.range(0, numMembers) - .mapToObj(__ -> new CompletableFuture()) - .collect(Collectors.toList()); - - IntStream.range(0, numMembers).forEach(i -> { - CoordinatorResult result = context.sendGenericGroupJoin( - request.setMemberId(memberIds.get(i)), secondJoinFutures.get(i), requireKnownMemberId); - - assertTrue(result.records().isEmpty()); - }); - secondJoinFutures.forEach(future -> assertFalse(future.isDone())); - - // Advance clock by initial rebalance delay. - assertNoOrEmptyResult(context.sleep(context.genericGroupInitialRebalanceDelayMs)); - secondJoinFutures.forEach(future -> assertFalse(future.isDone())); - // Advance clock by rebalance timeout to complete join phase. - assertNoOrEmptyResult(context.sleep(rebalanceTimeoutMs)); - - List joinResponses = secondJoinFutures.stream().map(future -> { - assertTrue(future.isDone()); - try { - assertEquals(Errors.NONE.code(), future.get().errorCode()); - return future.get(); - } catch (Exception e) { - fail("Unexpected exception: " + e.getMessage()); - } - return null; - }).collect(Collectors.toList()); - - assertEquals(numMembers, group.size()); - assertTrue(group.isInState(COMPLETING_REBALANCE)); - - return joinResponses; - } - - private void verifyHeartbeat( - GroupMetadataManagerTestContext context, - String groupId, - JoinGroupResponseData joinResponse, - Errors expectedError - ) { - HeartbeatRequestData request = new HeartbeatRequestData() - .setGroupId(groupId) - .setMemberId(joinResponse.memberId()) - .setGenerationId(joinResponse.generationId()); - - if (expectedError != Errors.NONE) { - assertThrows(expectedError.exception().getClass(), () -> context.sendGenericGroupHeartbeat(request)); - } else { - HeartbeatResponseData response = context.sendGenericGroupHeartbeat(request); - assertEquals(Errors.NONE.code(), response.errorCode()); - } - } } -