Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@

import java.nio.ByteBuffer;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;

Expand All @@ -49,6 +50,71 @@ public class ApiVersionsResponse extends AbstractResponse {

private final ApiVersionsResponseData data;

public static class Builder {
private Errors error = Errors.NONE;
private int throttleTimeMs = 0;
private ApiVersionCollection apiVersions = null;
private Features<SupportedVersionRange> supportedFeatures = null;
private Map<String, Short> finalizedFeatures = null;
private long finalizedFeaturesEpoch = 0;
private boolean zkMigrationEnabled = false;
private boolean suppressFeatureLevel0 = false;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

suppressFeatureLevel0 => alterFeatureLevel0 ?


public Builder setError(Errors error) {
this.error = error;
return this;
}

public Builder setThrottleTimeMs(int throttleTimeMs) {
this.throttleTimeMs = throttleTimeMs;
return this;
}

public Builder setApiVersions(ApiVersionCollection apiVersions) {
this.apiVersions = apiVersions;
return this;
}

public Builder setSupportedFeatures(Features<SupportedVersionRange> supportedFeatures) {
this.supportedFeatures = supportedFeatures;
return this;
}

public Builder setFinalizedFeatures(Map<String, Short> finalizedFeatures) {
this.finalizedFeatures = finalizedFeatures;
return this;
}

public Builder setFinalizedFeaturesEpoch(long finalizedFeaturesEpoch) {
this.finalizedFeaturesEpoch = finalizedFeaturesEpoch;
return this;
}

public Builder setZkMigrationEnabled(boolean zkMigrationEnabled) {
this.zkMigrationEnabled = zkMigrationEnabled;
return this;
}

public Builder setSuppressFeatureLevel0(boolean suppressFeatureLevel0) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

suppressFeatureLevel0 => alterFeatureLevel0 ?

this.suppressFeatureLevel0 = suppressFeatureLevel0;
return this;
}

public ApiVersionsResponse build() {
final ApiVersionsResponseData data = new ApiVersionsResponseData();
data.setErrorCode(error.code());
data.setApiKeys(Objects.requireNonNull(apiVersions));
data.setThrottleTimeMs(throttleTimeMs);
data.setSupportedFeatures(
createSupportedFeatureKeys(Objects.requireNonNull(supportedFeatures), suppressFeatureLevel0));
data.setFinalizedFeatures(
createFinalizedFeatureKeys(Objects.requireNonNull(finalizedFeatures)));
data.setFinalizedFeaturesEpoch(finalizedFeaturesEpoch);
data.setZkMigrationReady(zkMigrationEnabled);
return new ApiVersionsResponse(data);
}
}

public ApiVersionsResponse(ApiVersionsResponseData data) {
super(ApiKeys.API_VERSIONS);
this.data = data;
Expand Down Expand Up @@ -105,65 +171,32 @@ public static ApiVersionsResponse parse(ByteBuffer buffer, short version) {
}
}

public static ApiVersionsResponse createApiVersionsResponse(
int throttleTimeMs,
public static ApiVersionCollection controllerApiVersions(
RecordVersion minRecordVersion,
Features<SupportedVersionRange> latestSupportedFeatures,
Map<String, Short> finalizedFeatures,
long finalizedFeaturesEpoch,
NodeApiVersions controllerApiVersions,
ListenerType listenerType,
boolean enableUnstableLastVersion,
boolean zkMigrationEnabled,
boolean clientTelemetryEnabled
) {
ApiVersionCollection apiKeys;
if (controllerApiVersions != null) {
apiKeys = intersectForwardableApis(
listenerType,
minRecordVersion,
controllerApiVersions.allSupportedApiVersions(),
enableUnstableLastVersion,
clientTelemetryEnabled
);
} else {
apiKeys = filterApis(
minRecordVersion,
listenerType,
enableUnstableLastVersion,
clientTelemetryEnabled
);
}

return createApiVersionsResponse(
throttleTimeMs,
apiKeys,
latestSupportedFeatures,
finalizedFeatures,
finalizedFeaturesEpoch,
zkMigrationEnabled
);
return intersectForwardableApis(
listenerType,
minRecordVersion,
controllerApiVersions.allSupportedApiVersions(),
enableUnstableLastVersion,
clientTelemetryEnabled);
}

public static ApiVersionsResponse createApiVersionsResponse(
int throttleTimeMs,
ApiVersionCollection apiVersions,
Features<SupportedVersionRange> latestSupportedFeatures,
Map<String, Short> finalizedFeatures,
long finalizedFeaturesEpoch,
boolean zkMigrationEnabled
public static ApiVersionCollection brokerApiVersions(
RecordVersion minRecordVersion,
ListenerType listenerType,
boolean enableUnstableLastVersion,
boolean clientTelemetryEnabled
) {
return new ApiVersionsResponse(
createApiVersionsResponseData(
throttleTimeMs,
Errors.NONE,
apiVersions,
latestSupportedFeatures,
finalizedFeatures,
finalizedFeaturesEpoch,
zkMigrationEnabled
)
);
return filterApis(
minRecordVersion,
listenerType,
enableUnstableLastVersion,
clientTelemetryEnabled);
}

public static ApiVersionCollection filterApis(
Expand Down Expand Up @@ -249,37 +282,24 @@ public static ApiVersionCollection intersectForwardableApis(
return apiKeys;
}

private static ApiVersionsResponseData createApiVersionsResponseData(
final int throttleTimeMs,
final Errors error,
final ApiVersionCollection apiKeys,
final Features<SupportedVersionRange> latestSupportedFeatures,
final Map<String, Short> finalizedFeatures,
final long finalizedFeaturesEpoch,
final boolean zkMigrationEnabled
) {
final ApiVersionsResponseData data = new ApiVersionsResponseData();
data.setThrottleTimeMs(throttleTimeMs);
data.setErrorCode(error.code());
data.setApiKeys(apiKeys);
data.setSupportedFeatures(createSupportedFeatureKeys(latestSupportedFeatures));
data.setFinalizedFeatures(createFinalizedFeatureKeys(finalizedFeatures));
data.setFinalizedFeaturesEpoch(finalizedFeaturesEpoch);
data.setZkMigrationReady(zkMigrationEnabled);

return data;
}

private static SupportedFeatureKeyCollection createSupportedFeatureKeys(
Features<SupportedVersionRange> latestSupportedFeatures) {
Features<SupportedVersionRange> latestSupportedFeatures,
Comment thread
cmccabe marked this conversation as resolved.
boolean suppressV0

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

suppressV0 => alterV0 ?

) {
SupportedFeatureKeyCollection converted = new SupportedFeatureKeyCollection();
for (Map.Entry<String, SupportedVersionRange> feature : latestSupportedFeatures.features().entrySet()) {
final SupportedFeatureKey key = new SupportedFeatureKey();
final SupportedVersionRange versionRange = feature.getValue();
key.setName(feature.getKey());
key.setMinVersion(versionRange.min());
key.setMaxVersion(versionRange.max());
converted.add(key);
// Some older clients will have deserialization problems if a feature's
// minimum supported level is 0. Therefore, when preparing ApiVersionResponse
// at versions less than 4, we must leave those features out. See KAFKA-17011
// for details.
if (!(suppressV0 && versionRange.min() == 0)) {
final SupportedFeatureKey key = new SupportedFeatureKey();
key.setName(feature.getKey());
key.setMinVersion(versionRange.min());
key.setMaxVersion(versionRange.max());
converted.add(key);
}
}

return converted;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,9 @@
// Versions 0 through 2 of ApiVersionsRequest are the same.
//
// Version 3 is the first flexible version and adds ClientSoftwareName and ClientSoftwareVersion.
"validVersions": "0-3",
//
// Version 4 fixes KAFKA-17011, which blocked SupportedFeatures.MinVersion in the response from being 0.
"validVersions": "0-4",
"flexibleVersions": "3+",
"fields": [
{ "name": "ClientSoftwareName", "type": "string", "versions": "3+",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,9 @@
//
// Starting from Apache Kafka 2.4 (KIP-511), ApiKeys field is populated with the supported
// versions of the ApiVersionsRequest when an UNSUPPORTED_VERSION error is returned.
"validVersions": "0-3",
//
// Version 4 fixes KAFKA-17011, which blocked SupportedFeatures.MinVersion from being 0.
"validVersions": "0-4",
"flexibleVersions": "3+",
"fields": [
{ "name": "ErrorCode", "type": "int16", "versions": "0+",
Expand Down Expand Up @@ -67,7 +69,7 @@
{"name": "MaxVersionLevel", "type": "int16", "versions": "3+",
"about": "The cluster-wide finalized max version level for the feature."},
{"name": "MinVersionLevel", "type": "int16", "versions": "3+",
"about": "The cluster-wide finalized min version level for the feature."}
"about": "The cluster-wide finalized min version level for the feature. Due to KAFKA-17011, this must be 1 or greater in v3. It can be 0 in v4."}
]
},
{ "name": "ZkMigrationReady", "type": "bool", "versions": "3+", "taggedVersions": "3+",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -717,14 +717,16 @@ private static Features<org.apache.kafka.common.feature.SupportedVersionRange> c

private static ApiVersionsResponse prepareApiVersionsResponseForDescribeFeatures(Errors error) {
if (error == Errors.NONE) {
return ApiVersionsResponse.createApiVersionsResponse(
0,
ApiVersionsResponse.filterApis(RecordVersion.current(), ApiMessageType.ListenerType.ZK_BROKER, false, false),
convertSupportedFeaturesMap(defaultFeatureMetadata().supportedFeatures()),
Collections.singletonMap("test_feature_1", (short) 2),
defaultFeatureMetadata().finalizedFeaturesEpoch().get(),
false
);
return new ApiVersionsResponse.Builder().
setApiVersions(ApiVersionsResponse.filterApis(
RecordVersion.current(), ApiMessageType.ListenerType.ZK_BROKER, false, false)).
setSupportedFeatures(
convertSupportedFeaturesMap(defaultFeatureMetadata().supportedFeatures())).
setFinalizedFeatures(
Collections.singletonMap("test_feature_1", (short) 2)).
setFinalizedFeaturesEpoch(
defaultFeatureMetadata().finalizedFeaturesEpoch().get()).
build();
}
return new ApiVersionsResponse(
new ApiVersionsResponseData()
Expand Down
Loading