diff --git a/metadata/src/main/java/org/apache/kafka/controller/ClusterFeatureSupportDescriber.java b/metadata/src/main/java/org/apache/kafka/controller/ClusterFeatureSupportDescriber.java index 77dc11fb938bb..89b81a6312672 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterFeatureSupportDescriber.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterFeatureSupportDescriber.java @@ -22,9 +22,23 @@ import java.util.Iterator; import java.util.Map; import java.util.Map.Entry; +import java.util.Set; public interface ClusterFeatureSupportDescriber { Iterator>> brokerSupported(); Iterator>> controllerSupported(); + + /** + * The IDs of the controllers which are currently members of the quorum. + * + *

This is separate from {@link #controllerSupported()} because a controller + * registration can outlive the controller's membership in a dynamic quorum. + * + * @return the current quorum controller IDs, or an empty set when the caller + * should use its static quorum configuration + */ + default Set controllerIds() { + return Set.of(); + } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java index 01188bdd6d7d4..01a054e54d2f4 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java @@ -46,6 +46,7 @@ import java.util.Map; import java.util.Map.Entry; import java.util.Optional; +import java.util.Set; import java.util.function.Consumer; import static org.apache.kafka.common.metadata.MetadataRecordType.FEATURE_LEVEL_RECORD; @@ -341,10 +342,17 @@ private Optional reasonNotSupported( HashSet foundControllers = new HashSet<>(); foundControllers.add(quorumFeatures.nodeId()); if (metadataVersionOrThrow().isControllerRegistrationSupported()) { + Set controllerIds = clusterSupportDescriber.controllerIds(); + if (controllerIds.isEmpty()) { + controllerIds = Set.copyOf(quorumFeatures.quorumNodeIds()); + } for (Iterator>> iter = clusterSupportDescriber.controllerSupported(); iter.hasNext(); ) { Entry> entry = iter.next(); + if (!controllerIds.contains(entry.getKey())) { + continue; + } if (entry.getKey() == quorumFeatures.nodeId()) { // No need to re-check the features supported by this controller, since we // already checked that above. @@ -357,7 +365,7 @@ private Optional reasonNotSupported( foundControllers.add(entry.getKey()); numControllersChecked++; } - for (int id : quorumFeatures.quorumNodeIds()) { + for (int id : controllerIds) { if (!foundControllers.contains(id)) { return Optional.of("controller " + id + " has not registered, and may not " + "support this feature"); diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index 6620c7e6dc983..e95df9fdba120 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -527,6 +527,11 @@ public Iterator>> brokerSupported() { public Iterator>> controllerSupported() { return clusterControl.controllerSupportedFeatures(); } + + @Override + public Set controllerIds() { + return raftClient.voterIds(); + } } class PeriodicTaskControlManagerQueueAccessor implements PeriodicTaskControlManager.QueueAccessor { diff --git a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java index 941578cedcc2a..08fc5d211265c 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/FeatureControlManagerTest.java @@ -176,6 +176,14 @@ public void testReplayKraftVersionFeatureLevel() { static ClusterFeatureSupportDescriber createFakeClusterFeatureSupportDescriber( List>> brokerRanges, List>> controllerRanges + ) { + return createFakeClusterFeatureSupportDescriber(brokerRanges, controllerRanges, Set.of()); + } + + static ClusterFeatureSupportDescriber createFakeClusterFeatureSupportDescriber( + List>> brokerRanges, + List>> controllerRanges, + Set controllerIds ) { return new ClusterFeatureSupportDescriber() { @Override @@ -187,9 +195,76 @@ public Iterator>> brokerSupported() public Iterator>> controllerSupported() { return controllerRanges.iterator(); } + + @Override + public Set controllerIds() { + return controllerIds; + } }; } + @Test + public void testFeatureUpgradeIgnoresRemovedControllerRegistration() { + FeatureControlManager manager = new FeatureControlManager.Builder(). + setQuorumFeatures(new QuorumFeatures( + 0, + QuorumFeatures.defaultSupportedFeatureMap(true), + List.of(0, 1, 2))). + setClusterFeatureSupportDescriber(createFakeClusterFeatureSupportDescriber( + List.of(), + List.of( + new SimpleImmutableEntry<>(1, Map.of( + MetadataVersion.FEATURE_NAME, + VersionRange.of(MetadataVersion.IBP_3_7_IV0.featureLevel(), + MetadataVersion.IBP_4_3_IV0.featureLevel()))), + new SimpleImmutableEntry<>(2, Map.of( + MetadataVersion.FEATURE_NAME, + VersionRange.of(MetadataVersion.IBP_3_7_IV0.featureLevel(), + MetadataVersion.IBP_4_2_IV1.featureLevel()))) + ), + Set.of(0, 1))). + build(); + manager.replay(new FeatureLevelRecord().setName(MetadataVersion.FEATURE_NAME). + setFeatureLevel(MetadataVersion.IBP_3_7_IV0.featureLevel())); + + ControllerResult result = manager.updateFeatures( + Map.of(MetadataVersion.FEATURE_NAME, MetadataVersion.IBP_4_3_IV0.featureLevel()), + Map.of(), + false, + 0); + + assertEquals(ApiError.NONE, result.response()); + } + + @Test + public void testFeatureUpgradeRequiresCurrentControllerRegistration() { + FeatureControlManager manager = new FeatureControlManager.Builder(). + setQuorumFeatures(new QuorumFeatures( + 0, + QuorumFeatures.defaultSupportedFeatureMap(true), + List.of(0, 1))). + setClusterFeatureSupportDescriber(createFakeClusterFeatureSupportDescriber( + List.of(), + List.of(new SimpleImmutableEntry<>(1, Map.of( + MetadataVersion.FEATURE_NAME, + VersionRange.of(MetadataVersion.IBP_3_7_IV0.featureLevel(), + MetadataVersion.IBP_4_3_IV0.featureLevel())))), + Set.of(0, 1, 2))). + build(); + manager.replay(new FeatureLevelRecord().setName(MetadataVersion.FEATURE_NAME). + setFeatureLevel(MetadataVersion.IBP_3_7_IV0.featureLevel())); + + ControllerResult result = manager.updateFeatures( + Map.of(MetadataVersion.FEATURE_NAME, MetadataVersion.IBP_4_3_IV0.featureLevel()), + Map.of(), + false, + 0); + + assertEquals(new ApiError(Errors.INVALID_UPDATE_VERSION, + "Invalid update version 30 for feature metadata.version. controller 2 has not " + + "registered, and may not support this feature"), result.response()); + } + @Test public void testUpdateFeaturesErrorCases() { LogContext logContext = new LogContext(); diff --git a/metadata/src/test/java/org/apache/kafka/controller/MockRaftClient.java b/metadata/src/test/java/org/apache/kafka/controller/MockRaftClient.java index f817717a2b342..e6b15c14267f6 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/MockRaftClient.java +++ b/metadata/src/test/java/org/apache/kafka/controller/MockRaftClient.java @@ -61,6 +61,7 @@ import java.util.Optional; import java.util.OptionalInt; import java.util.OptionalLong; +import java.util.Set; import java.util.TreeMap; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; @@ -765,6 +766,11 @@ public Optional voterNode(int id, ListenerName listenerName) { return Optional.empty(); } + @Override + public Set voterIds() { + return Set.copyOf(shared.raftClients.keySet()); + } + public List> listeners() { final CompletableFuture>> future = new CompletableFuture<>(); eventQueue.append(() -> diff --git a/metadata/src/test/java/org/apache/kafka/image/publisher/SnapshotEmitterTest.java b/metadata/src/test/java/org/apache/kafka/image/publisher/SnapshotEmitterTest.java index 57fe845c98fce..8b1e6f261f6c7 100644 --- a/metadata/src/test/java/org/apache/kafka/image/publisher/SnapshotEmitterTest.java +++ b/metadata/src/test/java/org/apache/kafka/image/publisher/SnapshotEmitterTest.java @@ -37,6 +37,7 @@ import java.util.Optional; import java.util.OptionalInt; import java.util.OptionalLong; +import java.util.Set; import java.util.TreeMap; import java.util.concurrent.CompletableFuture; @@ -80,6 +81,11 @@ public Optional voterNode(int id, ListenerName listenerName) { return Optional.empty(); } + @Override + public Set voterIds() { + return Set.of(); + } + @Override public long prepareAppend(int epoch, List records) { return 0; diff --git a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java index 44f6013e10812..b2773130f0815 100644 --- a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java +++ b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java @@ -3900,6 +3900,11 @@ public Optional voterNode(int id, ListenerName listenerName) { return partitionState.lastVoterSet().voterNode(id, listenerName); } + @Override + public Set voterIds() { + return partitionState.lastVoterSet().voterIds(); + } + // Visible only for test QuorumState quorum() { // It's okay to return null since this method is only called by tests diff --git a/raft/src/main/java/org/apache/kafka/raft/RaftClient.java b/raft/src/main/java/org/apache/kafka/raft/RaftClient.java index d028fe848a338..74f0170755a7d 100644 --- a/raft/src/main/java/org/apache/kafka/raft/RaftClient.java +++ b/raft/src/main/java/org/apache/kafka/raft/RaftClient.java @@ -30,6 +30,7 @@ import java.util.Optional; import java.util.OptionalInt; import java.util.OptionalLong; +import java.util.Set; import java.util.concurrent.CompletionStage; public interface RaftClient extends AutoCloseable { @@ -152,6 +153,13 @@ default void beginShutdown() {} */ Optional voterNode(int id, ListenerName listenerName); + /** + * Return the IDs of the voters in the current quorum. + * + * @return the current voter IDs + */ + Set voterIds(); + /** * Prepare a list of records to be appended to the log. * diff --git a/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientReconfigTest.java b/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientReconfigTest.java index 08f3d63833b68..63116372512ad 100644 --- a/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientReconfigTest.java +++ b/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientReconfigTest.java @@ -1124,6 +1124,7 @@ public void testRemoveVoter() throws Exception { int epoch = context.currentEpoch(); assertTrue(context.client.quorum().isVoter(follower2)); + assertEquals(voters.voterIds(), context.client.voterIds()); checkLeaderMetricValues(3, 0, 0, context); @@ -1144,6 +1145,7 @@ public void testRemoveVoter() throws Exception { // follower2 should not be a voter in the latest voter set assertFalse(context.client.quorum().isVoter(follower2)); + assertEquals(Set.of(local.id(), follower1.id()), context.client.voterIds()); checkLeaderMetricValues(2, 1, 1, context); // Send a FETCH to increase the HWM and commit the new voter set