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 8c0894c79c1 NIFI-16228 - Versioned Flow Update Fails to Apply
Auto-Termination for Previously Connected Relationship (#11564)
8c0894c79c1 is described below
commit 8c0894c79c1ad26fc43c568a3f7c099c382d893a
Author: Pierre Villard <[email protected]>
AuthorDate: Wed Aug 19 21:19:22 2026 +0200
NIFI-16228 - Versioned Flow Update Fails to Apply Auto-Termination for
Previously Connected Relationship (#11564)
---
.../apache/nifi/util/FlowDifferenceFilters.java | 4 +
...StandardVersionedComponentSynchronizerTest.java | 105 +++++++++++++++++++++
.../nifi/util/TestFlowDifferenceFilters.java | 58 ++++++++++++
3 files changed, 167 insertions(+)
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/util/FlowDifferenceFilters.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/util/FlowDifferenceFilters.java
index 1ec0dc97ede..c4d0e955062 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/util/FlowDifferenceFilters.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/util/FlowDifferenceFilters.java
@@ -472,6 +472,10 @@ public class FlowDifferenceFilters {
return false;
}
+ if (!replaceNull(processorNode.getConnections(relationship),
Collections.emptySet()).isEmpty()) {
+ return false;
+ }
+
if (hasConnection(processGroup, processorA, relationshipName)) {
return false;
}
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java
index ec977b5453f..9c0953c1f45 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java
@@ -35,6 +35,7 @@ import org.apache.nifi.controller.ReloadComponent;
import org.apache.nifi.controller.ReportingTaskNode;
import org.apache.nifi.controller.flow.FlowManager;
import org.apache.nifi.controller.queue.FlowFileQueue;
+import org.apache.nifi.controller.queue.LoadBalanceCompression;
import org.apache.nifi.controller.queue.LoadBalanceStrategy;
import org.apache.nifi.controller.service.ControllerServiceNode;
import org.apache.nifi.controller.service.ControllerServiceProvider;
@@ -77,7 +78,13 @@ import org.apache.nifi.parameter.StandardParameterContext;
import org.apache.nifi.parameter.StandardParameterContextManager;
import org.apache.nifi.processor.ProcessContext;
import org.apache.nifi.processor.Relationship;
+import org.apache.nifi.registry.flow.diff.DifferenceType;
+import org.apache.nifi.registry.flow.diff.FlowComparatorVersionedStrategy;
+import org.apache.nifi.registry.flow.diff.StandardComparableDataFlow;
+import org.apache.nifi.registry.flow.diff.StandardFlowComparator;
+import org.apache.nifi.registry.flow.diff.StaticDifferenceDescriptor;
import org.apache.nifi.registry.flow.mapping.FlowMappingOptions;
+import org.apache.nifi.registry.flow.mapping.VersionedComponentFlowMapper;
import org.apache.nifi.remote.PublicPort;
import org.apache.nifi.reporting.InitializationException;
import org.apache.nifi.scheduling.ExecutionNode;
@@ -106,6 +113,7 @@ import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.stream.Collectors;
@@ -695,6 +703,103 @@ public class StandardVersionedComponentSynchronizerTest {
verify(processGroup, never()).startProcessor(runningProcessor, false);
}
+ @Test
+ public void
testGroupSynchronizeAppliesAutoTerminatedRelationshipMovedOffExistingConnection()
{
+ final ProcessGroup processGroup = createMockProcessGroup();
+ final String processGroupId = processGroup.getIdentifier();
+
when(processGroup.getVersionedComponentId()).thenReturn(Optional.of(processGroupId));
+ when(processGroup.getInputPorts()).thenReturn(Collections.emptySet());
+ when(processGroup.getOutputPorts()).thenReturn(Collections.emptySet());
+ when(processGroup.getFunnels()).thenReturn(Collections.emptySet());
+ when(processGroup.getLabels()).thenReturn(Collections.emptySet());
+
when(processGroup.getRemoteProcessGroups()).thenReturn(Collections.emptySet());
+
when(processGroup.getProcessGroups()).thenReturn(Collections.emptySet());
+
when(processGroup.getControllerServices(false)).thenReturn(Collections.emptySet());
+
+ final ProcessorNode sourceProcessor =
createMappableProcessor(processGroup);
+ final ProcessorNode destinationProcessor =
createMappableProcessor(processGroup);
+
when(sourceProcessor.getProcessGroupIdentifier()).thenReturn(processGroupId);
+
when(destinationProcessor.getProcessGroupIdentifier()).thenReturn(processGroupId);
+
when(sourceProcessor.getVersionedComponentId()).thenReturn(Optional.of("source-processor"));
+
when(destinationProcessor.getVersionedComponentId()).thenReturn(Optional.of("destination-processor"));
+
+ final Relationship discardRelationship = new Relationship.Builder()
+ .name("discard")
+ .autoTerminateDefault(true)
+ .build();
+ final Relationship successRelationship = new
Relationship.Builder().name("success").build();
+
+
when(sourceProcessor.getRelationships()).thenReturn(Set.of(discardRelationship,
successRelationship));
+
when(sourceProcessor.getRelationship("discard")).thenReturn(discardRelationship);
+
when(sourceProcessor.getRelationship("success")).thenReturn(successRelationship);
+ final AtomicReference<Set<Relationship>> autoTerminatedRelationships =
new AtomicReference<>(Collections.emptySet());
+
when(sourceProcessor.getAutoTerminatedRelationships()).thenAnswer(invocation ->
autoTerminatedRelationships.get());
+ doAnswer(invocation -> {
+ autoTerminatedRelationships.set(invocation.getArgument(0));
+ return null;
+ }).when(sourceProcessor).setAutoTerminatedRelationships(anySet());
+
when(destinationProcessor.getAutoTerminatedRelationships()).thenReturn(Collections.emptySet());
+
+ final Connection connection = createMockConnection(sourceProcessor,
destinationProcessor, processGroup);
+
when(connection.getVersionedComponentId()).thenReturn(Optional.of("connection-ab"));
+ when(connection.getName()).thenReturn("connection-ab");
+ when(connection.getLabelIndex()).thenReturn(0);
+ final AtomicReference<Collection<Relationship>> selectedRelationships
= new AtomicReference<>(Set.of(discardRelationship));
+ when(connection.getRelationships()).thenAnswer(invocation ->
selectedRelationships.get());
+ doAnswer(invocation -> {
+ selectedRelationships.set(invocation.getArgument(0));
+ return null;
+ }).when(connection).setRelationships(anyCollection());
+ when(connection.getZIndex()).thenReturn(0L);
+ when(connection.getBendPoints()).thenReturn(Collections.emptyList());
+
when(sourceProcessor.getConnections(discardRelationship)).thenReturn(Set.of(connection));
+
when(processGroup.getConnection(connection.getIdentifier())).thenReturn(connection);
+
+ final FlowFileQueue flowFileQueue = connection.getFlowFileQueue();
+ when(flowFileQueue.getBackPressureDataSizeThreshold()).thenReturn("1
GB");
+
when(flowFileQueue.getBackPressureObjectThreshold()).thenReturn(10000L);
+ when(flowFileQueue.getFlowFileExpiration()).thenReturn("0 sec");
+
when(flowFileQueue.getPriorities()).thenReturn(Collections.emptyList());
+
when(flowFileQueue.getLoadBalanceStrategy()).thenReturn(LoadBalanceStrategy.DO_NOT_LOAD_BALANCE);
+ when(flowFileQueue.getPartitioningAttribute()).thenReturn(null);
+
when(flowFileQueue.getLoadBalanceCompression()).thenReturn(LoadBalanceCompression.DO_NOT_COMPRESS);
+
+ when(processGroup.getProcessors()).thenReturn(List.of(sourceProcessor,
destinationProcessor));
+
when(flowManager.getProcessorNode(sourceProcessor.getIdentifier())).thenReturn(sourceProcessor);
+
when(flowManager.getProcessorNode(destinationProcessor.getIdentifier())).thenReturn(destinationProcessor);
+ when(flowManager.getGroup(processGroupId)).thenReturn(processGroup);
+
+ final VersionedComponentFlowMapper mapper = new
VersionedComponentFlowMapper(mock(ExtensionManager.class),
FlowMappingOptions.DEFAULT_OPTIONS);
+ final VersionedProcessGroup proposedGroup =
mapper.mapProcessGroup(processGroup, controllerServiceProvider, flowManager,
true);
+ final VersionedProcessor proposedSourceProcessor =
proposedGroup.getProcessors().stream()
+ .filter(processor ->
processor.getIdentifier().equals("source-processor"))
+ .findFirst()
+ .orElseThrow();
+ final VersionedConnection proposedConnection =
proposedGroup.getConnections().stream()
+ .filter(versionedConnection ->
versionedConnection.getIdentifier().equals("connection-ab"))
+ .findFirst()
+ .orElseThrow();
+
+
proposedSourceProcessor.setAutoTerminatedRelationships(Set.of("discard"));
+ proposedConnection.setSelectedRelationships(Collections.emptySet());
+
+ final VersionedExternalFlow externalFlow = new VersionedExternalFlow();
+ externalFlow.setFlowContents(proposedGroup);
+
+ assertDoesNotThrow(() -> synchronizer.synchronize(processGroup,
externalFlow, synchronizationOptions));
+
+ verify(connection).setRelationships(Collections.emptySet());
+
verify(sourceProcessor).setAutoTerminatedRelationships(Set.of(discardRelationship));
+
+ final VersionedProcessGroup synchronizedGroup =
mapper.mapProcessGroup(processGroup, controllerServiceProvider, flowManager,
true);
+ final StandardFlowComparator comparator = new StandardFlowComparator(
+ new StandardComparableDataFlow("Target Flow", proposedGroup),
+ new StandardComparableDataFlow("Synchronized Flow",
synchronizedGroup),
+ new StaticDifferenceDescriptor(), Function.identity(),
VersionedComponent::getIdentifier, FlowComparatorVersionedStrategy.DEEP);
+ assertFalse(comparator.compare().getDifferences().stream()
+ .anyMatch(difference -> difference.getDifferenceType() ==
DifferenceType.AUTO_TERMINATED_RELATIONSHIPS_CHANGED));
+ }
+
@Test
public void testUpdateConnectionWithSourceDestStopped() throws
FlowSynchronizationException, TimeoutException {
final VersionedConnection versionedConnection =
createMinimalVersionedConnection(processorA, processorB);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/util/TestFlowDifferenceFilters.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/util/TestFlowDifferenceFilters.java
index 5f942172bfd..3d587e77b0e 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/util/TestFlowDifferenceFilters.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/util/TestFlowDifferenceFilters.java
@@ -31,6 +31,7 @@ import org.apache.nifi.flow.ScheduledState;
import org.apache.nifi.flow.VersionedConnection;
import org.apache.nifi.flow.VersionedControllerService;
import org.apache.nifi.flow.VersionedPort;
+import org.apache.nifi.flow.VersionedProcessGroup;
import org.apache.nifi.flow.VersionedProcessor;
import org.apache.nifi.flow.VersionedPropertyDescriptor;
import org.apache.nifi.flow.VersionedRemoteGroupPort;
@@ -1093,6 +1094,63 @@ public class TestFlowDifferenceFilters {
assertTrue(FlowDifferenceFilters.isComponentUpdateRequired(scheduledStateDiff,
null, flowManager));
}
+ @Test
+ public void
testAutoTerminatedRelationshipChangeRequiresUpdateWhenCurrentlyConnected() {
+ final FlowManager flowManager = Mockito.mock(FlowManager.class);
+ final ProcessorNode processorNode = Mockito.mock(ProcessorNode.class);
+ final Connection connection = Mockito.mock(Connection.class);
+ final String processorInstanceId = "processor-instance";
+ final String relationshipName = "discard";
+
+ final VersionedProcessor currentProcessor = new VersionedProcessor();
+ currentProcessor.setAutoTerminatedRelationships(Set.of());
+ final InstantiatedVersionedProcessor proposedProcessor = new
InstantiatedVersionedProcessor(processorInstanceId, "group-id");
+
proposedProcessor.setAutoTerminatedRelationships(Set.of(relationshipName));
+ final Relationship relationship = new Relationship.Builder()
+ .name(relationshipName)
+ .autoTerminateDefault(true)
+ .build();
+
+
Mockito.when(flowManager.getProcessorNode(processorInstanceId)).thenReturn(processorNode);
+
Mockito.when(processorNode.getRelationship(relationshipName)).thenReturn(relationship);
+
Mockito.when(processorNode.getConnections(relationship)).thenReturn(Set.of(connection));
+
+ final FlowDifference difference = new StandardFlowDifference(
+ DifferenceType.AUTO_TERMINATED_RELATIONSHIPS_CHANGED,
currentProcessor, proposedProcessor,
+ Set.of(), Set.of(relationshipName), "Auto-terminated
relationships changed");
+
+ assertFalse(FlowDifferenceFilters.isEnvironmentalChange(difference,
new VersionedProcessGroup(), flowManager));
+ assertTrue(FlowDifferenceFilters.isComponentUpdateRequired(difference,
new VersionedProcessGroup(), flowManager));
+ }
+
+ @Test
+ public void
testAutoTerminatedRelationshipChangeDoesNotRequireUpdateWhenNotCurrentlyConnected()
{
+ final FlowManager flowManager = Mockito.mock(FlowManager.class);
+ final ProcessorNode processorNode = Mockito.mock(ProcessorNode.class);
+ final String processorInstanceId = "processor-instance";
+ final String relationshipName = "discard";
+
+ final VersionedProcessor currentProcessor = new VersionedProcessor();
+ currentProcessor.setAutoTerminatedRelationships(Set.of());
+ final InstantiatedVersionedProcessor proposedProcessor = new
InstantiatedVersionedProcessor(processorInstanceId, "group-id");
+
proposedProcessor.setAutoTerminatedRelationships(Set.of(relationshipName));
+ final Relationship relationship = new Relationship.Builder()
+ .name(relationshipName)
+ .autoTerminateDefault(true)
+ .build();
+
+
Mockito.when(flowManager.getProcessorNode(processorInstanceId)).thenReturn(processorNode);
+
Mockito.when(processorNode.getRelationship(relationshipName)).thenReturn(relationship);
+
Mockito.when(processorNode.getConnections(relationship)).thenReturn(Set.of());
+
+ final FlowDifference difference = new StandardFlowDifference(
+ DifferenceType.AUTO_TERMINATED_RELATIONSHIPS_CHANGED,
currentProcessor, proposedProcessor,
+ Set.of(), Set.of(relationshipName), "Auto-terminated
relationships changed");
+
+ assertTrue(FlowDifferenceFilters.isEnvironmentalChange(difference, new
VersionedProcessGroup(), flowManager));
+
assertFalse(FlowDifferenceFilters.isComponentUpdateRequired(difference, new
VersionedProcessGroup(), flowManager));
+ }
+
@Test
public void testIsComponentUpdateRequiredForPublicPortNameChange() {
// A public-port name change must still be reported as "update
required" by the shared filter. The preservation of a