This is an automated email from the ASF dual-hosted git repository.
pvillard31 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 02440436538 NIFI-16272 Aligned Process Group Authorization for
Connector Migration (#11606)
02440436538 is described below
commit 024404365389a469d8c5bdb45ee51455176be843
Author: David Handermann <[email protected]>
AuthorDate: Tue Sep 1 02:49:00 2026 -0500
NIFI-16272 Aligned Process Group Authorization for Connector Migration
(#11606)
---
.../apache/nifi/web/StandardNiFiServiceFacade.java | 19 +++++-
.../org/apache/nifi/web/api/ConnectorResource.java | 11 +++-
.../nifi/web/StandardNiFiServiceFacadeTest.java | 53 +++++++++++++++
.../apache/nifi/web/api/TestConnectorResource.java | 75 ++++++++++++++++++++++
4 files changed, 153 insertions(+), 5 deletions(-)
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java
index abf4c573f16..61295122ea2 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java
@@ -40,6 +40,7 @@ import
org.apache.nifi.authorization.AuthorizationResult.Result;
import org.apache.nifi.authorization.AuthorizeAccess;
import org.apache.nifi.authorization.Authorizer;
import org.apache.nifi.authorization.Group;
+import org.apache.nifi.authorization.ProcessGroupAuthorizable;
import org.apache.nifi.authorization.RequestAction;
import org.apache.nifi.authorization.Resource;
import org.apache.nifi.authorization.User;
@@ -4209,9 +4210,21 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
@Override
public VersionedFlowMigrationSourcesEntity
getConnectorMigrationSources(final String connectorId) {
- final List<VersionedFlowMigrationSourceDTO> migrationSources =
connectorDAO.getMigrationSources(connectorId).stream()
- .map(this::createVersionedFlowMigrationSourceDto)
- .toList();
+ final NiFiUser user = NiFiUserUtils.getNiFiUser();
+ final List<VersionedFlowMigrationSourceDTO> migrationSources = new
ArrayList<>();
+
+ for (final ConnectorMigrationSource migrationSource :
connectorDAO.getMigrationSources(connectorId)) {
+ final ProcessGroupAuthorizable processGroupAuthorizable;
+ try {
+ processGroupAuthorizable =
authorizableLookup.getProcessGroup(migrationSource.getProcessGroupId());
+ } catch (final ResourceNotFoundException e) {
+ continue;
+ }
+
+ if
(processGroupAuthorizable.getAuthorizable().isAuthorized(authorizer,
RequestAction.READ, user)) {
+
migrationSources.add(createVersionedFlowMigrationSourceDto(migrationSource));
+ }
+ }
final VersionedFlowMigrationSourcesEntity entity = new
VersionedFlowMigrationSourcesEntity();
entity.setMigrationSources(migrationSources);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java
index c0de37b3f76..29ef824a5fd 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java
@@ -49,6 +49,7 @@ import org.apache.commons.lang3.exception.ExceptionUtils;
import org.apache.nifi.asset.Asset;
import org.apache.nifi.authorization.AuthorizeConnectorConfigReferences;
import org.apache.nifi.authorization.Authorizer;
+import org.apache.nifi.authorization.ProcessGroupAuthorizable;
import org.apache.nifi.authorization.RequestAction;
import org.apache.nifi.authorization.resource.Authorizable;
import org.apache.nifi.authorization.resource.DataAuthorizable;
@@ -2429,7 +2430,8 @@ public class ConnectorResource extends
ApplicationResource {
@ApiResponse(responseCode = "409", description = "The
request was valid but NiFi was not in the appropriate state to process it.")
},
security = {
- @SecurityRequirement(name = "Read - /connectors/{uuid}")
+ @SecurityRequirement(name = "Read - /connectors/{uuid}"),
+ @SecurityRequirement(name = "Read - /process-groups/{uuid}
- For each Process Group returned")
}
)
public Response getMigrationSources(@PathParam("id") final String
connectorId) {
@@ -2528,7 +2530,8 @@ public class ConnectorResource extends
ApplicationResource {
@ApiResponse(responseCode = "409", description = "The
request was valid but NiFi was not in the appropriate state to process it.")
},
security = {
- @SecurityRequirement(name = "Write - /connectors/{uuid}")
+ @SecurityRequirement(name = "Write - /connectors/{uuid}"),
+ @SecurityRequirement(name = "Write -
/process-groups/{uuid} - For a migration from a local Process Group")
}
)
public Response createMigrationRequest(@PathParam("id") final String
connectorId, final MigrationRequestEntity requestEntity) {
@@ -2573,6 +2576,10 @@ public class ConnectorResource extends
ApplicationResource {
lookup -> {
final Authorizable connector =
lookup.getConnector(connectorId);
connector.authorize(authorizer, RequestAction.WRITE, user);
+ if (hasLocalSource) {
+ final ProcessGroupAuthorizable sourceGroupAuthorizable
= lookup.getProcessGroup(request.getLocalSource().getProcessGroupId());
+ authorizeProcessGroup(sourceGroupAuthorizable,
authorizer, lookup, RequestAction.WRITE, true, false, false, false, true);
+ }
},
() -> {
// Verify the target Connector is ready (stopped and
unmodified from its initial flow) for every
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java
index dab62bdaecc..97f7d20ca0b 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java
@@ -29,6 +29,8 @@ import org.apache.nifi.authorization.AuthorizationResult;
import org.apache.nifi.authorization.Authorizer;
import org.apache.nifi.authorization.ComponentAuthorizable;
import org.apache.nifi.authorization.Group;
+import org.apache.nifi.authorization.ProcessGroupAuthorizable;
+import org.apache.nifi.authorization.RequestAction;
import org.apache.nifi.authorization.Resource;
import org.apache.nifi.authorization.User;
import org.apache.nifi.authorization.resource.Authorizable;
@@ -43,6 +45,7 @@ import org.apache.nifi.components.BacklogReportingException;
import org.apache.nifi.components.PropertyDescriptor;
import org.apache.nifi.components.connector.BacklogReportingConnector;
import org.apache.nifi.components.connector.Connector;
+import org.apache.nifi.components.connector.ConnectorMigrationSource;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorState;
import org.apache.nifi.components.connector.ConnectorSyncMode;
@@ -139,6 +142,7 @@ import org.apache.nifi.web.api.dto.ProcessGroupDTO;
import org.apache.nifi.web.api.dto.RemoteProcessGroupDTO;
import org.apache.nifi.web.api.dto.RevisionDTO;
import org.apache.nifi.web.api.dto.VersionControlInformationDTO;
+import org.apache.nifi.web.api.dto.VersionedFlowMigrationSourceDTO;
import org.apache.nifi.web.api.dto.action.HistoryDTO;
import org.apache.nifi.web.api.dto.action.HistoryQueryDTO;
import org.apache.nifi.web.api.dto.search.SearchResultsDTO;
@@ -161,6 +165,7 @@ import org.apache.nifi.web.api.entity.StatusHistoryEntity;
import org.apache.nifi.web.api.entity.TenantEntity;
import org.apache.nifi.web.api.entity.TenantsEntity;
import org.apache.nifi.web.api.entity.VersionControlInformationEntity;
+import org.apache.nifi.web.api.entity.VersionedFlowMigrationSourcesEntity;
import org.apache.nifi.web.controller.ControllerFacade;
import org.apache.nifi.web.dao.ComponentStateDAO;
import org.apache.nifi.web.dao.ConnectorDAO;
@@ -2114,6 +2119,54 @@ public class StandardNiFiServiceFacadeTest {
assertThrows(ResourceNotFoundException.class, () ->
serviceFacade.getConnectorProcessorState(connectorId, processorId));
}
+ @Test
+ public void
testGetConnectorMigrationSourcesReturnsAuthorizedProcessGroups() {
+ final Authentication authentication = new NiFiAuthenticationToken(new
NiFiUserDetails(new Builder().identity(USER_1).build()));
+ SecurityContextHolder.getContext().setAuthentication(authentication);
+
+ final String connectorId = "connector-id";
+ final String authorizedGroupId = "authorized-group-id";
+ final String deniedGroupId = "denied-group-id";
+ final String removedGroupId = "removed-group-id";
+
+ final ConnectorDAO connectorDAO = mock(ConnectorDAO.class);
+ serviceFacade.setConnectorDAO(connectorDAO);
+ when(connectorDAO.getMigrationSources(connectorId)).thenReturn(List.of(
+ createMigrationSource(authorizedGroupId),
+ createMigrationSource(deniedGroupId),
+ createMigrationSource(removedGroupId)));
+
+ final ProcessGroupAuthorizable authorizedGroup =
createProcessGroupAuthorizable(true);
+ final ProcessGroupAuthorizable deniedGroup =
createProcessGroupAuthorizable(false);
+
+ final AuthorizableLookup migrationLookup =
mock(AuthorizableLookup.class);
+ serviceFacade.setAuthorizableLookup(migrationLookup);
+
when(migrationLookup.getProcessGroup(authorizedGroupId)).thenReturn(authorizedGroup);
+
when(migrationLookup.getProcessGroup(deniedGroupId)).thenReturn(deniedGroup);
+ when(migrationLookup.getProcessGroup(removedGroupId)).thenThrow(new
ResourceNotFoundException("Process Group was removed"));
+
+ final VersionedFlowMigrationSourcesEntity entity =
serviceFacade.getConnectorMigrationSources(connectorId);
+
+ final List<VersionedFlowMigrationSourceDTO> migrationSources =
entity.getMigrationSources();
+ assertEquals(1, migrationSources.size());
+ assertEquals(authorizedGroupId,
migrationSources.getFirst().getProcessGroupId());
+ }
+
+ private ConnectorMigrationSource createMigrationSource(final String
processGroupId) {
+ final ConnectorMigrationSource migrationSource = new
ConnectorMigrationSource();
+ migrationSource.setProcessGroupId(processGroupId);
+ return migrationSource;
+ }
+
+ private ProcessGroupAuthorizable createProcessGroupAuthorizable(final
boolean readAuthorized) {
+ final Authorizable authorizable = mock(Authorizable.class);
+ when(authorizable.isAuthorized(any(Authorizer.class),
eq(RequestAction.READ), any(NiFiUser.class))).thenReturn(readAuthorized);
+
+ final ProcessGroupAuthorizable processGroupAuthorizable =
mock(ProcessGroupAuthorizable.class);
+
when(processGroupAuthorizable.getAuthorizable()).thenReturn(authorizable);
+ return processGroupAuthorizable;
+ }
+
@Test
public void testVerifyCanClearConnectorProcessorState() {
final String connectorId = "connector-id";
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestConnectorResource.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestConnectorResource.java
index fb32e890f03..13eca09c528 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestConnectorResource.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestConnectorResource.java
@@ -26,6 +26,7 @@ import org.apache.nifi.authorization.AccessDeniedException;
import org.apache.nifi.authorization.AuthorizableLookup;
import org.apache.nifi.authorization.AuthorizeAccess;
import org.apache.nifi.authorization.Authorizer;
+import org.apache.nifi.authorization.ProcessGroupAuthorizable;
import org.apache.nifi.authorization.RequestAction;
import org.apache.nifi.authorization.resource.Authorizable;
import org.apache.nifi.authorization.user.NiFiUser;
@@ -60,6 +61,7 @@ import
org.apache.nifi.web.api.entity.ConnectorPropertyAllowableValuesEntity;
import org.apache.nifi.web.api.entity.ConnectorRunStatusEntity;
import org.apache.nifi.web.api.entity.ControllerServiceEntity;
import org.apache.nifi.web.api.entity.ControllerServicesEntity;
+import org.apache.nifi.web.api.entity.MigrationPayloadEntity;
import org.apache.nifi.web.api.entity.MigrationRequestEntity;
import org.apache.nifi.web.api.entity.ParameterContextEntity;
import org.apache.nifi.web.api.entity.ParameterEntity;
@@ -80,10 +82,12 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.security.core.Authentication;
import org.springframework.security.core.context.SecurityContextHolder;
+import java.io.ByteArrayInputStream;
import java.io.File;
import java.io.IOException;
import java.io.InputStream;
import java.net.URI;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.util.Arrays;
import java.util.List;
@@ -261,6 +265,77 @@ public class TestConnectorResource {
verify(serviceFacade, never()).verifyCanMigrateConnector(anyString(),
anyString());
}
+ @Test
+ public void
testCreateMigrationRequestDeniedWhenLocalSourceProcessGroupNotAuthorized() {
+ authenticate();
+ final ConnectorResource spyResource = spy(connectorResource);
+ doReturn(false).when(spyResource).isReplicateRequest();
+
+ final AuthorizableLookup lookup = wireMigrationAuthorization();
+ final ProcessGroupAuthorizable sourceGroup =
mock(ProcessGroupAuthorizable.class);
+ final Authorizable sourceGroupAuthorizable = mock(Authorizable.class);
+
when(sourceGroup.getAuthorizable()).thenReturn(sourceGroupAuthorizable);
+ when(lookup.getProcessGroup(PROCESS_GROUP_ID)).thenReturn(sourceGroup);
+ doThrow(new AccessDeniedException("Not authorized to write the source
Process Group"))
+ .when(sourceGroupAuthorizable).authorize(any(),
eq(RequestAction.WRITE), any());
+
+ final MigrationRequestDTO requestDto = new MigrationRequestDTO();
+ requestDto.setConnectorId(CONNECTOR_ID);
+
requestDto.setLocalSource(createLocalMigrationSource(PROCESS_GROUP_ID));
+
+ final MigrationRequestEntity requestEntity = new
MigrationRequestEntity();
+ requestEntity.setRequest(requestDto);
+
+ assertThrows(AccessDeniedException.class, () ->
spyResource.createMigrationRequest(CONNECTOR_ID, requestEntity));
+
+ verify(lookup).getProcessGroup(PROCESS_GROUP_ID);
+ verify(serviceFacade,
never()).verifyConnectorReadyForMigration(anyString());
+ verify(serviceFacade, never()).verifyCanMigrateConnector(anyString(),
anyString());
+ }
+
+ @Test
+ public void
testCreateMigrationRequestFromUploadedPayloadDoesNotAuthorizeProcessGroup()
throws IOException {
+ authenticate();
+ final ConnectorResource spyResource = spy(connectorResource);
+ doReturn(false).when(spyResource).isReplicateRequest();
+
+ final AuthorizableLookup lookup = wireMigrationAuthorization();
+
+ final String payloadId;
+ try (Response payloadResponse =
spyResource.createMigrationPayload(CONNECTOR_ID, new
ByteArrayInputStream("{}".getBytes(StandardCharsets.UTF_8)))) {
+ payloadId = ((MigrationPayloadEntity)
payloadResponse.getEntity()).getPayload().getPayloadId();
+ }
+
+ final MigrationRequestDTO requestDto = new MigrationRequestDTO();
+ requestDto.setConnectorId(CONNECTOR_ID);
+ requestDto.setPayloadId(payloadId);
+
+ final MigrationRequestEntity requestEntity = new
MigrationRequestEntity();
+ requestEntity.setRequest(requestDto);
+
+ // Stopping at the readiness check keeps the migration from being
submitted asynchronously while still
+ // exercising the authorization callback, which runs before
verification.
+ doThrow(new IllegalStateException("Connector must be stopped before it
can be migrated"))
+
.when(serviceFacade).verifyConnectorReadyForMigration(CONNECTOR_ID);
+
+ assertThrows(IllegalStateException.class, () ->
spyResource.createMigrationRequest(CONNECTOR_ID, requestEntity));
+
+ verify(lookup, never()).getProcessGroup(anyString());
+ }
+
+ private AuthorizableLookup wireMigrationAuthorization() {
+ final AuthorizableLookup lookup = mock(AuthorizableLookup.class);
+
when(lookup.getConnector(CONNECTOR_ID)).thenReturn(mock(Authorizable.class));
+
+ doAnswer(invocation -> {
+ final AuthorizeAccess authorizeAccess = invocation.getArgument(0);
+ authorizeAccess.authorize(lookup);
+ return null;
+ }).when(serviceFacade).authorizeAccess(any(AuthorizeAccess.class));
+
+ return lookup;
+ }
+
@Test
public void testCreateMigrationRequestRequiresLocalSourceOrPayload() {
final MigrationRequestDTO requestDto = new MigrationRequestDTO();