diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java index 17c0300f112b4..53efb4029ce75 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java @@ -644,7 +644,7 @@ public void startFunction( @Consumes(MediaType.MULTIPART_FORM_DATA) public void uploadFunction(final @FormDataParam("data") InputStream uploadedInputStream, final @FormDataParam("path") String path) { - functions.uploadFunction(uploadedInputStream, path, clientAppId()); + functions.uploadFunction(uploadedInputStream, path, clientAppId(), clientAuthData()); } @GET @@ -711,6 +711,6 @@ public void updateFunctionOnWorkerLeader(final @PathParam("tenant") String tenan final @FormDataParam("delete") boolean delete) { functions.updateFunctionOnWorkerLeader(tenant, namespace, functionName, uploadedInputStream, - delete, uri.getRequestUri(), clientAppId()); + delete, uri.getRequestUri(), clientAppId(), clientAuthData()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SinksBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SinksBase.java index 40ef5a4bda35b..a58f9ce715190 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SinksBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SinksBase.java @@ -517,6 +517,6 @@ public List getSinkConfigDefinition( }) @Path("/reloadBuiltInSinks") public void reloadSinks() { - sink.reloadConnectors(clientAppId()); + sink.reloadConnectors(clientAppId(), clientAuthData()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SourcesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SourcesBase.java index 50859b7129872..495b9fb2a41f7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SourcesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SourcesBase.java @@ -507,6 +507,6 @@ public List getSourceConfigDefinition( }) @Path("/reloadBuiltInSources") public void reloadSources() { - source.reloadConnectors(clientAppId()); + source.reloadConnectors(clientAppId(), clientAuthData()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Functions.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Functions.java index 6c4b367f0cea5..38b27160f094c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Functions.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Functions.java @@ -302,7 +302,7 @@ public Response stopFunction(final @PathParam("tenant") String tenant, @Consumes(MediaType.MULTIPART_FORM_DATA) public Response uploadFunction(final @FormDataParam("data") InputStream uploadedInputStream, final @FormDataParam("path") String path) { - return functions.uploadFunction(uploadedInputStream, path, clientAppId()); + return functions.uploadFunction(uploadedInputStream, path, clientAppId(), clientAuthData()); } @GET 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 3bcea1b90f463..27fe4a7dba088 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 @@ -878,13 +878,13 @@ public List getListOfConnectors() { return this.worker().getConnectorsManager().getConnectors(); } - public void reloadConnectors(String clientRole) { + public void reloadConnectors(String clientRole, AuthenticationDataSource authenticationData) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); } if (worker().getWorkerConfig().isAuthorizationEnabled()) { // Only superuser has permission to do this operation. - if (!isSuperUser(clientRole)) { + if (!isSuperUser(clientRole, authenticationData)) { throw new RestException(Status.UNAUTHORIZED, "This operation requires super-user access"); } } @@ -1178,13 +1178,14 @@ public void putFunctionState(final String tenant, } } - public void uploadFunction(final InputStream uploadedInputStream, final String path, String clientRole) { + public void uploadFunction(final InputStream uploadedInputStream, final String path, String clientRole, + AuthenticationDataSource authenticationData) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); } - if (worker().getWorkerConfig().isAuthorizationEnabled() && !isSuperUser(clientRole)) { + if (worker().getWorkerConfig().isAuthorizationEnabled() && !isSuperUser(clientRole, authenticationData)) { throw new RestException(Status.UNAUTHORIZED, "client is not authorize to perform operation"); } @@ -1280,7 +1281,7 @@ public StreamingOutput downloadFunction(final String path, String clientRole, Au throw new RestException(Status.INTERNAL_SERVER_ERROR, e.getMessage()); } } else { - if (!isSuperUser(clientRole)) { + if (!isSuperUser(clientRole, clientAuthenticationDataHttps)) { throw new RestException(Status.UNAUTHORIZED, "client is not authorize to perform operation"); } } @@ -1430,7 +1431,7 @@ public boolean isAuthorizedRole(String tenant, String namespace, String clientRo AuthenticationDataSource authenticationData) throws PulsarAdminException { if (worker().getWorkerConfig().isAuthorizationEnabled()) { // skip authorization if client role is super-user - if (isSuperUser(clientRole)) { + if (isSuperUser(clientRole, authenticationData)) { return true; } @@ -1513,14 +1514,14 @@ protected void componentInstanceStatusRequestValidate (final String tenant, } } - public boolean isSuperUser(String clientRole) { + public boolean isSuperUser(String clientRole, AuthenticationDataSource authenticationData) { if (clientRole != null) { try { if ((worker().getWorkerConfig().getSuperUserRoles() != null && worker().getWorkerConfig().getSuperUserRoles().contains(clientRole))) { return true; } - return worker().getAuthorizationService().isSuperUser(clientRole, null) + return worker().getAuthorizationService().isSuperUser(clientRole, authenticationData) .get(worker().getWorkerConfig().getZooKeeperOperationTimeoutSeconds(), SECONDS); } catch (InterruptedException e) { log.warn("Time-out {} sec while checking the role {} is a super user role ", 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 5b0ec90ad5c6e..a53e89125fc12 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 @@ -655,14 +655,15 @@ public void updateFunctionOnWorkerLeader(final String tenant, final InputStream uploadedInputStream, final boolean delete, URI uri, - final String clientRole) { + final String clientRole, + AuthenticationDataSource authenticationData) { if (!isWorkerServiceAvailable()) { throwUnavailableException(); } if (worker().getWorkerConfig().isAuthorizationEnabled()) { - if (!isSuperUser(clientRole)) { + if (!isSuperUser(clientRole, authenticationData)) { log.error("{}/{}/{} Client [{}] is not superuser to update on worker leader {}", tenant, namespace, functionName, clientRole, ComponentTypeUtils.toString(componentType)); throw new RestException(Response.Status.UNAUTHORIZED, "client is not authorize to perform operation"); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplV2.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplV2.java index 86118cea3b208..d99ca6445732c 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplV2.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplV2.java @@ -20,6 +20,7 @@ import com.google.gson.Gson; import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.common.functions.FunctionConfig; import org.apache.pulsar.common.functions.FunctionState; import org.apache.pulsar.common.io.ConnectorDefinition; @@ -177,8 +178,9 @@ public Response stopFunctionInstances(String tenant, String namespace, String fu return Response.ok().build(); } - public Response uploadFunction(InputStream uploadedInputStream, String path, String clientRole) { - delegate.uploadFunction(uploadedInputStream, path, clientRole); + public Response uploadFunction(InputStream uploadedInputStream, String path, String clientRole, + AuthenticationDataSource authenticationData) { + delegate.uploadFunction(uploadedInputStream, path, clientRole, authenticationData); return Response.ok().build(); } @@ -229,4 +231,4 @@ private InstanceCommunication.FunctionStatus toProto( return functionStatus; } -} \ No newline at end of file +} diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionsApiV2Resource.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionsApiV2Resource.java index cd8acda4a24b6..bcfdcaf758962 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionsApiV2Resource.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionsApiV2Resource.java @@ -296,7 +296,7 @@ public Response stopFunction(final @PathParam("tenant") String tenant, @Consumes(MediaType.MULTIPART_FORM_DATA) public Response uploadFunction(final @FormDataParam("data") InputStream uploadedInputStream, final @FormDataParam("path") String path) { - return functions.uploadFunction(uploadedInputStream, path, clientAppId()); + return functions.uploadFunction(uploadedInputStream, path, clientAppId(), clientAuthData()); } @GET diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java index 0fe8c480aadcd..2fad782850193 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java @@ -312,7 +312,7 @@ public void startFunction(final @PathParam("tenant") String tenant, @Consumes(MediaType.MULTIPART_FORM_DATA) public void uploadFunction(final @FormDataParam("data") InputStream uploadedInputStream, final @FormDataParam("path") String path) { - functions.uploadFunction(uploadedInputStream, path, clientAppId()); + functions.uploadFunction(uploadedInputStream, path, clientAppId(), clientAuthData()); } @GET @@ -386,6 +386,6 @@ public void updateFunctionOnWorkerLeader(final @PathParam("tenant") String tenan final @FormDataParam("delete") boolean delete) { functions.updateFunctionOnWorkerLeader(tenant, namespace, functionName, uploadedInputStream, - delete, uri.getRequestUri(), clientAppId()); + delete, uri.getRequestUri(), clientAppId(), clientAuthData()); } } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SinksApiV3Resource.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SinksApiV3Resource.java index 8ee166241cf91..7fd7ef85d6b60 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SinksApiV3Resource.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SinksApiV3Resource.java @@ -273,6 +273,6 @@ public List getSinkConfigDefinition( }) @Path("/reloadBuiltInSinks") public void reloadSinks() { - sink.reloadConnectors(clientAppId()); + sink.reloadConnectors(clientAppId(), clientAuthData()); } } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SourcesApiV3Resource.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SourcesApiV3Resource.java index 6b3daa98b5269..cb4356fe221d4 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SourcesApiV3Resource.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/SourcesApiV3Resource.java @@ -289,6 +289,6 @@ public List getSourceConfigDefinition( }) @Path("/reloadBuiltInSources") public void reloadSources() { - source.reloadConnectors(clientAppId()); + source.reloadConnectors(clientAppId(), clientAuthData()); } } diff --git a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java index e8618c254c9ea..7bf94d738a2ef 100644 --- a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java +++ b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java @@ -250,9 +250,9 @@ public void testIsAuthorizedRole() throws PulsarAdminException, InterruptedExcep // test pulsar super user final String pulsarSuperUser = "pulsarSuperUser"; - when(authorizationService.isSuperUser(pulsarSuperUser, null)).thenReturn(CompletableFuture.completedFuture(true)); + when(authorizationService.isSuperUser(eq(pulsarSuperUser), any())).thenReturn(CompletableFuture.completedFuture(true)); assertTrue(functionImpl.isAuthorizedRole("test-tenant", "test-ns", pulsarSuperUser, authenticationDataSource)); - assertTrue(functionImpl.isSuperUser(pulsarSuperUser)); + assertTrue(functionImpl.isSuperUser(pulsarSuperUser, null)); // test normal user functionImpl = spy(new FunctionsImpl(() -> mockedWorkerService)); @@ -263,7 +263,7 @@ public void testIsAuthorizedRole() throws PulsarAdminException, InterruptedExcep when(admin.tenants()).thenReturn(tenants); when(this.mockedWorkerService.getBrokerAdmin()).thenReturn(admin); when(authorizationService.isTenantAdmin("test-tenant", "test-user", tenantInfo, authenticationDataSource)).thenReturn(CompletableFuture.completedFuture(false)); - when(authorizationService.isSuperUser("test-user", null)).thenReturn(CompletableFuture.completedFuture(false)); + when(authorizationService.isSuperUser(eq("test-user"), any())).thenReturn(CompletableFuture.completedFuture(false)); assertFalse(functionImpl.isAuthorizedRole("test-tenant", "test-ns", "test-user", authenticationDataSource)); // if user is tenant admin @@ -322,10 +322,10 @@ public void testIsSuperUser() throws PulsarAdminException { }); AuthenticationDataSource authenticationDataSource = mock(AuthenticationDataSource.class); - assertTrue(functionImpl.isSuperUser(superUser)); + assertTrue(functionImpl.isSuperUser(superUser, null)); - assertFalse(functionImpl.isSuperUser("normal-user")); - assertFalse(functionImpl.isSuperUser( null)); + assertFalse(functionImpl.isSuperUser("normal-user", null)); + assertFalse(functionImpl.isSuperUser( null, null)); // test super roles is null and it's not a pulsar super user when(authorizationService.isSuperUser(superUser, null)) @@ -334,7 +334,7 @@ public void testIsSuperUser() throws PulsarAdminException { workerConfig = new WorkerConfig(); workerConfig.setAuthorizationEnabled(true); doReturn(workerConfig).when(mockedWorkerService).getWorkerConfig(); - assertFalse(functionImpl.isSuperUser(superUser)); + assertFalse(functionImpl.isSuperUser(superUser, null)); } public static FunctionConfig createDefaultFunctionConfig() {