Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,8 @@ public enum SaslState {
private RequestHeader currentRequestHeader;
// Version of SaslAuthenticate request/responses
private short saslAuthenticateVersion;
// Version of SaslHandshake request/responses
private short saslHandshakeVersion;

public SaslClientAuthenticator(Map<String, ?> configs,
AuthenticateCallbackHandler callbackHandler,
Expand Down Expand Up @@ -213,13 +215,13 @@ public void authenticate() throws IOException {
if (apiVersionsResponse == null)
break;
else {
saslAuthenticateVersion(apiVersionsResponse);
setSaslAuthenticateAndHandshakeVersions(apiVersionsResponse);
reauthInfo.apiVersionsResponseReceivedFromBroker = apiVersionsResponse;
setSaslState(SaslState.SEND_HANDSHAKE_REQUEST);
// Fall through to send handshake request with the latest supported version
}
case SEND_HANDSHAKE_REQUEST:
sendHandshakeRequest(reauthInfo.apiVersionsResponseReceivedFromBroker);
sendHandshakeRequest(saslHandshakeVersion);
setSaslState(SaslState.RECEIVE_HANDSHAKE_RESPONSE);
break;
case RECEIVE_HANDSHAKE_RESPONSE:
Expand All @@ -236,11 +238,11 @@ public void authenticate() throws IOException {
setSaslState(SaslState.INTERMEDIATE);
break;
case REAUTH_PROCESS_ORIG_APIVERSIONS_RESPONSE:
saslAuthenticateVersion(reauthInfo.apiVersionsResponseFromOriginalAuthentication);
setSaslAuthenticateAndHandshakeVersions(reauthInfo.apiVersionsResponseFromOriginalAuthentication);
setSaslState(SaslState.REAUTH_SEND_HANDSHAKE_REQUEST); // Will set immediately
// Fall through to send handshake request with the latest supported version
case REAUTH_SEND_HANDSHAKE_REQUEST:
sendHandshakeRequest(reauthInfo.apiVersionsResponseFromOriginalAuthentication);
sendHandshakeRequest(saslHandshakeVersion);
setSaslState(SaslState.REAUTH_RECEIVE_HANDSHAKE_OR_OTHER_RESPONSE);
break;
case REAUTH_RECEIVE_HANDSHAKE_OR_OTHER_RESPONSE:
Expand Down Expand Up @@ -285,9 +287,8 @@ public void authenticate() throws IOException {
}
}

private void sendHandshakeRequest(ApiVersionsResponse apiVersionsResponse) throws IOException {
SaslHandshakeRequest handshakeRequest = createSaslHandshakeRequest(
apiVersionsResponse.apiVersion(ApiKeys.SASL_HANDSHAKE.id).maxVersion());
private void sendHandshakeRequest(short version) throws IOException {
SaslHandshakeRequest handshakeRequest = createSaslHandshakeRequest(version);
send(handshakeRequest.toSend(node, nextRequestHeader(ApiKeys.SASL_HANDSHAKE, handshakeRequest.version())));
}

Expand Down Expand Up @@ -345,11 +346,17 @@ protected SaslHandshakeRequest createSaslHandshakeRequest(short version) {
}

// Visible to override for testing
protected void saslAuthenticateVersion(ApiVersionsResponse apiVersionsResponse) {
protected void setSaslAuthenticateAndHandshakeVersions(ApiVersionsResponse apiVersionsResponse) {
ApiVersionsResponseKey authenticateVersion = apiVersionsResponse.apiVersion(ApiKeys.SASL_AUTHENTICATE.id);
if (authenticateVersion != null)
if (authenticateVersion != null) {
this.saslAuthenticateVersion = (short) Math.min(authenticateVersion.maxVersion(),
ApiKeys.SASL_AUTHENTICATE.latestVersion());
}
ApiVersionsResponseKey handshakeVersion = apiVersionsResponse.apiVersion(ApiKeys.SASL_HANDSHAKE.id);
if (handshakeVersion != null) {
this.saslHandshakeVersion = (short) Math.min(handshakeVersion.maxVersion(),
ApiKeys.SASL_HANDSHAKE.latestVersion());
}
}

private void setSaslState(SaslState saslState) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,10 @@
"type": "request",
"name": "SaslHandshakeRequest",
// Version 1 supports SASL_AUTHENTICATE.
// Version 2 adds flexible version support
"validVersions": "0-2",
"flexibleVersions": "2+",
// NOTE: Version cannot be easily bumped due to incorrect

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Rollback version with comment

// client negotiation for clients <= 2.4.
// See https://issues.apache.org/jira/browse/KAFKA-9577
"validVersions": "0-1",
"fields": [
{ "name": "Mechanism", "type": "string", "versions": "0+",
"about": "The SASL mechanism chosen by the client." }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,10 @@
"type": "response",
"name": "SaslHandshakeResponse",
// Version 1 is the same as version 0.
// Version 2 adds flexible version support
"validVersions": "0-2",
"flexibleVersions": "2+",
// NOTE: Version cannot be easily bumped due to incorrect
// client negotiation for clients <= 2.4.
// See https://issues.apache.org/jira/browse/KAFKA-9577
"validVersions": "0-1",
"fields": [
{ "name": "ErrorCode", "type": "int16", "versions": "0+",
"about": "The error code, or 0 if there was no error." },
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -677,7 +677,7 @@ public void testUnauthenticatedApiVersionsRequestOverSslHandshakeVersion1() thro
* {@link SaslServerAuthenticator} of the server prior to client authentication.
*/
@Test
public void testApiVersionsRequestWithUnsupportedVersion() throws Exception {
public void testApiVersionsRequestWithServerUnsupportedVersion() throws Exception {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Rename as this only tested part of the version support.

short handshakeVersion = ApiKeys.SASL_HANDSHAKE.latestVersion();
SecurityProtocol securityProtocol = SecurityProtocol.SASL_PLAINTEXT;
configureMechanisms("PLAIN", Arrays.asList("PLAIN"));
Expand Down Expand Up @@ -714,6 +714,25 @@ public void testApiVersionsRequestWithUnsupportedVersion() throws Exception {
authenticateUsingSaslPlainAndCheckConnection(node, handshakeVersion > 0);
}

/**
* Tests correct negotiation of handshake and authenticate api versions by having the server
* return a higher version than supported on the client.
* Note, that due to KAFKA-9577 this will require a workaround to effectively bump
* SASL_HANDSHAKE in the future.
*/
@Test
public void testSaslUnsupportedClientVersions() throws Exception {
configureMechanisms("SCRAM-SHA-512", Arrays.asList("SCRAM-SHA-512"));

server = startServerApiVersionsUnsupportedByClient(SecurityProtocol.SASL_SSL, "SCRAM-SHA-512");
updateScramCredentialCache(TestJaasConfig.USERNAME, TestJaasConfig.PASSWORD);

String node = "0";

createClientConnection(SecurityProtocol.SASL_SSL, "SCRAM-SHA-512", node, true);
NetworkTestUtils.checkClientConnection(selector, "0", 100, 10);
}

/**
* Tests that invalid ApiVersionRequest is handled by the server correctly and
* returns an INVALID_REQUEST error.
Expand Down Expand Up @@ -748,6 +767,16 @@ public void testInvalidApiVersionsRequest() throws Exception {
authenticateUsingSaslPlainAndCheckConnection(node, handshakeVersion > 0);
}


@Test
public void testForBrokenSaslHandshakeVersionBump() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Block anyone from bumping the schema without realizing.

assertEquals("It is not possible to easily bump SASL_HANDSHAKE schema"
+ " due to improper version negotiation in clients < 2.5."
+ " Please see https://issues.apache.org/jira/browse/KAFKA-9577",
ApiKeys.SASL_HANDSHAKE.latestVersion(),
1);
}

/**
* Tests that valid ApiVersionRequest is handled by the server correctly and
* returns an NONE error.
Expand Down Expand Up @@ -1730,6 +1759,47 @@ private void createClientConnection(SecurityProtocol securityProtocol, String sa
createClientConnectionWithoutSaslAuthenticateHeader(securityProtocol, saslMechanism, node);
}

private NioEchoServer startServerApiVersionsUnsupportedByClient(final SecurityProtocol securityProtocol, String saslMechanism) throws Exception {
final ListenerName listenerName = ListenerName.forSecurityProtocol(securityProtocol);
final Map<String, ?> configs = Collections.emptyMap();
final JaasContext jaasContext = JaasContext.loadServerContext(listenerName, saslMechanism, configs);
final Map<String, JaasContext> jaasContexts = Collections.singletonMap(saslMechanism, jaasContext);

boolean isScram = ScramMechanism.isScram(saslMechanism);
if (isScram)
ScramCredentialUtils.createCache(credentialCache, Arrays.asList(saslMechanism));
SaslChannelBuilder serverChannelBuilder = new SaslChannelBuilder(Mode.SERVER, jaasContexts,
securityProtocol, listenerName, false, saslMechanism, true,
credentialCache, null, time, new LogContext()) {

@Override
protected SaslServerAuthenticator buildServerAuthenticator(Map<String, ?> configs,
Map<String, AuthenticateCallbackHandler> callbackHandlers,
String id,
TransportLayer transportLayer,
Map<String, Subject> subjects,
Map<String, Long> connectionsMaxReauthMsByMechanism,
ChannelMetadataRegistry metadataRegistry) {
return new SaslServerAuthenticator(configs, callbackHandlers, id, subjects, null, listenerName,
securityProtocol, transportLayer, connectionsMaxReauthMsByMechanism, metadataRegistry, time) {

@Override
protected ApiVersionsResponse apiVersionsResponse() {
ApiVersionsResponseKeyCollection versionCollection = new ApiVersionsResponseKeyCollection(2);
versionCollection.add(new ApiVersionsResponseKey().setApiKey(ApiKeys.SASL_HANDSHAKE.id).setMinVersion((short) 0).setMaxVersion((short) 100));
versionCollection.add(new ApiVersionsResponseKey().setApiKey(ApiKeys.SASL_AUTHENTICATE.id).setMinVersion((short) 0).setMaxVersion((short) 100));
return new ApiVersionsResponse(new ApiVersionsResponseData().setApiKeys(versionCollection));
}
};
}
};
serverChannelBuilder.configure(saslServerConfigs);
server = new NioEchoServer(listenerName, securityProtocol, new TestSecurityConfig(saslServerConfigs),
"localhost", serverChannelBuilder, credentialCache, time);
server.start();
return server;
}

private NioEchoServer startServerWithoutSaslAuthenticateHeader(final SecurityProtocol securityProtocol, String saslMechanism)
throws Exception {
final ListenerName listenerName = ListenerName.forSecurityProtocol(securityProtocol);
Expand Down Expand Up @@ -1821,7 +1891,7 @@ protected SaslHandshakeRequest createSaslHandshakeRequest(short version) {
return buildSaslHandshakeRequest(saslMechanism, (short) 0);
}
@Override
protected void saslAuthenticateVersion(ApiVersionsResponse apiVersionsResponse) {
protected void setSaslAuthenticateAndHandshakeVersions(ApiVersionsResponse apiVersionsResponse) {
// Don't set version so that headers are disabled
}
};
Expand Down