This is an automated email from the ASF dual-hosted git repository.

chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 58f63f448e3 MINOR: Complete nodeApiVersions future when 
describeFeatures fails (#22978)
58f63f448e3 is described below

commit 58f63f448e3b8f94e71bbf1d96935b268443f1e6
Author: Ming-Yen Chung <[email protected]>
AuthorDate: Wed Jul 29 10:50:18 2026 +0800

    MINOR: Complete nodeApiVersions future when describeFeatures fails (#22978)
    
    Follow-up to
    https://github.com/apache/kafka/pull/20598#discussion_r3667202261
    
    `describeFeatures()` now returns an `InternalDescribeFeaturesResult`
    carrying a `nodeApiVersions` future, but `handleFailure` only completes
    the `FeatureMetadata` future. When the ApiVersions request fails (for
    example on timeout), `nodeApiVersions()` is never completed and callers
    waiting on it hang indefinitely. The error-code branch in
    `handleResponse` already completes both.
    
    Extended the existing `testDescribeFeatures*` tests to assert both
    futures. `testDescribeFeaturesWithNodeFailure` is the one that covers
    `handleFailure`; it hangs until the class-level timeout without this
    fix.
    
    Reviewers: Chia-Ping Tsai <[email protected]>, Ken Huang
     <[email protected]>
---
 .../kafka/clients/admin/KafkaAdminClient.java      |  3 ++-
 .../kafka/clients/admin/KafkaAdminClientTest.java  | 22 ++++++++++++----------
 2 files changed, 14 insertions(+), 11 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java 
b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
index 312bd28e947..873c2a3e5c5 100644
--- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
+++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
@@ -4690,7 +4690,8 @@ public class KafkaAdminClient extends AdminClient {
 
             @Override
             void handleFailure(Throwable throwable) {
-                completeAllExceptionally(Collections.singletonList(future), 
throwable);
+                future.completeExceptionally(throwable);
+                nodeApiVersionsFuture.completeExceptionally(throwable);
             }
         };
 
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
index a61a660a59c..0162d2faa49 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java
@@ -24,6 +24,7 @@ import org.apache.kafka.clients.NodeApiVersions;
 import org.apache.kafka.clients.admin.DeleteAclsResult.FilterResults;
 import org.apache.kafka.clients.admin.ListOffsetsResult.ListOffsetsResultInfo;
 import org.apache.kafka.clients.admin.internals.AdminMetadataManager;
+import org.apache.kafka.clients.admin.internals.InternalDescribeFeaturesResult;
 import org.apache.kafka.clients.consumer.ConsumerPartitionAssignor;
 import org.apache.kafka.clients.consumer.OffsetAndMetadata;
 import org.apache.kafka.clients.consumer.internals.ConsumerProtocol;
@@ -9026,10 +9027,10 @@ public class KafkaAdminClientTest {
             env.kafkaClient().prepareResponse(
                 body -> body instanceof ApiVersionsRequest,
                 prepareApiVersionsResponseForDescribeFeatures(Errors.NONE));
-            final KafkaFuture<FeatureMetadata> future = 
env.adminClient().describeFeatures(
-                new 
DescribeFeaturesOptions().timeoutMs(10000)).featureMetadata();
-            final FeatureMetadata metadata = future.get();
-            assertEquals(defaultFeatureMetadata(), metadata);
+            final var result = (InternalDescribeFeaturesResult) 
env.adminClient().describeFeatures(
+                new DescribeFeaturesOptions().timeoutMs(10000));
+            assertEquals(defaultFeatureMetadata(), 
result.featureMetadata().get());
+            
assertNotNull(result.nodeApiVersions().get().apiVersion(ApiKeys.API_VERSIONS));
         }
     }
 
@@ -9041,9 +9042,9 @@ public class KafkaAdminClientTest {
                 
prepareApiVersionsResponseForDescribeFeatures(Errors.INVALID_REQUEST));
             final DescribeFeaturesOptions options = new 
DescribeFeaturesOptions();
             options.timeoutMs(10000);
-            final KafkaFuture<FeatureMetadata> future = 
env.adminClient().describeFeatures(options).featureMetadata();
-            final ExecutionException e = 
assertThrows(ExecutionException.class, future::get);
-            assertEquals(Errors.INVALID_REQUEST.exception().getClass(), 
e.getCause().getClass());
+            final var result = (InternalDescribeFeaturesResult) 
env.adminClient().describeFeatures(options);
+            TestUtils.assertFutureThrows(InvalidRequestException.class, 
result.featureMetadata());
+            TestUtils.assertFutureThrows(InvalidRequestException.class, 
result.nodeApiVersions());
         }
     }
 
@@ -9068,9 +9069,10 @@ public class KafkaAdminClientTest {
                 body -> body instanceof ApiVersionsRequest,
                 prepareApiVersionsResponseForDescribeFeatures(Errors.NONE),
                 env.cluster().nodeById(1));
-            final KafkaFuture<FeatureMetadata> future = 
env.adminClient().describeFeatures(
-                new 
DescribeFeaturesOptions().timeoutMs(1000).nodeId(0)).featureMetadata();
-            assertThrows(ExecutionException.class, future::get);
+            final var result = (InternalDescribeFeaturesResult) 
env.adminClient().describeFeatures(
+                new DescribeFeaturesOptions().timeoutMs(1000).nodeId(0));
+            TestUtils.assertFutureThrows(TimeoutException.class, 
result.featureMetadata());
+            TestUtils.assertFutureThrows(TimeoutException.class, 
result.nodeApiVersions());
         }
     }
 

Reply via email to