This is an automated email from the ASF dual-hosted git repository.
markap14 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 2597636c8e3 NIFI-16174 - Treat a stateless process group as a single
lifecycle unit when starting/stopping a controller service's referencing
components (#11515)
2597636c8e3 is described below
commit 2597636c8e363725d2cc0769f360bcf77ad8378c
Author: Noah <[email protected]>
AuthorDate: Fri Aug 28 13:15:00 2026 -0700
NIFI-16174 - Treat a stateless process group as a single lifecycle unit
when starting/stopping a controller service's referencing components (#11515)
* NIFI-16174 - Treat a stateless process group as a single lifecycle unit
when starting/stopping a controller service's referencing components
StandardControllerServiceProvider scheduled processors that reference a
controller service individually, even when they belong to a stateless
process
group. Because a stateless group is a single scheduling unit, this left the
group with a mixed running/stopped processor state and a group node stuck
RUNNING, from which it could not recover.
Resolve each referenced processor's owning stateless group (STATELESS ->
self,
INHERITED -> nearest explicit ancestor) and, for stateless members, stop the
group once via ProcessGroup.stopProcessing() / start it once via
ComponentScheduler.startStatelessGroup(), mapping the group's single future
to
every affected member. Standard processors are unchanged. Public
ControllerServiceProvider signatures are unchanged.
Adds unit coverage in StandardControllerServiceProviderTest and an
end-to-end
regression (ConnectorTroubleshootingIT) backed by a stateless
controller-service
reference in the ComponentLifecycleConnector test fixture.
* Update ComponentLifecycleConnector.java
* NIFI-16174 - Guard against null execution engine when resolving the
owning stateless group
A referenced processor's process group can report a null execution engine
(e.g. in unit-test fixtures backed by mock process groups). Treat a null
engine as non-stateless so getStatelessGroup returns null and the processor
is handled on the standard per-component path, rather than throwing an NPE
in the switch.
* NIFI-16174 - Resolve the top-most stateless group when starting/stopping
a controller service's referencing components
The previous helper returned the inner-most Process Group marked STATELESS.
Only
the top-most stateless group may be started or stopped directly, so on a
nested
stateless chain both paths became silent no-ops: startProcessing() logs a
warning
and returns, and stopProcessing() returns an already-completed Future
without
stopping anything.
Walk up while the group resolves to STATELESS and operate on the last one.
Document the rule on ProcessGroup.startProcessing()/stopProcessing(), and
pin the
no-op with a test, since it is what terminates the recursive
stopComponents() walk
during stateless shutdown.
The system test now carries a second stateless subtree whose only
referencing
processor lives in a nested stateless group, which is what makes it
discriminate:
when a referencing processor exists in the outer group too, the outer
group's
transition masks the nested no-op.
---
.../service/StandardControllerServiceProvider.java | 79 ++++++--
.../nifi/groups/StandardProcessGroupTest.java | 26 +++
.../service/ControllerServiceProvider.java | 15 ++
.../StandardControllerServiceProviderTest.java | 216 +++++++++++++++++++++
.../tests/system/ComponentLifecycleConnector.java | 81 ++++++--
.../connectors/ConnectorTroubleshootingIT.java | 109 +++++++++++
6 files changed, 496 insertions(+), 30 deletions(-)
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/service/StandardControllerServiceProvider.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/service/StandardControllerServiceProvider.java
index 4e0d5dadf43..fe1f98127fe 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/service/StandardControllerServiceProvider.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/controller/service/StandardControllerServiceProvider.java
@@ -27,6 +27,7 @@ import org.apache.nifi.controller.ReportingTaskNode;
import org.apache.nifi.controller.ScheduledState;
import org.apache.nifi.controller.flow.FlowManager;
import org.apache.nifi.events.BulletinFactory;
+import org.apache.nifi.flow.ExecutionEngine;
import org.apache.nifi.groups.ComponentScheduler;
import org.apache.nifi.groups.DefaultComponentScheduler;
import org.apache.nifi.groups.ProcessGroup;
@@ -148,17 +149,27 @@ public class StandardControllerServiceProvider implements
ControllerServiceProvi
}
}
- // start all of the components that are not disabled
+ // A processor in a stateless group is started through its group (the
group is a single scheduling unit), and
+ // each stateless group is started at most once.
final Set<ComponentNode> updated = new HashSet<>();
+ final Set<ProcessGroup> startedStatelessGroups = new HashSet<>();
for (final ProcessorNode node : processors) {
if (candidates != null && !candidates.contains(node)) {
continue;
}
- if (node.getScheduledState() != ScheduledState.DISABLED) {
+ if (node.getScheduledState() == ScheduledState.DISABLED) {
+ continue;
+ }
+
+ final ProcessGroup statelessGroup =
getTopMostStatelessGroup(node.getProcessGroup());
+ if (statelessGroup == null) {
componentScheduler.startComponent(node);
- updated.add(node);
+ } else if (startedStatelessGroups.add(statelessGroup)) {
+ componentScheduler.startStatelessGroup(statelessGroup);
}
+
+ updated.add(node);
}
for (final ReportingTaskNode node : reportingTasks) {
if (candidates != null && !candidates.contains(node)) {
@@ -194,15 +205,30 @@ public class StandardControllerServiceProvider implements
ControllerServiceProvi
final Map<ComponentNode, Future<Void>> updated = new HashMap<>();
+ // A processor in a stateless group is stopped through the group (a
single scheduling unit); standard processors
+ // are stopped individually.
+ final Map<ProcessGroup, List<ProcessorNode>> statelessMembersByGroup =
new HashMap<>();
+ final List<ProcessorNode> standardProcessors = new ArrayList<>();
+ for (final ProcessorNode node : processors) {
+ if (!isRunningOrStarting(node)) {
+ continue;
+ }
+
+ final ProcessGroup statelessGroup =
getTopMostStatelessGroup(node.getProcessGroup());
+ if (statelessGroup == null) {
+ standardProcessors.add(node);
+ } else {
+ statelessMembersByGroup.computeIfAbsent(statelessGroup, group
-> new ArrayList<>()).add(node);
+ }
+ }
+
// verify that we can stop all components (that are running or
starting) before doing anything
// Note: We check both RUNNING and STARTING states because a processor
might be stuck in STARTING
// state if it references an invalid controller service (e.g., after a
restart when the controller
// service configuration became invalid). Such processors need to be
stopped before the controller
- // service can be disabled.
- for (final ProcessorNode node : processors) {
- if (isRunningOrStarting(node)) {
- node.verifyCanStop();
- }
+ // service can be disabled. Stateless-group members are verified and
stopped through their group.
+ for (final ProcessorNode node : standardProcessors) {
+ node.verifyCanStop();
}
for (final ReportingTaskNode node : reportingTasks) {
if (isRunningOrStarting(node)) {
@@ -215,13 +241,19 @@ public class StandardControllerServiceProvider implements
ControllerServiceProvi
}
}
- // stop all of the components that are running or starting
- for (final ProcessorNode node : processors) {
- if (isRunningOrStarting(node)) {
- final Future<Void> future =
node.getProcessGroup().stopProcessor(node);
- updated.put(node, future);
+ // stop each stateless group once as a single unit, mapping the
group's single future to every affected member
+ for (final Map.Entry<ProcessGroup, List<ProcessorNode>> entry :
statelessMembersByGroup.entrySet()) {
+ final Future<Void> future = entry.getKey().stopProcessing();
+ for (final ProcessorNode member : entry.getValue()) {
+ updated.put(member, future);
}
}
+
+ // stop the standard processors
+ for (final ProcessorNode node : standardProcessors) {
+ final Future<Void> future =
node.getProcessGroup().stopProcessor(node);
+ updated.put(node, future);
+ }
for (final ReportingTaskNode node : reportingTasks) {
if (isRunningOrStarting(node)) {
final Future<Void> future = processScheduler.unschedule(node);
@@ -264,6 +296,27 @@ public class StandardControllerServiceProvider implements
ControllerServiceProvi
return scheduledState == ScheduledState.RUNNING || scheduledState ==
ScheduledState.STARTING;
}
+ /**
+ * Returns the top-most Process Group that runs using the Stateless Engine
and contains the given Process Group, or
+ * {@code null} if the given Process Group is {@code null} or does not run
using the Stateless Engine. Only the
+ * top-most stateless group may be started or stopped directly; {@link
ProcessGroup#startProcessing()} and
+ * {@link ProcessGroup#stopProcessing()} are no-ops on a nested stateless
group.
+ */
+ private ProcessGroup getTopMostStatelessGroup(final ProcessGroup start) {
+ ProcessGroup topMost = null;
+
+ // a Standard-Engine group cannot be a child of a stateless group, so
the stateless chain is contiguous
+ for (ProcessGroup group = start; group != null; group =
group.getParent()) {
+ if (group.resolveExecutionEngine() != ExecutionEngine.STATELESS) {
+ break;
+ }
+
+ topMost = group;
+ }
+
+ return topMost;
+ }
+
@Override
public CompletableFuture<Void> enableControllerService(final
ControllerServiceNode serviceNode) {
if (serviceNode.isActive()) {
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java
index 070bd5ad42a..0e46d63531f 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java
@@ -28,6 +28,7 @@ import org.apache.nifi.controller.ReloadComponent;
import org.apache.nifi.controller.flow.FlowManager;
import org.apache.nifi.controller.service.ControllerServiceProvider;
import org.apache.nifi.encrypt.PropertyEncryptor;
+import org.apache.nifi.flow.ExecutionEngine;
import org.apache.nifi.nar.ExtensionManager;
import org.apache.nifi.registry.flow.VersionControlInformation;
import org.apache.nifi.registry.flow.VersionedFlowStatus;
@@ -41,13 +42,18 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.io.IOException;
import java.util.Map;
import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -144,6 +150,26 @@ class StandardProcessGroupTest {
);
}
+ @Test
+ void testStartAndStopProcessingDoNothingWhenParentIsStateless() throws
Exception {
+ // Callers that resolve a component's owning stateless group rely on
this no-op: only the top-most stateless
+ // group may be started or stopped directly. It is also what
terminates the recursive stopComponents() walk,
+ // so making it throw or delegate upward would break stateless
shutdown.
+
when(parentProcessGroup.getExecutionEngine()).thenReturn(ExecutionEngine.STATELESS);
+
when(parentProcessGroup.resolveExecutionEngine()).thenReturn(ExecutionEngine.STATELESS);
+ processGroup.setParent(parentProcessGroup);
+ processGroup.setExecutionEngine(ExecutionEngine.STATELESS);
+
+ processGroup.startProcessing();
+
+ final CompletableFuture<Void> stopFuture =
processGroup.stopProcessing();
+
+ assertTrue(stopFuture.isDone());
+ assertNull(stopFuture.get());
+ verify(processScheduler, never()).startStatelessGroup(any());
+ verify(processScheduler, never()).stopStatelessGroup(any());
+ }
+
@Test
void testGetLoggingAttributesWithoutParentProcessGroup() {
processGroup.setName(NAME);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/service/ControllerServiceProvider.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/service/ControllerServiceProvider.java
index 6e375542796..362e1b814be 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/service/ControllerServiceProvider.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/service/ControllerServiceProvider.java
@@ -134,7 +134,15 @@ public interface ControllerServiceProvider extends
ControllerServiceLookup {
* Controller services that reference this one, its schedulable referencing
* components will also be unscheduled.
*
+ * A referencing processor that belongs to a stateless process group is not
+ * stopped individually; instead its owning stateless group is stopped as a
+ * single unit, and the group's single stop {@link Future} is mapped to
every
+ * affected processor in that group in the returned map.
+ *
* @param serviceNode the node
+ *
+ * @return a map of each affected component to the {@link Future} that
+ * completes when the component (or its stateless group) has stopped
*/
Map<ComponentNode, Future<Void>>
unscheduleReferencingComponents(ControllerServiceNode serviceNode);
@@ -199,12 +207,19 @@ public interface ControllerServiceProvider extends
ControllerServiceLookup {
* recursively, so if a Processor is referencing Service A, which is
* referencing serviceNode, then the Processor will also be started.
*
+ * A referencing processor that belongs to a stateless process group is not
+ * started individually; instead its owning stateless group is started as a
+ * single unit. Each affected stateless group is started at most once even
+ * when several of its processors reference the service.
+ *
* @param serviceNode the node
*/
Set<ComponentNode> scheduleReferencingComponents(ControllerServiceNode
serviceNode);
/**
* Schedules any of the candidate components that are currently
referencing the given Controller Service to run.
+ * A candidate processor that belongs to a stateless process group causes
its owning stateless group to be started
+ * as a single unit; the group is started once even if only one of its
processors is among the candidates.
* @return the components that were scheduled
*/
Set<ComponentNode> scheduleReferencingComponents(ControllerServiceNode
serviceNode, Set<ComponentNode> candidates, ComponentScheduler
componentScheduler);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/service/StandardControllerServiceProviderTest.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/service/StandardControllerServiceProviderTest.java
index 36be3f53163..82df40c22ae 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/service/StandardControllerServiceProviderTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/service/StandardControllerServiceProviderTest.java
@@ -21,12 +21,19 @@ import org.apache.nifi.bundle.BundleCoordinate;
import org.apache.nifi.components.state.StateManagerProvider;
import org.apache.nifi.components.validation.ValidationTrigger;
import org.apache.nifi.components.validation.VerifiableComponentFactory;
+import org.apache.nifi.controller.ComponentNode;
import org.apache.nifi.controller.ControllerService;
import org.apache.nifi.controller.ExtensionBuilder;
+import org.apache.nifi.controller.FlowAnalysisRuleNode;
import org.apache.nifi.controller.NodeTypeProvider;
import org.apache.nifi.controller.ProcessScheduler;
+import org.apache.nifi.controller.ProcessorNode;
import org.apache.nifi.controller.ReloadComponent;
+import org.apache.nifi.controller.ReportingTaskNode;
+import org.apache.nifi.controller.ScheduledState;
import org.apache.nifi.controller.flow.FlowManager;
+import org.apache.nifi.flow.ExecutionEngine;
+import org.apache.nifi.groups.ComponentScheduler;
import org.apache.nifi.groups.ProcessGroup;
import org.apache.nifi.nar.ExtensionDiscoveringManager;
import org.apache.nifi.nar.ExtensionManager;
@@ -50,6 +57,7 @@ import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -57,7 +65,9 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -236,6 +246,212 @@ class StandardControllerServiceProviderTest {
assertTrue(identifiers.contains(serviceId));
}
+ @Test
+ void testUnscheduleReferencingComponentsStopsStatelessGroupAsUnit() {
+ final ProcessGroup statelessGroup =
createGroup(ExecutionEngine.STATELESS, null);
+ final ProcessGroup standardGroup =
createGroup(ExecutionEngine.STANDARD, null);
+
+ final ProcessorNode statelessProcessorA =
createProcessor(statelessGroup, ScheduledState.RUNNING, ScheduledState.RUNNING);
+ final ProcessorNode statelessProcessorB =
createProcessor(statelessGroup, ScheduledState.RUNNING, ScheduledState.RUNNING);
+ final ProcessorNode standardProcessor = createProcessor(standardGroup,
ScheduledState.RUNNING, ScheduledState.RUNNING);
+
+
when(statelessGroup.stopProcessing()).thenReturn(CompletableFuture.completedFuture(null));
+
when(standardGroup.stopProcessor(standardProcessor)).thenReturn(CompletableFuture.completedFuture(null));
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(statelessProcessorA,
statelessProcessorB, standardProcessor));
+
+ final Map<ComponentNode, Future<Void>> result =
serviceProvider.unscheduleReferencingComponents(serviceNode);
+
+ // The stateless group is stopped once as a single unit; its member
processors are not stopped individually.
+ verify(statelessGroup, times(1)).stopProcessing();
+ verify(statelessGroup, never()).stopProcessor(any());
+ verify(standardGroup, never()).stopProcessing();
+
+ // The standard processor is stopped directly through its process
group.
+ verify(standardGroup, times(1)).stopProcessor(standardProcessor);
+
+ // Every affected processor is represented in the returned map, and
the stateless members share the group's single future.
+ assertTrue(result.containsKey(statelessProcessorA));
+ assertTrue(result.containsKey(statelessProcessorB));
+ assertTrue(result.containsKey(standardProcessor));
+ assertEquals(result.get(statelessProcessorA),
result.get(statelessProcessorB));
+ }
+
+ @Test
+ void
testUnscheduleReferencingComponentsResolvesInheritedChildToStatelessAncestor() {
+ final ProcessGroup statelessGroup =
createGroup(ExecutionEngine.STATELESS, null);
+ final ProcessGroup inheritedChild =
createGroup(ExecutionEngine.INHERITED, statelessGroup);
+ final ProcessGroup inheritedGrandchild =
createGroup(ExecutionEngine.INHERITED, inheritedChild);
+ final ProcessorNode childProcessor =
createProcessor(inheritedGrandchild, ScheduledState.RUNNING,
ScheduledState.RUNNING);
+
+
when(statelessGroup.stopProcessing()).thenReturn(CompletableFuture.completedFuture(null));
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(childProcessor));
+
+ final Map<ComponentNode, Future<Void>> result =
serviceProvider.unscheduleReferencingComponents(serviceNode);
+
+ // A processor in an INHERITED child resolves up to its explicit
stateless ancestor, which is stopped as a unit.
+ verify(statelessGroup, times(1)).stopProcessing();
+ verify(inheritedChild, never()).stopProcessing();
+ verify(inheritedGrandchild, never()).stopProcessing();
+ verify(inheritedGrandchild, never()).stopProcessor(any());
+ assertTrue(result.containsKey(childProcessor));
+ }
+
+ @Test
+ void testScheduleReferencingComponentsStartsStatelessGroupAsUnit() {
+ final ProcessGroup statelessGroup =
createGroup(ExecutionEngine.STATELESS, null);
+ final ProcessGroup standardGroup =
createGroup(ExecutionEngine.STANDARD, null);
+
+ final ProcessorNode statelessProcessorA =
createProcessor(statelessGroup, ScheduledState.STOPPED, ScheduledState.STOPPED);
+ final ProcessorNode statelessProcessorB =
createProcessor(statelessGroup, ScheduledState.STOPPED, ScheduledState.STOPPED);
+ final ProcessorNode standardProcessor = createProcessor(standardGroup,
ScheduledState.STOPPED, ScheduledState.STOPPED);
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(statelessProcessorA,
statelessProcessorB, standardProcessor));
+ final ComponentScheduler componentScheduler =
mock(ComponentScheduler.class);
+
+ final Set<ComponentNode> result =
serviceProvider.scheduleReferencingComponents(serviceNode, null,
componentScheduler);
+
+ // The stateless group is started once as a single unit; its members
are not started individually.
+ verify(componentScheduler,
times(1)).startStatelessGroup(statelessGroup);
+ verify(componentScheduler,
never()).startComponent(statelessProcessorA);
+ verify(componentScheduler,
never()).startComponent(statelessProcessorB);
+ verify(componentScheduler, never()).startStatelessGroup(standardGroup);
+
+ // The standard processor is started directly.
+ verify(componentScheduler, times(1)).startComponent(standardProcessor);
+
+ assertTrue(result.contains(statelessProcessorA));
+ assertTrue(result.contains(statelessProcessorB));
+ assertTrue(result.contains(standardProcessor));
+ }
+
+ @Test
+ void
testScheduleReferencingComponentsWithCandidateStartsOwningStatelessGroup() {
+ final ProcessGroup statelessGroup =
createGroup(ExecutionEngine.STATELESS, null);
+
+ final ProcessorNode statelessProcessorA =
createProcessor(statelessGroup, ScheduledState.STOPPED, ScheduledState.STOPPED);
+ final ProcessorNode statelessProcessorB =
createProcessor(statelessGroup, ScheduledState.STOPPED, ScheduledState.STOPPED);
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(statelessProcessorA,
statelessProcessorB));
+ final ComponentScheduler componentScheduler =
mock(ComponentScheduler.class);
+
+ // Only one member of the stateless group is a candidate, but the
group must still be started as a whole.
+ final Set<ComponentNode> result =
serviceProvider.scheduleReferencingComponents(serviceNode,
Set.of(statelessProcessorA), componentScheduler);
+
+ verify(componentScheduler,
times(1)).startStatelessGroup(statelessGroup);
+ verify(componentScheduler,
never()).startComponent(statelessProcessorA);
+ assertTrue(result.contains(statelessProcessorA));
+ }
+
+ @Test
+ void testUnscheduleReferencingComponentsStopsEachStatelessGroupOnce() {
+ final ProcessGroup firstGroup = createGroup(ExecutionEngine.STATELESS,
null);
+ final ProcessGroup secondGroup =
createGroup(ExecutionEngine.STATELESS, null);
+
+ final ProcessorNode firstMember = createProcessor(firstGroup,
ScheduledState.RUNNING, ScheduledState.RUNNING);
+ final ProcessorNode secondMember = createProcessor(secondGroup,
ScheduledState.RUNNING, ScheduledState.RUNNING);
+
+
when(firstGroup.stopProcessing()).thenReturn(CompletableFuture.completedFuture(null));
+
when(secondGroup.stopProcessing()).thenReturn(CompletableFuture.completedFuture(null));
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(firstMember, secondMember));
+
+ serviceProvider.unscheduleReferencingComponents(serviceNode);
+
+ // Two distinct stateless groups are each stopped exactly once.
+ verify(firstGroup, times(1)).stopProcessing();
+ verify(secondGroup, times(1)).stopProcessing();
+ }
+
+ @Test
+ void
testUnscheduleReferencingComponentsStopsTopMostStatelessGroupWhenNested() {
+ final ProcessGroup root = createGroup(ExecutionEngine.STANDARD, null);
+ final ProcessGroup inheritedChild =
createGroup(ExecutionEngine.INHERITED, root);
+ final ProcessGroup outer = createGroup(ExecutionEngine.STATELESS,
inheritedChild);
+ final ProcessGroup middle = createGroup(ExecutionEngine.STATELESS,
outer);
+ final ProcessGroup inner = createGroup(ExecutionEngine.STATELESS,
middle);
+
+ final ProcessorNode innerProcessor = createProcessor(inner,
ScheduledState.RUNNING, ScheduledState.RUNNING);
+
+
when(outer.stopProcessing()).thenReturn(CompletableFuture.completedFuture(null));
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(innerProcessor));
+
+ final Map<ComponentNode, Future<Void>> result =
serviceProvider.unscheduleReferencingComponents(serviceNode);
+
+ // stopProcessing() on a nested stateless group is a no-op that still
returns a completed future, so stopping the
+ // wrong group would leave it running while the caller believes the
stop succeeded.
+ verify(outer, times(1)).stopProcessing();
+ verify(middle, never()).stopProcessing();
+ verify(inner, never()).stopProcessing();
+ verify(inner, never()).stopProcessor(any());
+ assertTrue(result.containsKey(innerProcessor));
+ }
+
+ @Test
+ void
testScheduleReferencingComponentsStartsTopMostStatelessGroupWhenNested() {
+ final ProcessGroup root = createGroup(ExecutionEngine.STANDARD, null);
+ final ProcessGroup inheritedChild =
createGroup(ExecutionEngine.INHERITED, root);
+ final ProcessGroup outer = createGroup(ExecutionEngine.STATELESS,
inheritedChild);
+ final ProcessGroup middle = createGroup(ExecutionEngine.STATELESS,
outer);
+ final ProcessGroup inner = createGroup(ExecutionEngine.STATELESS,
middle);
+
+ final ProcessorNode innerProcessor = createProcessor(inner,
ScheduledState.STOPPED, ScheduledState.STOPPED);
+
+ final ControllerServiceNode serviceNode =
createServiceWithProcessorReferences(List.of(innerProcessor));
+ final ComponentScheduler componentScheduler =
mock(ComponentScheduler.class);
+
+ final Set<ComponentNode> result =
serviceProvider.scheduleReferencingComponents(serviceNode, null,
componentScheduler);
+
+ verify(componentScheduler, times(1)).startStatelessGroup(outer);
+ verify(componentScheduler, never()).startStatelessGroup(middle);
+ verify(componentScheduler, never()).startStatelessGroup(inner);
+ verify(componentScheduler, never()).startComponent(innerProcessor);
+ assertTrue(result.contains(innerProcessor));
+ }
+
+ private ProcessGroup createGroup(final ExecutionEngine executionEngine,
final ProcessGroup parent) {
+ // Resolve before stubbing: reading the parent mock inside an
in-progress when(...) trips Mockito's
+ // unfinished-stubbing detection.
+ final ExecutionEngine resolvedExecutionEngine =
resolveExecutionEngine(executionEngine, parent);
+
+ final ProcessGroup group = mock(ProcessGroup.class);
+ lenient().when(group.getParent()).thenReturn(parent);
+
lenient().when(group.resolveExecutionEngine()).thenReturn(resolvedExecutionEngine);
+ return group;
+ }
+
+ /**
+ * Mirrors {@code StandardProcessGroup.resolveExecutionEngine()}. Callers
must create a parent before its children
+ * so that the parent's stubbed resolution is already in place.
+ */
+ private ExecutionEngine resolveExecutionEngine(final ExecutionEngine
executionEngine, final ProcessGroup parent) {
+ if (executionEngine != ExecutionEngine.INHERITED) {
+ return executionEngine;
+ }
+
+ return parent == null ? ExecutionEngine.STANDARD :
parent.resolveExecutionEngine();
+ }
+
+ private ProcessorNode createProcessor(final ProcessGroup group, final
ScheduledState scheduledState, final ScheduledState physicalState) {
+ final ProcessorNode processor = mock(ProcessorNode.class);
+ lenient().when(processor.getProcessGroup()).thenReturn(group);
+
lenient().when(processor.getScheduledState()).thenReturn(scheduledState);
+
lenient().when(processor.getPhysicalScheduledState()).thenReturn(physicalState);
+ return processor;
+ }
+
+ private ControllerServiceNode createServiceWithProcessorReferences(final
List<ProcessorNode> processors) {
+ final ControllerServiceNode serviceNode =
mock(ControllerServiceNode.class);
+ final ControllerServiceReference reference =
mock(ControllerServiceReference.class);
+ when(serviceNode.getReferences()).thenReturn(reference);
+
when(reference.findRecursiveReferences(ProcessorNode.class)).thenReturn(processors);
+
when(reference.findRecursiveReferences(ReportingTaskNode.class)).thenReturn(Collections.emptyList());
+
when(reference.findRecursiveReferences(FlowAnalysisRuleNode.class)).thenReturn(Collections.emptyList());
+ return serviceNode;
+ }
+
private ControllerServiceNode createControllerService(
final String identifier,
final String type,
diff --git
a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/connectors/tests/system/ComponentLifecycleConnector.java
b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/connectors/tests/system/ComponentLifecycleConnector.java
index 9eddbf597e0..5021f99dd2a 100644
---
a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/connectors/tests/system/ComponentLifecycleConnector.java
+++
b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/connectors/tests/system/ComponentLifecycleConnector.java
@@ -47,6 +47,7 @@ import java.util.Set;
* - A child process group with input and output ports
* - A processor within the child group
* - A stateless group with a processor
+ * - A pair of nested stateless groups, where only the inner one holds a
processor
*
* This allows testing that start/stop operations properly handle all
component types recursively.
*/
@@ -95,7 +96,7 @@ public class ComponentLifecycleConnector extends
AbstractConnector {
final VersionedProcessor rootTerminateProcessor =
VersionedFlowUtils.addProcessor(rootGroup,
"org.apache.nifi.processors.tests.system.TerminateFlowFile",
SYSTEM_TEST_EXTENSIONS_BUNDLE, "Root TerminateFlowFile", new Position(300,
100));
- final VersionedProcessGroup childGroup =
createChildGroup(rootGroup.getIdentifier());
+ final VersionedProcessGroup childGroup =
createChildGroup(rootGroup.getIdentifier(),
rootControllerService.getIdentifier());
rootGroup.getProcessGroups().add(childGroup);
final VersionedPort childInputPort =
childGroup.getInputPorts().iterator().next();
@@ -109,7 +110,7 @@ public class ComponentLifecycleConnector extends
AbstractConnector {
return rootGroup;
}
- private VersionedProcessGroup createChildGroup(final String parentGroupId)
{
+ private VersionedProcessGroup createChildGroup(final String parentGroupId,
final String rootCountServiceId) {
final VersionedProcessGroup childGroup =
VersionedFlowUtils.createProcessGroup("child-group-id", "Child Group");
childGroup.setPosition(new Position(100, 300));
childGroup.setRemoteProcessGroups(new HashSet<>());
@@ -131,39 +132,85 @@ public class ComponentLifecycleConnector extends
AbstractConnector {
final VersionedProcessor childProcessor =
VersionedFlowUtils.addProcessor(childGroup,
"org.apache.nifi.processors.tests.system.PassThrough",
SYSTEM_TEST_EXTENSIONS_BUNDLE, "Child Terminate", new Position(100, 100));
- final VersionedProcessGroup statelessGroup =
createStatelessGroup(childGroup.getIdentifier());
+ final VersionedProcessGroup statelessGroup =
createStatelessGroup(childGroup.getIdentifier(), rootCountServiceId);
childGroup.getProcessGroups().add(statelessGroup);
final VersionedPort statelessInputPort =
statelessGroup.getInputPorts().iterator().next();
+ final VersionedProcessGroup nestedOuterGroup =
createNestedStatelessGroups(childGroup.getIdentifier(), rootCountServiceId);
+ childGroup.getProcessGroups().add(nestedOuterGroup);
+
+ final VersionedPort nestedOuterInputPort =
nestedOuterGroup.getInputPorts().iterator().next();
+
VersionedFlowUtils.addConnection(childGroup,
VersionedFlowUtils.createConnectableComponent(inputPort),
VersionedFlowUtils.createConnectableComponent(childProcessor),
Set.of(""));
VersionedFlowUtils.addConnection(childGroup,
VersionedFlowUtils.createConnectableComponent(inputPort),
VersionedFlowUtils.createConnectableComponent(statelessInputPort),
Set.of(""));
+ VersionedFlowUtils.addConnection(childGroup,
VersionedFlowUtils.createConnectableComponent(inputPort),
+
VersionedFlowUtils.createConnectableComponent(nestedOuterInputPort),
Set.of(""));
VersionedFlowUtils.addConnection(childGroup,
VersionedFlowUtils.createConnectableComponent(childProcessor),
VersionedFlowUtils.createConnectableComponent(outputPort),
Set.of("success"));
return childGroup;
}
- private VersionedProcessGroup createStatelessGroup(final String
parentGroupId) {
- final VersionedProcessGroup statelessGroup =
VersionedFlowUtils.createProcessGroup("stateless-group-id", "Stateless Group");
- statelessGroup.setPosition(new Position(400, 100));
- statelessGroup.setRemoteProcessGroups(new HashSet<>());
- statelessGroup.setScheduledState(ScheduledState.ENABLED);
- statelessGroup.setExecutionEngine(ExecutionEngine.STATELESS);
- statelessGroup.setStatelessFlowTimeout("1 min");
- statelessGroup.setGroupIdentifier(parentGroupId);
+ private VersionedProcessGroup createStatelessGroup(final String
parentGroupId, final String rootCountServiceId) {
+ final VersionedProcessGroup statelessGroup =
createStatelessGroupShell("stateless-group-id", "Stateless", new Position(400,
100), parentGroupId);
+ addCountingFlow(statelessGroup, "Stateless", rootCountServiceId);
+ return statelessGroup;
+ }
- final VersionedPort statelessInput =
VersionedFlowUtils.addInputPort(statelessGroup, "Stateless Input", new
Position(0, 0));
+ private VersionedProcessGroup createNestedStatelessGroups(final String
parentGroupId, final String rootCountServiceId) {
+ final VersionedProcessGroup outerGroup =
createStatelessGroupShell("nested-stateless-outer-group-id", "Nested Outer",
+ new Position(400, 300), parentGroupId);
- final VersionedProcessor statelessProcessor =
VersionedFlowUtils.addProcessor(statelessGroup,
- "org.apache.nifi.processors.tests.system.TerminateFlowFile",
SYSTEM_TEST_EXTENSIONS_BUNDLE, "Stateless Terminate", new Position(100, 100));
+ // Only the inner group holds a processor referencing the root
service, so a resolver that stops at the inner
+ // group silently no-ops instead of transitioning the subtree.
+ final VersionedProcessGroup innerGroup =
createStatelessGroupShell("nested-stateless-inner-group-id", "Nested Inner",
+ new Position(200, 0), outerGroup.getIdentifier());
+ addCountingFlow(innerGroup, "Nested Inner", rootCountServiceId);
+ outerGroup.getProcessGroups().add(innerGroup);
- VersionedFlowUtils.addConnection(statelessGroup,
VersionedFlowUtils.createConnectableComponent(statelessInput),
- VersionedFlowUtils.createConnectableComponent(statelessProcessor),
Set.of(""));
+ VersionedFlowUtils.addConnection(outerGroup,
VersionedFlowUtils.createConnectableComponent(getInputPort(outerGroup)),
+
VersionedFlowUtils.createConnectableComponent(getInputPort(innerGroup)),
Set.of(""));
- return statelessGroup;
+ return outerGroup;
+ }
+
+ private VersionedProcessGroup createStatelessGroupShell(final String
identifier, final String namePrefix, final Position position,
+ final String
parentGroupId) {
+ final VersionedProcessGroup group =
VersionedFlowUtils.createProcessGroup(identifier, namePrefix + " Group");
+ group.setPosition(position);
+ group.setRemoteProcessGroups(new HashSet<>());
+ group.setScheduledState(ScheduledState.ENABLED);
+ group.setExecutionEngine(ExecutionEngine.STATELESS);
+ group.setStatelessFlowTimeout("1 min");
+ group.setGroupIdentifier(parentGroupId);
+
+ VersionedFlowUtils.addInputPort(group, namePrefix + " Input", new
Position(0, 0));
+ return group;
+ }
+
+ /**
+ * Wires {@code input port -> CountFlowFiles -> TerminateFlowFile} inside
the given group, with the CountFlowFiles
+ * processor referencing a Controller Service that lives outside the group.
+ */
+ private void addCountingFlow(final VersionedProcessGroup group, final
String namePrefix, final String rootCountServiceId) {
+ final VersionedProcessor countProcessor =
VersionedFlowUtils.addProcessor(group,
+ "org.apache.nifi.processors.tests.system.CountFlowFiles",
SYSTEM_TEST_EXTENSIONS_BUNDLE, namePrefix + " Count", new Position(100, 50));
+ countProcessor.getProperties().put("Count Service",
rootCountServiceId);
+
+ final VersionedProcessor terminateProcessor =
VersionedFlowUtils.addProcessor(group,
+ "org.apache.nifi.processors.tests.system.TerminateFlowFile",
SYSTEM_TEST_EXTENSIONS_BUNDLE, namePrefix + " Terminate", new Position(100,
100));
+
+ VersionedFlowUtils.addConnection(group,
VersionedFlowUtils.createConnectableComponent(getInputPort(group)),
+ VersionedFlowUtils.createConnectableComponent(countProcessor),
Set.of(""));
+ VersionedFlowUtils.addConnection(group,
VersionedFlowUtils.createConnectableComponent(countProcessor),
+ VersionedFlowUtils.createConnectableComponent(terminateProcessor),
Set.of("success"));
+ }
+
+ private VersionedPort getInputPort(final VersionedProcessGroup group) {
+ return group.getInputPorts().iterator().next();
}
@Override
diff --git
a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/connectors/ConnectorTroubleshootingIT.java
b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/connectors/ConnectorTroubleshootingIT.java
index 1bbe86c3f21..0bc80db4391 100644
---
a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/connectors/ConnectorTroubleshootingIT.java
+++
b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/connectors/ConnectorTroubleshootingIT.java
@@ -36,6 +36,8 @@ import org.apache.nifi.web.api.entity.AssetsEntity;
import org.apache.nifi.web.api.entity.ConnectionEntity;
import org.apache.nifi.web.api.entity.ConnectorEntity;
import org.apache.nifi.web.api.entity.ControllerServiceEntity;
+import
org.apache.nifi.web.api.entity.ControllerServiceReferencingComponentEntity;
+import
org.apache.nifi.web.api.entity.ControllerServiceReferencingComponentsEntity;
import org.apache.nifi.web.api.entity.HistoryEntity;
import org.apache.nifi.web.api.entity.ParameterProviderEntity;
import org.apache.nifi.web.api.entity.PortEntity;
@@ -43,6 +45,7 @@ import org.apache.nifi.web.api.entity.ProcessGroupEntity;
import org.apache.nifi.web.api.entity.ProcessGroupFlowEntity;
import org.apache.nifi.web.api.entity.ProcessorEntity;
import org.apache.nifi.web.api.entity.ScheduleComponentsEntity;
+import
org.apache.nifi.web.api.entity.UpdateControllerServiceReferenceRequestEntity;
import org.junit.jupiter.api.Test;
import java.io.File;
@@ -66,6 +69,74 @@ import static org.junit.jupiter.api.Assertions.fail;
*/
public class ConnectorTroubleshootingIT extends NiFiSystemIT {
+ /**
+ * Regression test for stateless-group handling in the controller-service
reference lifecycle. The managed flow has
+ * two stateless subtrees that both reference a controller service defined
at the connector root: a flat one, and one
+ * where the referencing processor sits in a nested group that also
declares the Stateless Engine. Stopping and
+ * starting the service's referencing components must transition each
stateless subtree as a single unit, resolving
+ * up to the top-most stateless group rather than leaving a mixed
running/stopped state.
+ */
+ @Test
+ public void
testControllerServiceReferenceLifecycleTransitionsStatelessGroupAsUnit() throws
NiFiClientException, IOException, InterruptedException {
+ final ConnectorEntity connector =
getClientUtil().createConnector("ComponentLifecycleConnector");
+ final String connectorId = connector.getId();
+
+ getClientUtil().applyConnectorUpdate(connector);
+ getClientUtil().waitForValidConnector(connectorId);
+
+ getClientUtil().enterTroubleshooting(connectorId);
+ assertConnectorState(connectorId, ConnectorState.TROUBLESHOOTING);
+
+ final List<ProcessorEntity> statelessProcessors =
findStatelessProcessors(connectorId);
+ assertEquals(4, statelessProcessors.size(), "Stateless groups should
contain exactly four processors");
+ final List<String> statelessProcessorIds =
statelessProcessors.stream().map(ProcessorEntity::getId).toList();
+
+ final String rootServiceId = findFirstControllerServiceId(connectorId);
+ assertNotNull(rootServiceId, "Managed flow should contain the root
controller service");
+
+ // Enable the root controller service so its referencing components
can be scheduled.
+ enableControllerService(rootServiceId);
+
+ final ControllerServiceReferencingComponentsEntity references =
+
getNifiClient().getControllerServicesClient().getControllerServiceReferences(rootServiceId);
+ final List<String> referencingProcessorIds =
references.getControllerServiceReferencingComponents().stream()
+ .map(ControllerServiceReferencingComponentEntity::getId)
+ .toList();
+ assertFalse(referencingProcessorIds.isEmpty(), "Root controller
service should have at least one referencing processor");
+ assertTrue(statelessProcessorIds.containsAll(referencingProcessorIds),
+ "Every referencing processor should be inside a stateless
group");
+ assertFalse(referencingProcessorIds.containsAll(statelessProcessorIds),
+ "At least one stateless processor should not directly
reference the service, to prove group-as-unit behavior");
+
+ final ProcessorEntity nestedCountProcessor =
findProcessorByName(connectorId, "Nested Inner Count");
+ assertNotNull(nestedCountProcessor, "Managed flow should contain the
nested stateless Count processor");
+
assertTrue(referencingProcessorIds.contains(nestedCountProcessor.getId()),
+ "The nested-group Count processor should reference the root
service, so top-most-group resolution is covered");
+
+ // Start the referencing components; every stateless subtree must
start whole.
+ updateReferenceState(rootServiceId, references,
ScheduledState.RUNNING.name());
+ for (final String processorId : statelessProcessorIds) {
+ waitForProcessorState(processorId, ScheduledState.RUNNING);
+ }
+
+ final ControllerServiceReferencingComponentsEntity runningReferences =
+
getNifiClient().getControllerServicesClient().getControllerServiceReferences(rootServiceId);
+ updateReferenceState(rootServiceId, runningReferences,
ScheduledState.STOPPED.name());
+ for (final String processorId : statelessProcessorIds) {
+ waitForProcessorState(processorId, ScheduledState.STOPPED);
+ }
+
+ // Disable services and exit Troubleshooting, then confirm the
Connector starts cleanly with no mixed state.
+ final String managedGroupId =
getNifiClient().getConnectorClient().getConnector(connectorId).getComponent().getManagedProcessGroupId();
+ getClientUtil().disableControllerServices(managedGroupId, true);
+
+ getClientUtil().endTroubleshooting(connectorId);
+ assertConnectorState(connectorId, ConnectorState.STOPPED);
+
+ getClientUtil().startConnector(connectorId);
+ assertConnectorState(connectorId, ConnectorState.RUNNING);
+ }
+
/**
* Transition a Connector into Troubleshooting, modify a processor inside
the managed flow, then transition back
* out. The Connector's authoritative flow should be restored on exit and
the Connector should start smoothly.
@@ -559,6 +630,44 @@ public class ConnectorTroubleshootingIT extends
NiFiSystemIT {
assertEquals(expected.name(), entity.getComponent().getState());
}
+ private List<ProcessorEntity> findStatelessProcessors(final String
connectorId) throws NiFiClientException, IOException {
+ final List<ProcessorEntity> result = new ArrayList<>();
+ final Map<String, Boolean> statelessByGroupId = new HashMap<>();
+ for (final ProcessorEntity processor : findAllProcessors(connectorId))
{
+ final String parentGroupId =
processor.getComponent().getParentGroupId();
+ Boolean stateless = statelessByGroupId.get(parentGroupId);
+ if (stateless == null) {
+ final ProcessGroupEntity parentGroup =
getNifiClient().getProcessGroupClient().getProcessGroup(parentGroupId);
+ stateless =
"STATELESS".equals(parentGroup.getComponent().getExecutionEngine());
+ statelessByGroupId.put(parentGroupId, stateless);
+ }
+
+ if (stateless) {
+ result.add(processor);
+ }
+ }
+ return result;
+ }
+
+ private void enableControllerService(final String serviceId) throws
NiFiClientException, IOException, InterruptedException {
+ final ControllerServiceEntity service =
getNifiClient().getControllerServicesClient().getControllerService(serviceId);
+ getClientUtil().enableControllerService(service);
+ getClientUtil().waitForControllerServiceRunStatus(serviceId,
"ENABLED");
+ }
+
+ private void updateReferenceState(final String serviceId, final
ControllerServiceReferencingComponentsEntity references, final String state)
throws NiFiClientException, IOException {
+ final Map<String, RevisionDTO> revisions = new HashMap<>();
+ for (final ControllerServiceReferencingComponentEntity component :
references.getControllerServiceReferencingComponents()) {
+ revisions.put(component.getId(), component.getRevision());
+ }
+
+ final UpdateControllerServiceReferenceRequestEntity request = new
UpdateControllerServiceReferenceRequestEntity();
+ request.setId(serviceId);
+ request.setReferencingComponentRevisions(revisions);
+ request.setState(state);
+
getNifiClient().getControllerServicesClient().updateControllerServiceReferences(request);
+ }
+
/**
* Stop and restart the NiFi instance, then wait for all nodes to
reconnect when running in a clustered
* environment. Subsequent flow-modifying requests (such as {@code
endTroubleshooting}) would otherwise be rejected