This is an automated email from the ASF dual-hosted git repository.
omkreddy pushed a commit to branch 4.1
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.1 by this push:
new bbc04ea298a MINOR: Validate OAuthBearer server callback handler
configuration at startup
bbc04ea298a is described below
commit bbc04ea298a46b0197da6ade5d4f6ae5f6c1f8bb
Author: Evan Zhou <[email protected]>
AuthorDate: Wed Jun 10 16:42:51 2026 -0700
MINOR: Validate OAuthBearer server callback handler configuration at startup
The OAuthBearer server callback handler now validates its configuration at
startup and fails fast with a ConfigException when the validator is not
fully
configured, rather than continuing with an incomplete configuration.
Co-Authored-By: MarkLee131 <[email protected]>
---
.../OAuthBearerValidatorCallbackHandler.java | 17 ++-
.../OAuthBearerValidatorCallbackHandlerTest.java | 116 ++++++++++++++++++++-
2 files changed, 130 insertions(+), 3 deletions(-)
diff --git
a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandler.java
b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandler.java
index 60fa8cdb678..2d03fa469e1 100644
---
a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandler.java
+++
b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandler.java
@@ -17,6 +17,7 @@
package org.apache.kafka.common.security.oauthbearer;
+import org.apache.kafka.common.config.ConfigException;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.security.auth.AuthenticateCallbackHandler;
import
org.apache.kafka.common.security.oauthbearer.internals.secured.CloseableVerificationKeyResolver;
@@ -104,13 +105,27 @@ public class OAuthBearerValidatorCallbackHandler
implements AuthenticateCallback
@Override
public void configure(Map<String, ?> configs, String saslMechanism,
List<AppConfigurationEntry> jaasConfigEntries) {
- jwtValidator = getConfiguredInstance(
+ JwtValidator validator = getConfiguredInstance(
configs,
saslMechanism,
jaasConfigEntries,
SaslConfigs.SASL_OAUTHBEARER_JWT_VALIDATOR_CLASS,
JwtValidator.class
);
+
+ if (validator instanceof DefaultJwtValidator
+ && ((DefaultJwtValidator) validator).delegate() instanceof
ClientJwtValidator) {
+ Utils.closeQuietly(validator, "JWT validator");
+ throw new ConfigException(String.format(
+ "The OAuth validator for the broker requires \"%s\" to be
configured so that JWT signatures" +
+ " can be verified, but it was not set. Set \"%s\" to the
OAuth/OIDC provider's JWKS endpoint URL" +
+ " for this listener, or configure a signature-verifying
\"%s\".",
+ SaslConfigs.SASL_OAUTHBEARER_JWKS_ENDPOINT_URL,
+ SaslConfigs.SASL_OAUTHBEARER_JWKS_ENDPOINT_URL,
+ SaslConfigs.SASL_OAUTHBEARER_JWT_VALIDATOR_CLASS));
+ }
+
+ jwtValidator = validator;
}
/*
diff --git
a/clients/src/test/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandlerTest.java
b/clients/src/test/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandlerTest.java
index adabec6bc95..96ef6aecc6e 100644
---
a/clients/src/test/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandlerTest.java
+++
b/clients/src/test/java/org/apache/kafka/common/security/oauthbearer/OAuthBearerValidatorCallbackHandlerTest.java
@@ -18,15 +18,24 @@
package org.apache.kafka.common.security.oauthbearer;
import org.apache.kafka.common.KafkaException;
+import org.apache.kafka.common.config.ConfigException;
+import org.apache.kafka.common.config.SaslConfigs;
import
org.apache.kafka.common.security.oauthbearer.internals.secured.AccessTokenBuilder;
import
org.apache.kafka.common.security.oauthbearer.internals.secured.CloseableVerificationKeyResolver;
import
org.apache.kafka.common.security.oauthbearer.internals.secured.OAuthBearerTest;
+import org.apache.kafka.common.utils.Time;
+import org.apache.kafka.common.utils.Utils;
+import org.jose4j.jwk.JsonWebKey;
+import org.jose4j.jwk.JsonWebKeySet;
+import org.jose4j.jwk.PublicJsonWebKey;
import org.jose4j.jws.AlgorithmIdentifiers;
+import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.util.Arrays;
+import java.util.Base64;
import java.util.List;
import java.util.Map;
@@ -34,7 +43,9 @@ import javax.security.auth.callback.Callback;
import javax.security.auth.login.AppConfigurationEntry;
import static
org.apache.kafka.common.config.SaslConfigs.SASL_OAUTHBEARER_EXPECTED_AUDIENCE;
+import static
org.apache.kafka.common.config.internals.BrokerSecurityConfigs.ALLOWED_SASL_OAUTHBEARER_URLS_CONFIG;
import static
org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule.OAUTHBEARER_MECHANISM;
+import static org.apache.kafka.test.TestUtils.tempFile;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
@@ -44,6 +55,15 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
public class OAuthBearerValidatorCallbackHandlerTest extends OAuthBearerTest {
+ private static final String ATTACKER_PRINCIPAL = "kafka-admin";
+
+ private static final String EXPECTED_ISSUER = "https://idp.legit.example/";
+
+ @AfterEach
+ public void tearDown() {
+ System.clearProperty(ALLOWED_SASL_OAUTHBEARER_URLS_CONFIG);
+ }
+
@Test
public void testBasic() throws Exception {
String expectedAudience = "a";
@@ -51,10 +71,13 @@ public class OAuthBearerValidatorCallbackHandlerTest
extends OAuthBearerTest {
AccessTokenBuilder builder = new AccessTokenBuilder()
.audience(expectedAudience)
.jwk(createRsaJwk())
- .alg(AlgorithmIdentifiers.RSA_USING_SHA256);
+ .alg(AlgorithmIdentifiers.RSA_USING_SHA256)
+ .addCustomClaim("iss", EXPECTED_ISSUER);
String accessToken = builder.build();
- Map<String, ?> configs =
getSaslConfigs(SASL_OAUTHBEARER_EXPECTED_AUDIENCE, allAudiences);
+ Map<String, ?> configs = getSaslConfigs(Map.of(
+ SASL_OAUTHBEARER_EXPECTED_AUDIENCE, allAudiences,
+ SaslConfigs.SASL_OAUTHBEARER_EXPECTED_ISSUER, EXPECTED_ISSUER));
CloseableVerificationKeyResolver verificationKeyResolver =
createVerificationKeyResolver(builder);
JwtValidator jwtValidator =
createJwtValidator(verificationKeyResolver);
OAuthBearerValidatorCallbackHandler handler = new
OAuthBearerValidatorCallbackHandler();
@@ -157,6 +180,67 @@ public class OAuthBearerValidatorCallbackHandlerTest
extends OAuthBearerTest {
assertDoesNotThrow(handler::close);
}
+ @Test
+ public void testFailsFastWhenJwksUrlAbsent() {
+ // Default validator, no JWKS endpoint URL and no verification key
resolver: the broker
+ // validator handler must fail fast rather than silently using a
validator that does not
+ // verify token signatures.
+ Map<String, ?> configs = getSaslConfigs();
+
assertNull(configs.get(SaslConfigs.SASL_OAUTHBEARER_JWKS_ENDPOINT_URL));
+ assertEquals(SaslConfigs.DEFAULT_SASL_OAUTHBEARER_JWT_VALIDATOR_CLASS,
+ ((Class<?>)
configs.get(SaslConfigs.SASL_OAUTHBEARER_JWT_VALIDATOR_CLASS)).getName());
+
+ OAuthBearerValidatorCallbackHandler handler = new
OAuthBearerValidatorCallbackHandler();
+ ConfigException e = assertThrows(ConfigException.class,
+ () -> handler.configure(configs, OAUTHBEARER_MECHANISM,
getJaasConfigEntries()));
+
assertTrue(e.getMessage().contains(SaslConfigs.SASL_OAUTHBEARER_JWKS_ENDPOINT_URL),
+ "Expected the failure to reference " +
SaslConfigs.SASL_OAUTHBEARER_JWKS_ENDPOINT_URL
+ + ", but was: " + e.getMessage());
+ }
+
+ @Test
+ public void testRejectsForgedUnsignedTokenWhenJwksUrlSet() throws
Exception {
+ // Build a real JWKS file so BrokerJwtValidator can configure
successfully.
+ PublicJsonWebKey jwk = createRsaJwk();
+ JsonWebKeySet jwks = new JsonWebKeySet(jwk);
+ String jwksJson =
jwks.toJson(JsonWebKey.OutputControlLevel.PUBLIC_ONLY);
+ String fileUrl = tempFile(jwksJson).toURI().toString();
+ System.setProperty(ALLOWED_SASL_OAUTHBEARER_URLS_CONFIG, fileUrl);
+
+ Map<String, ?> configs = getSaslConfigs(Map.of(
+ SaslConfigs.SASL_OAUTHBEARER_JWKS_ENDPOINT_URL, fileUrl,
+ SaslConfigs.SASL_OAUTHBEARER_EXPECTED_ISSUER, EXPECTED_ISSUER));
+
+ OAuthBearerValidatorCallbackHandler handler = new
OAuthBearerValidatorCallbackHandler();
+ assertDoesNotThrow(() -> handler.configure(configs,
OAUTHBEARER_MECHANISM, getJaasConfigEntries()));
+
+ try {
+ String forgedToken = forgeUnsignedJwt(ATTACKER_PRINCIPAL);
+ OAuthBearerValidatorCallback callback = new
OAuthBearerValidatorCallback(forgedToken);
+ handler.handle(new Callback[]{callback});
+
+ assertNull(callback.token(), "BrokerJwtValidator must not accept
the forged unsigned token");
+ assertNotNull(callback.errorStatus(), "Expected an invalid_token
error");
+ assertEquals("invalid_token", callback.errorStatus());
+ } finally {
+ handler.close();
+ }
+ }
+
+ @Test
+ public void testConfigureAcceptsCustomValidatorClass() {
+ // A custom (non-default) validator class must be accepted; the
startup check only applies
+ // when the default validator is in use.
+ Map<String, ?> configs =
getSaslConfigs(SaslConfigs.SASL_OAUTHBEARER_JWT_VALIDATOR_CLASS,
+ NoOpVerifyingJwtValidator.class.getName());
+ OAuthBearerValidatorCallbackHandler handler = new
OAuthBearerValidatorCallbackHandler();
+ try {
+ assertDoesNotThrow(() -> handler.configure(configs,
OAUTHBEARER_MECHANISM, getJaasConfigEntries()));
+ } finally {
+ handler.close();
+ }
+ }
+
private void assertInvalidAccessTokenFails(String accessToken, String
expectedMessageSubstring) throws Exception {
AccessTokenBuilder builder = new AccessTokenBuilder()
.alg(AlgorithmIdentifiers.RSA_USING_SHA256);
@@ -193,4 +277,32 @@ public class OAuthBearerValidatorCallbackHandlerTest
extends OAuthBearerTest {
private CloseableVerificationKeyResolver
createVerificationKeyResolver(AccessTokenBuilder builder) {
return (jws, nestingContext) -> builder.jwk().getPublicKey();
}
+
+ private String forgeUnsignedJwt(String subject) {
+ long nowSeconds = Time.SYSTEM.milliseconds() / 1000;
+ Base64.Encoder enc = Base64.getUrlEncoder().withoutPadding();
+
+ String header =
enc.encodeToString(Utils.utf8("{\"alg\":\"none\",\"typ\":\"JWT\"}"));
+ String payload = enc.encodeToString(Utils.utf8(String.format(
+
"{\"sub\":\"%s\",\"scope\":\"engineering\",\"iat\":%d,\"exp\":%d}",
+ subject,
+ nowSeconds,
+ nowSeconds + 3600)));
+ String garbageSignature = enc.encodeToString(Utils.utf8("forged"));
+ return String.format("%s.%s.%s", header, payload, garbageSignature);
+ }
+
+ /**
+ * Stand-in for a custom operator-supplied {@link JwtValidator}, used to
verify that the
+ * startup configuration check does not apply when a non-default validator
class is set.
+ */
+ public static class NoOpVerifyingJwtValidator implements JwtValidator {
+ @Override
+ public void configure(Map<String, ?> configs, String saslMechanism,
List<AppConfigurationEntry> jaasConfigEntries) { }
+
+ @Override
+ public OAuthBearerToken validate(String accessToken) {
+ return null;
+ }
+ }
}