jeffkbkim commented on code in PR #13663:
URL: https://github.com/apache/kafka/pull/13663#discussion_r1191340188


##########
group-coordinator/src/main/java/org/apache/kafka/coordinator/group/generic/GenericGroup.java:
##########
@@ -0,0 +1,998 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.coordinator.group.generic;
+
+import org.apache.kafka.clients.consumer.internals.ConsumerProtocol;
+import org.apache.kafka.common.message.JoinGroupResponseData;
+import org.apache.kafka.common.message.ListGroupsResponseData;
+import org.apache.kafka.common.message.SyncGroupResponseData;
+import 
org.apache.kafka.common.message.JoinGroupResponseData.JoinGroupResponseMember;
+import org.apache.kafka.common.protocol.Errors;
+import org.apache.kafka.common.protocol.types.SchemaException;
+import org.apache.kafka.common.utils.LogContext;
+import org.apache.kafka.common.utils.Time;
+import org.apache.kafka.coordinator.group.Group;
+import org.slf4j.Logger;
+
+import java.nio.ByteBuffer;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.Set;
+import java.util.UUID;
+import java.util.concurrent.CompletableFuture;
+import java.util.stream.Collectors;
+
+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;
+import static 
org.apache.kafka.coordinator.group.generic.GenericGroupState.PREPARING_REBALANCE;
+
+/**
+ *
+ * This class holds metadata for a generic group where the
+ * member assignment is driven solely from the client side.
+ *
+ * The APIs members use to make changes to the group membership
+ * consist of JoinGroup, SyncGroup, and LeaveGroup.
+ */
+public class GenericGroup implements Group {
+
+    /**
+     * Empty generation.
+     */
+    public static final int NO_GENERATION = -1;
+
+    /**
+     * Protocol with empty name.
+     */
+    public static final String NO_PROTOCOL_NAME = "";
+
+    /**
+     * No leader.
+     */
+    public static final String NO_LEADER = "";
+
+    /**
+     * Delimiter used to join a randomly generated UUID
+     * with a client id or a group instance id.
+     */
+    private static final String MEMBER_ID_DELIMITER = "-";
+
+    /**
+     * The slf4j log context, used to create new loggers.
+     */
+    private final LogContext logContext;
+
+    /**
+     * The slf4j logger.
+     */
+    private final Logger log;
+
+    /**
+     * The group id.
+     */
+    private final String groupId;
+
+    /**
+     * The time.
+     */
+    private final Time time;
+
+    /**
+     * The current group state.
+     */
+    private GenericGroupState state;
+
+    /**
+     * The timestamp of when the group transitioned
+     * to its current state.
+     */
+    private Optional<Long> currentStateTimestamp;
+
+    /**
+     * The protocol type used for rebalance.
+     */
+    private Optional<String> protocolType = Optional.empty();
+
+    /**
+     * The protocol name used for rebalance.
+     */
+    private Optional<String> protocolName = Optional.empty();
+
+    /**
+     * The generation id.
+     */
+    private int generationId = 0;
+
+    /**
+     * The id of the group's leader.
+     */
+    private Optional<String> leaderId = Optional.empty();
+
+    /**
+     * The members of the group.
+     */
+    private final Map<String, GenericGroupMember> members = new HashMap<>();
+
+    /**
+     * The static members of the group.
+     */
+    private final Map<String, String> staticMembers = new HashMap<>();
+
+    /**
+     * Members who have yet to (re)join the group
+     * during the join group phase.
+     */
+    private final Set<String> pendingJoinMembers = new HashSet<>();
+
+    /**
+     * The number of members awaiting a join response.
+     */
+    private int numMembersAwaitingJoinResponse = 0;
+
+    /**
+     * Map of protocol names to the number of members that support them.
+     */
+    private final Map<String, Integer> supportedProtocols = new HashMap<>();
+
+    /**
+     * Members who have yet to sync with the group
+     * during the sync group phase.
+     */
+    private final Set<String> pendingSyncMembers = new HashSet<>();
+
+    /**
+     * The list of topics to which the group members are subscribed.
+     */
+    private Optional<Set<String>> subscribedTopics = Optional.empty();
+
+    /**
+     * A flag to indiciate whether a new member was added. Used
+     * to further delay initial joins (new group).
+     */
+    private boolean newMemberAdded = false;
+
+    public GenericGroup(
+        LogContext logContext,
+        String groupId,
+        GenericGroupState initialState,
+        Time time
+    ) {
+        this.logContext = Objects.requireNonNull(logContext);
+        this.log = logContext.logger(GenericGroup.class);
+        this.groupId = Objects.requireNonNull(groupId);
+        this.state = Objects.requireNonNull(initialState);
+        this.time = Objects.requireNonNull(time);
+        this.currentStateTimestamp = Optional.of(time.milliseconds());
+    }
+
+    /**
+     * The type of this group.
+     *
+     * @return The group type (Generic).
+     */
+    @Override
+    public GroupType type() {
+        return GroupType.GENERIC;
+    }
+
+    /**
+     * The state of this group.
+     *
+     * @return The current state as a String.
+     */
+    @Override
+    public String stateAsString() {
+        return this.state.toString();
+    }
+
+    /**
+     * @return the group id.
+     */
+    public String groupId() {
+        return this.groupId;
+    }
+
+    /**
+     * @return the generation id.
+     */
+    public int generationId() {
+        return this.generationId;
+    }
+
+    /**
+     * @return the protocol name.
+     */
+    public Optional<String> protocolName() {
+        return this.protocolName;
+    }
+
+    /**
+     * @return the protocol type.
+     */
+    public Optional<String> protocolType() {
+        return this.protocolType;
+    }
+
+    /**
+     * @return the current group state.
+     */
+    public GenericGroupState currentState() {
+        return state;
+    }
+
+    /**
+     * Compares the group's current state with the given state.
+     *
+     * @param groupState the state to match against.
+     * @return true if the state matches, false otherwise.
+     */
+    public boolean isState(GenericGroupState groupState) {
+        return this.state == groupState;
+    }
+
+    /**
+     * To identify whether the given member id is part of this group.
+     *
+     * @param memberId the given member id.
+     * @return true if the member is part of this group, false otherwise.
+     */
+    public boolean hasMemberId(String memberId) {
+        return members.containsKey(memberId);
+    }
+
+    /**
+     * Get the member metadata associated with the provided member id.
+     *
+     * @param memberId the member id.
+     * @return the member metadata if it exists, null otherwise.
+     */
+    public GenericGroupMember member(String memberId) {
+        return members.get(memberId);
+    }
+
+    /**
+     * @return the total number of members in this group.
+     */
+    public int size() {
+        return members.size();
+    }
+
+    /**
+     * Used to identify whether the given member is the leader of this group.
+     *
+     * @param memberId the member id.
+     * @return true if the member is the leader, false otherwise.
+     */
+    public boolean isLeader(String memberId) {
+        return leaderId.map(id -> id.equals(memberId)).orElse(false);
+    }
+
+    /**
+     * @return the leader id or null if a leader does not exist.
+     */
+    public String leaderOrNull() {
+        return leaderId.orElse(null);
+    }
+
+    /**
+     * @return the current state timestamp.
+     */
+    public long currentStateTimestampOrDefault() {
+        return currentStateTimestamp.orElse(-1L);
+    }
+
+    /**
+     * @return whether the group is using the consumer protocol.
+     */
+    public boolean isGenericGroup() {
+        return protocolType.map(type ->
+            type.equals(ConsumerProtocol.PROTOCOL_TYPE)
+        ).orElse(false);
+    }
+
+    public void add(GenericGroupMember member) {
+        add(member, null);
+    }
+
+    /**
+     * Add a member to this group.
+     *
+     * @param member  the member to add.
+     * @param future  the future to complete once the join group phase 
completes.
+     */
+    public void add(GenericGroupMember member, 
CompletableFuture<JoinGroupResponseData> future) {
+        member.groupInstanceId().ifPresent(instanceId -> {
+            if (staticMembers.containsKey(instanceId)) {
+                throw new IllegalStateException("Static member with 
groupInstanceId=" +
+                    instanceId + " cannot be added to group " + groupId + " 
since" +
+                    " it is already a member.");
+            }
+            staticMembers.put(instanceId, member.memberId());
+        });
+
+        if (members.isEmpty()) {
+            this.protocolType = Optional.of(member.protocolType());
+        }
+
+        if (!Objects.equals(this.protocolType.orElse(null), 
member.protocolType())) {
+            throw new IllegalStateException("The group and member's protocol 
type must be the same.");
+        }
+
+        if (!supportsProtocols(
+            member.protocolType(),
+            GenericGroupMember.plainProtocolSet(member.supportedProtocols()))) 
{
+
+            throw new IllegalStateException("None of the member's protocols 
can be supported.");
+        }
+
+        if (!leaderId.isPresent()) {
+            leaderId = Optional.of(member.memberId());
+        }
+
+        members.put(member.memberId(), member);
+        incrementSupportedProtocols(member);
+        member.setAwaitingJoinFuture(future);
+
+        if (member.isAwaitingJoin()) {
+            numMembersAwaitingJoinResponse++;
+        }
+
+        pendingJoinMembers.remove(member.memberId());
+    }
+
+    /**
+     * Remove a member from the group.
+     *
+     * @param memberId the member id to remove.
+     */
+    public void remove(String memberId) {
+        GenericGroupMember removedMember = members.remove(memberId);
+        if (removedMember != null) {
+            decrementSupportedProtocols(removedMember);
+            if (removedMember.isAwaitingJoin()) {
+                numMembersAwaitingJoinResponse--;
+            }
+
+            removedMember.groupInstanceId().ifPresent(staticMembers::remove);
+        }
+
+        if (isLeader(memberId)) {
+            Iterator<String> iter = members.keySet().iterator();
+            String newLeader = iter.hasNext() ? iter.next() : null;
+            leaderId = Optional.ofNullable(newLeader);
+        }
+
+        pendingJoinMembers.remove(memberId);
+        pendingSyncMembers.remove(memberId);
+    }
+
+    /**
+     * Check whether current leader has rejoined. If not, choose a
+     * new leader from one of the joined members.
+     *
+     * Return false if
+     *   1. the group is currently empty (has no designated leader)
+     *   2. no member rejoined
+     *
+     * @return whether a new leader was elected.
+     */
+    public boolean maybeElectNewJoinedLeader() {
+        if (leaderId.isPresent()) {
+            GenericGroupMember currentLeader = member(leaderId.get());
+            if (!currentLeader.isAwaitingJoin()) {
+                for (GenericGroupMember member : members.values()) {
+                    if (member.isAwaitingJoin()) {
+                        leaderId = Optional.of(member.memberId());
+                        log.info("Group leader [memberId: {}, groupInstanceId: 
{}] " +
+                            "failed to join before the rebalance timeout. 
Member {} " +
+                            "was elected as the new leader.",
+                            currentLeader.memberId(),
+                            currentLeader.groupInstanceId().orElse("None"),
+                            member
+                        );
+                        return true;
+                    }
+                }
+                log.info("Group leader [memberId: {}, groupInstanceId: {}] " +
+                        "failed to join before the rebalance timeout and the " 
+
+                        "group couldn't proceed to the next generation because 
" +
+                        "no member joined.",
+                    currentLeader.memberId(),
+                    currentLeader.groupInstanceId().orElse("None")
+                );
+                return false;
+            }
+            return false;

Review Comment:
   thanks for the catch. added a test case, 
`testMaybeElectNewJoinedLeaderChooseExisting`



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

To unsubscribe, e-mail: jira-unsubscr...@kafka.apache.org

For queries about this service, please contact Infrastructure at:
us...@infra.apache.org

Reply via email to