From 77b049f847db25ffb19154e8943cd3f61c05f6c8 Mon Sep 17 00:00:00 2001 From: guangning Date: Sat, 29 Jan 2022 17:16:16 +0800 Subject: [PATCH 01/26] Support pass http auth status --- .../AuthenticationProvider.java | 8 +++++ .../AuthenticationProviderList.java | 33 ++++++++++++++++++- .../authentication/AuthenticationService.java | 5 ++- .../OneStageAuthenticationState.java | 8 +++++ .../broker/web/AuthenticationFilter.java | 12 +++++-- .../AuthenticationProviderListTest.java | 26 +++++++++++++++ .../impl/auth/AuthenticationDataBasic.java | 4 ++- .../impl/auth/AuthenticationDataTls.java | 10 ++++++ .../impl/auth/AuthenticationDataToken.java | 9 +++-- 9 files changed, 108 insertions(+), 7 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProvider.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProvider.java index c0716198a71af..5d8cdc51c9fa9 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProvider.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProvider.java @@ -104,6 +104,14 @@ default AuthenticationState newAuthState(AuthData authData, return new OneStageAuthenticationState(authData, remoteAddress, sslSession, this); } + /** + * Create an http authentication data State use passed in AuthenticationDataSource. + */ + default AuthenticationState newHttpAuthState(HttpServletRequest request) + throws AuthenticationException { + return new OneStageAuthenticationState(request, this); + } + /** * Validate the authentication for the given credentials with the specified authentication data. * diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java index 9ec1c2eb706cc..d87ad94581566 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java @@ -156,7 +156,9 @@ public String getAuthMethodName() { public String authenticate(AuthenticationDataSource authData) throws AuthenticationException { return applyAuthProcessor( providers, - provider -> provider.authenticate(authData) + provider -> { + return provider.authenticate(authData); + } ); } @@ -190,6 +192,35 @@ public AuthenticationState newAuthState(AuthData authData, SocketAddress remoteA } } + @Override + public AuthenticationState newHttpAuthState(HttpServletRequest request) throws AuthenticationException { + final List states = new ArrayList<>(providers.size()); + + AuthenticationException authenticationException = null; + try { + applyAuthProcessor( + providers, + provider -> { + AuthenticationState state = provider.newHttpAuthState(request); + states.add(state); + return state; + } + ); + } catch (AuthenticationException ae) { + authenticationException = ae; + } + if (states.isEmpty()) { + log.debug("Failed to initialize a new auth http state from {}", request.getRemoteHost(), authenticationException); + if (authenticationException != null) { + throw authenticationException; + } else { + throw new AuthenticationException("Failed to initialize a new http auth state from " + request.getRemoteHost()); + } + } else { + return new AuthenticationListState(states); + } + } + @Override public boolean authenticateHttpRequest(HttpServletRequest request, HttpServletResponse response) throws Exception { Boolean authenticated = applyAuthProcessor( diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java index 620dee3fb159a..51ef0b79e0447 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java @@ -87,7 +87,6 @@ public AuthenticationService(ServiceConfiguration conf) throws PulsarServerExcep public String authenticateHttpRequest(HttpServletRequest request) throws AuthenticationException { AuthenticationException authenticationException = null; - AuthenticationDataSource authData = new AuthenticationDataHttps(request); String authMethodName = request.getHeader("X-Pulsar-Auth-Method-Name"); if (authMethodName != null) { @@ -96,6 +95,8 @@ public String authenticateHttpRequest(HttpServletRequest request) throws Authent throw new AuthenticationException( String.format("Unsupported authentication method: [%s].", authMethodName)); } + AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); + AuthenticationDataSource authData = authenticationState.getAuthDataSource(); try { return providerToUse.authenticate(authData); } catch (AuthenticationException e) { @@ -109,6 +110,8 @@ public String authenticateHttpRequest(HttpServletRequest request) throws Authent } else { for (AuthenticationProvider provider : providers.values()) { try { + AuthenticationState authenticationState = provider.newHttpAuthState(request); + AuthenticationDataSource authData = authenticationState.getAuthDataSource(); return provider.authenticate(authData); } catch (AuthenticationException e) { if (LOG.isDebugEnabled()) { diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java index f2667c3b46c52..23c821aaf1476 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java @@ -22,6 +22,8 @@ import java.net.SocketAddress; import javax.naming.AuthenticationException; import javax.net.ssl.SSLSession; +import javax.servlet.http.HttpServlet; +import javax.servlet.http.HttpServletRequest; import org.apache.pulsar.common.api.AuthData; import static java.nio.charset.StandardCharsets.UTF_8; @@ -46,6 +48,12 @@ public OneStageAuthenticationState(AuthData authData, this.authRole = provider.authenticate(authenticationDataSource); } + public OneStageAuthenticationState(HttpServletRequest request, AuthenticationProvider provider) + throws AuthenticationException { + this.authenticationDataSource = new AuthenticationDataHttps(request); + this.authRole = provider.authenticate(authenticationDataSource); + } + @Override public String getAuthRole() { return authRole; diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java index 5d3cae49394d5..fadfafaef7e8a 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java @@ -33,6 +33,7 @@ import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationService; +import org.apache.pulsar.broker.authentication.AuthenticationState; import org.apache.pulsar.common.sasl.SaslConstants; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -76,8 +77,15 @@ public void doFilter(ServletRequest request, ServletResponse response, FilterCha // not sasl type, return role directly. String role = authenticationService.authenticateHttpRequest((HttpServletRequest) request); request.setAttribute(AuthenticatedRoleAttributeName, role); - request.setAttribute(AuthenticatedDataAttributeName, - new AuthenticationDataHttps((HttpServletRequest) request)); + String authMethodName = httpRequest.getHeader("X-Pulsar-Auth-Method-Name"); + if (authMethodName != null && authenticationService.getAuthenticationProvider(authMethodName) != null) { + AuthenticationState authenticationState = authenticationService + .getAuthenticationProvider(authMethodName).newHttpAuthState(httpRequest); + request.setAttribute(AuthenticatedDataAttributeName, authenticationState.getAuthDataSource()); + } else { + request.setAttribute(AuthenticatedDataAttributeName, + new AuthenticationDataHttps((HttpServletRequest) request)); + } if (LOG.isDebugEnabled()) { LOG.debug("[{}] Authenticated HTTP request with role {}", request.getRemoteAddr(), role); } diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java index b4bf974ff4f54..365b324df76f0 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java @@ -19,6 +19,8 @@ package org.apache.pulsar.broker.authentication; import static java.nio.charset.StandardCharsets.UTF_8; +import javax.servlet.http.HttpServletRequest; +import static org.mockito.Mockito.mock; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; @@ -165,6 +167,14 @@ private AuthenticationState newAuthState(String token, String expectedSubject) t return authState; } + private AuthenticationState newHttpAuthState(HttpServletRequest request, String expectedSubject) throws Exception { + AuthenticationState authState = authProvider.newHttpAuthState(request); + assertEquals(authState.getAuthRole(), expectedSubject); + assertTrue(authState.isComplete()); + assertFalse(authState.isExpired()); + return authState; + } + private void verifyAuthStateExpired(AuthenticationState authState, String expectedSubject) throws Exception { assertEquals(authState.getAuthRole(), expectedSubject); @@ -188,4 +198,20 @@ public void testNewAuthState() throws Exception { } + @Test + public void testNewHttpAuthState() throws Exception { + HttpServletRequest request = mock(HttpServletRequest.class); + AuthenticationState authStateAA = newHttpAuthState(request, SUBJECT_A); + AuthenticationState authStateAB = newHttpAuthState(request, SUBJECT_B); + AuthenticationState authStateBA = newHttpAuthState(request, SUBJECT_A); + AuthenticationState authStateBB = newHttpAuthState(request, SUBJECT_B); + + Thread.sleep(TimeUnit.SECONDS.toMillis(6)); + + verifyAuthStateExpired(authStateAA, SUBJECT_A); + verifyAuthStateExpired(authStateAB, SUBJECT_B); + verifyAuthStateExpired(authStateBA, SUBJECT_A); + verifyAuthStateExpired(authStateBB, SUBJECT_B); + } + } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index 9ba928c968034..a31382b9da2bc 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -29,12 +29,13 @@ public class AuthenticationDataBasic implements AuthenticationDataProvider { private static final String HTTP_HEADER_NAME = "Authorization"; + public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private String httpAuthToken; private String commandAuthToken; public AuthenticationDataBasic(String userId, String password) { httpAuthToken = "Basic " + Base64.getEncoder().encodeToString((userId + ":" + password).getBytes()); - commandAuthToken = userId+":"+password; + commandAuthToken = userId + ":" + password; } @Override @@ -46,6 +47,7 @@ public boolean hasDataForHttp() { public Set> getHttpHeaders() { Map headers = new HashMap<>(); headers.put(HTTP_HEADER_NAME, httpAuthToken); + headers.put(PULSAR_AUTH_METHOD_NAME, "basic"); return headers.entrySet(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java index 28b520c52246f..28d6854e49217 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java @@ -24,6 +24,9 @@ import java.security.PrivateKey; import java.security.cert.Certificate; import java.security.cert.X509Certificate; +import java.util.Collections; +import java.util.Map; +import java.util.Set; import java.util.function.Supplier; import org.apache.pulsar.client.api.AuthenticationDataProvider; @@ -35,6 +38,8 @@ import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; public class AuthenticationDataTls implements AuthenticationDataProvider { + public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + private static final long serialVersionUID = 1L; protected X509Certificate[] tlsCertificates; protected PrivateKey tlsPrivateKey; @@ -88,6 +93,11 @@ public boolean hasDataForTls() { return true; } + @Override + public Set> getHttpHeaders() { + return Collections.singletonMap(PULSAR_AUTH_METHOD_NAME, "tls").entrySet(); + } + @Override public Certificate[] getTlsCertificates() { if (certFile != null && certFile.checkAndRefresh()) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index 39436495253a8..b078c58703c75 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -19,7 +19,7 @@ package org.apache.pulsar.client.impl.auth; -import java.util.Collections; +import java.util.HashMap; import java.util.Map; import java.util.Set; import java.util.function.Supplier; @@ -29,6 +29,8 @@ public class AuthenticationDataToken implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; + public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + private final Supplier tokenSupplier; public AuthenticationDataToken(Supplier tokenSupplier) { @@ -42,7 +44,10 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return Collections.singletonMap(HTTP_HEADER_NAME, "Bearer " + getToken()).entrySet(); + Map headers = new HashMap<>(); + headers.put(PULSAR_AUTH_METHOD_NAME, "token"); + headers.put(HTTP_HEADER_NAME, "Bearer " + getToken()); + return headers.entrySet(); } @Override From ce3c5d3486837307f481f094ef341b86dccee8c0 Mon Sep 17 00:00:00 2001 From: guangning Date: Sat, 29 Jan 2022 17:43:24 +0800 Subject: [PATCH 02/26] Fixed comment and style --- .../broker/authentication/AuthenticationProviderList.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java index d87ad94581566..8c0cc7d34a39e 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java @@ -156,9 +156,7 @@ public String getAuthMethodName() { public String authenticate(AuthenticationDataSource authData) throws AuthenticationException { return applyAuthProcessor( providers, - provider -> { - return provider.authenticate(authData); - } + provider -> provider.authenticate(authData) ); } @@ -210,7 +208,7 @@ public AuthenticationState newHttpAuthState(HttpServletRequest request) throws A authenticationException = ae; } if (states.isEmpty()) { - log.debug("Failed to initialize a new auth http state from {}", request.getRemoteHost(), authenticationException); + log.debug("Failed to initialize a new http auth state from {}", request.getRemoteHost(), authenticationException); if (authenticationException != null) { throw authenticationException; } else { From 1badfaa1b50bb89929cf654203e6e4ed35c44293 Mon Sep 17 00:00:00 2001 From: guangning Date: Sat, 29 Jan 2022 17:49:38 +0800 Subject: [PATCH 03/26] Add method name for oauth2 auth --- .../impl/auth/oauth2/AuthenticationDataOAuth2.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java index 59810f50a62db..8e704aa97e82f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java @@ -18,7 +18,7 @@ */ package org.apache.pulsar.client.impl.auth.oauth2; -import java.util.Collections; +import java.util.HashMap; import java.util.Map; import java.util.Set; import org.apache.pulsar.client.api.AuthenticationDataProvider; @@ -28,13 +28,12 @@ */ class AuthenticationDataOAuth2 implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; + public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private final String accessToken; - private final Set> headers; public AuthenticationDataOAuth2(String accessToken) { this.accessToken = accessToken; - this.headers = Collections.singletonMap(HTTP_HEADER_NAME, "Bearer " + accessToken).entrySet(); } @Override @@ -44,7 +43,10 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return this.headers; + Map headers = new HashMap<>(); + headers.put(HTTP_HEADER_NAME, "Bearer " + accessToken); + headers.put(PULSAR_AUTH_METHOD_NAME, "token"); + return headers.entrySet(); } @Override From cbc79e235c3b572af2227ed1b173649b06e2dcad Mon Sep 17 00:00:00 2001 From: guangning Date: Sat, 29 Jan 2022 17:54:06 +0800 Subject: [PATCH 04/26] Change var type to private from public --- .../apache/pulsar/client/impl/auth/AuthenticationDataBasic.java | 2 +- .../apache/pulsar/client/impl/auth/AuthenticationDataTls.java | 2 +- .../apache/pulsar/client/impl/auth/AuthenticationDataToken.java | 2 +- .../client/impl/auth/oauth2/AuthenticationDataOAuth2.java | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index a31382b9da2bc..cbd6526ce0960 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -29,7 +29,7 @@ public class AuthenticationDataBasic implements AuthenticationDataProvider { private static final String HTTP_HEADER_NAME = "Authorization"; - public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private String httpAuthToken; private String commandAuthToken; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java index 28d6854e49217..70ee0787057c9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java @@ -38,7 +38,7 @@ import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; public class AuthenticationDataTls implements AuthenticationDataProvider { - public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private static final long serialVersionUID = 1L; protected X509Certificate[] tlsCertificates; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index b078c58703c75..b639d86f4588f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -29,7 +29,7 @@ public class AuthenticationDataToken implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; - public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private final Supplier tokenSupplier; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java index 8e704aa97e82f..a3d4e717c8c11 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java @@ -28,7 +28,7 @@ */ class AuthenticationDataOAuth2 implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; - public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private final String accessToken; From 8924aa67be7e1b46cdd9b809022680faa516947c Mon Sep 17 00:00:00 2001 From: guangning Date: Sat, 29 Jan 2022 18:00:12 +0800 Subject: [PATCH 05/26] Fixed style check --- .../broker/authentication/AuthenticationProviderList.java | 6 ++++-- .../broker/authentication/OneStageAuthenticationState.java | 1 - 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java index 8c0cc7d34a39e..3aff9be936194 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderList.java @@ -208,11 +208,13 @@ public AuthenticationState newHttpAuthState(HttpServletRequest request) throws A authenticationException = ae; } if (states.isEmpty()) { - log.debug("Failed to initialize a new http auth state from {}", request.getRemoteHost(), authenticationException); + log.debug("Failed to initialize a new http auth state from {}", + request.getRemoteHost(), authenticationException); if (authenticationException != null) { throw authenticationException; } else { - throw new AuthenticationException("Failed to initialize a new http auth state from " + request.getRemoteHost()); + throw new AuthenticationException( + "Failed to initialize a new http auth state from " + request.getRemoteHost()); } } else { return new AuthenticationListState(states); diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java index 23c821aaf1476..c1bc8d76f1efd 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/OneStageAuthenticationState.java @@ -22,7 +22,6 @@ import java.net.SocketAddress; import javax.naming.AuthenticationException; import javax.net.ssl.SSLSession; -import javax.servlet.http.HttpServlet; import javax.servlet.http.HttpServletRequest; import org.apache.pulsar.common.api.AuthData; From 3cc0a9dd6eae57b05dbf571e130dfdf62e655176 Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 22:09:10 +0800 Subject: [PATCH 06/26] Fixed function interface and client api var --- .../authentication/AuthenticationService.java | 18 ++++--- .../broker/web/AuthenticationFilter.java | 17 +++++-- .../pulsar/broker/web/PulsarWebResource.java | 7 ++- .../api/AuthenticationDataProvider.java | 2 + .../impl/auth/AuthenticationDataBasic.java | 1 - .../impl/auth/AuthenticationDataTls.java | 2 - .../impl/auth/AuthenticationDataToken.java | 3 -- .../auth/oauth2/AuthenticationDataOAuth2.java | 3 +- .../worker/rest/FunctionApiResource.java | 6 +-- .../worker/rest/api/ComponentImpl.java | 7 ++- .../worker/rest/api/FunctionsImpl.java | 5 +- .../functions/worker/rest/api/SinksImpl.java | 5 +- .../worker/rest/api/SourcesImpl.java | 5 +- .../worker/service/api/Component.java | 41 +++++++++++++-- .../worker/service/api/Functions.java | 50 ++++++++++++++++++- .../functions/worker/service/api/Sinks.java | 50 ++++++++++++++++++- .../functions/worker/service/api/Sources.java | 50 ++++++++++++++++++- .../websocket/admin/WebSocketWebResource.java | 11 ++-- 18 files changed, 232 insertions(+), 51 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java index 51ef0b79e0447..1ad09855cc507 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java @@ -32,6 +32,7 @@ import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.web.AuthenticationFilter; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -85,9 +86,9 @@ public AuthenticationService(ServiceConfiguration conf) throws PulsarServerExcep } } - public String authenticateHttpRequest(HttpServletRequest request) throws AuthenticationException { + public String authenticateHttpRequest(HttpServletRequest request, AuthenticationDataSource authData) throws AuthenticationException { AuthenticationException authenticationException = null; - String authMethodName = request.getHeader("X-Pulsar-Auth-Method-Name"); + String authMethodName = request.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); if (authMethodName != null) { AuthenticationProvider providerToUse = providers.get(authMethodName); @@ -95,8 +96,10 @@ public String authenticateHttpRequest(HttpServletRequest request) throws Authent throw new AuthenticationException( String.format("Unsupported authentication method: [%s].", authMethodName)); } - AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); - AuthenticationDataSource authData = authenticationState.getAuthDataSource(); + if (authData == null) { + AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); + authData = authenticationState.getAuthDataSource(); + } try { return providerToUse.authenticate(authData); } catch (AuthenticationException e) { @@ -111,8 +114,7 @@ public String authenticateHttpRequest(HttpServletRequest request) throws Authent for (AuthenticationProvider provider : providers.values()) { try { AuthenticationState authenticationState = provider.newHttpAuthState(request); - AuthenticationDataSource authData = authenticationState.getAuthDataSource(); - return provider.authenticate(authData); + return provider.authenticate(authenticationState.getAuthDataSource()); } catch (AuthenticationException e) { if (LOG.isDebugEnabled()) { LOG.debug("Authentication failed for provider " + provider.getAuthMethodName() + ": " + e.getMessage(), e); @@ -139,6 +141,10 @@ public String authenticateHttpRequest(HttpServletRequest request) throws Authent } } + public String authenticateHttpRequest(HttpServletRequest request) throws AuthenticationException { + return authenticateHttpRequest(request, null); + } + public AuthenticationProvider getAuthenticationProvider(String authMethodName) { return providers.get(authMethodName); } diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java index fadfafaef7e8a..b9ebd5bcda980 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java @@ -48,6 +48,8 @@ public class AuthenticationFilter implements Filter { public static final String AuthenticatedRoleAttributeName = AuthenticationFilter.class.getName() + "-role"; public static final String AuthenticatedDataAttributeName = AuthenticationFilter.class.getName() + "-data"; + public static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; + public AuthenticationFilter(AuthenticationService authenticationService) { this.authenticationService = authenticationService; @@ -75,17 +77,24 @@ public void doFilter(ServletRequest request, ServletResponse response, FilterCha if (!isSaslRequest(httpRequest)) { // not sasl type, return role directly. - String role = authenticationService.authenticateHttpRequest((HttpServletRequest) request); - request.setAttribute(AuthenticatedRoleAttributeName, role); - String authMethodName = httpRequest.getHeader("X-Pulsar-Auth-Method-Name"); + String authMethodName = httpRequest.getHeader(PULSAR_AUTH_METHOD_NAME); + AuthenticationState authenticationState = null; if (authMethodName != null && authenticationService.getAuthenticationProvider(authMethodName) != null) { - AuthenticationState authenticationState = authenticationService + authenticationState = authenticationService .getAuthenticationProvider(authMethodName).newHttpAuthState(httpRequest); request.setAttribute(AuthenticatedDataAttributeName, authenticationState.getAuthDataSource()); } else { request.setAttribute(AuthenticatedDataAttributeName, new AuthenticationDataHttps((HttpServletRequest) request)); } + String role; + if (authenticationState != null ) { + role = authenticationService.authenticateHttpRequest((HttpServletRequest) request, authenticationState.getAuthDataSource()); + } else { + role = authenticationService.authenticateHttpRequest((HttpServletRequest) request); + } + request.setAttribute(AuthenticatedRoleAttributeName, role); + if (LOG.isDebugEnabled()) { LOG.debug("[{}] Authenticated HTTP request with role {}", request.getRemoteAddr(), role); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java index e3ef2b3852838..9c917998b06b7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java @@ -48,7 +48,6 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.broker.authorization.AuthorizationService; import org.apache.pulsar.broker.namespace.LookupOptions; @@ -135,8 +134,8 @@ public String originalPrincipal() { return httpRequest.getHeader(ORIGINAL_PRINCIPAL_HEADER); } - public AuthenticationDataHttps clientAuthData() { - return (AuthenticationDataHttps) httpRequest.getAttribute(AuthenticationFilter.AuthenticatedDataAttributeName); + public AuthenticationDataSource clientAuthData() { + return (AuthenticationDataSource) httpRequest.getAttribute(AuthenticationFilter.AuthenticatedDataAttributeName); } public boolean isRequestHttps() { @@ -1051,7 +1050,7 @@ && pulsar().getBrokerService().isAuthorizationEnabled()) { throw new RestException(Status.UNAUTHORIZED, "Need to authenticate to perform the request"); } - AuthenticationDataHttps authData = clientAuthData(); + AuthenticationDataSource authData = clientAuthData(); authData.setSubscription(subscription); Boolean isAuthorized = pulsar().getBrokerService().getAuthorizationService() diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationDataProvider.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationDataProvider.java index 4624821f5e0c0..c0155fe66c05f 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationDataProvider.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationDataProvider.java @@ -36,6 +36,8 @@ @InterfaceAudience.LimitedPrivate @InterfaceStability.Stable public interface AuthenticationDataProvider extends Serializable { + + String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; /* * TLS */ diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index cbd6526ce0960..fe4a5f3aa46a9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -29,7 +29,6 @@ public class AuthenticationDataBasic implements AuthenticationDataProvider { private static final String HTTP_HEADER_NAME = "Authorization"; - private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private String httpAuthToken; private String commandAuthToken; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java index 70ee0787057c9..30b96cfc5ff52 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java @@ -38,8 +38,6 @@ import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; public class AuthenticationDataTls implements AuthenticationDataProvider { - private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; - private static final long serialVersionUID = 1L; protected X509Certificate[] tlsCertificates; protected PrivateKey tlsPrivateKey; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index b639d86f4588f..f97a251a77dce 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -28,9 +28,6 @@ public class AuthenticationDataToken implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; - - private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; - private final Supplier tokenSupplier; public AuthenticationDataToken(Supplier tokenSupplier) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java index a3d4e717c8c11..227443f0dc248 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java @@ -28,7 +28,6 @@ */ class AuthenticationDataOAuth2 implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; - private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; private final String accessToken; @@ -45,7 +44,7 @@ public boolean hasDataForHttp() { public Set> getHttpHeaders() { Map headers = new HashMap<>(); headers.put(HTTP_HEADER_NAME, "Bearer " + accessToken); - headers.put(PULSAR_AUTH_METHOD_NAME, "token"); + headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationOAuth2.AUTH_METHOD_NAME); return headers.entrySet(); } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/FunctionApiResource.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/FunctionApiResource.java index 18ac5b201368c..ce08c2e8c0cd7 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/FunctionApiResource.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/FunctionApiResource.java @@ -24,7 +24,7 @@ import javax.ws.rs.core.Context; import javax.ws.rs.core.UriInfo; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; +import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.broker.web.AuthenticationFilter; import org.apache.pulsar.functions.worker.WorkerService; @@ -54,7 +54,7 @@ public String clientAppId() { : null; } - public AuthenticationDataHttps clientAuthData() { - return (AuthenticationDataHttps) httpRequest.getAttribute(AuthenticationFilter.AuthenticatedDataAttributeName); + public AuthenticationDataSource clientAuthData() { + return (AuthenticationDataSource) httpRequest.getAttribute(AuthenticationFilter.AuthenticatedDataAttributeName); } } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java index af5d5e333a9f6..5511199df735c 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java @@ -45,7 +45,6 @@ import org.apache.bookkeeper.clients.exceptions.NamespaceNotFoundException; import org.apache.bookkeeper.clients.exceptions.StreamNotFoundException; import org.apache.commons.io.IOUtils; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.admin.internal.FunctionsImpl; @@ -372,7 +371,7 @@ public void deregisterFunction(final String tenant, final String namespace, final String componentName, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps) { + AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); @@ -1244,7 +1243,7 @@ public void uploadFunction(final InputStream uploadedInputStream, final String p @Override public StreamingOutput downloadFunction(String tenant, String namespace, String componentName, - String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps) { + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); } @@ -1305,7 +1304,7 @@ private StreamingOutput getStreamingOutput(String pkgPath) { } @Override - public StreamingOutput downloadFunction(final String path, String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps) { + public StreamingOutput downloadFunction(final String path, String clientRole, AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java index b00846de3f027..ba85cd5767b09 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java @@ -42,7 +42,6 @@ import javax.ws.rs.core.UriBuilder; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.functions.FunctionConfig; @@ -81,7 +80,7 @@ public void registerFunction(final String tenant, final String functionPkgUrl, final FunctionConfig functionConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps) { + AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); @@ -252,7 +251,7 @@ public void updateFunction(final String tenant, final String functionPkgUrl, final FunctionConfig functionConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps, + AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions) { if (!isWorkerServiceAvailable()) { diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java index 6d867b37d9307..c0c0effab94a2 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java @@ -39,7 +39,6 @@ import javax.ws.rs.core.Response; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.functions.UpdateOptionsImpl; @@ -81,7 +80,7 @@ public void registerSink(final String tenant, final String sinkPkgUrl, final SinkConfig sinkConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps) { + AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); @@ -251,7 +250,7 @@ public void updateSink(final String tenant, final String sinkPkgUrl, final SinkConfig sinkConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps, + AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions) { if (!isWorkerServiceAvailable()) { diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java index e3946ac3642ae..0db5871e4f82f 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java @@ -39,7 +39,6 @@ import javax.ws.rs.core.Response; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.functions.UpdateOptionsImpl; @@ -81,7 +80,7 @@ public void registerSource(final String tenant, final String sourcePkgUrl, final SourceConfig sourceConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps) { + AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); @@ -251,7 +250,7 @@ public void updateSource(final String tenant, final String sourcePkgUrl, final SourceConfig sourceConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps, + AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions) { if (!isWorkerServiceAvailable()) { diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java index 5f0c1a9670396..51fa693b042d6 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java @@ -43,7 +43,21 @@ void deregisterFunction(final String tenant, final String namespace, final String componentName, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps); + AuthenticationDataSource clientAuthenticationDataHttps); + + @Deprecated + default void deregisterFunction(final String tenant, + final String namespace, + final String componentName, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps) { + deregisterFunction( + tenant, + namespace, + componentName, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps); + } FunctionConfig getFunctionInfo(final String tenant, final String namespace, @@ -144,13 +158,34 @@ void uploadFunction(final InputStream uploadedInputStream, StreamingOutput downloadFunction(String path, String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps); + final AuthenticationDataSource clientAuthenticationDataHttps); + + @Deprecated + default StreamingOutput downloadFunction(String path, + String clientRole, + final AuthenticationDataHttps clientAuthenticationDataHttps) { + return downloadFunction(path, clientRole, (AuthenticationDataSource) clientAuthenticationDataHttps); + } StreamingOutput downloadFunction(String tenant, String namespace, String componentName, String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps); + final AuthenticationDataSource clientAuthenticationDataHttps); + + @Deprecated + default StreamingOutput downloadFunction(String tenant, + String namespace, + String componentName, + String clientRole, + final AuthenticationDataHttps clientAuthenticationDataHttps) { + return downloadFunction( + tenant, + namespace, + componentName, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps); + } List getListOfConnectors(); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java index 7a155d51fe2f3..897811ec700de 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java @@ -42,7 +42,29 @@ void registerFunction(final String tenant, final String functionPkgUrl, final FunctionConfig functionConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps); + AuthenticationDataSource clientAuthenticationDataHttps); + + @Deprecated + default void registerFunction(final String tenant, + final String namespace, + final String functionName, + final InputStream uploadedInputStream, + final FormDataContentDisposition fileDetail, + final String functionPkgUrl, + final FunctionConfig functionConfig, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps) { + registerFunction( + tenant, + namespace, + functionName, + uploadedInputStream, + fileDetail, + functionPkgUrl, + functionConfig, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps); + } void updateFunction(final String tenant, final String namespace, @@ -52,9 +74,33 @@ void updateFunction(final String tenant, final String functionPkgUrl, final FunctionConfig functionConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps, + AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); + @Deprecated + default void updateFunction(final String tenant, + final String namespace, + final String functionName, + final InputStream uploadedInputStream, + final FormDataContentDisposition fileDetail, + final String functionPkgUrl, + final FunctionConfig functionConfig, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps, + UpdateOptionsImpl updateOptions) { + updateFunction( + tenant, + namespace, + functionName, + uploadedInputStream, + fileDetail, + functionPkgUrl, + functionConfig, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps, + updateOptions); + } + void updateFunctionOnWorkerLeader(final String tenant, final String namespace, final String functionName, diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java index fc6bdf7bf6d32..cdcf437a1437d 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java @@ -45,7 +45,29 @@ void registerSink(final String tenant, final String sinkPkgUrl, final SinkConfig sinkConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps); + AuthenticationDataSource clientAuthenticationDataHttps); + + @Deprecated + default void registerSink(final String tenant, + final String namespace, + final String sinkName, + final InputStream uploadedInputStream, + final FormDataContentDisposition fileDetail, + final String sinkPkgUrl, + final SinkConfig sinkConfig, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps) { + registerSink( + tenant, + namespace, + sinkName, + uploadedInputStream, + fileDetail, + sinkPkgUrl, + sinkConfig, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps); + } void updateSink(final String tenant, final String namespace, @@ -55,9 +77,33 @@ void updateSink(final String tenant, final String sinkPkgUrl, final SinkConfig sinkConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps, + AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); + @Deprecated + default void updateSink(final String tenant, + final String namespace, + final String sinkName, + final InputStream uploadedInputStream, + final FormDataContentDisposition fileDetail, + final String sinkPkgUrl, + final SinkConfig sinkConfig, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps, + UpdateOptionsImpl updateOptions) { + updateSink( + tenant, + namespace, + sinkName, + uploadedInputStream, + fileDetail, + sinkPkgUrl, + sinkConfig, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps, + updateOptions); + } + SinkInstanceStatusData getSinkInstanceStatus(final String tenant, final String namespace, final String sinkName, diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java index cc372490f4bb2..ab57233aef013 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java @@ -45,7 +45,29 @@ void registerSource(final String tenant, final String sourcePkgUrl, final SourceConfig sourceConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps); + AuthenticationDataSource clientAuthenticationDataHttps); + + @Deprecated + default void registerSource(final String tenant, + final String namespace, + final String sourceName, + final InputStream uploadedInputStream, + final FormDataContentDisposition fileDetail, + final String sourcePkgUrl, + final SourceConfig sourceConfig, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps) { + registerSource( + tenant, + namespace, + sourceName, + uploadedInputStream, + fileDetail, + sourcePkgUrl, + sourceConfig, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps); + } void updateSource(final String tenant, final String namespace, @@ -55,9 +77,33 @@ void updateSource(final String tenant, final String sourcePkgUrl, final SourceConfig sourceConfig, final String clientRole, - AuthenticationDataHttps clientAuthenticationDataHttps, + AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); + @Deprecated + default void updateSource(final String tenant, + final String namespace, + final String sourceName, + final InputStream uploadedInputStream, + final FormDataContentDisposition fileDetail, + final String sourcePkgUrl, + final SourceConfig sourceConfig, + final String clientRole, + AuthenticationDataHttps clientAuthenticationDataHttps, + UpdateOptionsImpl updateOptions) { + updateSource( + tenant, + namespace, + sourceName, + uploadedInputStream, + fileDetail, + sourcePkgUrl, + sourceConfig, + clientRole, + (AuthenticationDataSource) clientAuthenticationDataHttps, + updateOptions); + } + SourceStatus getSourceStatus(final String tenant, final String namespace, diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java index 9e7388971e4ca..c58e67f84de91 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java @@ -26,7 +26,8 @@ import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.UriInfo; -import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; +import org.apache.pulsar.broker.authentication.AuthenticationDataSource; +import org.apache.pulsar.broker.web.AuthenticationFilter; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.RestException; import org.apache.pulsar.websocket.WebSocketService; @@ -51,7 +52,7 @@ public class WebSocketWebResource { private WebSocketService socketService; private String clientId; - private AuthenticationDataHttps authData; + private AuthenticationDataSource authData; protected WebSocketService service() { if (socketService == null) { @@ -82,9 +83,11 @@ public String clientAppId() { return clientId; } - public AuthenticationDataHttps authData() { + public AuthenticationDataSource authData() throws AuthenticationException { if (authData == null) { - authData = new AuthenticationDataHttps(httpRequest); + String authMethodName = httpRequest.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); + authData = service().getAuthenticationService().getAuthenticationProvider(authMethodName) + .newHttpAuthState(httpRequest).getAuthDataSource(); } return authData; } From 3407283538fe3e6fbd2882190e188c187fe22d1f Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 22:21:47 +0800 Subject: [PATCH 07/26] Update auth method name --- .../apache/pulsar/client/impl/auth/AuthenticationBasic.java | 3 ++- .../pulsar/client/impl/auth/AuthenticationDataBasic.java | 2 +- .../apache/pulsar/client/impl/auth/AuthenticationDataTls.java | 2 +- .../pulsar/client/impl/auth/AuthenticationDataToken.java | 2 +- .../org/apache/pulsar/client/impl/auth/AuthenticationTls.java | 4 ++-- .../apache/pulsar/client/impl/auth/AuthenticationToken.java | 3 ++- 6 files changed, 9 insertions(+), 7 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationBasic.java index 4aecb10da3303..52dedc4f16e43 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationBasic.java @@ -30,6 +30,7 @@ import java.util.Map; public class AuthenticationBasic implements Authentication, EncodedAuthenticationParameterSupport { + static final String AUTH_METHOD_NAME = "basic"; private String userId; private String password; @@ -40,7 +41,7 @@ public void close() throws IOException { @Override public String getAuthMethodName() { - return "basic"; + return AUTH_METHOD_NAME; } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index fe4a5f3aa46a9..99a649a67b2e2 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -46,7 +46,7 @@ public boolean hasDataForHttp() { public Set> getHttpHeaders() { Map headers = new HashMap<>(); headers.put(HTTP_HEADER_NAME, httpAuthToken); - headers.put(PULSAR_AUTH_METHOD_NAME, "basic"); + headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationBasic.AUTH_METHOD_NAME); return headers.entrySet(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java index 30b96cfc5ff52..937922ae999f4 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java @@ -93,7 +93,7 @@ public boolean hasDataForTls() { @Override public Set> getHttpHeaders() { - return Collections.singletonMap(PULSAR_AUTH_METHOD_NAME, "tls").entrySet(); + return Collections.singletonMap(PULSAR_AUTH_METHOD_NAME, AuthenticationTls.AUTH_METHOD_NAME).entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index f97a251a77dce..42cd030044223 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -42,7 +42,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { Map headers = new HashMap<>(); - headers.put(PULSAR_AUTH_METHOD_NAME, "token"); + headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationToken.AUTH_METHOD_NAME); headers.put(HTTP_HEADER_NAME, "Bearer " + getToken()); return headers.entrySet(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationTls.java index 51464572d799a..d0414c7c0d217 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationTls.java @@ -41,7 +41,7 @@ * */ public class AuthenticationTls implements Authentication, EncodedAuthenticationParameterSupport { - private static final String AUTH_NAME = "tls"; + static final String AUTH_METHOD_NAME = "tls"; private static final long serialVersionUID = 1L; private String certFilePath; @@ -76,7 +76,7 @@ public void close() throws IOException { @Override public String getAuthMethodName() { - return AUTH_NAME; + return AUTH_METHOD_NAME; } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationToken.java index b85a40c343ce2..e71c438db6596 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationToken.java @@ -40,6 +40,7 @@ * Token based authentication provider. */ public class AuthenticationToken implements Authentication, EncodedAuthenticationParameterSupport { + static final String AUTH_METHOD_NAME = "token"; private static final long serialVersionUID = 1L; private Supplier tokenSupplier = null; @@ -60,7 +61,7 @@ public void close() throws IOException { @Override public String getAuthMethodName() { - return "token"; + return AUTH_METHOD_NAME; } @Override From a4d8718e633af42fdda3b819b15d50b8a7badbfc Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 22:27:41 +0800 Subject: [PATCH 08/26] Backward compatible --- .../pulsar/websocket/admin/WebSocketWebResource.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java index c58e67f84de91..3598621fe4759 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java @@ -26,6 +26,7 @@ import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.UriInfo; +import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.broker.web.AuthenticationFilter; import org.apache.pulsar.common.naming.TopicName; @@ -86,8 +87,12 @@ public String clientAppId() { public AuthenticationDataSource authData() throws AuthenticationException { if (authData == null) { String authMethodName = httpRequest.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); - authData = service().getAuthenticationService().getAuthenticationProvider(authMethodName) - .newHttpAuthState(httpRequest).getAuthDataSource(); + if (authMethodName != null && service().getAuthenticationService().getAuthenticationProvider(authMethodName) != null) { + authData = service().getAuthenticationService().getAuthenticationProvider(authMethodName) + .newHttpAuthState(httpRequest).getAuthDataSource(); + } else { + authData = new AuthenticationDataHttps(httpRequest); + } } return authData; } From 713b1d1e35ac81c6b5e09d3411041014c40b6217 Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 22:38:14 +0800 Subject: [PATCH 09/26] Fixed style --- .../pulsar/broker/authentication/AuthenticationService.java | 3 ++- .../org/apache/pulsar/broker/web/AuthenticationFilter.java | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java index 2919fccbcf317..e270e6c311b18 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java @@ -85,7 +85,8 @@ public AuthenticationService(ServiceConfiguration conf) throws PulsarServerExcep } } - public String authenticateHttpRequest(HttpServletRequest request, AuthenticationDataSource authData) throws AuthenticationException { + public String authenticateHttpRequest(HttpServletRequest request, AuthenticationDataSource authData) + throws AuthenticationException { AuthenticationException authenticationException = null; String authMethodName = request.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java index 8bfb1724cec65..6152a81d512cb 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java @@ -86,7 +86,8 @@ public void doFilter(ServletRequest request, ServletResponse response, FilterCha } String role; if (authenticationState != null ) { - role = authenticationService.authenticateHttpRequest((HttpServletRequest) request, authenticationState.getAuthDataSource()); + role = authenticationService.authenticateHttpRequest( + (HttpServletRequest) request, authenticationState.getAuthDataSource()); } else { role = authenticationService.authenticateHttpRequest((HttpServletRequest) request); } From 6c34f62b2392501278201874471514cf968c8353 Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 22:43:53 +0800 Subject: [PATCH 10/26] Delete whitespace --- .../java/org/apache/pulsar/broker/web/AuthenticationFilter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java index 6152a81d512cb..29ff9b331bbdf 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java @@ -85,7 +85,7 @@ public void doFilter(ServletRequest request, ServletResponse response, FilterCha new AuthenticationDataHttps((HttpServletRequest) request)); } String role; - if (authenticationState != null ) { + if (authenticationState != null) { role = authenticationService.authenticateHttpRequest( (HttpServletRequest) request, authenticationState.getAuthDataSource()); } else { From 8c2ef67fab3dd24eec99af22041f07a83e0dc434 Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 22:52:30 +0800 Subject: [PATCH 11/26] Reduce length --- .../apache/pulsar/websocket/admin/WebSocketWebResource.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java index 1d096fcbb2374..16a8402743916 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java @@ -86,7 +86,8 @@ public String clientAppId() { public AuthenticationDataSource authData() throws AuthenticationException { if (authData == null) { String authMethodName = httpRequest.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); - if (authMethodName != null && service().getAuthenticationService().getAuthenticationProvider(authMethodName) != null) { + if (authMethodName != null + && service().getAuthenticationService().getAuthenticationProvider(authMethodName) != null) { authData = service().getAuthenticationService().getAuthenticationProvider(authMethodName) .newHttpAuthState(httpRequest).getAuthDataSource(); } else { From c9f800091a6fa1f53e8640b6b448521b2a5cfa3c Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 3 Feb 2022 23:55:54 +0800 Subject: [PATCH 12/26] Update test --- .../broker/authentication/AuthenticationProviderListTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java index 365b324df76f0..1901673712eaf 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java @@ -21,6 +21,7 @@ import static java.nio.charset.StandardCharsets.UTF_8; import javax.servlet.http.HttpServletRequest; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; @@ -201,6 +202,8 @@ public void testNewAuthState() throws Exception { @Test public void testNewHttpAuthState() throws Exception { HttpServletRequest request = mock(HttpServletRequest.class); + when(request.getRemoteHost()).thenReturn("127.0.0.1"); + when(request.getRemotePort()).thenReturn(8080); AuthenticationState authStateAA = newHttpAuthState(request, SUBJECT_A); AuthenticationState authStateAB = newHttpAuthState(request, SUBJECT_B); AuthenticationState authStateBA = newHttpAuthState(request, SUBJECT_A); From 3fb6e3b4315a3c7d659eef9f6beb6902ae62ca45 Mon Sep 17 00:00:00 2001 From: guangning Date: Fri, 4 Feb 2022 00:21:23 +0800 Subject: [PATCH 13/26] Update remote addr --- .../broker/authentication/AuthenticationProviderListTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java index 1901673712eaf..6a32e4604a2a8 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java @@ -202,7 +202,7 @@ public void testNewAuthState() throws Exception { @Test public void testNewHttpAuthState() throws Exception { HttpServletRequest request = mock(HttpServletRequest.class); - when(request.getRemoteHost()).thenReturn("127.0.0.1"); + when(request.getRemoteAddr()).thenReturn("127.0.0.1"); when(request.getRemotePort()).thenReturn(8080); AuthenticationState authStateAA = newHttpAuthState(request, SUBJECT_A); AuthenticationState authStateAB = newHttpAuthState(request, SUBJECT_B); From 97ae44d4f8c8afdd6a1b7171d5c8f3f0bc384921 Mon Sep 17 00:00:00 2001 From: guangning Date: Fri, 4 Feb 2022 10:20:37 +0800 Subject: [PATCH 14/26] Fixed authentication token test --- .../AuthenticationProviderToken.java | 61 +++++++++++++++++++ .../AuthenticationProviderListTest.java | 30 ++++++--- 2 files changed, 84 insertions(+), 7 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index 451c63fb807da..67dd0492e66cf 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -38,6 +38,7 @@ import java.util.List; import javax.naming.AuthenticationException; import javax.net.ssl.SSLSession; +import javax.servlet.http.HttpServletRequest; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.authentication.metrics.AuthenticationMetrics; @@ -165,6 +166,11 @@ public AuthenticationState newAuthState(AuthData authData, SocketAddress remoteA return new TokenAuthenticationState(this, authData, remoteAddress, sslSession); } + @Override + public AuthenticationState newHttpAuthState(HttpServletRequest request) throws AuthenticationException { + return new TokenAuthenticationHttpState(this, request); + } + public static String getToken(AuthenticationDataSource authData) throws AuthenticationException { if (authData.hasDataFromCommand()) { // Authenticate Pulsar binary connection @@ -363,4 +369,59 @@ public boolean isExpired() { return expiration < System.currentTimeMillis(); } } + + private static final class TokenAuthenticationHttpState implements AuthenticationState { + + private final AuthenticationProviderToken provider; + private AuthenticationDataSource authenticationDataSource; + private Jwt jwt; + private long expiration; + + TokenAuthenticationHttpState(AuthenticationProviderToken provider, HttpServletRequest request) + throws AuthenticationException { + this.provider = provider; + String httpHeaderValue = request.getHeader(HTTP_HEADER_NAME); + if (httpHeaderValue == null || !httpHeaderValue.startsWith(HTTP_HEADER_VALUE_PREFIX)) { + throw new AuthenticationException("Invalid HTTP Authorization header"); + } + + // Remove prefix + String token = httpHeaderValue.substring(HTTP_HEADER_VALUE_PREFIX.length()); + this.jwt = provider.authenticateToken(token); + this.authenticationDataSource = new AuthenticationDataHttps(request); + if (jwt.getBody().getExpiration() != null) { + this.expiration = jwt.getBody().getExpiration().getTime(); + } else { + // Disable expiration + this.expiration = Long.MAX_VALUE; + } + } + + @Override + public String getAuthRole() throws AuthenticationException { + return provider.getPrincipal(jwt); + } + + @Override + public AuthenticationDataSource getAuthDataSource() { + return authenticationDataSource; + } + + @Override + public AuthData authenticate(AuthData authData) throws AuthenticationException { + return null; + } + + @Override + public boolean isComplete() { + // The authentication of tokens is always done in one single stage + return true; + } + + @Override + public boolean isExpired() { + return expiration < System.currentTimeMillis(); + } + + } } diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java index 6a32e4604a2a8..7d7e0ca92f61a 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/authentication/AuthenticationProviderListTest.java @@ -201,13 +201,29 @@ public void testNewAuthState() throws Exception { @Test public void testNewHttpAuthState() throws Exception { - HttpServletRequest request = mock(HttpServletRequest.class); - when(request.getRemoteAddr()).thenReturn("127.0.0.1"); - when(request.getRemotePort()).thenReturn(8080); - AuthenticationState authStateAA = newHttpAuthState(request, SUBJECT_A); - AuthenticationState authStateAB = newHttpAuthState(request, SUBJECT_B); - AuthenticationState authStateBA = newHttpAuthState(request, SUBJECT_A); - AuthenticationState authStateBB = newHttpAuthState(request, SUBJECT_B); + HttpServletRequest requestAA = mock(HttpServletRequest.class); + when(requestAA.getRemoteAddr()).thenReturn("127.0.0.1"); + when(requestAA.getRemotePort()).thenReturn(8080); + when(requestAA.getHeader("Authorization")).thenReturn("Bearer " + expiringTokenAA); + AuthenticationState authStateAA = newHttpAuthState(requestAA, SUBJECT_A); + + HttpServletRequest requestAB = mock(HttpServletRequest.class); + when(requestAB.getRemoteAddr()).thenReturn("127.0.0.1"); + when(requestAB.getRemotePort()).thenReturn(8080); + when(requestAB.getHeader("Authorization")).thenReturn("Bearer " + expiringTokenAB); + AuthenticationState authStateAB = newHttpAuthState(requestAB, SUBJECT_B); + + HttpServletRequest requestBA = mock(HttpServletRequest.class); + when(requestBA.getRemoteAddr()).thenReturn("127.0.0.1"); + when(requestBA.getRemotePort()).thenReturn(8080); + when(requestBA.getHeader("Authorization")).thenReturn("Bearer " + expiringTokenBA); + AuthenticationState authStateBA = newHttpAuthState(requestBA, SUBJECT_A); + + HttpServletRequest requestBB = mock(HttpServletRequest.class); + when(requestBB.getRemoteAddr()).thenReturn("127.0.0.1"); + when(requestBB.getRemotePort()).thenReturn(8080); + when(requestBB.getHeader("Authorization")).thenReturn("Bearer " + expiringTokenBB); + AuthenticationState authStateBB = newHttpAuthState(requestBB, SUBJECT_B); Thread.sleep(TimeUnit.SECONDS.toMillis(6)); From 197411fff53e29da7d0045a5660f1044501deb9a Mon Sep 17 00:00:00 2001 From: guangning Date: Fri, 4 Feb 2022 10:50:17 +0800 Subject: [PATCH 15/26] Fixed authentication token test --- .../client/impl/auth/AuthenticationTokenTest.java | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/auth/AuthenticationTokenTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/auth/AuthenticationTokenTest.java index d5d42c9686b11..d5174a0880bf7 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/auth/AuthenticationTokenTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/auth/AuthenticationTokenTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.client.impl.auth; +import java.util.Map; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNull; @@ -25,7 +26,6 @@ import java.io.*; import java.nio.charset.StandardCharsets; -import java.util.Collections; import java.util.function.Supplier; import org.apache.commons.io.FileUtils; @@ -34,6 +34,7 @@ import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.testng.annotations.Test; +import org.testng.collections.Maps; public class AuthenticationTokenTest { @@ -51,8 +52,10 @@ public void testAuthToken() throws Exception { assertNull(authData.getTlsPrivateKey()); assertTrue(authData.hasDataForHttp()); - assertEquals(authData.getHttpHeaders(), - Collections.singletonMap("Authorization", "Bearer token-xyz").entrySet()); + Map headers = Maps.newHashMap(); + headers.put("Authorization", "Bearer token-xyz"); + headers.put("X-Pulsar-Auth-Method-Name", "token"); + assertEquals(authData.getHttpHeaders(), headers.entrySet()); authToken.close(); } @@ -78,8 +81,10 @@ public void testAuthTokenClientConfig() throws Exception { assertNull(authData.getTlsPrivateKey()); assertTrue(authData.hasDataForHttp()); - assertEquals(authData.getHttpHeaders(), - Collections.singletonMap("Authorization", "Bearer token-xyz").entrySet()); + Map headers = Maps.newHashMap(); + headers.put("Authorization", "Bearer token-xyz"); + headers.put("X-Pulsar-Auth-Method-Name", "token"); + assertEquals(authData.getHttpHeaders(), headers.entrySet()); authToken.close(); } From b20a7d5f13f5608d71cffaa51dfaf61d28a6b122 Mon Sep 17 00:00:00 2001 From: guangning Date: Wed, 9 Feb 2022 11:31:46 +0800 Subject: [PATCH 16/26] Update var type, add final field, add comment for authentication method --- .../authentication/AuthenticationProviderToken.java | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index 67dd0492e66cf..32fe78e5ff67c 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -373,9 +373,9 @@ public boolean isExpired() { private static final class TokenAuthenticationHttpState implements AuthenticationState { private final AuthenticationProviderToken provider; - private AuthenticationDataSource authenticationDataSource; - private Jwt jwt; - private long expiration; + private final AuthenticationDataSource authenticationDataSource; + private final Jwt jwt; + private final long expiration; TokenAuthenticationHttpState(AuthenticationProviderToken provider, HttpServletRequest request) throws AuthenticationException { @@ -407,6 +407,10 @@ public AuthenticationDataSource getAuthDataSource() { return authenticationDataSource; } + /** + * Here is an explanation of why the null value is returned. + * pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationState.java#L49 + */ @Override public AuthData authenticate(AuthData authData) throws AuthenticationException { return null; From 0467c965860341aa01275f9dbc5830cbae0a1543 Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 10 Feb 2022 18:05:40 +0800 Subject: [PATCH 17/26] Update comment --- .../AuthenticationProviderToken.java | 92 ++++++------------- .../authentication/AuthenticationService.java | 17 +++- .../broker/web/AuthenticationFilter.java | 12 +-- .../impl/auth/AuthenticationDataBasic.java | 6 +- .../impl/auth/AuthenticationDataTls.java | 4 +- .../impl/auth/AuthenticationDataToken.java | 8 +- .../auth/oauth2/AuthenticationDataOAuth2.java | 6 +- .../worker/service/api/Functions.java | 35 +++++++ .../functions/worker/service/api/Sinks.java | 35 +++++++ .../functions/worker/service/api/Sources.java | 35 +++++++ .../websocket/admin/WebSocketWebResource.java | 27 +++--- 11 files changed, 178 insertions(+), 99 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index 32fe78e5ff67c..e06d142f405c7 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -168,7 +168,7 @@ public AuthenticationState newAuthState(AuthData authData, SocketAddress remoteA @Override public AuthenticationState newHttpAuthState(HttpServletRequest request) throws AuthenticationException { - return new TokenAuthenticationHttpState(this, request); + return new TokenAuthenticationState(this, request); } public static String getToken(AuthenticationDataSource authData) throws AuthenticationException { @@ -316,9 +316,10 @@ private static final class TokenAuthenticationState implements AuthenticationSta private final AuthenticationProviderToken provider; private AuthenticationDataSource authenticationDataSource; private Jwt jwt; - private final SocketAddress remoteAddress; - private final SSLSession sslSession; + private SocketAddress remoteAddress; + private SSLSession sslSession; private long expiration; + private HttpServletRequest request; TokenAuthenticationState( AuthenticationProviderToken provider, @@ -331,75 +332,50 @@ private static final class TokenAuthenticationState implements AuthenticationSta this.authenticate(authData); } + TokenAuthenticationState( + AuthenticationProviderToken provider, + HttpServletRequest request) throws AuthenticationException { + this.provider = provider; + this.request = request; + String httpHeaderValue = this.request.getHeader(HTTP_HEADER_NAME); + if (httpHeaderValue == null || !httpHeaderValue.startsWith(HTTP_HEADER_VALUE_PREFIX)) { + throw new AuthenticationException("Invalid HTTP Authorization header"); + } + + // Remove prefix + String token = httpHeaderValue.substring(HTTP_HEADER_VALUE_PREFIX.length()); + AuthData authData = AuthData.of(token.getBytes()); + this.authenticate(authData); + } + @Override public String getAuthRole() throws AuthenticationException { return provider.getPrincipal(jwt); } + /** + * Here is an explanation of why the null value is returned. + * pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationState.java#L49 + */ @Override public AuthData authenticate(AuthData authData) throws AuthenticationException { String token = new String(authData.getBytes(), UTF_8); this.jwt = provider.authenticateToken(token); - this.authenticationDataSource = new AuthenticationDataCommand(token, remoteAddress, sslSession); - if (jwt.getBody().getExpiration() != null) { - this.expiration = jwt.getBody().getExpiration().getTime(); + if (remoteAddress == null || sslSession == null) { + this.authenticationDataSource = new AuthenticationDataHttps(this.request); } else { - // Disable expiration - this.expiration = Long.MAX_VALUE; + this.authenticationDataSource = new AuthenticationDataCommand(token, remoteAddress, sslSession); } - - // There's no additional auth stage required - return null; - } - - @Override - public AuthenticationDataSource getAuthDataSource() { - return authenticationDataSource; - } - - @Override - public boolean isComplete() { - // The authentication of tokens is always done in one single stage - return true; - } - - @Override - public boolean isExpired() { - return expiration < System.currentTimeMillis(); - } - } - - private static final class TokenAuthenticationHttpState implements AuthenticationState { - - private final AuthenticationProviderToken provider; - private final AuthenticationDataSource authenticationDataSource; - private final Jwt jwt; - private final long expiration; - - TokenAuthenticationHttpState(AuthenticationProviderToken provider, HttpServletRequest request) - throws AuthenticationException { - this.provider = provider; - String httpHeaderValue = request.getHeader(HTTP_HEADER_NAME); - if (httpHeaderValue == null || !httpHeaderValue.startsWith(HTTP_HEADER_VALUE_PREFIX)) { - throw new AuthenticationException("Invalid HTTP Authorization header"); - } - - // Remove prefix - String token = httpHeaderValue.substring(HTTP_HEADER_VALUE_PREFIX.length()); - this.jwt = provider.authenticateToken(token); - this.authenticationDataSource = new AuthenticationDataHttps(request); if (jwt.getBody().getExpiration() != null) { this.expiration = jwt.getBody().getExpiration().getTime(); } else { // Disable expiration this.expiration = Long.MAX_VALUE; } - } - @Override - public String getAuthRole() throws AuthenticationException { - return provider.getPrincipal(jwt); + // There's no additional auth stage required + return null; } @Override @@ -407,15 +383,6 @@ public AuthenticationDataSource getAuthDataSource() { return authenticationDataSource; } - /** - * Here is an explanation of why the null value is returned. - * pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationState.java#L49 - */ - @Override - public AuthData authenticate(AuthData authData) throws AuthenticationException { - return null; - } - @Override public boolean isComplete() { // The authentication of tokens is always done in one single stage @@ -426,6 +393,5 @@ public boolean isComplete() { public boolean isExpired() { return expiration < System.currentTimeMillis(); } - } } diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java index e270e6c311b18..23d40bdd8eb36 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java @@ -96,11 +96,15 @@ public String authenticateHttpRequest(HttpServletRequest request, Authentication throw new AuthenticationException( String.format("Unsupported authentication method: [%s].", authMethodName)); } - if (authData == null) { - AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); - authData = authenticationState.getAuthDataSource(); - } try { + if (authData == null) { + // In the default implementation clas OneStageAuthenticationStates, the authentication has been + // done, in order to avoid secondary authentication here directly return the role, if the user + // custom implementation, should consider adding authentication in their own implementation class. + AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); + return authenticationState.getAuthRole(); + } + // Backward compatible, the authData value was null in the previous implementation return providerToUse.authenticate(authData); } catch (AuthenticationException e) { if (LOG.isDebugEnabled()) { @@ -143,6 +147,11 @@ public String authenticateHttpRequest(HttpServletRequest request, Authentication } } + /** + * Mark this function as deprecated, it is recommended to use a method with the AuthenticationDataSource + * signature to implement it. + */ + @Deprecated public String authenticateHttpRequest(HttpServletRequest request) throws AuthenticationException { return authenticateHttpRequest(request, null); } diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java index 29ff9b331bbdf..6c69bce42685e 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/web/AuthenticationFilter.java @@ -75,20 +75,16 @@ public void doFilter(ServletRequest request, ServletResponse response, FilterCha if (!isSaslRequest(httpRequest)) { // not sasl type, return role directly. String authMethodName = httpRequest.getHeader(PULSAR_AUTH_METHOD_NAME); - AuthenticationState authenticationState = null; + String role; if (authMethodName != null && authenticationService.getAuthenticationProvider(authMethodName) != null) { - authenticationState = authenticationService + AuthenticationState authenticationState = authenticationService .getAuthenticationProvider(authMethodName).newHttpAuthState(httpRequest); request.setAttribute(AuthenticatedDataAttributeName, authenticationState.getAuthDataSource()); - } else { - request.setAttribute(AuthenticatedDataAttributeName, - new AuthenticationDataHttps((HttpServletRequest) request)); - } - String role; - if (authenticationState != null) { role = authenticationService.authenticateHttpRequest( (HttpServletRequest) request, authenticationState.getAuthDataSource()); } else { + request.setAttribute(AuthenticatedDataAttributeName, + new AuthenticationDataHttps((HttpServletRequest) request)); role = authenticationService.authenticateHttpRequest((HttpServletRequest) request); } request.setAttribute(AuthenticatedRoleAttributeName, role); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index 983b238de7e03..332f17b407ed1 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -29,10 +29,13 @@ public class AuthenticationDataBasic implements AuthenticationDataProvider { private static final String HTTP_HEADER_NAME = "Authorization"; private String httpAuthToken; private String commandAuthToken; + private final Map headers = new HashMap<>(); public AuthenticationDataBasic(String userId, String password) { httpAuthToken = "Basic " + Base64.getEncoder().encodeToString((userId + ":" + password).getBytes()); commandAuthToken = userId + ":" + password; + headers.put(HTTP_HEADER_NAME, httpAuthToken); + headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationBasic.AUTH_METHOD_NAME); } @Override @@ -42,9 +45,6 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - Map headers = new HashMap<>(); - headers.put(HTTP_HEADER_NAME, httpAuthToken); - headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationBasic.AUTH_METHOD_NAME); return headers.entrySet(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java index 423f402bef18e..f6397237a4ab6 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java @@ -45,6 +45,8 @@ public class AuthenticationDataTls implements AuthenticationDataProvider { @SuppressFBWarnings(value = "SE_TRANSIENT_FIELD_NOT_RESTORED", justification = "Using custom serializer which Findbugs can't detect") private transient Supplier certStreamProvider, keyStreamProvider, trustStoreStreamProvider; + private final Map headers = Collections.singletonMap( + PULSAR_AUTH_METHOD_NAME, AuthenticationTls.AUTH_METHOD_NAME); public AuthenticationDataTls(String certFilePath, String keyFilePath) throws KeyManagementException { if (certFilePath == null) { @@ -92,7 +94,7 @@ public boolean hasDataForTls() { @Override public Set> getHttpHeaders() { - return Collections.singletonMap(PULSAR_AUTH_METHOD_NAME, AuthenticationTls.AUTH_METHOD_NAME).entrySet(); + return headers.entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index 1008d2a90738c..00def8f4260c3 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -28,9 +28,12 @@ public class AuthenticationDataToken implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; private final Supplier tokenSupplier; + private final Map headers = new HashMap<>(); public AuthenticationDataToken(Supplier tokenSupplier) { this.tokenSupplier = tokenSupplier; + headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationToken.AUTH_METHOD_NAME); + headers.put(HTTP_HEADER_NAME, "Bearer " + getToken()); } @Override @@ -40,10 +43,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - Map headers = new HashMap<>(); - headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationToken.AUTH_METHOD_NAME); - headers.put(HTTP_HEADER_NAME, "Bearer " + getToken()); - return headers.entrySet(); + return this.headers.entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java index 227443f0dc248..61ea75f04e732 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java @@ -30,9 +30,12 @@ class AuthenticationDataOAuth2 implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; private final String accessToken; + private final Map headers = new HashMap<>(); public AuthenticationDataOAuth2(String accessToken) { this.accessToken = accessToken; + headers.put(HTTP_HEADER_NAME, "Bearer " + accessToken); + headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationOAuth2.AUTH_METHOD_NAME); } @Override @@ -42,9 +45,6 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - Map headers = new HashMap<>(); - headers.put(HTTP_HEADER_NAME, "Bearer " + accessToken); - headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationOAuth2.AUTH_METHOD_NAME); return headers.entrySet(); } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java index 897811ec700de..0f4deb3cf37a1 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java @@ -34,6 +34,18 @@ */ public interface Functions extends Component { + /** + * Register a new function + * @param tenant The tenant of a Pulsar Function + * @param namespace The namespace of a Pulsar Function + * @param functionName The name of a Pulsar Function + * @param uploadedInputStream Input stream of bytes + * @param fileDetail A form-data content disposition header + * @param functionPkgUrl URL path of the Pulsar Function package + * @param functionConfig Configuration of Pulsar Function + * @param clientRole Client role for running the pulsar function + * @param clientAuthenticationDataHttps Authentication status of the http client + */ void registerFunction(final String tenant, final String namespace, final String functionName, @@ -44,6 +56,11 @@ void registerFunction(final String tenant, final String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); + /** + * This method uses an incorrect signature 'AuthenticationDataHttps' that prevents the extension of auth status, + * so it is marked as deprecated and kept here only for backward compatibility. Please use the method that accepts + * the signature of the AuthenticationDataSource. + */ @Deprecated default void registerFunction(final String tenant, final String namespace, @@ -66,6 +83,19 @@ default void registerFunction(final String tenant, (AuthenticationDataSource) clientAuthenticationDataHttps); } + /** + * Update a function + * @param tenant The tenant of a Pulsar Function + * @param namespace The namespace of a Pulsar Function + * @param functionName The name of a Pulsar Function + * @param uploadedInputStream Input stream of bytes + * @param fileDetail A form-data content disposition header + * @param functionPkgUrl URL path of the Pulsar Function package + * @param functionConfig Configuration of Pulsar Function + * @param clientRole Client role for running the Pulsar Function + * @param clientAuthenticationDataHttps Authentication status of the http client + * @param updateOptions Options while updating the function + */ void updateFunction(final String tenant, final String namespace, final String functionName, @@ -77,6 +107,11 @@ void updateFunction(final String tenant, AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); + /** + * This method uses an incorrect signature 'AuthenticationDataHttps' that prevents the extension of auth status, + * so it is marked as deprecated and kept here only for backward compatibility. Please use the method that accepts + * the signature of the AuthenticationDataSource. + */ @Deprecated default void updateFunction(final String tenant, final String namespace, diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java index cdcf437a1437d..83c303870f700 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java @@ -37,6 +37,18 @@ */ public interface Sinks extends Component { + /** + * Update a function + * @param tenant The tenant of a Pulsar Sink + * @param namespace The namespace of a Pulsar Sink + * @param sinkName The name of a Pulsar Sink + * @param uploadedInputStream Input stream of bytes + * @param fileDetail A form-data content disposition header + * @param sinkPkgUrl URL path of the Pulsar Sink package + * @param sinkConfig Configuration of Pulsar Sink + * @param clientRole Client role for running the Pulsar Sink + * @param clientAuthenticationDataHttps Authentication status of the http client + */ void registerSink(final String tenant, final String namespace, final String sinkName, @@ -47,6 +59,11 @@ void registerSink(final String tenant, final String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); + /** + * This method uses an incorrect signature 'AuthenticationDataHttps' that prevents the extension of auth status, + * so it is marked as deprecated and kept here only for backward compatibility. Please use the method that accepts + * the signature of the AuthenticationDataSource. + */ @Deprecated default void registerSink(final String tenant, final String namespace, @@ -69,6 +86,19 @@ default void registerSink(final String tenant, (AuthenticationDataSource) clientAuthenticationDataHttps); } + /** + * Update a function + * @param tenant The tenant of a Pulsar Sink + * @param namespace The namespace of a Pulsar Sink + * @param sinkName The name of a Pulsar Sink + * @param uploadedInputStream Input stream of bytes + * @param fileDetail A form-data content disposition header + * @param sinkPkgUrl URL path of the Pulsar Sink package + * @param sinkConfig Configuration of Pulsar Sink + * @param clientRole Client role for running the Pulsar Sink + * @param clientAuthenticationDataHttps Authentication status of the http client + * @param updateOptions Options while updating the sink + */ void updateSink(final String tenant, final String namespace, final String sinkName, @@ -80,6 +110,11 @@ void updateSink(final String tenant, AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); + /** + * This method uses an incorrect signature 'AuthenticationDataHttps' that prevents the extension of auth status, + * so it is marked as deprecated and kept here only for backward compatibility. Please use the method that accepts + * the signature of the AuthenticationDataSource. + */ @Deprecated default void updateSink(final String tenant, final String namespace, diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java index ab57233aef013..ded37c0978a48 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java @@ -37,6 +37,18 @@ */ public interface Sources extends Component { + /** + * Update a function + * @param tenant The tenant of a Pulsar Source + * @param namespace The namespace of a Pulsar Source + * @param sourceName The name of a Pulsar Source + * @param uploadedInputStream Input stream of bytes + * @param fileDetail A form-data content disposition header + * @param sourcePkgUrl URL path of the Pulsar Source package + * @param sourceConfig Configuration of Pulsar Source + * @param clientRole Client role for running the Pulsar Source + * @param clientAuthenticationDataHttps Authentication status of the http client + */ void registerSource(final String tenant, final String namespace, final String sourceName, @@ -47,6 +59,11 @@ void registerSource(final String tenant, final String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); + /** + * This method uses an incorrect signature 'AuthenticationDataHttps' that prevents the extension of auth status, + * so it is marked as deprecated and kept here only for backward compatibility. Please use the method that accepts + * the signature of the AuthenticationDataSource. + */ @Deprecated default void registerSource(final String tenant, final String namespace, @@ -69,6 +86,19 @@ default void registerSource(final String tenant, (AuthenticationDataSource) clientAuthenticationDataHttps); } + /** + * Update a function + * @param tenant The tenant of a Pulsar Source + * @param namespace The namespace of a Pulsar Source + * @param sourceName The name of a Pulsar Source + * @param uploadedInputStream Input stream of bytes + * @param fileDetail A form-data content disposition header + * @param sourcePkgUrl URL path of the Pulsar Source package + * @param sourceConfig Configuration of Pulsar Source + * @param clientRole Client role for running the Pulsar Source + * @param clientAuthenticationDataHttps Authentication status of the http client + * @param updateOptions Options while updating the source + */ void updateSource(final String tenant, final String namespace, final String sourceName, @@ -80,6 +110,11 @@ void updateSource(final String tenant, AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); + /** + * This method uses an incorrect signature 'AuthenticationDataHttps' that prevents the extension of auth status, + * so it is marked as deprecated and kept here only for backward compatibility. Please use the method that accepts + * the signature of the AuthenticationDataSource. + */ @Deprecated default void updateSource(final String tenant, final String namespace, diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java index 16a8402743916..b6a0f43a01aae 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/admin/WebSocketWebResource.java @@ -52,7 +52,7 @@ public class WebSocketWebResource { private WebSocketService socketService; private String clientId; - private AuthenticationDataSource authData; + private AuthenticationDataSource authenticationDataSource; protected WebSocketService service() { if (socketService == null) { @@ -69,7 +69,18 @@ protected WebSocketService service() { public String clientAppId() { if (isBlank(clientId)) { try { - clientId = service().getAuthenticationService().authenticateHttpRequest(httpRequest); + String authMethodName = httpRequest.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); + if (authMethodName != null + && service().getAuthenticationService().getAuthenticationProvider(authMethodName) != null) { + authenticationDataSource = service().getAuthenticationService() + .getAuthenticationProvider(authMethodName) + .newHttpAuthState(httpRequest).getAuthDataSource(); + clientId = service().getAuthenticationService().authenticateHttpRequest( + httpRequest, authenticationDataSource); + } else { + clientId = service().getAuthenticationService().authenticateHttpRequest(httpRequest); + authenticationDataSource = new AuthenticationDataHttps(httpRequest); + } } catch (AuthenticationException e) { if (service().getConfig().isAuthenticationEnabled()) { throw new RestException(Status.UNAUTHORIZED, "Failed to get clientId from request"); @@ -84,17 +95,7 @@ public String clientAppId() { } public AuthenticationDataSource authData() throws AuthenticationException { - if (authData == null) { - String authMethodName = httpRequest.getHeader(AuthenticationFilter.PULSAR_AUTH_METHOD_NAME); - if (authMethodName != null - && service().getAuthenticationService().getAuthenticationProvider(authMethodName) != null) { - authData = service().getAuthenticationService().getAuthenticationProvider(authMethodName) - .newHttpAuthState(httpRequest).getAuthDataSource(); - } else { - authData = new AuthenticationDataHttps(httpRequest); - } - } - return authData; + return authenticationDataSource; } /** From 3fc7d57f53e983acdba7aca9ab208b8fe94b5f4b Mon Sep 17 00:00:00 2001 From: guangning Date: Thu, 10 Feb 2022 21:51:16 +0800 Subject: [PATCH 18/26] Fixed test --- .../broker/authentication/AuthenticationProviderToken.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index e06d142f405c7..a53829828b41c 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -320,6 +320,7 @@ private static final class TokenAuthenticationState implements AuthenticationSta private SSLSession sslSession; private long expiration; private HttpServletRequest request; + private boolean authHttpState; TokenAuthenticationState( AuthenticationProviderToken provider, @@ -329,6 +330,7 @@ private static final class TokenAuthenticationState implements AuthenticationSta this.provider = provider; this.remoteAddress = remoteAddress; this.sslSession = sslSession; + this.authHttpState = false; this.authenticate(authData); } @@ -337,6 +339,7 @@ private static final class TokenAuthenticationState implements AuthenticationSta HttpServletRequest request) throws AuthenticationException { this.provider = provider; this.request = request; + this.authHttpState = true; String httpHeaderValue = this.request.getHeader(HTTP_HEADER_NAME); if (httpHeaderValue == null || !httpHeaderValue.startsWith(HTTP_HEADER_VALUE_PREFIX)) { throw new AuthenticationException("Invalid HTTP Authorization header"); @@ -362,7 +365,7 @@ public AuthData authenticate(AuthData authData) throws AuthenticationException { String token = new String(authData.getBytes(), UTF_8); this.jwt = provider.authenticateToken(token); - if (remoteAddress == null || sslSession == null) { + if (authHttpState) { this.authenticationDataSource = new AuthenticationDataHttps(this.request); } else { this.authenticationDataSource = new AuthenticationDataCommand(token, remoteAddress, sslSession); From 5b08617167a0f4faa6d637d58eddeb1d447970c1 Mon Sep 17 00:00:00 2001 From: guangning Date: Fri, 11 Feb 2022 09:02:00 +0800 Subject: [PATCH 19/26] Fixed auth test --- .../pulsar/websocket/admin/WebSocketWebResourceTest.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/admin/WebSocketWebResourceTest.java b/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/admin/WebSocketWebResourceTest.java index 59e681de51322..6c8156ef79fe9 100644 --- a/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/admin/WebSocketWebResourceTest.java +++ b/pulsar-websocket/src/test/java/org/apache/pulsar/websocket/admin/WebSocketWebResourceTest.java @@ -121,6 +121,10 @@ public void setup(Method method) throws Exception { // Mock ServletContext when(servletContext.getAttribute(anyString())).thenReturn(socketService); + // Mock HttpServletRequest + when(httpRequest.getRemoteAddr()).thenReturn("127.0.0.1"); + when(httpRequest.getRemotePort()).thenReturn(8080); + // Mock UriInfo when(uri.getRequestUri()).thenReturn(null); From ae5b4477e16669e3224f028670035e06927ec088 Mon Sep 17 00:00:00 2001 From: guangning Date: Wed, 16 Feb 2022 17:55:35 +0800 Subject: [PATCH 20/26] Update comment --- .../AuthenticationProviderToken.java | 24 +++++++----------- .../authentication/AuthenticationService.java | 2 +- .../websocket/AbstractWebSocketHandler.java | 25 ++++++++++++++++--- 3 files changed, 32 insertions(+), 19 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index a53829828b41c..a9e0a529a9fa9 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -320,7 +320,6 @@ private static final class TokenAuthenticationState implements AuthenticationSta private SSLSession sslSession; private long expiration; private HttpServletRequest request; - private boolean authHttpState; TokenAuthenticationState( AuthenticationProviderToken provider, @@ -330,8 +329,9 @@ private static final class TokenAuthenticationState implements AuthenticationSta this.provider = provider; this.remoteAddress = remoteAddress; this.sslSession = sslSession; - this.authHttpState = false; - this.authenticate(authData); + String token = new String(authData.getBytes(), UTF_8); + this.authenticationDataSource = new AuthenticationDataCommand(token, remoteAddress, sslSession); + this.checkExpiration(token); } TokenAuthenticationState( @@ -339,7 +339,6 @@ private static final class TokenAuthenticationState implements AuthenticationSta HttpServletRequest request) throws AuthenticationException { this.provider = provider; this.request = request; - this.authHttpState = true; String httpHeaderValue = this.request.getHeader(HTTP_HEADER_NAME); if (httpHeaderValue == null || !httpHeaderValue.startsWith(HTTP_HEADER_VALUE_PREFIX)) { throw new AuthenticationException("Invalid HTTP Authorization header"); @@ -347,8 +346,8 @@ private static final class TokenAuthenticationState implements AuthenticationSta // Remove prefix String token = httpHeaderValue.substring(HTTP_HEADER_VALUE_PREFIX.length()); - AuthData authData = AuthData.of(token.getBytes()); - this.authenticate(authData); + this.authenticationDataSource = new AuthenticationDataHttps(this.request); + this.checkExpiration(token); } @Override @@ -362,23 +361,18 @@ public String getAuthRole() throws AuthenticationException { */ @Override public AuthData authenticate(AuthData authData) throws AuthenticationException { - String token = new String(authData.getBytes(), UTF_8); + // There's no additional auth stage required + return null; + } + private void checkExpiration(String token) throws AuthenticationException { this.jwt = provider.authenticateToken(token); - if (authHttpState) { - this.authenticationDataSource = new AuthenticationDataHttps(this.request); - } else { - this.authenticationDataSource = new AuthenticationDataCommand(token, remoteAddress, sslSession); - } if (jwt.getBody().getExpiration() != null) { this.expiration = jwt.getBody().getExpiration().getTime(); } else { // Disable expiration this.expiration = Long.MAX_VALUE; } - - // There's no additional auth stage required - return null; } @Override diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java index 23d40bdd8eb36..0803c5dc5d76c 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java @@ -98,7 +98,7 @@ public String authenticateHttpRequest(HttpServletRequest request, Authentication } try { if (authData == null) { - // In the default implementation clas OneStageAuthenticationStates, the authentication has been + // In the default implementation class OneStageAuthenticationStates, the authentication has been // done, in order to avoid secondary authentication here directly return the role, if the user // custom implementation, should consider adding authentication in their own implementation class. AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java index 5217091b18fb4..8262782bf64f0 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java @@ -30,6 +30,8 @@ import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; +import org.apache.pulsar.broker.authentication.AuthenticationProvider; +import org.apache.pulsar.broker.authentication.AuthenticationState; import org.apache.pulsar.client.api.PulsarClientException.AuthenticationException; import org.apache.pulsar.client.api.PulsarClientException.AuthorizationException; import org.apache.pulsar.client.api.PulsarClientException.ConsumerBusyException; @@ -59,7 +61,7 @@ public abstract class AbstractWebSocketHandler extends WebSocketAdapter implemen protected final TopicName topic; protected final Map queryParams; - + private static final String PULSAR_AUTH_METHOD_NAME = "X-Pulsar-Auth-Method-Name"; public AbstractWebSocketHandler(WebSocketService service, HttpServletRequest request, @@ -76,9 +78,21 @@ public AbstractWebSocketHandler(WebSocketService service, protected boolean checkAuth(ServletUpgradeResponse response) { String authRole = ""; + String authMethodName = request.getHeader(PULSAR_AUTH_METHOD_NAME); + AuthenticationState authenticationState = null; if (service.isAuthenticationEnabled()) { try { - authRole = service.getAuthenticationService().authenticateHttpRequest(request); + if (authMethodName != null + && service.getAuthenticationService().getAuthenticationProvider(authMethodName) != null) { + authenticationState = service.getAuthenticationService() + .getAuthenticationProvider(authMethodName).newHttpAuthState(request); + } + if (authenticationState != null) { + authRole = service.getAuthenticationService() + .authenticateHttpRequest(request, authenticationState.getAuthDataSource()); + } else { + authRole = service.getAuthenticationService().authenticateHttpRequest(request); + } log.info("[{}:{}] Authenticated WebSocket client {} on topic {}", request.getRemoteAddr(), request.getRemotePort(), authRole, topic); @@ -96,7 +110,12 @@ protected boolean checkAuth(ServletUpgradeResponse response) { } if (service.isAuthorizationEnabled()) { - AuthenticationDataSource authenticationData = new AuthenticationDataHttps(request); + AuthenticationDataSource authenticationData; + if (authenticationState != null) { + authenticationData = authenticationState.getAuthDataSource(); + } else { + authenticationData = new AuthenticationDataHttps(request); + } try { if (!isAuthorized(authRole, authenticationData)) { log.warn("[{}:{}] WebSocket Client [{}] is not authorized on topic {}", request.getRemoteAddr(), From e4efa86aabf259c57a8f97869e0553114b44c6e2 Mon Sep 17 00:00:00 2001 From: guangning Date: Wed, 16 Feb 2022 18:42:31 +0800 Subject: [PATCH 21/26] Revert code --- .../pulsar/broker/authentication/AuthenticationService.java | 5 +---- .../apache/pulsar/websocket/AbstractWebSocketHandler.java | 1 - 2 files changed, 1 insertion(+), 5 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java index 0803c5dc5d76c..6fe6c5f8e7ca7 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationService.java @@ -98,11 +98,8 @@ public String authenticateHttpRequest(HttpServletRequest request, Authentication } try { if (authData == null) { - // In the default implementation class OneStageAuthenticationStates, the authentication has been - // done, in order to avoid secondary authentication here directly return the role, if the user - // custom implementation, should consider adding authentication in their own implementation class. AuthenticationState authenticationState = providerToUse.newHttpAuthState(request); - return authenticationState.getAuthRole(); + authData = authenticationState.getAuthDataSource(); } // Backward compatible, the authData value was null in the previous implementation return providerToUse.authenticate(authData); diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java index 8262782bf64f0..ea4cc128a764c 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/AbstractWebSocketHandler.java @@ -30,7 +30,6 @@ import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.authentication.AuthenticationDataHttps; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; -import org.apache.pulsar.broker.authentication.AuthenticationProvider; import org.apache.pulsar.broker.authentication.AuthenticationState; import org.apache.pulsar.client.api.PulsarClientException.AuthenticationException; import org.apache.pulsar.client.api.PulsarClientException.AuthorizationException; From bb1ebfe2ee95dd2c3458b1dd4d52e7881184ab18 Mon Sep 17 00:00:00 2001 From: guangning Date: Wed, 23 Feb 2022 10:35:29 +0800 Subject: [PATCH 22/26] Remove no used var --- .../authentication/AuthenticationProviderToken.java | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index a9e0a529a9fa9..7eea1ee2fcbd0 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -316,10 +316,7 @@ private static final class TokenAuthenticationState implements AuthenticationSta private final AuthenticationProviderToken provider; private AuthenticationDataSource authenticationDataSource; private Jwt jwt; - private SocketAddress remoteAddress; - private SSLSession sslSession; private long expiration; - private HttpServletRequest request; TokenAuthenticationState( AuthenticationProviderToken provider, @@ -327,8 +324,6 @@ private static final class TokenAuthenticationState implements AuthenticationSta SocketAddress remoteAddress, SSLSession sslSession) throws AuthenticationException { this.provider = provider; - this.remoteAddress = remoteAddress; - this.sslSession = sslSession; String token = new String(authData.getBytes(), UTF_8); this.authenticationDataSource = new AuthenticationDataCommand(token, remoteAddress, sslSession); this.checkExpiration(token); @@ -338,15 +333,14 @@ private static final class TokenAuthenticationState implements AuthenticationSta AuthenticationProviderToken provider, HttpServletRequest request) throws AuthenticationException { this.provider = provider; - this.request = request; - String httpHeaderValue = this.request.getHeader(HTTP_HEADER_NAME); + String httpHeaderValue = request.getHeader(HTTP_HEADER_NAME); if (httpHeaderValue == null || !httpHeaderValue.startsWith(HTTP_HEADER_VALUE_PREFIX)) { throw new AuthenticationException("Invalid HTTP Authorization header"); } // Remove prefix String token = httpHeaderValue.substring(HTTP_HEADER_VALUE_PREFIX.length()); - this.authenticationDataSource = new AuthenticationDataHttps(this.request); + this.authenticationDataSource = new AuthenticationDataHttps(request); this.checkExpiration(token); } From 376412160a4fa0c0b5673227ed3ddf49081659e9 Mon Sep 17 00:00:00 2001 From: guangning Date: Wed, 23 Feb 2022 16:45:53 +0800 Subject: [PATCH 23/26] Update map headers to unmodifiableMap type --- .../broker/authentication/AuthenticationProviderToken.java | 5 +++-- .../pulsar/client/impl/auth/AuthenticationDataBasic.java | 3 ++- .../pulsar/client/impl/auth/AuthenticationDataTls.java | 2 +- .../pulsar/client/impl/auth/AuthenticationDataToken.java | 3 ++- .../client/impl/auth/oauth2/AuthenticationDataOAuth2.java | 3 ++- 5 files changed, 10 insertions(+), 6 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java index 7eea1ee2fcbd0..164d5ee672ce6 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationProviderToken.java @@ -350,8 +350,9 @@ public String getAuthRole() throws AuthenticationException { } /** - * Here is an explanation of why the null value is returned. - * pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authentication/AuthenticationState.java#L49 + * @param authData Authentication data. + * @return null. Explanation of returning null values, {@link AuthenticationState#authenticateAsync(AuthData)} + * @throws AuthenticationException */ @Override public AuthData authenticate(AuthData authData) throws AuthenticationException { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index 332f17b407ed1..3ee5374210cc9 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -20,6 +20,7 @@ package org.apache.pulsar.client.impl.auth; import java.util.Base64; +import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.Set; @@ -45,7 +46,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return headers.entrySet(); + return Collections.unmodifiableMap(this.headers).entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java index f6397237a4ab6..14e67ba4ddf65 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataTls.java @@ -45,7 +45,7 @@ public class AuthenticationDataTls implements AuthenticationDataProvider { @SuppressFBWarnings(value = "SE_TRANSIENT_FIELD_NOT_RESTORED", justification = "Using custom serializer which Findbugs can't detect") private transient Supplier certStreamProvider, keyStreamProvider, trustStoreStreamProvider; - private final Map headers = Collections.singletonMap( + private static final Map headers = Collections.singletonMap( PULSAR_AUTH_METHOD_NAME, AuthenticationTls.AUTH_METHOD_NAME); public AuthenticationDataTls(String certFilePath, String keyFilePath) throws KeyManagementException { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index 00def8f4260c3..bf105df32bf2a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -19,6 +19,7 @@ package org.apache.pulsar.client.impl.auth; +import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.Set; @@ -43,7 +44,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return this.headers.entrySet(); + return Collections.unmodifiableMap(this.headers).entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java index 61ea75f04e732..7af587cd9ca1f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.client.impl.auth.oauth2; +import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.Set; @@ -45,7 +46,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return headers.entrySet(); + return Collections.unmodifiableMap(this.headers).entrySet(); } @Override From 49c1080cb686708d94e143254662118a5c133ced Mon Sep 17 00:00:00 2001 From: guangning Date: Sun, 27 Feb 2022 21:11:36 +0800 Subject: [PATCH 24/26] Remove redundant final --- .../worker/service/api/Functions.java | 78 +++++++++---------- .../functions/worker/service/api/Sinks.java | 34 ++++---- .../functions/worker/service/api/Sources.java | 68 ++++++++-------- 3 files changed, 90 insertions(+), 90 deletions(-) diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java index f99f9443910f8..ac77b76ec2e30 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java @@ -35,7 +35,7 @@ public interface Functions extends Component { /** - * Register a new function + * Register a new function. * @param tenant The tenant of a Pulsar Function * @param namespace The namespace of a Pulsar Function * @param functionName The name of a Pulsar Function @@ -46,14 +46,14 @@ public interface Functions extends Component { * @param clientRole Client role for running the pulsar function * @param clientAuthenticationDataHttps Authentication status of the http client */ - void registerFunction(final String tenant, - final String namespace, - final String functionName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String functionPkgUrl, - final FunctionConfig functionConfig, - final String clientRole, + void registerFunction(String tenant, + String namespace, + String functionName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String functionPkgUrl, + FunctionConfig functionConfig, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); /** @@ -62,14 +62,14 @@ void registerFunction(final String tenant, * the signature of the AuthenticationDataSource. */ @Deprecated - default void registerFunction(final String tenant, - final String namespace, - final String functionName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String functionPkgUrl, - final FunctionConfig functionConfig, - final String clientRole, + default void registerFunction(String tenant, + String namespace, + String functionName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String functionPkgUrl, + FunctionConfig functionConfig, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps) { registerFunction( tenant, @@ -84,7 +84,7 @@ default void registerFunction(final String tenant, } /** - * Update a function + * Update a function. * @param tenant The tenant of a Pulsar Function * @param namespace The namespace of a Pulsar Function * @param functionName The name of a Pulsar Function @@ -96,14 +96,14 @@ default void registerFunction(final String tenant, * @param clientAuthenticationDataHttps Authentication status of the http client * @param updateOptions Options while updating the function */ - void updateFunction(final String tenant, - final String namespace, - final String functionName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String functionPkgUrl, - final FunctionConfig functionConfig, - final String clientRole, + void updateFunction(String tenant, + String namespace, + String functionName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String functionPkgUrl, + FunctionConfig functionConfig, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); @@ -113,14 +113,14 @@ void updateFunction(final String tenant, * the signature of the AuthenticationDataSource. */ @Deprecated - default void updateFunction(final String tenant, - final String namespace, - final String functionName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String functionPkgUrl, - final FunctionConfig functionConfig, - final String clientRole, + default void updateFunction(String tenant, + String namespace, + String functionName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String functionPkgUrl, + FunctionConfig functionConfig, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions) { updateFunction( @@ -136,11 +136,11 @@ default void updateFunction(final String tenant, updateOptions); } - void updateFunctionOnWorkerLeader(final String tenant, - final String namespace, - final String functionName, - final InputStream uploadedInputStream, - final boolean delete, + void updateFunctionOnWorkerLeader(String tenant, + String namespace, + String functionName, + InputStream uploadedInputStream, + boolean delete, URI uri, String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java index 3a023c8eb4c56..db8411151ed7a 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java @@ -38,7 +38,7 @@ public interface Sinks extends Component { /** - * Update a function + * Update a function. * @param tenant The tenant of a Pulsar Sink * @param namespace The namespace of a Pulsar Sink * @param sinkName The name of a Pulsar Sink @@ -87,7 +87,7 @@ default void registerSink(final String tenant, } /** - * Update a function + * Update a function. * @param tenant The tenant of a Pulsar Sink * @param namespace The namespace of a Pulsar Sink * @param sinkName The name of a Pulsar Sink @@ -116,14 +116,14 @@ void updateSink(final String tenant, * the signature of the AuthenticationDataSource. */ @Deprecated - default void updateSink(final String tenant, - final String namespace, - final String sinkName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sinkPkgUrl, - final SinkConfig sinkConfig, - final String clientRole, + default void updateSink(String tenant, + String namespace, + String sinkName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sinkPkgUrl, + SinkConfig sinkConfig, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions) { updateSink( @@ -139,13 +139,13 @@ default void updateSink(final String tenant, updateOptions); } - SinkInstanceStatusData getSinkInstanceStatus(final String tenant, - final String namespace, - final String sinkName, - final String instanceId, - final URI uri, - final String clientRole, - final AuthenticationDataSource clientAuthenticationDataHttps); + SinkInstanceStatusData getSinkInstanceStatus(String tenant, + String namespace, + String sinkName, + String instanceId, + URI uri, + String clientRole, + AuthenticationDataSource clientAuthenticationDataHttps); SinkStatus getSinkStatus(String tenant, String namespace, diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java index 92f43d1b72865..089115228b2d9 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sources.java @@ -38,7 +38,7 @@ public interface Sources extends Component { /** - * Update a function + * Update a function. * @param tenant The tenant of a Pulsar Source * @param namespace The namespace of a Pulsar Source * @param sourceName The name of a Pulsar Source @@ -49,14 +49,14 @@ public interface Sources extends Component { * @param clientRole Client role for running the Pulsar Source * @param clientAuthenticationDataHttps Authentication status of the http client */ - void registerSource(final String tenant, - final String namespace, - final String sourceName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sourcePkgUrl, - final SourceConfig sourceConfig, - final String clientRole, + void registerSource(String tenant, + String namespace, + String sourceName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sourcePkgUrl, + SourceConfig sourceConfig, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); /** @@ -65,14 +65,14 @@ void registerSource(final String tenant, * the signature of the AuthenticationDataSource. */ @Deprecated - default void registerSource(final String tenant, - final String namespace, - final String sourceName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sourcePkgUrl, - final SourceConfig sourceConfig, - final String clientRole, + default void registerSource(String tenant, + String namespace, + String sourceName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sourcePkgUrl, + SourceConfig sourceConfig, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps) { registerSource( tenant, @@ -87,7 +87,7 @@ default void registerSource(final String tenant, } /** - * Update a function + * Update a function. * @param tenant The tenant of a Pulsar Source * @param namespace The namespace of a Pulsar Source * @param sourceName The name of a Pulsar Source @@ -99,14 +99,14 @@ default void registerSource(final String tenant, * @param clientAuthenticationDataHttps Authentication status of the http client * @param updateOptions Options while updating the source */ - void updateSource(final String tenant, - final String namespace, - final String sourceName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sourcePkgUrl, - final SourceConfig sourceConfig, - final String clientRole, + void updateSource(String tenant, + String namespace, + String sourceName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sourcePkgUrl, + SourceConfig sourceConfig, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); @@ -116,14 +116,14 @@ void updateSource(final String tenant, * the signature of the AuthenticationDataSource. */ @Deprecated - default void updateSource(final String tenant, - final String namespace, - final String sourceName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sourcePkgUrl, - final SourceConfig sourceConfig, - final String clientRole, + default void updateSource(String tenant, + String namespace, + String sourceName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sourcePkgUrl, + SourceConfig sourceConfig, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions) { updateSource( From b1def58156f9db86161f847224a9624b13ffec46 Mon Sep 17 00:00:00 2001 From: guangning Date: Sun, 27 Feb 2022 21:21:54 +0800 Subject: [PATCH 25/26] Fixed style --- .../worker/rest/api/ComponentImpl.java | 6 ++- .../worker/service/api/Component.java | 24 +++++----- .../functions/worker/service/api/Sinks.java | 48 +++++++++---------- 3 files changed, 40 insertions(+), 38 deletions(-) diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java index 61c72faf5b8fe..d0db047a02048 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/ComponentImpl.java @@ -1210,7 +1210,8 @@ public FunctionState getFunctionState(final String tenant, .getBytes(kv.value(), kv.value().readerIndex(), kv.value().readableBytes()), UTF_8), null, null, kv.version()); } catch (Exception e) { - value = new FunctionState(key, null, ByteBufUtil.getBytes(kv.value()), null, kv.version()); + value = new FunctionState( + key, null, ByteBufUtil.getBytes(kv.value()), null, kv.version()); } } } @@ -1409,7 +1410,8 @@ private StreamingOutput getStreamingOutput(String pkgPath) { } @Override - public StreamingOutput downloadFunction(final String path, String clientRole, AuthenticationDataSource clientAuthenticationDataHttps) { + public StreamingOutput downloadFunction( + final String path, String clientRole, AuthenticationDataSource clientAuthenticationDataHttps) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java index 8dae5055d3fe9..c305d64b9f3cf 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Component.java @@ -40,17 +40,17 @@ public interface Component { W worker(); - void deregisterFunction(final String tenant, - final String namespace, - final String componentName, - final String clientRole, + void deregisterFunction(String tenant, + String namespace, + String componentName, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); @Deprecated - default void deregisterFunction(final String tenant, - final String namespace, - final String componentName, - final String clientRole, + default void deregisterFunction(String tenant, + String namespace, + String componentName, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps) { deregisterFunction( tenant, @@ -159,12 +159,12 @@ void uploadFunction(InputStream uploadedInputStream, StreamingOutput downloadFunction(String path, String clientRole, - final AuthenticationDataSource clientAuthenticationDataHttps); + AuthenticationDataSource clientAuthenticationDataHttps); @Deprecated default StreamingOutput downloadFunction(String path, String clientRole, - final AuthenticationDataHttps clientAuthenticationDataHttps) { + AuthenticationDataHttps clientAuthenticationDataHttps) { return downloadFunction(path, clientRole, (AuthenticationDataSource) clientAuthenticationDataHttps); } @@ -172,14 +172,14 @@ StreamingOutput downloadFunction(String tenant, String namespace, String componentName, String clientRole, - final AuthenticationDataSource clientAuthenticationDataHttps); + AuthenticationDataSource clientAuthenticationDataHttps); @Deprecated default StreamingOutput downloadFunction(String tenant, String namespace, String componentName, String clientRole, - final AuthenticationDataHttps clientAuthenticationDataHttps) { + AuthenticationDataHttps clientAuthenticationDataHttps) { return downloadFunction( tenant, namespace, diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java index db8411151ed7a..d97a2856cd2c6 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Sinks.java @@ -49,14 +49,14 @@ public interface Sinks extends Component { * @param clientRole Client role for running the Pulsar Sink * @param clientAuthenticationDataHttps Authentication status of the http client */ - void registerSink(final String tenant, - final String namespace, - final String sinkName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sinkPkgUrl, - final SinkConfig sinkConfig, - final String clientRole, + void registerSink(String tenant, + String namespace, + String sinkName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sinkPkgUrl, + SinkConfig sinkConfig, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps); /** @@ -65,14 +65,14 @@ void registerSink(final String tenant, * the signature of the AuthenticationDataSource. */ @Deprecated - default void registerSink(final String tenant, - final String namespace, - final String sinkName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sinkPkgUrl, - final SinkConfig sinkConfig, - final String clientRole, + default void registerSink(String tenant, + String namespace, + String sinkName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sinkPkgUrl, + SinkConfig sinkConfig, + String clientRole, AuthenticationDataHttps clientAuthenticationDataHttps) { registerSink( tenant, @@ -99,14 +99,14 @@ default void registerSink(final String tenant, * @param clientAuthenticationDataHttps Authentication status of the http client * @param updateOptions Options while updating the sink */ - void updateSink(final String tenant, - final String namespace, - final String sinkName, - final InputStream uploadedInputStream, - final FormDataContentDisposition fileDetail, - final String sinkPkgUrl, - final SinkConfig sinkConfig, - final String clientRole, + void updateSink(String tenant, + String namespace, + String sinkName, + InputStream uploadedInputStream, + FormDataContentDisposition fileDetail, + String sinkPkgUrl, + SinkConfig sinkConfig, + String clientRole, AuthenticationDataSource clientAuthenticationDataHttps, UpdateOptionsImpl updateOptions); From 485e9618f8d021ee4e24d2e686704632ef92876c Mon Sep 17 00:00:00 2001 From: guangning Date: Mon, 28 Feb 2022 10:05:48 +0800 Subject: [PATCH 26/26] Fixed comment --- .../pulsar/client/impl/auth/AuthenticationDataBasic.java | 5 +++-- .../pulsar/client/impl/auth/AuthenticationDataToken.java | 5 +++-- .../client/impl/auth/oauth2/AuthenticationDataOAuth2.java | 5 +++-- 3 files changed, 9 insertions(+), 6 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java index 3ee5374210cc9..ba9e728a2d739 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataBasic.java @@ -30,13 +30,14 @@ public class AuthenticationDataBasic implements AuthenticationDataProvider { private static final String HTTP_HEADER_NAME = "Authorization"; private String httpAuthToken; private String commandAuthToken; - private final Map headers = new HashMap<>(); + private Map headers = new HashMap<>(); public AuthenticationDataBasic(String userId, String password) { httpAuthToken = "Basic " + Base64.getEncoder().encodeToString((userId + ":" + password).getBytes()); commandAuthToken = userId + ":" + password; headers.put(HTTP_HEADER_NAME, httpAuthToken); headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationBasic.AUTH_METHOD_NAME); + this.headers = Collections.unmodifiableMap(this.headers); } @Override @@ -46,7 +47,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return Collections.unmodifiableMap(this.headers).entrySet(); + return this.headers.entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java index bf105df32bf2a..b69222a57c9b0 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDataToken.java @@ -29,12 +29,13 @@ public class AuthenticationDataToken implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; private final Supplier tokenSupplier; - private final Map headers = new HashMap<>(); + private Map headers = new HashMap<>(); public AuthenticationDataToken(Supplier tokenSupplier) { this.tokenSupplier = tokenSupplier; headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationToken.AUTH_METHOD_NAME); headers.put(HTTP_HEADER_NAME, "Bearer " + getToken()); + this.headers = Collections.unmodifiableMap(this.headers); } @Override @@ -44,7 +45,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return Collections.unmodifiableMap(this.headers).entrySet(); + return this.headers.entrySet(); } @Override diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java index 7af587cd9ca1f..788f2d5ba251f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/oauth2/AuthenticationDataOAuth2.java @@ -31,12 +31,13 @@ class AuthenticationDataOAuth2 implements AuthenticationDataProvider { public static final String HTTP_HEADER_NAME = "Authorization"; private final String accessToken; - private final Map headers = new HashMap<>(); + private Map headers = new HashMap<>(); public AuthenticationDataOAuth2(String accessToken) { this.accessToken = accessToken; headers.put(HTTP_HEADER_NAME, "Bearer " + accessToken); headers.put(PULSAR_AUTH_METHOD_NAME, AuthenticationOAuth2.AUTH_METHOD_NAME); + this.headers = Collections.unmodifiableMap(this.headers); } @Override @@ -46,7 +47,7 @@ public boolean hasDataForHttp() { @Override public Set> getHttpHeaders() { - return Collections.unmodifiableMap(this.headers).entrySet(); + return this.headers.entrySet(); } @Override