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());
}
}