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;
+        }
+    }
 }

Reply via email to