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..8d8f79668949
--- /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..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
@@ -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..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;
@@ -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,139 @@ 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 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({
+ 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..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
@@ -3656,6 +3656,64 @@ export class FlowEffects {
)
);
+ stopSources$ = createEffect(() =>
+ this.actions$.pipe(
+ ofType(FlowActions.stopSources),
+ map((action) => action.request),
+ mergeMap((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 {