From e7947c8674a870b85145777cc755fec282dff550 Mon Sep 17 00:00:00 2001 From: Matt Gilman Date: Mon, 14 Sep 2026 16:39:29 -0400 Subject: [PATCH 1/2] NIFI-16337: Add support for stopping Process Group sources - add the Stop sources REST endpoint and cluster response handling - expose resolved execution engines on Process Group payloads - add Flow Designer actions for stopping source components - reject STATELESS targets and skip nested STATELESS groups - add backend and frontend test coverage --- .../nifi/web/api/dto/ProcessGroupDTO.java | 12 + .../web/api/dto/flow/ProcessGroupFlowDTO.java | 12 + .../web/api/entity/ProcessGroupEntity.java | 13 + .../http/StandardHttpResponseMapper.java | 2 + .../endpoints/StopSourcesEndpointMerger.java | 60 ++++ .../StopSourcesEndpointMergerTest.java | 92 ++++++ .../apache/nifi/web/NiFiServiceFacade.java | 30 ++ .../nifi/web/StandardNiFiServiceFacade.java | 61 ++++ .../org/apache/nifi/web/api/FlowResource.java | 119 +++++++ .../apache/nifi/web/api/dto/DtoFactory.java | 4 + .../nifi/web/api/dto/EntityFactory.java | 1 + .../web/StandardNiFiServiceFacadeTest.java | 290 ++++++++++++++++++ .../apache/nifi/web/api/TestFlowResource.java | 251 +++++++++++++++ .../nifi/web/api/dto/DtoFactoryTest.java | 106 +++++++ .../nifi/web/api/dto/EntityFactoryTest.java | 58 ++++ .../canvas-context-menu.service.spec.ts | 175 +++++++++++ .../service/canvas-context-menu.service.ts | 31 ++ .../service/canvas-utils.service.ts | 21 +- .../flow-designer/service/flow.service.ts | 14 + .../flow-designer/state/flow/flow.actions.ts | 8 + .../state/flow/flow.effects.spec.ts | 104 ++++++- .../flow-designer/state/flow/flow.effects.ts | 58 ++++ .../flow-designer/state/flow/flow.reducer.ts | 5 + .../state/flow/flow.selectors.ts | 5 + .../pages/flow-designer/state/flow/index.ts | 14 + .../system/pg/ClusteredStopSourcesIT.java | 260 ++++++++++++++++ .../nifi/toolkit/client/FlowClient.java | 10 + .../toolkit/client/impl/JerseyFlowClient.java | 26 ++ 28 files changed, 1840 insertions(+), 2 deletions(-) create mode 100644 nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/http/endpoints/StopSourcesEndpointMerger.java create mode 100644 nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/test/java/org/apache/nifi/cluster/coordination/http/endpoints/StopSourcesEndpointMergerTest.java create mode 100644 nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/EntityFactoryTest.java create mode 100644 nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/pg/ClusteredStopSourcesIT.java diff --git a/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/ProcessGroupDTO.java b/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/ProcessGroupDTO.java index 1a830d5420e2..17cfa451c51b 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/ProcessGroupDTO.java +++ b/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/ProcessGroupDTO.java @@ -38,6 +38,7 @@ public class ProcessGroupDTO extends ComponentDTO { private String defaultBackPressureDataSizeThreshold; private String logFileSuffix; private String executionEngine; + private String resolvedExecutionEngine; private Integer maxConcurrentTasks; private String statelessFlowTimeout; @@ -395,6 +396,17 @@ public void setExecutionEngine(final String executionEngine) { this.executionEngine = executionEngine; } + @Schema(description = "The Execution Engine that will actually run this Process Group after resolving INHERITED. Never INHERITED; a root group configured as INHERITED resolves to STANDARD.", + allowableValues = {"STATELESS", "STANDARD"}, + accessMode = Schema.AccessMode.READ_ONLY) + public String getResolvedExecutionEngine() { + return resolvedExecutionEngine; + } + + public void setResolvedExecutionEngine(final String resolvedExecutionEngine) { + this.resolvedExecutionEngine = resolvedExecutionEngine; + } + @Schema(description = "If the Process Group is configured to run in using the Stateless Engine, represents the current state. Otherwise, will be STOPPED.", allowableValues = {"STOPPED", "RUNNING"}) public String getStatelessGroupScheduledState() { diff --git a/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/flow/ProcessGroupFlowDTO.java b/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/flow/ProcessGroupFlowDTO.java index f2bdaea9db9e..a430d88c9d20 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/flow/ProcessGroupFlowDTO.java +++ b/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/dto/flow/ProcessGroupFlowDTO.java @@ -38,6 +38,7 @@ public class ProcessGroupFlowDTO { private FlowBreadcrumbEntity breadcrumb; private FlowDTO flow; private Date lastRefreshed; + private String resolvedExecutionEngine; /** * @return contents of this process group. This field will be populated if the request is marked verbose @@ -130,4 +131,15 @@ public ParameterContextReferenceEntity getParameterContext() { public void setParameterContext(ParameterContextReferenceEntity parameterContext) { this.parameterContext = parameterContext; } + + @Schema(description = "The Execution Engine that will actually run this Process Group after resolving INHERITED. Never INHERITED; a root group configured as INHERITED resolves to STANDARD.", + allowableValues = {"STATELESS", "STANDARD"}, + accessMode = Schema.AccessMode.READ_ONLY) + public String getResolvedExecutionEngine() { + return resolvedExecutionEngine; + } + + public void setResolvedExecutionEngine(final String resolvedExecutionEngine) { + this.resolvedExecutionEngine = resolvedExecutionEngine; + } } diff --git a/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/entity/ProcessGroupEntity.java b/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/entity/ProcessGroupEntity.java index 70fddd5eaf82..21eb63d51ffd 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/entity/ProcessGroupEntity.java +++ b/nifi-framework-bundle/nifi-framework/nifi-client-dto/src/main/java/org/apache/nifi/web/api/entity/ProcessGroupEntity.java @@ -56,6 +56,7 @@ public class ProcessGroupEntity extends ComponentEntity implements Permissible

{ + public static final Pattern STOP_SOURCES_URI_PATTERN = Pattern.compile("/nifi-api/flow/process-groups/(?:(?:root)|(?:[a-f0-9\\-]{36}))/sources"); + + @Override + public boolean canHandle(final URI uri, final String method) { + return "PUT".equalsIgnoreCase(method) && STOP_SOURCES_URI_PATTERN.matcher(uri.getPath()).matches(); + } + + @Override + protected Class getEntityClass() { + return ScheduleComponentsEntity.class; + } + + @Override + protected void mergeResponses(final ScheduleComponentsEntity clientEntity, final Map entityMap, + final Set successfulResponses, final Set problematicResponses) { + if (clientEntity.getComponents() == null) { + clientEntity.setComponents(new HashMap<>()); + } + + for (final ScheduleComponentsEntity nodeEntity : entityMap.values()) { + if (nodeEntity.getComponents() == null) { + continue; + } + + for (final Map.Entry entry : nodeEntity.getComponents().entrySet()) { + clientEntity.getComponents().putIfAbsent(entry.getKey(), entry.getValue()); + } + } + } +} diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/test/java/org/apache/nifi/cluster/coordination/http/endpoints/StopSourcesEndpointMergerTest.java b/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/test/java/org/apache/nifi/cluster/coordination/http/endpoints/StopSourcesEndpointMergerTest.java new file mode 100644 index 000000000000..3f0f48b3b6f8 --- /dev/null +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/test/java/org/apache/nifi/cluster/coordination/http/endpoints/StopSourcesEndpointMergerTest.java @@ -0,0 +1,92 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.nifi.cluster.coordination.http.endpoints; + +import org.apache.nifi.cluster.protocol.NodeIdentifier; +import org.apache.nifi.web.api.dto.RevisionDTO; +import org.apache.nifi.web.api.entity.ScheduleComponentsEntity; +import org.junit.jupiter.api.Test; + +import java.net.URI; +import java.util.HashMap; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class StopSourcesEndpointMergerTest { + + private static final String GROUP_ID = "12345678-1234-1234-1234-123456789012"; + + @Test + public void testCanHandle() { + final StopSourcesEndpointMerger merger = new StopSourcesEndpointMerger(); + + assertTrue(merger.canHandle(URI.create("/nifi-api/flow/process-groups/" + GROUP_ID + "/sources"), "PUT")); + assertTrue(merger.canHandle(URI.create("/nifi-api/flow/process-groups/root/sources"), "PUT")); + assertTrue(merger.canHandle(URI.create("/nifi-api/flow/process-groups/" + GROUP_ID + "/sources"), "put")); + + assertFalse(merger.canHandle(URI.create("/nifi-api/flow/process-groups/" + GROUP_ID + "/sources"), "GET")); + assertFalse(merger.canHandle(URI.create("/nifi-api/flow/process-groups/" + GROUP_ID), "PUT")); + assertFalse(merger.canHandle(URI.create("/nifi-api/flow/process-groups/" + GROUP_ID), "GET")); + assertFalse(merger.canHandle(URI.create("/nifi-api/process-groups/" + GROUP_ID + "/sources"), "PUT")); + } + + @Test + public void testMergeResponsesUnionsComponentsAndKeepsClientRevision() { + final StopSourcesEndpointMerger merger = new StopSourcesEndpointMerger(); + + final RevisionDTO clientRevision = revision(1L, "client"); + final RevisionDTO nodeRevision = revision(2L, "node-2"); + final RevisionDTO extraRevision = revision(3L, "node-2-extra"); + + final ScheduleComponentsEntity clientEntity = new ScheduleComponentsEntity(); + clientEntity.setId(GROUP_ID); + clientEntity.setState("STOPPED"); + clientEntity.setComponents(new HashMap<>(Map.of("source-1", clientRevision))); + + final NodeIdentifier node1 = new NodeIdentifier("node1", "localhost", 8080, "localhost", 8081, "localhost", 8082, 8083, false); + final NodeIdentifier node2 = new NodeIdentifier("node2", "localhost", 8090, "localhost", 8091, "localhost", 8092, 8093, false); + + final ScheduleComponentsEntity node1Entity = new ScheduleComponentsEntity(); + node1Entity.setId(GROUP_ID); + node1Entity.setState("STOPPED"); + node1Entity.setComponents(Map.of("source-1", clientRevision)); + + final ScheduleComponentsEntity node2Entity = new ScheduleComponentsEntity(); + node2Entity.setId(GROUP_ID); + node2Entity.setState("STOPPED"); + node2Entity.setComponents(Map.of("source-1", nodeRevision, "source-2", extraRevision)); + + merger.mergeResponses(clientEntity, Map.of(node1, node1Entity, node2, node2Entity), null, null); + + assertEquals(GROUP_ID, clientEntity.getId()); + assertEquals("STOPPED", clientEntity.getState()); + assertEquals(2, clientEntity.getComponents().size()); + assertSame(clientRevision, clientEntity.getComponents().get("source-1")); + assertSame(extraRevision, clientEntity.getComponents().get("source-2")); + } + + private static RevisionDTO revision(final long version, final String clientId) { + final RevisionDTO revision = new RevisionDTO(); + revision.setVersion(version); + revision.setClientId(clientId); + return revision; + } +} diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/NiFiServiceFacade.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/NiFiServiceFacade.java index 2313b34f977c..2f9b8538d250 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/NiFiServiceFacade.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/NiFiServiceFacade.java @@ -387,6 +387,36 @@ Set getConnectorControllerServices(String connectorId, */ Set getRevisionsFromGroup(String groupId, Function> getComponents); + /** + * Identifies source component identifiers in the given Process Group and its Standard-engine descendants. + * Sources are processors with no non-loop incoming connection, Remote Process Group output ports, + * and public input ports with no non-loop incoming connection. Components owned by Process Groups + * that resolve to the Stateless Execution Engine are excluded. + * + * @param group the process group to search + * @return identifiers of source components + */ + Set findSourceComponentIds(ProcessGroup group); + + /** + * Verifies that source components can be stopped in the specified Process Group. + * + * @param groupId process group identifier + * @throws IllegalStateException when the Process Group resolves to the Stateless Execution Engine + */ + void verifyStopSources(String groupId); + + /** + * Verifies that the supplied component identifiers exactly match the source components currently identified + * in the specified Process Group. This ensures that every cluster node operates on the same source components. + * + * @param groupId process group identifier + * @param componentIds source component identifiers + * @throws IllegalStateException when the Process Group resolves to the Stateless Execution Engine or the supplied + * component identifiers do not match the source components currently identified + */ + void verifyStopSources(String groupId, Set componentIds); + /** * Gets the revisions from the specified snippet. * 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 b178d39d6769..8e2236e34ce0 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 @@ -234,6 +234,7 @@ import org.apache.nifi.registry.flow.mapping.InstantiatedVersionedRemoteGroupPort; import org.apache.nifi.registry.flow.mapping.VersionedComponentFlowMapper; import org.apache.nifi.registry.flow.mapping.VersionedComponentStateLookup; +import org.apache.nifi.remote.PublicPort; import org.apache.nifi.remote.RemoteGroupPort; import org.apache.nifi.reporting.Bulletin; import org.apache.nifi.reporting.BulletinQuery; @@ -242,6 +243,7 @@ import org.apache.nifi.reporting.ReportingTask; import org.apache.nifi.reporting.VerifiableReportingTask; import org.apache.nifi.util.BundleUtils; +import org.apache.nifi.util.Connectables; import org.apache.nifi.util.FlowDifferenceFilters; import org.apache.nifi.util.NiFiProperties; import org.apache.nifi.util.StringUtils; @@ -595,6 +597,65 @@ public Set getRevisionsFromGroup(final String groupId, final Function< return componentIds.stream().map(id -> revisionManager.getRevision(id)).collect(Collectors.toSet()); } + @Override + public Set findSourceComponentIds(final ProcessGroup group) { + final Set sourceIds = new LinkedHashSet<>(); + + for (final ProcessorNode processor : group.findAllProcessors()) { + if (processor.getProcessGroup().resolveExecutionEngine() == ExecutionEngine.STANDARD + && ProcessGroup.STOP_PROCESSORS_FILTER.test(processor) + && !Connectables.hasNonLoopConnection(processor)) { + sourceIds.add(processor.getIdentifier()); + } + } + + for (final RemoteProcessGroup remoteProcessGroup : group.findAllRemoteProcessGroups()) { + if (remoteProcessGroup.getProcessGroup().resolveExecutionEngine() == ExecutionEngine.STATELESS) { + continue; + } + + for (final RemoteGroupPort remotePort : remoteProcessGroup.getOutputPorts()) { + if (ProcessGroup.STOP_PORTS_FILTER.test(remotePort)) { + sourceIds.add(remotePort.getIdentifier()); + } + } + } + + for (final Port port : group.findAllInputPorts()) { + if (port.getProcessGroup().resolveExecutionEngine() == ExecutionEngine.STANDARD + && port instanceof PublicPort + && ProcessGroup.STOP_PORTS_FILTER.test(port) + && !Connectables.hasNonLoopConnection(port)) { + sourceIds.add(port.getIdentifier()); + } + } + + return sourceIds; + } + + @Override + public void verifyStopSources(final String groupId) { + final ProcessGroup group = processGroupDAO.getProcessGroup(groupId); + verifyStopSources(group); + } + + @Override + public void verifyStopSources(final String groupId, final Set componentIds) { + final ProcessGroup group = processGroupDAO.getProcessGroup(groupId); + verifyStopSources(group); + + final Set sourceComponentIds = findSourceComponentIds(group); + if (!sourceComponentIds.equals(componentIds)) { + throw new IllegalStateException("Source components changed while processing the request; refresh and retry"); + } + } + + private static void verifyStopSources(final ProcessGroup group) { + if (group.resolveExecutionEngine() == ExecutionEngine.STATELESS) { + throw new IllegalStateException("Cannot stop sources in a Process Group that resolves to the Stateless Execution Engine"); + } + } + @Override public Set getRevisionsFromSnippet(final String snippetId) { final Snippet snippet = snippetDAO.getSnippet(snippetId); diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/FlowResource.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/FlowResource.java index 753c6c66a97d..a77f79fdba72 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/FlowResource.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/FlowResource.java @@ -1217,6 +1217,125 @@ public Response scheduleComponents( ); } + /** + * Stops source components in the specified process group and its Standard-engine descendants. Sources are + * processors with no non-loop incoming connection, Remote Process Group output ports, and public input ports + * with no non-loop incoming connection. The operation is rejected when the specified group resolves to the + * Stateless Execution Engine. + * + * @param id The id of the process group. + * @param requestScheduleComponentsEntity A scheduleComponentsEntity with state STOPPED. + * @return A scheduleComponentsEntity identifying the components that were stopped. + */ + @PUT + @Consumes(MediaType.APPLICATION_JSON) + @Produces(MediaType.APPLICATION_JSON) + @Path("process-groups/{id}/sources") + @Operation( + summary = "Stop source components in the specified Process Group.", + description = "All source components included in the operation must be authorized. If any source is unauthorized, no components are stopped.", + responses = { + @ApiResponse(responseCode = "200", content = @Content(schema = @Schema(implementation = ScheduleComponentsEntity.class))), + @ApiResponse(responseCode = "400", description = "NiFi was unable to complete the request because it was invalid. The request should not be retried without modification."), + @ApiResponse(responseCode = "401", description = "Client could not be authenticated."), + @ApiResponse(responseCode = "403", description = "Client is not authorized to make this request."), + @ApiResponse(responseCode = "404", description = "The specified resource could not be found."), + @ApiResponse(responseCode = "409", description = "The request was valid but NiFi was not in the appropriate state to process it.") + }, + security = { + @SecurityRequirement(name = "Read - /flow"), + @SecurityRequirement(name = "Write - /{component-type}/{uuid} or /operation/{component-type}/{uuid} - For every source component being stopped") + } + ) + public Response stopSources( + @Parameter( + description = "The process group id.", + required = true + ) + @PathParam("id") final String id, + @Parameter( + description = "The request to stop sources. If the components in the request are not specified, all source components are identified.", + required = true + ) final ScheduleComponentsEntity requestScheduleComponentsEntity) { + + if (requestScheduleComponentsEntity == null) { + throw new IllegalArgumentException("Schedule Component must be specified."); + } + + if (!id.equals(requestScheduleComponentsEntity.getId())) { + throw new IllegalArgumentException(String.format("The process group id (%s) in the request body does " + + "not equal the process group id of the requested resource (%s).", requestScheduleComponentsEntity.getId(), id)); + } + + if (!ScheduledState.STOPPED.name().equals(requestScheduleComponentsEntity.getState())) { + throw new IllegalArgumentException("The scheduled state must be STOPPED."); + } + + authorizeFlow(); + serviceFacade.verifyStopSources(id); + + if (requestScheduleComponentsEntity.getComponents() == null) { + final Set revisions = serviceFacade.getRevisionsFromGroup(id, serviceFacade::findSourceComponentIds); + final Map componentsToStop = new HashMap<>(); + for (final Revision revision : revisions) { + final RevisionDTO dto = new RevisionDTO(); + dto.setClientId(revision.getClientId()); + dto.setVersion(revision.getVersion()); + componentsToStop.put(revision.getComponentId(), dto); + } + + requestScheduleComponentsEntity.setComponents(componentsToStop); + } + + final Map requestComponentsToStop = requestScheduleComponentsEntity.getComponents(); + final Set requestComponentIds = Set.copyOf(requestComponentsToStop.keySet()); + + serviceFacade.authorizeAccess(lookup -> { + for (final String componentId : requestComponentIds) { + final Authorizable connectable = lookup.getLocalConnectable(componentId); + OperationAuthorizable.authorizeOperation(connectable, authorizer, NiFiUserUtils.getNiFiUser()); + } + }); + + serviceFacade.verifyStopSources(id, requestComponentIds); + + if (isReplicateRequest()) { + return replicate(HttpMethod.PUT, requestScheduleComponentsEntity); + } else if (isDisconnectedFromCluster()) { + verifyDisconnectedNodeModification(requestScheduleComponentsEntity.isDisconnectedNodeAcknowledged()); + } + + final Map requestComponentRevisions = + requestComponentsToStop.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, e -> getRevision(e.getValue(), e.getKey()))); + final Set requestRevisions = new HashSet<>(requestComponentRevisions.values()); + + return withWriteLock( + serviceFacade, + requestScheduleComponentsEntity, + requestRevisions, + lookup -> { + authorizeFlow(); + + requestComponentsToStop.keySet().forEach(componentId -> { + final Authorizable connectable = lookup.getLocalConnectable(componentId); + OperationAuthorizable.authorizeOperation(connectable, authorizer, NiFiUserUtils.getNiFiUser()); + }); + }, + () -> { + serviceFacade.verifyStopSources(id, requestComponentIds); + serviceFacade.verifyScheduleComponents(id, ScheduledState.STOPPED, requestComponentIds); + }, + (revisions, scheduleComponentsEntity) -> { + final Map componentsToStop = scheduleComponentsEntity.getComponents(); + final Map componentRevisions = + componentsToStop.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, e -> getRevision(e.getValue(), e.getKey()))); + final ScheduleComponentsEntity entity = serviceFacade.scheduleComponents(id, ScheduledState.STOPPED, componentRevisions); + entity.setComponents(componentsToStop); + return generateOkResponse(entity).build(); + } + ); + } + @PUT @Consumes(MediaType.APPLICATION_JSON) @Produces(MediaType.APPLICATION_JSON) diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java index 46e8ab35d7ad..bc24e38ecda9 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java @@ -2568,6 +2568,8 @@ public ProcessGroupFlowDTO createProcessGroupFlowDto(final ProcessGroup group, f dto.setParentGroupId(parent.getIdentifier()); } + dto.setResolvedExecutionEngine(group.resolveExecutionEngine().name()); + final ParameterContext parameterContext = group.getParameterContext(); if (parameterContext != null) { dto.setParameterContext(entityFactory.createParameterReferenceEntity(createParameterContextReference(parameterContext), createPermissionsDto(parameterContext))); @@ -2834,6 +2836,7 @@ private ProcessGroupDTO createConciseProcessGroupDto(final ProcessGroup group) { dto.setLogFileSuffix(group.getLogFileSuffix()); dto.setStatelessGroupScheduledState(group.getStatelessScheduledState().name()); dto.setExecutionEngine(group.getExecutionEngine().name()); + dto.setResolvedExecutionEngine(group.resolveExecutionEngine().name()); dto.setMaxConcurrentTasks(group.getMaxConcurrentTasks()); dto.setStatelessFlowTimeout(group.getStatelessFlowTimeout()); @@ -4825,6 +4828,7 @@ public ProcessGroupDTO copy(final ProcessGroupDTO original, final boolean deep) copy.setDefaultBackPressureDataSizeThreshold(original.getDefaultBackPressureDataSizeThreshold()); copy.setLogFileSuffix(original.getLogFileSuffix()); copy.setExecutionEngine(original.getExecutionEngine()); + copy.setResolvedExecutionEngine(original.getResolvedExecutionEngine()); copy.setMaxConcurrentTasks(original.getMaxConcurrentTasks()); copy.setStatelessFlowTimeout(original.getStatelessFlowTimeout()); diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/EntityFactory.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/EntityFactory.java index be52686de488..f37331420822 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/EntityFactory.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/EntityFactory.java @@ -320,6 +320,7 @@ public ProcessGroupEntity createProcessGroupEntity(final ProcessGroupDTO dto, fi entity.setStaleCount(dto.getStaleCount()); entity.setLocallyModifiedAndStaleCount(dto.getLocallyModifiedAndStaleCount()); entity.setSyncFailureCount(dto.getSyncFailureCount()); + entity.setResolvedExecutionEngine(dto.getResolvedExecutionEngine()); final ParameterContextReferenceEntity parameterContextReference = dto.getParameterContext(); if (parameterContextReference != null) { 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 079f8cf8629a..47f729e17630 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 @@ -56,6 +56,9 @@ import org.apache.nifi.components.state.StateManagerProvider; import org.apache.nifi.components.state.StateMap; import org.apache.nifi.components.validation.ValidationStatus; +import org.apache.nifi.connectable.Connectable; +import org.apache.nifi.connectable.Connection; +import org.apache.nifi.connectable.Port; import org.apache.nifi.controller.ClusterTopologyProvider; import org.apache.nifi.controller.ControllerService; import org.apache.nifi.controller.Counter; @@ -114,6 +117,8 @@ import org.apache.nifi.registry.flow.mapping.FlowMappingOptions; import org.apache.nifi.registry.flow.mapping.InstantiatedVersionedProcessGroup; import org.apache.nifi.registry.flow.mapping.VersionedComponentFlowMapper; +import org.apache.nifi.remote.PublicPort; +import org.apache.nifi.remote.RemoteGroupPort; import org.apache.nifi.reporting.Bulletin; import org.apache.nifi.reporting.BulletinFactory; import org.apache.nifi.reporting.BulletinQuery; @@ -3260,6 +3265,291 @@ public void testDeleteAssetOwnedByContextRemovesAsset() { verify(assetManager).deleteAsset(ASSET_ID); } + @Test + public void testFindSourceComponentIdsIncludesProcessorWithNoIncomingConnections() { + final ProcessorNode processor = sourceProcessor("source-processor"); + final ProcessGroup group = sourceProcessGroup(List.of(processor), List.of(), List.of()); + + assertEquals(Set.of("source-processor"), serviceFacade.findSourceComponentIds(group)); + } + + @Test + public void testFindSourceComponentIdsIncludesProcessorWithSelfLoop() { + final ProcessorNode processor = sourceProcessor("self-loop-processor"); + final Connection selfLoop = connectionFrom(processor); + when(processor.getIncomingConnections()).thenReturn(List.of(selfLoop)); + final ProcessGroup group = sourceProcessGroup(List.of(processor), List.of(), List.of()); + + assertEquals(Set.of("self-loop-processor"), serviceFacade.findSourceComponentIds(group)); + } + + @Test + public void testFindSourceComponentIdsExcludesProcessorWithIncomingConnectionFromOtherComponent() { + final ProcessorNode upstream = sourceProcessor("upstream"); + final ProcessorNode downstream = sourceProcessor("downstream"); + final Connection incomingConnection = connectionFrom(upstream); + when(downstream.getIncomingConnections()).thenReturn(List.of(incomingConnection)); + final ProcessGroup group = sourceProcessGroup(List.of(upstream, downstream), List.of(), List.of()); + + assertEquals(Set.of("upstream"), serviceFacade.findSourceComponentIds(group)); + } + + @Test + public void testFindSourceComponentIdsIncludesNestedProcessGroupProcessors() { + final ProcessorNode parentSource = sourceProcessor("parent-source"); + final ProcessorNode nestedSource = sourceProcessor("nested-source"); + final ProcessGroup group = sourceProcessGroup(List.of(parentSource, nestedSource), List.of(), List.of()); + + assertEquals(Set.of("parent-source", "nested-source"), serviceFacade.findSourceComponentIds(group)); + } + + @Test + public void testFindSourceComponentIdsExcludesSourcesInStatelessProcessGroups() { + final ProcessorNode statelessProcessor = sourceProcessor("stateless-processor"); + final RemoteGroupPort statelessRemoteOutput = remotePort("stateless-remote-output"); + final RemoteProcessGroup statelessRemoteProcessGroup = mock(RemoteProcessGroup.class); + when(statelessRemoteProcessGroup.getOutputPorts()).thenReturn(Set.of(statelessRemoteOutput)); + final PublicPort statelessPublicInput = sourcePublicPort("stateless-public-input"); + final ProcessGroup group = sourceProcessGroup( + List.of(statelessProcessor), + List.of(statelessRemoteProcessGroup), + List.of(statelessPublicInput) + ); + final ProcessGroup statelessGroup = mock(ProcessGroup.class); + when(statelessGroup.resolveExecutionEngine()).thenReturn(ExecutionEngine.STATELESS); + when(statelessProcessor.getProcessGroup()).thenReturn(statelessGroup); + when(statelessRemoteProcessGroup.getProcessGroup()).thenReturn(statelessGroup); + when(statelessPublicInput.getProcessGroup()).thenReturn(statelessGroup); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testVerifyStopSourcesRejectsStatelessProcessGroup() { + final ProcessGroup group = mock(ProcessGroup.class); + when(group.resolveExecutionEngine()).thenReturn(ExecutionEngine.STATELESS); + when(processGroupDAO.getProcessGroup("stateless-group")).thenReturn(group); + + assertThrows(IllegalStateException.class, () -> serviceFacade.verifyStopSources("stateless-group")); + } + + @Test + public void testVerifyStopSourcesAllowsStandardProcessGroup() { + final ProcessGroup group = mock(ProcessGroup.class); + when(group.resolveExecutionEngine()).thenReturn(ExecutionEngine.STANDARD); + when(processGroupDAO.getProcessGroup("standard-group")).thenReturn(group); + + serviceFacade.verifyStopSources("standard-group"); + } + + @Test + public void testVerifyStopSourcesAllowsMatchingSourceComponents() { + final ProcessorNode sourceProcessor = sourceProcessor("source-processor"); + final ProcessGroup group = sourceProcessGroup(List.of(sourceProcessor), List.of(), List.of()); + when(processGroupDAO.getProcessGroup("standard-group")).thenReturn(group); + + serviceFacade.verifyStopSources("standard-group", Set.of("source-processor")); + } + + @Test + public void testVerifyStopSourcesRejectsMissingSourceComponent() { + final ProcessorNode sourceProcessor = sourceProcessor("source-processor"); + final ProcessGroup group = sourceProcessGroup(List.of(sourceProcessor), List.of(), List.of()); + when(processGroupDAO.getProcessGroup("standard-group")).thenReturn(group); + + assertThrows(IllegalStateException.class, () -> serviceFacade.verifyStopSources("standard-group", Set.of())); + } + + @Test + public void testVerifyStopSourcesRejectsNonSourceComponent() { + final ProcessorNode sourceProcessor = sourceProcessor("source-processor"); + final ProcessorNode downstreamProcessor = sourceProcessor("downstream-processor"); + final Connection incomingConnection = connectionFrom(sourceProcessor); + when(downstreamProcessor.getIncomingConnections()).thenReturn(List.of(incomingConnection)); + final ProcessGroup group = sourceProcessGroup(List.of(sourceProcessor, downstreamProcessor), List.of(), List.of()); + when(processGroupDAO.getProcessGroup("standard-group")).thenReturn(group); + + assertThrows(IllegalStateException.class, + () -> serviceFacade.verifyStopSources("standard-group", Set.of("source-processor", "downstream-processor"))); + } + + @Test + public void testVerifyStopSourcesRejectsComponentInStatelessDescendant() { + final ProcessorNode statelessProcessor = sourceProcessor("stateless-processor"); + final ProcessGroup group = sourceProcessGroup(List.of(statelessProcessor), List.of(), List.of()); + final ProcessGroup statelessGroup = mock(ProcessGroup.class); + when(statelessGroup.resolveExecutionEngine()).thenReturn(ExecutionEngine.STATELESS); + when(statelessProcessor.getProcessGroup()).thenReturn(statelessGroup); + when(processGroupDAO.getProcessGroup("standard-group")).thenReturn(group); + + assertThrows(IllegalStateException.class, + () -> serviceFacade.verifyStopSources("standard-group", Set.of("stateless-processor"))); + } + + @Test + public void testFindSourceComponentIdsExcludesNonRunningSourceProcessor() { + final ProcessorNode nonRunningSource = sourceProcessor("non-running-source", false); + final ProcessGroup group = sourceProcessGroup(List.of(nonRunningSource), List.of(), List.of()); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testFindSourceComponentIdsExcludesNonTransmittingRemoteOutputPort() { + final RemoteGroupPort remoteOutput = remotePort("remote-output", ScheduledState.STOPPED); + final RemoteProcessGroup remoteProcessGroup = mock(RemoteProcessGroup.class); + when(remoteProcessGroup.getOutputPorts()).thenReturn(Set.of(remoteOutput)); + when(remoteProcessGroup.getInputPorts()).thenReturn(Set.of()); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(remoteProcessGroup), List.of()); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testFindSourceComponentIdsExcludesStoppedPublicInputPort() { + final PublicPort publicInput = sourcePublicPort("public-input", ScheduledState.STOPPED); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(), List.of(publicInput)); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testFindSourceComponentIdsIncludesRemoteOutputPortAndExcludesRemoteInputPort() { + final RemoteGroupPort remoteOutput = remotePort("remote-output"); + final RemoteGroupPort remoteInput = remotePort("remote-input"); + final RemoteProcessGroup remoteProcessGroup = mock(RemoteProcessGroup.class); + when(remoteProcessGroup.getIdentifier()).thenReturn("rpg"); + when(remoteProcessGroup.getOutputPorts()).thenReturn(Set.of(remoteOutput)); + when(remoteProcessGroup.getInputPorts()).thenReturn(Set.of(remoteInput)); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(remoteProcessGroup), List.of()); + + final Set sourceIds = serviceFacade.findSourceComponentIds(group); + assertEquals(Set.of("remote-output"), sourceIds); + assertFalse(sourceIds.contains("rpg")); + assertFalse(sourceIds.contains("remote-input")); + } + + @Test + public void testFindSourceComponentIdsExcludesLocalPorts() { + final Port localInput = sourcePort("local-input"); + final Port localOutput = sourcePort("local-output"); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(), List.of(localInput)); + when(group.findAllOutputPorts()).thenReturn(List.of(localOutput)); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testFindSourceComponentIdsIncludesPublicInputPortWithNoIncomingConnections() { + final PublicPort publicInput = sourcePublicPort("public-input"); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(), List.of(publicInput)); + + assertEquals(Set.of("public-input"), serviceFacade.findSourceComponentIds(group)); + } + + @Test + public void testFindSourceComponentIdsExcludesPublicInputPortWithIncomingConnectionFromParent() { + final Connectable parentOutput = sourceProcessor("parent-output"); + final PublicPort publicInput = sourcePublicPort("nested-public-input"); + final Connection incomingConnection = connectionFrom(parentOutput); + when(publicInput.getIncomingConnections()).thenReturn(List.of(incomingConnection)); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(), List.of(publicInput)); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testFindSourceComponentIdsExcludesPublicOutputPort() { + final PublicPort publicOutput = sourcePublicPort("public-output"); + final ProcessGroup group = sourceProcessGroup(List.of(), List.of(), List.of()); + when(group.findAllOutputPorts()).thenReturn(List.of(publicOutput)); + + assertTrue(serviceFacade.findSourceComponentIds(group).isEmpty()); + } + + @Test + public void testFindSourceComponentIdsMixOfSourcesAndNonSources() { + final ProcessorNode sourceProcessor = sourceProcessor("source-processor"); + final ProcessorNode downstream = sourceProcessor("downstream"); + final Connection incomingConnection = connectionFrom(sourceProcessor); + when(downstream.getIncomingConnections()).thenReturn(List.of(incomingConnection)); + + final RemoteGroupPort remoteOutput = remotePort("remote-output"); + final RemoteGroupPort remoteInput = remotePort("remote-input"); + final RemoteProcessGroup remoteProcessGroup = mock(RemoteProcessGroup.class); + when(remoteProcessGroup.getIdentifier()).thenReturn("rpg"); + when(remoteProcessGroup.getOutputPorts()).thenReturn(Set.of(remoteOutput)); + when(remoteProcessGroup.getInputPorts()).thenReturn(Set.of(remoteInput)); + + final PublicPort publicInput = sourcePublicPort("public-input"); + final Port localInput = sourcePort("local-input"); + + final ProcessGroup group = sourceProcessGroup(List.of(sourceProcessor, downstream), List.of(remoteProcessGroup), List.of(publicInput, localInput)); + + assertEquals(Set.of("source-processor", "remote-output", "public-input"), serviceFacade.findSourceComponentIds(group)); + } + + private static ProcessGroup sourceProcessGroup(final List processors, final List remoteProcessGroups, final List inputPorts) { + final ProcessGroup group = mock(ProcessGroup.class); + when(group.resolveExecutionEngine()).thenReturn(ExecutionEngine.STANDARD); + when(group.findAllProcessors()).thenReturn(processors); + when(group.findAllRemoteProcessGroups()).thenReturn(remoteProcessGroups); + when(group.findAllInputPorts()).thenReturn(inputPorts); + when(group.findAllOutputPorts()).thenReturn(Collections.emptyList()); + processors.forEach(processor -> when(processor.getProcessGroup()).thenReturn(group)); + remoteProcessGroups.forEach(remoteProcessGroup -> when(remoteProcessGroup.getProcessGroup()).thenReturn(group)); + inputPorts.forEach(inputPort -> when(inputPort.getProcessGroup()).thenReturn(group)); + return group; + } + + private static ProcessorNode sourceProcessor(final String identifier) { + return sourceProcessor(identifier, true); + } + + private static ProcessorNode sourceProcessor(final String identifier, final boolean running) { + final ProcessorNode processor = mock(ProcessorNode.class); + when(processor.getIdentifier()).thenReturn(identifier); + when(processor.getIncomingConnections()).thenReturn(Collections.emptyList()); + when(processor.isRunning()).thenReturn(running); + return processor; + } + + private static Port sourcePort(final String identifier) { + final Port port = mock(Port.class); + when(port.getIdentifier()).thenReturn(identifier); + when(port.getIncomingConnections()).thenReturn(Collections.emptyList()); + return port; + } + + private static PublicPort sourcePublicPort(final String identifier) { + return sourcePublicPort(identifier, ScheduledState.RUNNING); + } + + private static PublicPort sourcePublicPort(final String identifier, final ScheduledState scheduledState) { + final PublicPort port = mock(PublicPort.class); + when(port.getIdentifier()).thenReturn(identifier); + when(port.getIncomingConnections()).thenReturn(Collections.emptyList()); + when(port.getScheduledState()).thenReturn(scheduledState); + return port; + } + + private static RemoteGroupPort remotePort(final String identifier) { + return remotePort(identifier, ScheduledState.RUNNING); + } + + private static RemoteGroupPort remotePort(final String identifier, final ScheduledState scheduledState) { + final RemoteGroupPort port = mock(RemoteGroupPort.class); + when(port.getIdentifier()).thenReturn(identifier); + when(port.getScheduledState()).thenReturn(scheduledState); + return port; + } + + private static Connection connectionFrom(final Connectable source) { + final Connection connection = mock(Connection.class); + when(connection.getSource()).thenReturn(source); + return connection; + } + private Asset createAsset(final String assetId, final String ownerId) { final Asset asset = mock(Asset.class); when(asset.getIdentifier()).thenReturn(assetId); diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java index 8356406c31ab..3841ed70afec 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java @@ -26,15 +26,22 @@ import io.prometheus.client.Collector.MetricFamilySamples.Sample; import io.prometheus.client.CollectorRegistry; import io.prometheus.client.exporter.common.TextFormat; +import jakarta.ws.rs.HttpMethod; import jakarta.ws.rs.core.MediaType; import jakarta.ws.rs.core.Response; import jakarta.ws.rs.core.StreamingOutput; 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.RequestAction; +import org.apache.nifi.authorization.resource.Authorizable; +import org.apache.nifi.authorization.user.NiFiUser; import org.apache.nifi.components.ValidationResult; import org.apache.nifi.components.validation.DisabledServiceValidationResult; import org.apache.nifi.connectable.Port; import org.apache.nifi.controller.ProcessorNode; +import org.apache.nifi.controller.ScheduledState; import org.apache.nifi.controller.service.ControllerServiceNode; import org.apache.nifi.controller.status.ProcessGroupStatus; import org.apache.nifi.controller.status.ProcessingPerformanceStatus; @@ -52,12 +59,15 @@ import org.apache.nifi.util.NiFiProperties; import org.apache.nifi.web.NiFiServiceFacade; import org.apache.nifi.web.ResourceNotFoundException; +import org.apache.nifi.web.Revision; import org.apache.nifi.web.api.dto.ComponentDifferenceDTO; import org.apache.nifi.web.api.dto.DifferenceDTO; +import org.apache.nifi.web.api.dto.RevisionDTO; import org.apache.nifi.web.api.entity.ActivateControllerServicesEntity; import org.apache.nifi.web.api.entity.ClearBulletinsForGroupRequestEntity; import org.apache.nifi.web.api.entity.ConnectorEntity; import org.apache.nifi.web.api.entity.FlowComparisonEntity; +import org.apache.nifi.web.api.entity.ScheduleComponentsEntity; import org.apache.nifi.web.api.request.FlowMetricsProducer; import org.apache.nifi.web.api.request.FlowMetricsReportingStrategy; import org.jetbrains.annotations.NotNull; @@ -95,13 +105,18 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.anySet; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.nullable; import static org.mockito.Mockito.any; import static org.mockito.Mockito.anyString; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -145,6 +160,9 @@ public class TestFlowResource { @Mock private NiFiProperties properties; + @Mock + private Authorizer authorizer; + @Mock private ConnectorResource connectorResource; @@ -152,6 +170,7 @@ public class TestFlowResource { public void setUp() { lenient().when(properties.isNode()).thenReturn(Boolean.FALSE); resource.properties = properties; + resource.setAuthorizer(authorizer); resource.setConnectorResource(connectorResource); } @@ -616,6 +635,238 @@ public void testClearBulletinsIncludesControllerServices() { assertEquals(5, componentIds.size(), "Should have exactly 5 authorized components"); } + @Test + public void testStopSourcesUsesFacadeToIdentifySources() { + final ScheduleComponentsEntity entity = new ScheduleComponentsEntity(); + entity.setId(PROCESS_GROUP_ID); + entity.setState(ScheduledState.STOPPED.name()); + + when(properties.isNode()).thenReturn(false); + resource.httpServletRequest = new MockHttpServletRequest(); + + final ProcessGroup processGroup = mock(ProcessGroup.class); + final Set identifiedSources = Set.of("source-processor", "remote-output", "public-input"); + when(serviceFacade.findSourceComponentIds(processGroup)).thenReturn(identifiedSources); + + final ArgumentCaptor>> revisionsCaptor = ArgumentCaptor.captor(); + when(serviceFacade.getRevisionsFromGroup(eq(PROCESS_GROUP_ID), revisionsCaptor.capture())).thenReturn(Set.of()); + when(serviceFacade.scheduleComponents(eq(PROCESS_GROUP_ID), eq(ScheduledState.STOPPED), any())).thenReturn(entity); + + final Response response = resource.stopSources(PROCESS_GROUP_ID, entity); + + assertNotNull(response); + assertEquals(HttpURLConnection.HTTP_OK, response.getStatus()); + assertTrue(entity.getComponents().isEmpty()); + assertEquals(identifiedSources, revisionsCaptor.getValue().apply(processGroup)); + } + + @Test + public void testStopSourcesUnauthorizedFlowDoesNotDiscoverOrSchedule() { + final ScheduleComponentsEntity entity = new ScheduleComponentsEntity(); + entity.setId(PROCESS_GROUP_ID); + entity.setState(ScheduledState.STOPPED.name()); + + resource.httpServletRequest = new MockHttpServletRequest(); + + doThrow(new AccessDeniedException("denied")).when(serviceFacade).authorizeAccess(any()); + + assertThrows(AccessDeniedException.class, () -> resource.stopSources(PROCESS_GROUP_ID, entity)); + + verify(serviceFacade, never()).verifyStopSources(anyString()); + verify(serviceFacade, never()).getRevisionsFromGroup(anyString(), any()); + verify(serviceFacade, never()).scheduleComponents(anyString(), any(), any()); + verify(serviceFacade, never()).verifyScheduleComponents(anyString(), any(), any()); + } + + @Test + public void testStopSourcesWithNoSourcesReturnsEmptyComponents() { + final ScheduleComponentsEntity entity = new ScheduleComponentsEntity(); + entity.setId(PROCESS_GROUP_ID); + entity.setState(ScheduledState.STOPPED.name()); + + when(properties.isNode()).thenReturn(false); + resource.httpServletRequest = new MockHttpServletRequest(); + + when(serviceFacade.getRevisionsFromGroup(eq(PROCESS_GROUP_ID), any())).thenReturn(Set.of()); + when(serviceFacade.scheduleComponents(eq(PROCESS_GROUP_ID), eq(ScheduledState.STOPPED), any())).thenReturn(entity); + + final Response response = resource.stopSources(PROCESS_GROUP_ID, entity); + + assertEquals(HttpURLConnection.HTTP_OK, response.getStatus()); + assertTrue(entity.getComponents().isEmpty()); + verify(serviceFacade).scheduleComponents(eq(PROCESS_GROUP_ID), eq(ScheduledState.STOPPED), eq(Map.of())); + } + + @Test + public void testStopSourcesStatelessProcessGroupDoesNotDiscoverOrSchedule() { + final RevisionDTO revision = new RevisionDTO(); + revision.setVersion(1L); + final ScheduleComponentsEntity entity = new ScheduleComponentsEntity(); + entity.setId(PROCESS_GROUP_ID); + entity.setState(ScheduledState.STOPPED.name()); + entity.setComponents(Map.of("source-processor", revision)); + + doThrow(new IllegalStateException("Stateless Process Group")).when(serviceFacade).verifyStopSources(PROCESS_GROUP_ID); + + assertThrows(IllegalStateException.class, () -> resource.stopSources(PROCESS_GROUP_ID, entity)); + + verify(serviceFacade, never()).getRevisionsFromGroup(anyString(), any()); + verify(serviceFacade).authorizeAccess(any()); + verify(serviceFacade, never()).verifyScheduleComponents(anyString(), any(), any()); + verify(serviceFacade, never()).scheduleComponents(anyString(), any(), any()); + } + + @Test + public void testStopSourcesResponseIncludesStoppedComponentIds() { + final ScheduleComponentsEntity request = new ScheduleComponentsEntity(); + request.setId(PROCESS_GROUP_ID); + request.setState(ScheduledState.STOPPED.name()); + + when(properties.isNode()).thenReturn(false); + resource.httpServletRequest = new MockHttpServletRequest(); + + when(serviceFacade.getRevisionsFromGroup(eq(PROCESS_GROUP_ID), any())) + .thenReturn(Set.of(new Revision(1L, "client", "source-processor"))); + + final ScheduleComponentsEntity facadeResponse = new ScheduleComponentsEntity(); + facadeResponse.setId(PROCESS_GROUP_ID); + facadeResponse.setState(ScheduledState.STOPPED.name()); + when(serviceFacade.scheduleComponents(eq(PROCESS_GROUP_ID), eq(ScheduledState.STOPPED), any())).thenReturn(facadeResponse); + + final Response response = resource.stopSources(PROCESS_GROUP_ID, request); + + assertEquals(HttpURLConnection.HTTP_OK, response.getStatus()); + final ScheduleComponentsEntity responseEntity = (ScheduleComponentsEntity) response.getEntity(); + assertNotNull(responseEntity.getComponents()); + assertEquals(Set.of("source-processor"), responseEntity.getComponents().keySet()); + } + + @Test + public void testStopSourcesIdentifiesSourcesBeforeReplicate() { + final FlowResource spyResource = spy(resource); + doReturn(true).when(spyResource).isReplicateRequest(); + doReturn(Response.ok().build()).when(spyResource).replicate(anyString(), any(ScheduleComponentsEntity.class)); + + final ScheduleComponentsEntity request = new ScheduleComponentsEntity(); + request.setId(PROCESS_GROUP_ID); + request.setState(ScheduledState.STOPPED.name()); + when(serviceFacade.getRevisionsFromGroup(eq(PROCESS_GROUP_ID), any())).thenReturn(Set.of( + new Revision(1L, "client", "first-source"), + new Revision(2L, "client", "second-source") + )); + + final Response response = spyResource.stopSources(PROCESS_GROUP_ID, request); + + assertEquals(HttpURLConnection.HTTP_OK, response.getStatus()); + assertEquals(Set.of("first-source", "second-source"), request.getComponents().keySet()); + verify(serviceFacade).verifyStopSources(PROCESS_GROUP_ID, Set.of("first-source", "second-source")); + verify(serviceFacade, times(2)).authorizeAccess(any()); + verify(spyResource).replicate(HttpMethod.PUT, request); + verify(serviceFacade, never()).scheduleComponents(anyString(), any(), any()); + } + + @Test + public void testStopSourcesRechecksIdentifiedSourcesBeforeScheduling() { + final RevisionDTO revisionDto = new RevisionDTO(); + revisionDto.setClientId("client"); + revisionDto.setVersion(1L); + + final ScheduleComponentsEntity request = new ScheduleComponentsEntity(); + request.setId(PROCESS_GROUP_ID); + request.setState(ScheduledState.STOPPED.name()); + request.setComponents(Map.of("source-processor", revisionDto)); + resource.httpServletRequest = new MockHttpServletRequest(); + + doNothing().when(serviceFacade).verifyStopSources(PROCESS_GROUP_ID); + doNothing() + .doThrow(new IllegalStateException("Source components changed")) + .when(serviceFacade).verifyStopSources(PROCESS_GROUP_ID, Set.of("source-processor")); + + assertThrows(IllegalStateException.class, () -> resource.stopSources(PROCESS_GROUP_ID, request)); + + verify(serviceFacade, times(2)).verifyStopSources(PROCESS_GROUP_ID, Set.of("source-processor")); + verify(serviceFacade, never()).scheduleComponents(anyString(), any(), any()); + } + + @Test + public void testStopSourcesAuthorizesSuppliedComponentsBeforeReplicate() { + final FlowResource spyResource = spy(resource); + doReturn(true).when(spyResource).isReplicateRequest(); + doReturn(Response.ok().build()).when(spyResource).replicate(anyString(), any(ScheduleComponentsEntity.class)); + + final RevisionDTO firstRevision = new RevisionDTO(); + firstRevision.setClientId("client"); + firstRevision.setVersion(1L); + final RevisionDTO secondRevision = new RevisionDTO(); + secondRevision.setClientId("client"); + secondRevision.setVersion(2L); + + final ScheduleComponentsEntity request = new ScheduleComponentsEntity(); + request.setId(PROCESS_GROUP_ID); + request.setState(ScheduledState.STOPPED.name()); + request.setComponents(Map.of("first-source", firstRevision, "second-source", secondRevision)); + + final Response response = spyResource.stopSources(PROCESS_GROUP_ID, request); + + assertEquals(HttpURLConnection.HTTP_OK, response.getStatus()); + final ArgumentCaptor authorizeAccessCaptor = ArgumentCaptor.captor(); + verify(serviceFacade, times(2)).authorizeAccess(authorizeAccessCaptor.capture()); + final AuthorizableLookup lookup = mock(AuthorizableLookup.class); + final Authorizable flow = mock(Authorizable.class); + final Authorizable firstSource = mock(Authorizable.class); + final Authorizable secondSource = mock(Authorizable.class); + when(lookup.getFlow()).thenReturn(flow); + when(lookup.getLocalConnectable("first-source")).thenReturn(firstSource); + when(lookup.getLocalConnectable("second-source")).thenReturn(secondSource); + + authorizeAccessCaptor.getAllValues().forEach(authorizeAccess -> authorizeAccess.authorize(lookup)); + + verify(flow).authorize(eq(authorizer), eq(RequestAction.READ), nullable(NiFiUser.class)); + verify(firstSource).authorize(eq(authorizer), eq(RequestAction.WRITE), nullable(NiFiUser.class)); + verify(secondSource).authorize(eq(authorizer), eq(RequestAction.WRITE), nullable(NiFiUser.class)); + verify(serviceFacade).verifyStopSources(PROCESS_GROUP_ID, Set.of("first-source", "second-source")); + verify(serviceFacade, never()).getRevisionsFromGroup(anyString(), any()); + verify(spyResource).replicate(eq(HttpMethod.PUT), eq(request)); + verify(serviceFacade, never()).scheduleComponents(anyString(), any(), any()); + assertEquals(Set.of("first-source", "second-source"), request.getComponents().keySet()); + } + + @Test + public void testStopSourcesUnauthorizedSuppliedComponentsDoesNotReplicate() { + final FlowResource spyResource = spy(resource); + + final RevisionDTO revisionDto = new RevisionDTO(); + revisionDto.setClientId("client"); + revisionDto.setVersion(1L); + + final ScheduleComponentsEntity request = new ScheduleComponentsEntity(); + request.setId(PROCESS_GROUP_ID); + request.setState(ScheduledState.STOPPED.name()); + request.setComponents(Map.of("source-processor", revisionDto)); + + final AuthorizableLookup lookup = mock(AuthorizableLookup.class); + final Authorizable flow = mock(Authorizable.class); + final Authorizable sourceProcessor = mock(Authorizable.class); + when(lookup.getFlow()).thenReturn(flow); + when(lookup.getLocalConnectable("source-processor")).thenReturn(sourceProcessor); + doThrow(new AccessDeniedException("denied")).when(sourceProcessor) + .authorize(eq(authorizer), eq(RequestAction.WRITE), nullable(NiFiUser.class)); + doAnswer(invocation -> { + invocation.getArgument(0, AuthorizeAccess.class).authorize(lookup); + return null; + }).when(serviceFacade).authorizeAccess(any()); + + assertThrows(AccessDeniedException.class, () -> spyResource.stopSources(PROCESS_GROUP_ID, request)); + + verify(serviceFacade, times(2)).authorizeAccess(any()); + verify(flow).authorize(eq(authorizer), eq(RequestAction.READ), nullable(NiFiUser.class)); + verify(sourceProcessor).authorize(eq(authorizer), eq(RequestAction.WRITE), nullable(NiFiUser.class)); + verify(serviceFacade, never()).verifyStopSources(PROCESS_GROUP_ID, Set.of("source-processor")); + verify(spyResource, never()).replicate(anyString(), any(ScheduleComponentsEntity.class)); + verify(serviceFacade, never()).scheduleComponents(anyString(), any(), any()); + verify(serviceFacade, never()).getRevisionsFromGroup(anyString(), any()); + } + @Test public void testGetConnectors() { final ConnectorEntity connectorEntity = new ConnectorEntity(); diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java index a6f4fb1cd812..5bd149c33f79 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java @@ -26,6 +26,7 @@ import org.apache.nifi.connectable.Connectable; import org.apache.nifi.connectable.ConnectableType; import org.apache.nifi.connectable.Connection; +import org.apache.nifi.connectable.Position; import org.apache.nifi.controller.ControllerService; import org.apache.nifi.controller.queue.FlowFileQueue; import org.apache.nifi.controller.queue.LoadBalanceCompression; @@ -33,7 +34,13 @@ import org.apache.nifi.controller.service.ControllerServiceNode; import org.apache.nifi.controller.service.ControllerServiceProvider; import org.apache.nifi.controller.service.ControllerServiceState; +import org.apache.nifi.controller.status.ProcessGroupStatus; +import org.apache.nifi.flow.ExecutionEngine; +import org.apache.nifi.groups.FlowFileConcurrency; +import org.apache.nifi.groups.FlowFileOutboundPolicy; import org.apache.nifi.groups.ProcessGroup; +import org.apache.nifi.groups.ProcessGroupCounts; +import org.apache.nifi.groups.StatelessGroupScheduledState; import org.apache.nifi.logging.LogLevel; import org.apache.nifi.nar.ExtensionManager; import org.apache.nifi.nar.NarManifest; @@ -51,6 +58,7 @@ import org.apache.nifi.registry.flow.diff.DifferenceType; import org.apache.nifi.registry.flow.diff.FlowDifference; import org.apache.nifi.web.ResourceNotFoundException; +import org.apache.nifi.web.api.dto.flow.ProcessGroupFlowDTO; import org.apache.nifi.web.api.entity.AllowableValueEntity; import org.apache.nifi.web.api.entity.ParameterContextReferenceEntity; import org.apache.nifi.web.revision.RevisionManager; @@ -1062,4 +1070,102 @@ private static void configureBaseParameterContext(final ParameterContext context when(context.getParameterReferenceManager()).thenReturn(ParameterReferenceManager.EMPTY); } + @Test + void testCreateProcessGroupDtoSetsResolvedExecutionEngineForInheritedUnderStatelessParent() { + final ProcessGroup group = stubProcessGroup(ExecutionEngine.INHERITED, ExecutionEngine.STATELESS); + + final ProcessGroupDTO dto = newDtoFactoryForParameters().createProcessGroupDto(group, false); + + assertEquals("INHERITED", dto.getExecutionEngine()); + assertEquals("STATELESS", dto.getResolvedExecutionEngine()); + } + + @Test + void testCreateProcessGroupDtoSetsResolvedExecutionEngineStandardForRootInherited() { + final ProcessGroup group = stubProcessGroup(ExecutionEngine.INHERITED, ExecutionEngine.STANDARD); + + final ProcessGroupDTO dto = newDtoFactoryForParameters().createProcessGroupDto(group, false); + + assertEquals("INHERITED", dto.getExecutionEngine()); + assertEquals("STANDARD", dto.getResolvedExecutionEngine()); + } + + @Test + void testCreateProcessGroupDtoSetsResolvedExecutionEngineForConfiguredStateless() { + final ProcessGroup group = stubProcessGroup(ExecutionEngine.STATELESS, ExecutionEngine.STATELESS); + + final ProcessGroupDTO dto = newDtoFactoryForParameters().createProcessGroupDto(group, false); + + assertEquals("STATELESS", dto.getExecutionEngine()); + assertEquals("STATELESS", dto.getResolvedExecutionEngine()); + } + + @Test + void testCreateProcessGroupFlowDtoSetsResolvedExecutionEngine() { + final ProcessGroup group = stubProcessGroup(ExecutionEngine.INHERITED, ExecutionEngine.STATELESS); + final ProcessGroupStatus groupStatus = mock(ProcessGroupStatus.class); + when(groupStatus.getProcessorStatus()).thenReturn(Collections.emptyList()); + when(groupStatus.getConnectionStatus()).thenReturn(Collections.emptyList()); + when(groupStatus.getProcessGroupStatus()).thenReturn(Collections.emptyList()); + when(groupStatus.getRemoteProcessGroupStatus()).thenReturn(Collections.emptyList()); + when(groupStatus.getInputPortStatus()).thenReturn(Collections.emptyList()); + when(groupStatus.getOutputPortStatus()).thenReturn(Collections.emptyList()); + + final ProcessGroupFlowDTO dto = newDtoFactoryForParameters().createProcessGroupFlowDto( + group, + groupStatus, + mock(RevisionManager.class), + ignored -> Collections.emptyList(), + false + ); + + assertEquals("STATELESS", dto.getResolvedExecutionEngine()); + } + + @Test + void testCopyProcessGroupDtoCopiesResolvedExecutionEngine() { + final ProcessGroupDTO original = new ProcessGroupDTO(); + original.setContents(new FlowSnippetDTO()); + original.setExecutionEngine("INHERITED"); + original.setResolvedExecutionEngine("STATELESS"); + + final ProcessGroupDTO copy = newDtoFactoryForParameters().copy(original, false); + + assertEquals("INHERITED", copy.getExecutionEngine()); + assertEquals("STATELESS", copy.getResolvedExecutionEngine()); + } + + private static ProcessGroup stubProcessGroup(final ExecutionEngine configured, final ExecutionEngine resolved) { + final ProcessGroup group = mock(ProcessGroup.class); + when(group.getIdentifier()).thenReturn("pg-1"); + when(group.getPosition()).thenReturn(new Position(0, 0)); + when(group.getComments()).thenReturn(""); + when(group.getName()).thenReturn("group"); + when(group.getVersionedComponentId()).thenReturn(Optional.empty()); + when(group.getVersionControlInformation()).thenReturn(null); + when(group.getFlowFileConcurrency()).thenReturn(FlowFileConcurrency.UNBOUNDED); + when(group.getFlowFileOutboundPolicy()).thenReturn(FlowFileOutboundPolicy.STREAM_WHEN_AVAILABLE); + when(group.getDefaultFlowFileExpiration()).thenReturn("0 sec"); + when(group.getDefaultBackPressureObjectThreshold()).thenReturn(10000L); + when(group.getDefaultBackPressureDataSizeThreshold()).thenReturn("1 GB"); + when(group.getLogFileSuffix()).thenReturn(null); + when(group.getStatelessScheduledState()).thenReturn(StatelessGroupScheduledState.STOPPED); + when(group.getExecutionEngine()).thenReturn(configured); + when(group.resolveExecutionEngine()).thenReturn(resolved); + when(group.getMaxConcurrentTasks()).thenReturn(1); + when(group.getStatelessFlowTimeout()).thenReturn("1 min"); + when(group.getParameterContext()).thenReturn(null); + when(group.getParent()).thenReturn(null); + when(group.getCounts()).thenReturn(new ProcessGroupCounts(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0)); + when(group.getProcessors()).thenReturn(Collections.emptySet()); + when(group.getConnections()).thenReturn(Collections.emptySet()); + when(group.getLabels()).thenReturn(Collections.emptySet()); + when(group.getFunnels()).thenReturn(Collections.emptySet()); + when(group.getProcessGroups()).thenReturn(Collections.emptySet()); + when(group.getRemoteProcessGroups()).thenReturn(Collections.emptySet()); + when(group.getInputPorts()).thenReturn(Collections.emptySet()); + when(group.getOutputPorts()).thenReturn(Collections.emptySet()); + return group; + } + } diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/EntityFactoryTest.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/EntityFactoryTest.java new file mode 100644 index 000000000000..2edf44b21225 --- /dev/null +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/EntityFactoryTest.java @@ -0,0 +1,58 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.nifi.web.api.dto; + +import org.apache.nifi.web.api.entity.ProcessGroupEntity; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; + +public class EntityFactoryTest { + + @Test + void testCreateProcessGroupEntityPromotesResolvedExecutionEngineWhenUnauthorized() { + final ProcessGroupDTO dto = new ProcessGroupDTO(); + dto.setId("pg-1"); + dto.setResolvedExecutionEngine("STATELESS"); + + final PermissionsDTO permissions = new PermissionsDTO(); + permissions.setCanRead(false); + permissions.setCanWrite(false); + + final ProcessGroupEntity entity = new EntityFactory().createProcessGroupEntity(dto, null, permissions, null, null); + + assertEquals("STATELESS", entity.getResolvedExecutionEngine()); + assertNull(entity.getComponent()); + } + + @Test + void testCreateProcessGroupEntityIncludesComponentWhenAuthorized() { + final ProcessGroupDTO dto = new ProcessGroupDTO(); + dto.setId("pg-1"); + dto.setResolvedExecutionEngine("STANDARD"); + + final PermissionsDTO permissions = new PermissionsDTO(); + permissions.setCanRead(true); + permissions.setCanWrite(true); + + final ProcessGroupEntity entity = new EntityFactory().createProcessGroupEntity(dto, null, permissions, null, null); + + assertEquals("STANDARD", entity.getResolvedExecutionEngine()); + assertEquals(dto, entity.getComponent()); + } +} diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts new file mode 100644 index 000000000000..a4c30864e3ee --- /dev/null +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts @@ -0,0 +1,175 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { TestBed } from '@angular/core/testing'; +import { MockStore, provideMockStore } from '@ngrx/store/testing'; +import type { Mock } from 'vitest'; + +import { CanvasContextMenu } from './canvas-context-menu.service'; +import { CanvasUtils } from './canvas-utils.service'; +import { Client } from '../../../service/client.service'; +import { CanvasView } from './canvas-view.service'; +import { CanvasActionsService } from './canvas-actions.service'; +import { DraggableBehavior } from './behavior/draggable-behavior.service'; +import * as FlowActions from '../state/flow/flow.actions'; +import type { ResolvedExecutionEngine } from '../state/flow'; +import type { ContextMenuItemDefinition } from '../../../ui/common/context-menu/context-menu.component'; + +interface SetupOptions { + currentProcessGroupId?: string; + isProcessGroup?: boolean; + resolvedExecutionEngine?: ResolvedExecutionEngine; +} + +function menuItem(menuItems: ContextMenuItemDefinition[], text: string): ContextMenuItemDefinition { + const item = menuItems.find((candidate) => candidate.text === text); + if (!item) { + throw new Error(`Expected menu item "${text}"`); + } + return item; +} + +async function setup(options: SetupOptions = {}) { + const currentProcessGroupId = options.currentProcessGroupId ?? 'current-pg'; + const canvasUtils = { + getProcessGroupId: vi.fn().mockReturnValue(currentProcessGroupId), + isProcessGroup: vi.fn().mockReturnValue(options.isProcessGroup ?? false), + getResolvedExecutionEngine: vi.fn().mockReturnValue(options.resolvedExecutionEngine ?? 'STANDARD') + }; + + await TestBed.configureTestingModule({ + providers: [ + CanvasContextMenu, + provideMockStore(), + { provide: CanvasUtils, useValue: canvasUtils }, + { provide: Client, useValue: {} }, + { provide: CanvasView, useValue: {} }, + { + provide: CanvasActionsService, + useValue: { + getConditionFunction: () => () => false, + getActionFunction: () => () => undefined + } + }, + { provide: DraggableBehavior, useValue: {} } + ] + }).compileComponents(); + + const service = TestBed.inject(CanvasContextMenu); + const store = TestBed.inject(MockStore); + const dispatchSpy = vi.spyOn(store, 'dispatch') as Mock; + + return { service, canvasUtils, dispatchSpy }; +} + +describe('CanvasContextMenu', () => { + describe('Stop sources', () => { + it('is immediately after Stop in the root menu', async () => { + const { service } = await setup(); + const menuItems = service.getMenu('root')!.menuItems; + const texts = menuItems.map((item) => item.text); + const stopIndex = texts.indexOf('Stop'); + const stopSourcesIndex = texts.indexOf('Stop sources'); + + expect(stopIndex).toBeGreaterThan(-1); + expect(stopSourcesIndex).toBe(stopIndex + 1); + }); + + it('is visible on an empty canvas and dispatches stopSources for the current group', async () => { + const { service, canvasUtils, dispatchSpy } = await setup({ currentProcessGroupId: 'current-pg' }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { empty: () => true }; + + expect(stopSources.condition!(selection as never)).toBe(true); + stopSources.action!(selection as never); + + expect(canvasUtils.getProcessGroupId).toHaveBeenCalled(); + expect(dispatchSpy).toHaveBeenCalledWith(FlowActions.stopSources({ request: { id: 'current-pg' } })); + }); + + it('is visible for a selected process group and dispatches stopSources for that group', async () => { + const { service, dispatchSpy } = await setup({ isProcessGroup: true }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { + empty: () => false, + datum: () => ({ id: 'pg-child', resolvedExecutionEngine: 'STANDARD' }) + }; + + expect(stopSources.condition!(selection as never)).toBe(true); + stopSources.action!(selection as never); + + expect(dispatchSpy).toHaveBeenCalledWith(FlowActions.stopSources({ request: { id: 'pg-child' } })); + }); + + it('is hidden when the selection is not a process group', async () => { + const { service } = await setup({ isProcessGroup: false }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { empty: () => false }; + + expect(stopSources.condition!(selection as never)).toBe(false); + }); + + it('is hidden on an empty canvas when the current group resolves to STATELESS', async () => { + const { service } = await setup({ resolvedExecutionEngine: 'STATELESS' }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { empty: () => true }; + + expect(stopSources.condition!(selection as never)).toBe(false); + }); + + it('is hidden for a selected process group when the resolved engine is missing', async () => { + const { service } = await setup({ isProcessGroup: true }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { + empty: () => false, + datum: () => ({ id: 'pg-child' }) + }; + + expect(stopSources.condition!(selection as never)).toBe(false); + }); + + it('is visible on an empty canvas when the current group is configured INHERITED but resolves to STANDARD', async () => { + const { service } = await setup({ resolvedExecutionEngine: 'STANDARD' }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { empty: () => true }; + + expect(stopSources.condition!(selection as never)).toBe(true); + }); + + it('is visible for a selected process group with no component when the entity resolves to STANDARD', async () => { + const { service } = await setup({ isProcessGroup: true }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { + empty: () => false, + datum: () => ({ id: 'pg-child', resolvedExecutionEngine: 'STANDARD' }) + }; + + expect(stopSources.condition!(selection as never)).toBe(true); + }); + + it('is hidden for a selected process group with no component when the entity resolves to STATELESS', async () => { + const { service } = await setup({ isProcessGroup: true }); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const selection = { + empty: () => false, + datum: () => ({ id: 'pg-child', resolvedExecutionEngine: 'STATELESS' }) + }; + + expect(stopSources.condition!(selection as never)).toBe(false); + }); + }); +}); diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts index 5eb220748949..73f365bf7fef 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts @@ -47,6 +47,7 @@ import { replayLastProvenanceEvent, requestRefreshRemoteProcessGroup, runOnce, + stopSources, stopVersionControlRequest, terminateThreads, updatePositions @@ -773,6 +774,36 @@ export class CanvasContextMenu implements ContextMenuDefinitionProvider { text: 'Stop', action: this.canvasActionsService.getActionFunction('stop') }, + { + condition: (selection: any) => { + if (!(selection.empty() || this.canvasUtils.isProcessGroup(selection))) { + return false; + } + const resolved = selection.empty() + ? this.canvasUtils.getResolvedExecutionEngine() + : selection.datum().resolvedExecutionEngine; + return resolved === 'STANDARD'; + }, + clazz: 'fa fa-stop-circle-o', + text: 'Stop sources', + action: (selection: any) => { + let processGroupId: string; + if (selection.empty()) { + processGroupId = this.canvasUtils.getProcessGroupId(); + } else { + const selectionData = selection.datum(); + processGroupId = selectionData.id; + } + + this.store.dispatch( + stopSources({ + request: { + id: processGroupId + } + }) + ); + } + }, { condition: (selection: any) => { if (selection.size() !== 1) { diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-utils.service.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-utils.service.ts index b8847fc7271e..8da73babd8e5 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-utils.service.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-utils.service.ts @@ -20,13 +20,15 @@ import * as d3 from 'd3'; import { humanizer, Humanizer } from 'humanize-duration'; import { Store } from '@ngrx/store'; import { CanvasState } from '../state'; +import type { ResolvedExecutionEngine } from '../state/flow'; import { selectBreadcrumbs, selectCanvasPermissions, selectConnections, selectCurrentParameterContext, selectCurrentProcessGroupId, - selectParentProcessGroupId + selectParentProcessGroupId, + selectResolvedExecutionEngine } from '../state/flow/flow.selectors'; import { initialState as initialFlowState } from '../state/flow/flow.reducer'; import { takeUntilDestroyed } from '@angular/core/rxjs-interop'; @@ -80,6 +82,8 @@ export class CanvasUtils { private trimLengthCaches: Map>> = new Map(); private currentProcessGroupId: string = initialFlowState.id; private parentProcessGroupId: string | null = initialFlowState.flow.processGroupFlow.parentGroupId; + private currentResolvedExecutionEngine: ResolvedExecutionEngine = + initialFlowState.flow.processGroupFlow.resolvedExecutionEngine; private canvasPermissions: Permissions = initialFlowState.flow.permissions; private currentUser: CurrentUser = initialUserState.user; private currentParameterContext: ParameterContextReferenceEntity | null = @@ -110,6 +114,13 @@ export class CanvasUtils { this.parentProcessGroupId = parentProcessGroupId; }); + this.store + .select(selectResolvedExecutionEngine) + .pipe(takeUntilDestroyed(this.destroyRef)) + .subscribe((resolvedExecutionEngine) => { + this.currentResolvedExecutionEngine = resolvedExecutionEngine; + }); + this.store .select(selectCanvasPermissions) .pipe(takeUntilDestroyed(this.destroyRef)) @@ -229,6 +240,14 @@ export class CanvasUtils { return this.currentProcessGroupId; } + /** + * The Execution Engine that will actually run the current Process Group + * after resolving INHERITED. Always STANDARD or STATELESS. + */ + public getResolvedExecutionEngine(): ResolvedExecutionEngine { + return this.currentResolvedExecutionEngine; + } + /** * Returns the parent group id or null if current is root. */ diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/flow.service.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/flow.service.ts index b605afeda3f2..1f5f469ab300 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/flow.service.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/flow.service.ts @@ -43,6 +43,8 @@ import { SaveToVersionControlRequest, StartProcessGroupRequest, StopProcessGroupRequest, + StopSourcesRequest, + StopSourcesResponse, StopVersionControlRequest, TerminateThreadsRequest, UploadProcessGroupRequest, @@ -403,6 +405,18 @@ export class FlowService implements PropertyDescriptorRetriever { return this.httpClient.put(`${FlowService.API}/flow/process-groups/${request.id}`, stopRequest); } + stopSources(request: StopSourcesRequest): Observable { + const stopRequest: ProcessGroupRunStatusRequest = { + id: request.id, + disconnectedNodeAcknowledged: this.clusterConnectionService.isDisconnectionAcknowledged(), + state: 'STOPPED' + }; + return this.httpClient.put( + `${FlowService.API}/flow/process-groups/${request.id}/sources`, + stopRequest + ); + } + stopRemoteProcessGroupsInProcessGroup(request: StopProcessGroupRequest): Observable { const stopRequest = { disconnectedNodeAcknowledged: this.clusterConnectionService.isDisconnectionAcknowledged(), diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.actions.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.actions.ts index 76587c6a0cc4..441f5b5c5e8d 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.actions.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.actions.ts @@ -92,6 +92,7 @@ import { StopComponentsRequest, StopProcessGroupRequest, StopProcessGroupResponse, + StopSourcesRequest, StopVersionControlRequest, StopVersionControlResponse, TerminateThreadsRequest, @@ -749,6 +750,13 @@ export const startCurrentProcessGroup = createAction(`${CANVAS_PREFIX} Start Cur export const stopCurrentProcessGroup = createAction(`${CANVAS_PREFIX} Stop Current Process Group`); +export const stopSources = createAction(`${CANVAS_PREFIX} Stop Sources`, props<{ request: StopSourcesRequest }>()); + +export const stopSourcesSuccess = createAction( + `${CANVAS_PREFIX} Stop Sources Success`, + props<{ response: StopProcessGroupResponse }>() +); + export const enableControllerServicesInCurrentProcessGroup = createAction( `${CANVAS_PREFIX} Enable Controller Services In Current Process Group` ); diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts index f8ffd720ac88..fb83f4b0606b 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts @@ -834,7 +834,8 @@ describe('FlowEffects', () => { createConnection: vi.fn(), createLabel: vi.fn(), clearBulletinsForProcessGroup: vi.fn(), - submitProcessorBacklogRequest: vi.fn() + submitProcessorBacklogRequest: vi.fn(), + stopSources: vi.fn() } }, { @@ -1127,6 +1128,107 @@ describe('FlowEffects', () => { }); }); + describe('stopSources$', () => { + beforeEach(() => { + effects = TestBed.inject(FlowEffects); + }); + + it('should call flowService.stopSources and dispatch stopSourcesSuccess', async () => { + vi.spyOn(flowService, 'stopSources').mockReturnValue( + of({ id: 'test-group-id', state: 'STOPPED', components: {} }) + ); + + action$.next(FlowActions.stopSources({ request: { id: 'test-group-id' } })); + + const result = await new Promise((resolve) => effects.stopSources$.pipe(take(1)).subscribe(resolve)); + + expect(result).toEqual( + FlowActions.stopSourcesSuccess({ + response: { + type: ComponentType.ProcessGroup, + component: { id: 'test-group-id', state: 'STOPPED' } + } + }) + ); + expect(flowService.stopSources).toHaveBeenCalledWith({ id: 'test-group-id' }); + }); + + it('should dispatch flowSnackbarError and not stopSourcesSuccess when stopSources fails', async () => { + const errorHelper = TestBed.inject(ErrorHelper); + const errorResponse = new HttpErrorResponse({ + error: 'stop sources failed', + status: 409, + statusText: 'Conflict' + }); + vi.spyOn(flowService, 'stopSources').mockReturnValue(throwError(() => errorResponse)); + vi.spyOn(errorHelper, 'getErrorString').mockReturnValue('Formatted error message'); + + action$.next(FlowActions.stopSources({ request: { id: 'test-group-id' } })); + + const result = await new Promise((resolve) => effects.stopSources$.pipe(take(1)).subscribe(resolve)); + + expect(result).toEqual(FlowActions.flowSnackbarError({ error: 'Formatted error message' })); + expect((result as Action).type).not.toBe(FlowActions.stopSourcesSuccess.type); + expect(flowService.stopSources).toHaveBeenCalledWith({ id: 'test-group-id' }); + }); + }); + + describe('stopSourcesCurrentProcessGroupSuccess$', () => { + beforeEach(() => { + effects = TestBed.inject(FlowEffects); + }); + + it('should dispatch reloadFlow when sources are stopped in the current process group', async () => { + store.overrideSelector(selectCurrentProcessGroupId, 'test-group-id'); + store.refreshState(); + + action$.next( + FlowActions.stopSourcesSuccess({ + response: { + type: ComponentType.ProcessGroup, + component: { id: 'test-group-id', state: 'STOPPED' } + } + }) + ); + + const result = await new Promise((resolve) => + effects.stopSourcesCurrentProcessGroupSuccess$.pipe(take(1)).subscribe(resolve) + ); + + expect(result).toEqual(FlowActions.reloadFlow()); + }); + }); + + describe('stopSourcesSuccess$', () => { + beforeEach(() => { + effects = TestBed.inject(FlowEffects); + }); + + it('should dispatch loadChildProcessGroup when sources are stopped in a child process group', async () => { + store.overrideSelector(selectCurrentProcessGroupId, 'current-pg'); + store.refreshState(); + + action$.next( + FlowActions.stopSourcesSuccess({ + response: { + type: ComponentType.ProcessGroup, + component: { id: 'pg-child', state: 'STOPPED' } + } + }) + ); + + const result = await new Promise((resolve) => effects.stopSourcesSuccess$.pipe(take(1)).subscribe(resolve)); + + expect(result).toEqual( + FlowActions.loadChildProcessGroup({ + request: { + id: 'pg-child' + } + }) + ); + }); + }); + describe('clearBulletinsForProcessGroup$', () => { beforeEach(() => { effects = TestBed.inject(FlowEffects); diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts index 5413002cecfb..ec40ed54e953 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts @@ -3656,6 +3656,64 @@ export class FlowEffects { ) ); + stopSources$ = createEffect(() => + this.actions$.pipe( + ofType(FlowActions.stopSources), + map((action) => action.request), + switchMap((request) => + from(this.flowService.stopSources(request)).pipe( + map((response) => + FlowActions.stopSourcesSuccess({ + response: { + type: ComponentType.ProcessGroup, + component: { + id: response.id, + state: response.state + } + } + }) + ), + catchError((errorResponse: HttpErrorResponse) => of(this.snackBarOrFullScreenError(errorResponse))) + ) + ) + ) + ); + + /** + * If sources were stopped in the current process group, reload the flow + */ + stopSourcesCurrentProcessGroupSuccess$ = createEffect(() => + this.actions$.pipe( + ofType(FlowActions.stopSourcesSuccess), + map((action) => action.response), + concatLatestFrom(() => this.store.select(selectCurrentProcessGroupId)), + filter(([response, currentPg]) => response.component.id === currentPg), + switchMap(() => of(FlowActions.reloadFlow())) + ) + ); + + /** + * If sources were stopped in a child ProcessGroup, reload that row; the + * schedule response does not contain all the displayed info + */ + stopSourcesSuccess$ = createEffect(() => + this.actions$.pipe( + ofType(FlowActions.stopSourcesSuccess), + map((action) => action.response), + concatLatestFrom(() => this.store.select(selectCurrentProcessGroupId)), + filter(([response, currentPg]) => response.component.id !== currentPg), + switchMap(([response]) => + of( + FlowActions.loadChildProcessGroup({ + request: { + id: response.component.id + } + }) + ) + ) + ) + ); + enableControllerServicesInCurrentProcessGroup$ = createEffect(() => this.actions$.pipe( ofType(FlowActions.enableControllerServicesInCurrentProcessGroup), diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.reducer.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.reducer.ts index f9d0122704c1..6e6f602ef428 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.reducer.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.reducer.ts @@ -73,6 +73,8 @@ import { startComponentSuccess, startPollingProcessorUntilStopped, startProcessGroupSuccess, + stopSources, + stopSourcesSuccess, startRemoteProcessGroupPolling, stopComponent, stopComponentSuccess, @@ -123,6 +125,7 @@ export const initialState: FlowState = { } }, parameterContext: null, + resolvedExecutionEngine: 'STANDARD', flow: { processGroups: [], remoteProcessGroups: [], @@ -432,6 +435,7 @@ export const flowReducer = createReducer( disableComponent, startComponent, stopComponent, + stopSources, runOnce, (state) => ({ ...state, @@ -443,6 +447,7 @@ export const flowReducer = createReducer( disableProcessGroupSuccess, startProcessGroupSuccess, stopProcessGroupSuccess, + stopSourcesSuccess, (state) => ({ ...state, saving: false diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.selectors.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.selectors.ts index 27b84bddbb0a..a699a24460a9 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.selectors.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.selectors.ts @@ -70,6 +70,11 @@ export const selectParentProcessGroupId = createSelector( (state: FlowState) => state.flow.processGroupFlow.parentGroupId ); +export const selectResolvedExecutionEngine = createSelector( + selectFlowState, + (state: FlowState) => state.flow.processGroupFlow.resolvedExecutionEngine +); + export const selectProcessGroupIdFromRoute = createSelector(selectCurrentRoute, (route) => { if (route) { // always select the process group from the route diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/index.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/index.ts index c4e821d01524..1b7b830ec22b 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/index.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/index.ts @@ -470,6 +470,8 @@ export interface Flow { funnels: ComponentEntity[]; } +export type ResolvedExecutionEngine = 'STANDARD' | 'STATELESS'; + export interface ProcessGroupFlow { id: string; uri: string; @@ -478,6 +480,7 @@ export interface ProcessGroupFlow { parameterContext: ParameterContextReferenceEntity | null; flow: Flow; lastRefreshed: string; + resolvedExecutionEngine: ResolvedExecutionEngine; } export interface ProcessGroupFlowEntity { @@ -648,6 +651,17 @@ export interface StopProcessGroupRequest { errorStrategy: 'snackbar' | 'banner'; } +export interface StopSourcesRequest { + id: string; +} + +export interface StopSourcesResponse { + id: string; + state: 'STOPPED'; + components: Record; + disconnectedNodeAcknowledged?: boolean; +} + export interface StopComponentResponse { type: ComponentType; component: ComponentEntity; diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/pg/ClusteredStopSourcesIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/pg/ClusteredStopSourcesIT.java new file mode 100644 index 000000000000..3b5b2c9e1ec1 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/pg/ClusteredStopSourcesIT.java @@ -0,0 +1,260 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.nifi.tests.system.pg; + +import jakarta.ws.rs.WebApplicationException; +import org.apache.nifi.controller.ScheduledState; +import org.apache.nifi.groups.StatelessGroupScheduledState; +import org.apache.nifi.tests.system.NiFiInstanceFactory; +import org.apache.nifi.tests.system.NiFiSystemIT; +import org.apache.nifi.toolkit.client.NiFiClientException; +import org.apache.nifi.web.api.entity.ProcessGroupEntity; +import org.apache.nifi.web.api.entity.ProcessorEntity; +import org.apache.nifi.web.api.entity.ScheduleComponentsEntity; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.util.Map; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.fail; + +public class ClusteredStopSourcesIT extends NiFiSystemIT { + + @Override + public NiFiInstanceFactory getInstanceFactory() { + return createTwoNodeInstanceFactory(); + } + + @Test + public void testStopSourcesSkipsStatelessDescendants() throws NiFiClientException, IOException, InterruptedException { + final MixedEngineFlow flow = createMixedEngineFlow(); + startMixedEngineFlow(flow); + + final ScheduleComponentsEntity response = getNifiClient().getFlowClient() + .stopProcessGroupSources(flow.parentGroup().getId(), stopSourcesRequest(flow.parentGroup().getId())); + + assertEquals(Set.of(flow.standardSource().getId()), response.getComponents().keySet()); + assertProcessorStateOnAllNodes(flow.standardSource().getId(), ScheduledState.STOPPED); + assertProcessorStateOnAllNodes(flow.standardDownstream().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessSource().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessDownstream().getId(), ScheduledState.RUNNING); + assertStatelessGroupStateOnAllNodes(flow.statelessGroup().getId(), StatelessGroupScheduledState.RUNNING); + } + + @Test + public void testStopSourcesRejectsStatelessProcessGroup() throws NiFiClientException, IOException, InterruptedException { + final MixedEngineFlow flow = createMixedEngineFlow(); + startMixedEngineFlow(flow); + + final NiFiClientException exception = assertThrows(NiFiClientException.class, () -> getNifiClient().getFlowClient() + .stopProcessGroupSources(flow.statelessGroup().getId(), stopSourcesRequest(flow.statelessGroup().getId()))); + + assertConflict(exception); + assertProcessorStateOnAllNodes(flow.statelessSource().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessDownstream().getId(), ScheduledState.RUNNING); + assertStatelessGroupStateOnAllNodes(flow.statelessGroup().getId(), StatelessGroupScheduledState.RUNNING); + } + + @Test + public void testStopSourcesRejectsMismatchedComponentIds() throws NiFiClientException, IOException, InterruptedException { + final MixedEngineFlow flow = createMixedEngineFlow(); + startMixedEngineFlow(flow); + final ScheduleComponentsEntity request = stopSourcesRequest(flow.parentGroup().getId()); + request.setComponents(Map.of()); + + final NiFiClientException exception = assertThrows(NiFiClientException.class, () -> getNifiClient().getFlowClient() + .stopProcessGroupSources(flow.parentGroup().getId(), request)); + + assertConflict(exception); + assertProcessorStateOnAllNodes(flow.standardSource().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.standardDownstream().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessSource().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessDownstream().getId(), ScheduledState.RUNNING); + assertStatelessGroupStateOnAllNodes(flow.statelessGroup().getId(), StatelessGroupScheduledState.RUNNING); + } + + @Test + public void testStopSourcesRejectsWhenNodeIdentifiesDifferentSources() throws NiFiClientException, IOException, InterruptedException { + final MixedEngineFlow flow = createMixedEngineFlow(); + startMixedEngineFlow(flow); + final long node1Revision = stopAndRestartProcessorOnNode(flow.standardSource().getId(), 1); + final long node2Revision = stopAndRenameProcessorOnNode(flow.standardSource().getId(), 2); + + assertEquals(node1Revision, node2Revision, "Processor revisions must match so source-set verification rejects the request"); + assertProcessorStateOnNode(flow.standardSource().getId(), ScheduledState.RUNNING, 1); + assertProcessorStateOnNode(flow.standardSource().getId(), ScheduledState.STOPPED, 2); + + final NiFiClientException exception = assertThrows(NiFiClientException.class, () -> getNifiClient().getFlowClient() + .stopProcessGroupSources(flow.parentGroup().getId(), stopSourcesRequest(flow.parentGroup().getId()))); + + assertConflict(exception); + assertProcessorStateOnNode(flow.standardSource().getId(), ScheduledState.RUNNING, 1); + assertProcessorStateOnNode(flow.standardSource().getId(), ScheduledState.STOPPED, 2); + assertProcessorStateOnAllNodes(flow.standardDownstream().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessSource().getId(), ScheduledState.RUNNING); + assertProcessorStateOnAllNodes(flow.statelessDownstream().getId(), ScheduledState.RUNNING); + assertStatelessGroupStateOnAllNodes(flow.statelessGroup().getId(), StatelessGroupScheduledState.RUNNING); + } + + private MixedEngineFlow createMixedEngineFlow() throws NiFiClientException, IOException, InterruptedException { + final ProcessGroupEntity parentGroup = getClientUtil().createProcessGroup("Parent", "root"); + final ProcessGroupEntity standardGroup = getClientUtil().createProcessGroup("Standard", parentGroup.getId()); + final ProcessGroupEntity statelessGroup = getClientUtil().createProcessGroup("Stateless", parentGroup.getId()); + getClientUtil().markStateless(statelessGroup, "1 min"); + + final ProcessorEntity standardSource = getClientUtil().createProcessor(GENERATE_FLOWFILE, standardGroup.getId()); + final ProcessorEntity standardDownstream = getClientUtil().createProcessor(TERMINATE_FLOWFILE, standardGroup.getId()); + getClientUtil().createConnection(standardSource, standardDownstream, SUCCESS, standardGroup.getId()); + + final ProcessorEntity statelessSource = getClientUtil().createProcessor(GENERATE_FLOWFILE, statelessGroup.getId()); + final ProcessorEntity statelessDownstream = getClientUtil().createProcessor(TERMINATE_FLOWFILE, statelessGroup.getId()); + getClientUtil().createConnection(statelessSource, statelessDownstream, SUCCESS, statelessGroup.getId()); + + getClientUtil().waitForValidProcessor(standardSource.getId()); + getClientUtil().waitForValidProcessor(standardDownstream.getId()); + getClientUtil().waitForValidProcessor(statelessSource.getId()); + getClientUtil().waitForValidProcessor(statelessDownstream.getId()); + + return new MixedEngineFlow(parentGroup, standardGroup, statelessGroup, + standardSource, standardDownstream, statelessSource, statelessDownstream); + } + + private void startMixedEngineFlow(final MixedEngineFlow flow) throws NiFiClientException, IOException, InterruptedException { + getClientUtil().startProcessGroupComponents(flow.standardGroup().getId()); + getClientUtil().startProcessGroupComponents(flow.statelessGroup().getId()); + getClientUtil().waitForRunningProcessor(flow.standardSource().getId()); + getClientUtil().waitForRunningProcessor(flow.standardDownstream().getId()); + getClientUtil().waitForRunningProcessor(flow.statelessSource().getId()); + getClientUtil().waitForRunningProcessor(flow.statelessDownstream().getId()); + } + + private ScheduleComponentsEntity stopSourcesRequest(final String groupId) { + final ScheduleComponentsEntity request = new ScheduleComponentsEntity(); + request.setId(groupId); + request.setState(ScheduledState.STOPPED.name()); + request.setDisconnectedNodeAcknowledged(true); + return request; + } + + private long stopAndRestartProcessorOnNode(final String processorId, final int nodeIndex) + throws NiFiClientException, IOException, InterruptedException { + try { + switchClientToNode(nodeIndex); + final ProcessorEntity processor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + processor.setDisconnectedNodeAcknowledged(true); + getNifiClient().getProcessorClient(DO_NOT_REPLICATE).stopProcessor(processor); + waitForProcessorStateOnCurrentNode(processorId, ScheduledState.STOPPED); + + final ProcessorEntity stoppedProcessor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + stoppedProcessor.setDisconnectedNodeAcknowledged(true); + final ProcessorEntity restartedProcessor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).startProcessor(stoppedProcessor); + return restartedProcessor.getRevision().getVersion(); + } finally { + switchClientToNode(1); + } + } + + private long stopAndRenameProcessorOnNode(final String processorId, final int nodeIndex) + throws NiFiClientException, IOException, InterruptedException { + try { + switchClientToNode(nodeIndex); + final ProcessorEntity processor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + processor.setDisconnectedNodeAcknowledged(true); + getNifiClient().getProcessorClient(DO_NOT_REPLICATE).stopProcessor(processor); + waitForProcessorStateOnCurrentNode(processorId, ScheduledState.STOPPED); + + final ProcessorEntity stoppedProcessor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + stoppedProcessor.getComponent().setName(stoppedProcessor.getComponent().getName() + " Node " + nodeIndex); + stoppedProcessor.setDisconnectedNodeAcknowledged(true); + final ProcessorEntity renamedProcessor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).updateProcessor(stoppedProcessor); + return renamedProcessor.getRevision().getVersion(); + } finally { + switchClientToNode(1); + } + } + + private void waitForProcessorStateOnCurrentNode(final String processorId, final ScheduledState expectedState) throws InterruptedException { + waitFor(() -> { + final ProcessorEntity processor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + return expectedState.name().equals(processor.getComponent().getState()) + && expectedState.name().equals(processor.getComponent().getPhysicalState()); + }); + } + + private void assertProcessorStateOnNode(final String processorId, final ScheduledState expectedState, final int nodeIndex) + throws NiFiClientException, IOException { + try { + switchClientToNode(nodeIndex); + final ProcessorEntity processor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + assertEquals(expectedState.name(), processor.getComponent().getState(), + "Unexpected state for Processor %s on Node %d".formatted(processorId, nodeIndex)); + } finally { + switchClientToNode(1); + } + } + + private void assertProcessorStateOnAllNodes(final String processorId, final ScheduledState expectedState) + throws NiFiClientException, IOException { + try { + for (int nodeIndex = 1; nodeIndex <= 2; nodeIndex++) { + switchClientToNode(nodeIndex); + final ProcessorEntity processor = getNifiClient().getProcessorClient(DO_NOT_REPLICATE).getProcessor(processorId); + assertEquals(expectedState.name(), processor.getComponent().getState(), + "Unexpected state for Processor %s on Node %d".formatted(processorId, nodeIndex)); + } + } finally { + switchClientToNode(1); + } + } + + private void assertStatelessGroupStateOnAllNodes(final String groupId, final StatelessGroupScheduledState expectedState) + throws NiFiClientException, IOException { + try { + for (int nodeIndex = 1; nodeIndex <= 2; nodeIndex++) { + switchClientToNode(nodeIndex); + final ProcessGroupEntity group = getNifiClient().getProcessGroupClient(DO_NOT_REPLICATE).getProcessGroup(groupId); + assertEquals(expectedState.name(), group.getComponent().getStatelessGroupScheduledState(), + "Unexpected state for Stateless Process Group %s on Node %d".formatted(groupId, nodeIndex)); + } + } finally { + switchClientToNode(1); + } + } + + private void assertConflict(final NiFiClientException exception) { + final Throwable cause = exception.getCause(); + if (cause instanceof final WebApplicationException webApplicationException) { + assertEquals(409, webApplicationException.getResponse().getStatus()); + return; + } + + fail("Expected WebApplicationException 409, got: " + cause); + } + + private record MixedEngineFlow( + ProcessGroupEntity parentGroup, + ProcessGroupEntity standardGroup, + ProcessGroupEntity statelessGroup, + ProcessorEntity standardSource, + ProcessorEntity standardDownstream, + ProcessorEntity statelessSource, + ProcessorEntity statelessDownstream) { + } +} diff --git a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/FlowClient.java b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/FlowClient.java index 5218f17bf3c3..a2c0818f4f35 100644 --- a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/FlowClient.java +++ b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/FlowClient.java @@ -85,6 +85,16 @@ public interface FlowClient { ScheduleComponentsEntity scheduleProcessGroupComponents( String processGroupId, ScheduleComponentsEntity scheduleComponentsEntity) throws NiFiClientException, IOException; + /** + * Stops source components in a process group. + * + * @param processGroupId the id of a process group + * @param scheduleComponentsEntity the scheduled state to update to + * @return the entity representing the stopped source components + */ + ScheduleComponentsEntity stopProcessGroupSources( + String processGroupId, ScheduleComponentsEntity scheduleComponentsEntity) throws NiFiClientException, IOException; + /** * Gets the possible versions for the given flow in the given bucket in the * given registry. diff --git a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyFlowClient.java b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyFlowClient.java index acf0a700a534..8a684cbfa9e5 100644 --- a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyFlowClient.java +++ b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyFlowClient.java @@ -166,6 +166,32 @@ public ScheduleComponentsEntity scheduleProcessGroupComponents( }); } + @Override + public ScheduleComponentsEntity stopProcessGroupSources( + final String processGroupId, final ScheduleComponentsEntity scheduleComponentsEntity) + throws NiFiClientException, IOException { + + if (StringUtils.isBlank(processGroupId)) { + throw new IllegalArgumentException("Process group id cannot be null"); + } + + if (scheduleComponentsEntity == null) { + throw new IllegalArgumentException("ScheduleComponentsEntity cannot be null"); + } + + scheduleComponentsEntity.setId(processGroupId); + + return executeAction("Error stopping process group sources", () -> { + final WebTarget target = flowTarget + .path("process-groups/{id}/sources") + .resolveTemplate("id", processGroupId); + + return getRequestBuilder(target).put( + Entity.entity(scheduleComponentsEntity, MediaType.APPLICATION_JSON_TYPE), + ScheduleComponentsEntity.class); + }); + } + @Override public VersionedFlowSnapshotMetadataSetEntity getVersions(final String registryId, final String bucketId, final String flowId) throws NiFiClientException, IOException { From 73ec93a6b4428f1276756846cbf18e6c12ab379c Mon Sep 17 00:00:00 2001 From: Matt Gilman Date: Tue, 15 Sep 2026 13:12:26 -0400 Subject: [PATCH 2/2] NIFI-16337: Addressing review feedback. --- .../canvas-context-menu.service.spec.ts | 20 +++++----- .../service/canvas-context-menu.service.ts | 2 +- .../state/flow/flow.effects.spec.ts | 38 +++++++++++++++++-- .../flow-designer/state/flow/flow.effects.ts | 2 +- 4 files changed, 47 insertions(+), 15 deletions(-) diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts index a4c30864e3ee..8d8f79668949 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.spec.ts @@ -77,13 +77,13 @@ async function setup(options: SetupOptions = {}) { } describe('CanvasContextMenu', () => { - describe('Stop sources', () => { + describe('Stop Sources', () => { it('is immediately after Stop in the root menu', async () => { const { service } = await setup(); const menuItems = service.getMenu('root')!.menuItems; const texts = menuItems.map((item) => item.text); const stopIndex = texts.indexOf('Stop'); - const stopSourcesIndex = texts.indexOf('Stop sources'); + const stopSourcesIndex = texts.indexOf('Stop Sources'); expect(stopIndex).toBeGreaterThan(-1); expect(stopSourcesIndex).toBe(stopIndex + 1); @@ -91,7 +91,7 @@ describe('CanvasContextMenu', () => { it('is visible on an empty canvas and dispatches stopSources for the current group', async () => { const { service, canvasUtils, dispatchSpy } = await setup({ currentProcessGroupId: 'current-pg' }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => true }; expect(stopSources.condition!(selection as never)).toBe(true); @@ -103,7 +103,7 @@ describe('CanvasContextMenu', () => { it('is visible for a selected process group and dispatches stopSources for that group', async () => { const { service, dispatchSpy } = await setup({ isProcessGroup: true }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => false, datum: () => ({ id: 'pg-child', resolvedExecutionEngine: 'STANDARD' }) @@ -117,7 +117,7 @@ describe('CanvasContextMenu', () => { it('is hidden when the selection is not a process group', async () => { const { service } = await setup({ isProcessGroup: false }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => false }; expect(stopSources.condition!(selection as never)).toBe(false); @@ -125,7 +125,7 @@ describe('CanvasContextMenu', () => { it('is hidden on an empty canvas when the current group resolves to STATELESS', async () => { const { service } = await setup({ resolvedExecutionEngine: 'STATELESS' }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => true }; expect(stopSources.condition!(selection as never)).toBe(false); @@ -133,7 +133,7 @@ describe('CanvasContextMenu', () => { it('is hidden for a selected process group when the resolved engine is missing', async () => { const { service } = await setup({ isProcessGroup: true }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => false, datum: () => ({ id: 'pg-child' }) @@ -144,7 +144,7 @@ describe('CanvasContextMenu', () => { it('is visible on an empty canvas when the current group is configured INHERITED but resolves to STANDARD', async () => { const { service } = await setup({ resolvedExecutionEngine: 'STANDARD' }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => true }; expect(stopSources.condition!(selection as never)).toBe(true); @@ -152,7 +152,7 @@ describe('CanvasContextMenu', () => { it('is visible for a selected process group with no component when the entity resolves to STANDARD', async () => { const { service } = await setup({ isProcessGroup: true }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => false, datum: () => ({ id: 'pg-child', resolvedExecutionEngine: 'STANDARD' }) @@ -163,7 +163,7 @@ describe('CanvasContextMenu', () => { it('is hidden for a selected process group with no component when the entity resolves to STATELESS', async () => { const { service } = await setup({ isProcessGroup: true }); - const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop sources'); + const stopSources = menuItem(service.getMenu('root')!.menuItems, 'Stop Sources'); const selection = { empty: () => false, datum: () => ({ id: 'pg-child', resolvedExecutionEngine: 'STATELESS' }) diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts index 73f365bf7fef..6267e71468fc 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/service/canvas-context-menu.service.ts @@ -785,7 +785,7 @@ export class CanvasContextMenu implements ContextMenuDefinitionProvider { return resolved === 'STANDARD'; }, clazz: 'fa fa-stop-circle-o', - text: 'Stop sources', + text: 'Stop Sources', action: (selection: any) => { let processGroupId: string; if (selection.empty()) { diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts index fb83f4b0606b..4dc36a34dd7a 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.spec.ts @@ -17,7 +17,7 @@ import { FlowService } from '../../service/flow.service'; import * as FlowActions from './flow.actions'; -import { of, ReplaySubject, take, throwError } from 'rxjs'; +import { firstValueFrom, of, ReplaySubject, Subject, take, throwError, toArray } from 'rxjs'; import { MatDialog, MatDialogRef } from '@angular/material/dialog'; import { ComponentHistoryEntity } from '../../../../state/shared'; import { EditProcessor } from '../../../../ui/common/component-dialogs/edit-processor/edit-processor.component'; @@ -35,7 +35,8 @@ import { CreateComponentResponse, CreateConnection, flowFeatureKey, - MoveToFrontRequest + MoveToFrontRequest, + StopSourcesResponse } from './index'; import { BacklogRequestEntity, @@ -81,7 +82,6 @@ import * as fromParameter from '../parameter/parameter.reducer'; import { flowAnalysisFeatureKey } from '../flow-analysis'; import * as fromFlowAnalysis from '../flow-analysis/flow-analysis.reducer'; import * as EmptyQueueActions from '../../../../state/empty-queue/empty-queue.actions'; -import { firstValueFrom } from 'rxjs'; describe('FlowEffects', () => { let action$: ReplaySubject; @@ -1153,6 +1153,38 @@ describe('FlowEffects', () => { expect(flowService.stopSources).toHaveBeenCalledWith({ id: 'test-group-id' }); }); + it('should keep overlapping stop sources requests active', async () => { + const firstResponse$ = new Subject(); + const secondResponse$ = new Subject(); + vi.spyOn(flowService, 'stopSources').mockImplementation((request) => + request.id === 'first-group-id' ? firstResponse$ : secondResponse$ + ); + const resultsPromise = firstValueFrom(effects.stopSources$.pipe(take(2), toArray())); + + action$.next(FlowActions.stopSources({ request: { id: 'first-group-id' } })); + action$.next(FlowActions.stopSources({ request: { id: 'second-group-id' } })); + secondResponse$.next({ id: 'second-group-id', state: 'STOPPED', components: {} }); + secondResponse$.complete(); + firstResponse$.next({ id: 'first-group-id', state: 'STOPPED', components: {} }); + firstResponse$.complete(); + + expect(await resultsPromise).toEqual([ + FlowActions.stopSourcesSuccess({ + response: { + type: ComponentType.ProcessGroup, + component: { id: 'second-group-id', state: 'STOPPED' } + } + }), + FlowActions.stopSourcesSuccess({ + response: { + type: ComponentType.ProcessGroup, + component: { id: 'first-group-id', state: 'STOPPED' } + } + }) + ]); + expect(flowService.stopSources).toHaveBeenCalledTimes(2); + }); + it('should dispatch flowSnackbarError and not stopSourcesSuccess when stopSources fails', async () => { const errorHelper = TestBed.inject(ErrorHelper); const errorResponse = new HttpErrorResponse({ diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts index ec40ed54e953..5d9220f4699a 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/flow-designer/state/flow/flow.effects.ts @@ -3660,7 +3660,7 @@ export class FlowEffects { this.actions$.pipe( ofType(FlowActions.stopSources), map((action) => action.request), - switchMap((request) => + mergeMap((request) => from(this.flowService.stopSources(request)).pipe( map((response) => FlowActions.stopSourcesSuccess({