From 74b37a25e48f31264e3def71c336d4b73954ef6a Mon Sep 17 00:00:00 2001 From: Ron Dagostino Date: Wed, 30 Sep 2020 11:41:49 -0400 Subject: [PATCH 1/4] KAFKA-10556: NPE if sasl.mechanism is unrecognized --- .../security/authenticator/SaslClientAuthenticator.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslClientAuthenticator.java b/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslClientAuthenticator.java index 6d279ac9f4f2b..bba1c43cc3f02 100644 --- a/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslClientAuthenticator.java +++ b/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslClientAuthenticator.java @@ -214,7 +214,11 @@ SaslClient createSaslClient() { String[] mechs = {mechanism}; log.debug("Creating SaslClient: client={};service={};serviceHostname={};mechs={}", clientPrincipalName, servicePrincipal, host, Arrays.toString(mechs)); - return Sasl.createSaslClient(mechs, clientPrincipalName, servicePrincipal, host, configs, callbackHandler); + SaslClient retvalSaslClient = Sasl.createSaslClient(mechs, clientPrincipalName, servicePrincipal, host, configs, callbackHandler); + if (retvalSaslClient == null) { + throw new SaslAuthenticationException("Failed to create SaslClient with mechanism " + mechanism); + } + return retvalSaslClient; }); } catch (PrivilegedActionException e) { throw new SaslAuthenticationException("Failed to create SaslClient with mechanism " + mechanism, e.getCause()); From 0523c0c5b25db33f4203ccfd3da1527141c35952 Mon Sep 17 00:00:00 2001 From: Ron Dagostino Date: Wed, 30 Sep 2020 12:15:41 -0400 Subject: [PATCH 2/4] Adjust test now that channel is never created --- .../authenticator/SaslAuthenticatorTest.java | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/common/security/authenticator/SaslAuthenticatorTest.java b/clients/src/test/java/org/apache/kafka/common/security/authenticator/SaslAuthenticatorTest.java index ed922b18f6f36..5c1ce3cfc0653 100644 --- a/clients/src/test/java/org/apache/kafka/common/security/authenticator/SaslAuthenticatorTest.java +++ b/clients/src/test/java/org/apache/kafka/common/security/authenticator/SaslAuthenticatorTest.java @@ -1236,9 +1236,18 @@ public void testInvalidMechanism() throws Exception { saslClientConfigs.put(SaslConfigs.SASL_MECHANISM, "INVALID"); server = createEchoServer(securityProtocol); - createAndCheckClientConnectionFailure(securityProtocol, node); - server.verifyAuthenticationMetrics(0, 1); - server.verifyReauthenticationMetrics(0, 0); + try { + createAndCheckClientConnectionFailure(securityProtocol, node); + fail("Did not generate exception prior to creating channel"); + } catch (IOException expected) { + server.verifyAuthenticationMetrics(0, 0); + server.verifyReauthenticationMetrics(0, 0); + Throwable underlyingCause = expected.getCause().getCause().getCause(); + assertEquals(SaslAuthenticationException.class, underlyingCause.getClass()); + assertEquals("Failed to create SaslClient with mechanism INVALID", underlyingCause.getMessage()); + } finally { + closeClientConnectionIfNecessary(); + } } /** From 7a94ebbf829a394e10986aabbfe8e28687a2d6c3 Mon Sep 17 00:00:00 2001 From: Ron Dagostino Date: Wed, 30 Sep 2020 12:42:01 -0400 Subject: [PATCH 3/4] Fix SaslServerAuthenticator as well --- .../security/authenticator/SaslServerAuthenticator.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java b/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java index b959d68896426..dce46102ebf08 100644 --- a/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java +++ b/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java @@ -190,11 +190,15 @@ private void createSaslServer(String mechanism) throws IOException { if (mechanism.equals(SaslConfigs.GSSAPI_MECHANISM)) { saslServer = createSaslKerberosServer(callbackHandler, configs, subject); } else { + String failureErrorMessage = "Kafka Server failed to create a SaslServer to interact with a client during session authentication"; try { saslServer = Subject.doAs(subject, (PrivilegedExceptionAction) () -> Sasl.createSaslServer(saslMechanism, "kafka", serverAddress().getHostName(), configs, callbackHandler)); + if (saslServer == null) { + throw new SaslException(failureErrorMessage); + } } catch (PrivilegedActionException e) { - throw new SaslException("Kafka Server failed to create a SaslServer to interact with a client during session authentication", e.getCause()); + throw new SaslException(failureErrorMessage, e.getCause()); } } } From 85bde668e623fc79df628a7b9f03f0d698069cab Mon Sep 17 00:00:00 2001 From: Ron Dagostino Date: Wed, 30 Sep 2020 12:45:28 -0400 Subject: [PATCH 4/4] Better error message in SaslServerAuthenticator --- .../security/authenticator/SaslServerAuthenticator.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java b/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java index dce46102ebf08..20dbf7b0d9d08 100644 --- a/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java +++ b/clients/src/main/java/org/apache/kafka/common/security/authenticator/SaslServerAuthenticator.java @@ -190,15 +190,14 @@ private void createSaslServer(String mechanism) throws IOException { if (mechanism.equals(SaslConfigs.GSSAPI_MECHANISM)) { saslServer = createSaslKerberosServer(callbackHandler, configs, subject); } else { - String failureErrorMessage = "Kafka Server failed to create a SaslServer to interact with a client during session authentication"; try { saslServer = Subject.doAs(subject, (PrivilegedExceptionAction) () -> Sasl.createSaslServer(saslMechanism, "kafka", serverAddress().getHostName(), configs, callbackHandler)); if (saslServer == null) { - throw new SaslException(failureErrorMessage); + throw new SaslException("Kafka Server failed to create a SaslServer to interact with a client during session authentication with server mechanism " + saslMechanism); } } catch (PrivilegedActionException e) { - throw new SaslException(failureErrorMessage, e.getCause()); + throw new SaslException("Kafka Server failed to create a SaslServer to interact with a client during session authentication with server mechanism " + saslMechanism, e.getCause()); } } }