Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.function.Supplier;
import java.util.function.Function;

public class ChannelBuilders {
private static final Logger log = LoggerFactory.getLogger(ChannelBuilders.class);
Expand Down Expand Up @@ -104,7 +104,7 @@ public static ChannelBuilder serverChannelBuilder(ListenerName listenerName,
DelegationTokenCache tokenCache,
Time time,
LogContext logContext,
Supplier<ApiVersionsResponse> apiVersionSupplier) {
Function<Short, ApiVersionsResponse> apiVersionSupplier) {
return create(securityProtocol, Mode.SERVER, JaasContext.Type.SERVER, config, listenerName,
isInterBrokerListener, null, true, credentialCache,
tokenCache, time, logContext, apiVersionSupplier);
Expand All @@ -122,7 +122,7 @@ private static ChannelBuilder create(SecurityProtocol securityProtocol,
DelegationTokenCache tokenCache,
Time time,
LogContext logContext,
Supplier<ApiVersionsResponse> apiVersionSupplier) {
Function<Short, ApiVersionsResponse> apiVersionSupplier) {
Map<String, Object> configs = channelBuilderConfigs(config, listenerName);

ChannelBuilder channelBuilder;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Function;
import java.util.function.Supplier;

import javax.security.auth.Subject;
Expand All @@ -88,7 +89,7 @@ public class SaslChannelBuilder implements ChannelBuilder, ListenerReconfigurabl
private final DelegationTokenCache tokenCache;
private final Map<String, LoginManager> loginManagers;
private final Map<String, Subject> subjects;
private final Supplier<ApiVersionsResponse> apiVersionSupplier;
private final Function<Short, ApiVersionsResponse> apiVersionSupplier;
private final String sslClientAuthOverride;
private final Map<String, AuthenticateCallbackHandler> saslCallbackHandlers;
private final Map<String, Long> connectionsMaxReauthMsByMechanism;
Expand All @@ -112,7 +113,7 @@ public SaslChannelBuilder(Mode mode,
String sslClientAuthOverride,
Time time,
LogContext logContext,
Supplier<ApiVersionsResponse> apiVersionSupplier) {
Function<Short, ApiVersionsResponse> apiVersionSupplier) {
this.mode = mode;
this.jaasContexts = jaasContexts;
this.loginManagers = new HashMap<>(jaasContexts.size());
Expand Down
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 @@ -82,7 +82,7 @@
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.function.Supplier;
import java.util.function.Function;

import javax.net.ssl.SSLSession;
import javax.security.auth.Subject;
Expand Down Expand Up @@ -129,7 +129,7 @@ private enum SaslState {
private final Time time;
private final ReauthInfo reauthInfo;
private final ChannelMetadataRegistry metadataRegistry;
private final Supplier<ApiVersionsResponse> apiVersionSupplier;
private final Function<Short, ApiVersionsResponse> apiVersionSupplier;

// Current SASL state
private SaslState saslState = SaslState.INITIAL_REQUEST;
Expand Down Expand Up @@ -159,7 +159,7 @@ public SaslServerAuthenticator(Map<String, ?> configs,
Map<String, Long> connectionsMaxReauthMsByMechanism,
ChannelMetadataRegistry metadataRegistry,
Time time,
Supplier<ApiVersionsResponse> apiVersionSupplier) {
Function<Short, ApiVersionsResponse> apiVersionSupplier) {
this.callbackHandlers = callbackHandlers;
this.connectionId = connectionId;
this.subjects = subjects;
Expand Down Expand Up @@ -596,7 +596,7 @@ else if (!apiVersionsRequest.isValid())
else {
metadataRegistry.registerClientInformation(new ClientInformation(apiVersionsRequest.data().clientSoftwareName(),
apiVersionsRequest.data().clientSoftwareVersion()));
sendKafkaResponse(context, apiVersionSupplier.get());
sendKafkaResponse(context, apiVersionSupplier.apply(apiVersionsRequest.version()));
setSaslState(SaslState.HANDSHAKE_REQUEST);
}
}
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 @@ -358,7 +358,7 @@ public void testUnsupportedApiVersionsRequestWithVersionProvidedByTheBroker() {
ByteBuffer buffer = selector.completedSendBuffers().get(0).buffer();
RequestHeader header = parseHeader(buffer);
assertEquals(ApiKeys.API_VERSIONS, header.apiKey());
assertEquals(3, header.apiVersion());
assertEquals(4, header.apiVersion());

// prepare response
ApiVersionCollection apiKeys = new ApiVersionCollection();
Expand Down Expand Up @@ -430,7 +430,7 @@ public void testUnsupportedApiVersionsRequestWithoutVersionProvidedByTheBroker()
ByteBuffer buffer = selector.completedSendBuffers().get(0).buffer();
RequestHeader header = parseHeader(buffer);
assertEquals(ApiKeys.API_VERSIONS, header.apiKey());
assertEquals(3, header.apiVersion());
assertEquals(4, header.apiVersion());

// prepare response
delayedApiVersionsResponse(0, (short) 0,
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
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,7 @@ public NioEchoServer(ListenerName listenerName, SecurityProtocol securityProtoco
if (channelBuilder == null)
channelBuilder = ChannelBuilders.serverChannelBuilder(listenerName, false,
securityProtocol, config, credentialCache, tokenCache, time, logContext,
() -> TestUtils.defaultApiVersionsResponse(ApiMessageType.ListenerType.ZK_BROKER));
version -> TestUtils.defaultApiVersionsResponse(ApiMessageType.ListenerType.ZK_BROKER));
this.metrics = new Metrics();
this.selector = new Selector(10000, failedAuthenticationDelayMs, metrics, time,
"MetricGroup", channelBuilder, logContext);
Expand Down
Loading