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 f29ef183ec8 NIFI-16274 - Expose node connection state from 
FlowController (#11607)
f29ef183ec8 is described below

commit f29ef183ec811813aae943e316864c91d3368aee
Author: Pierre Villard <[email protected]>
AuthorDate: Wed Sep 2 21:55:26 2026 +0200

    NIFI-16274 - Expose node connection state from FlowController (#11607)
    
    * NIFI-16274 Expose node connection state from FlowController
    
    * NIFI-16274 Address node connection state review feedback
---
 .../org/apache/nifi/controller/FlowController.java | 39 ++++++++++++--
 .../nifi/controller/StandardFlowService.java       | 25 +++++----
 .../nifi/controller/TestStandardFlowService.java   | 62 ++++++++++++++++++++--
 3 files changed, 110 insertions(+), 16 deletions(-)

diff --git 
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
 
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
index 7bea91db241..8044d9abb66 100644
--- 
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
+++ 
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
@@ -39,7 +39,6 @@ import 
org.apache.nifi.cluster.coordination.ClusterCoordinator;
 import org.apache.nifi.cluster.coordination.heartbeat.HeartbeatMonitor;
 import org.apache.nifi.cluster.coordination.node.ClusterRoles;
 import org.apache.nifi.cluster.coordination.node.DisconnectionCode;
-import org.apache.nifi.cluster.coordination.node.NodeConnectionState;
 import org.apache.nifi.cluster.coordination.node.NodeConnectionStatus;
 import org.apache.nifi.cluster.protocol.DataFlow;
 import org.apache.nifi.cluster.protocol.Heartbeat;
@@ -3060,7 +3059,7 @@ public class FlowController implements 
ReportingTaskProvider, FlowAnalysisRulePr
             return Collections.emptyList();
         }
 
-        return 
clusterCoordinator.getNodeIdentifiers(NodeConnectionState.CONNECTED).stream()
+        return 
clusterCoordinator.getNodeIdentifiers(org.apache.nifi.cluster.coordination.node.NodeConnectionState.CONNECTED).stream()
                 
.sorted(Comparator.comparing(NodeIdentifier::getApiAddress).thenComparingInt(NodeIdentifier::getApiPort))
                 .toList();
     }
@@ -3070,6 +3069,39 @@ public class FlowController implements 
ReportingTaskProvider, FlowAnalysisRulePr
         return configuredForClustering;
     }
 
+    @Override
+    public NodeConnectionState getNodeConnectionState() {
+        final NodeConnectionState nodeConnectionState;
+
+        readLock.lock();
+        try {
+            if (!configuredForClustering) {
+                nodeConnectionState = NodeConnectionState.STANDALONE;
+            } else if (connectionStatus == null) {
+                nodeConnectionState = NodeConnectionState.DISCONNECTED;
+            } else {
+                nodeConnectionState = 
mapNodeConnectionState(connectionStatus.getState());
+            }
+        } finally {
+            readLock.unlock("getNodeConnectionState");
+        }
+
+        return nodeConnectionState;
+    }
+
+    private static NodeConnectionState mapNodeConnectionState(
+            final 
org.apache.nifi.cluster.coordination.node.NodeConnectionState connectionState) {
+        return switch (connectionState) {
+            case CONNECTING -> NodeConnectionState.CONNECTING;
+            case CONNECTED -> NodeConnectionState.CONNECTED;
+            case OFFLOADING -> NodeConnectionState.OFFLOADING;
+            case DISCONNECTING -> NodeConnectionState.DISCONNECTING;
+            case OFFLOADED -> NodeConnectionState.OFFLOADED;
+            case DISCONNECTED -> NodeConnectionState.DISCONNECTED;
+            case REMOVED -> NodeConnectionState.REMOVED;
+        };
+    }
+
     void registerForClusterCoordinator(final boolean participate) {
         final String participantId = participate ? 
heartbeatMonitor.getHeartbeatAddress() : null;
 
@@ -3648,7 +3680,8 @@ public class FlowController implements 
ReportingTaskProvider, FlowAnalysisRulePr
     public boolean isConnected() {
         rwLock.readLock().lock();
         try {
-            return connectionStatus != null && connectionStatus.getState() == 
NodeConnectionState.CONNECTED;
+            return connectionStatus != null
+                    && connectionStatus.getState() == 
org.apache.nifi.cluster.coordination.node.NodeConnectionState.CONNECTED;
         } finally {
             rwLock.readLock().unlock();
         }
diff --git 
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardFlowService.java
 
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardFlowService.java
index dc9bba00d8c..91c1197b271 100644
--- 
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardFlowService.java
+++ 
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardFlowService.java
@@ -762,7 +762,7 @@ public class StandardFlowService implements FlowService, 
ProtocolHandler {
     }
 
     // write lock must already be acquired
-    private void loadFromBytes(final DataFlow proposedFlow, final boolean 
allowEmptyFlow, final BundleUpdateStrategy bundleUpdateStrategy)
+    void loadFromBytes(final DataFlow proposedFlow, final boolean 
allowEmptyFlow, final BundleUpdateStrategy bundleUpdateStrategy)
             throws IOException, FlowSerializationException, 
FlowSynchronizationException, UninheritableFlowException, 
MissingBundleException {
         logger.trace("Loading flow from bytes");
 
@@ -920,15 +920,12 @@ public class StandardFlowService implements FlowService, 
ProtocolHandler {
         }
     }
 
-    private void loadFromConnectionResponse(final ConnectionResponse response) 
throws ConnectionException {
+    void loadFromConnectionResponse(final ConnectionResponse response) throws 
ConnectionException {
         writeLock.lock();
         try {
             connectionGeneration.incrementAndGet();
 
-            if (response.getNodeConnectionStatuses() != null) {
-                
clusterCoordinator.resetNodeStatuses(response.getNodeConnectionStatuses().stream()
-                    
.collect(Collectors.toMap(NodeConnectionStatus::getNodeIdentifier, status -> 
status)));
-            }
+            initializeNodeConnectionState(response);
 
             // get the dataflow from the response
             final DataFlow dataFlow = response.getDataFlow();
@@ -936,10 +933,6 @@ public class StandardFlowService implements FlowService, 
ProtocolHandler {
                 logger.trace("ResponseFlow = {}", new 
String(dataFlow.getFlow(), StandardCharsets.UTF_8));
             }
 
-            logger.info("Setting Flow Controller's Node ID: {}", nodeId);
-            nodeId = response.getNodeIdentifier();
-            controller.setNodeId(nodeId);
-
             // sync NARs before loading flow, otherwise components could be 
ghosted and fail to join the cluster
             narManager.syncWithClusterCoordinator();
 
@@ -994,6 +987,18 @@ public class StandardFlowService implements FlowService, 
ProtocolHandler {
 
     }
 
+    private void initializeNodeConnectionState(final ConnectionResponse 
response) {
+        if (response.getNodeConnectionStatuses() != null) {
+            
clusterCoordinator.resetNodeStatuses(response.getNodeConnectionStatuses().stream()
+                    
.collect(Collectors.toMap(NodeConnectionStatus::getNodeIdentifier, status -> 
status)));
+        }
+
+        nodeId = response.getNodeIdentifier();
+        logger.info("Setting Flow Controller's Node ID: {}", nodeId);
+        controller.setNodeId(nodeId);
+        controller.setConnectionStatus(new NodeConnectionStatus(nodeId, 
NodeConnectionState.CONNECTING));
+    }
+
     private void initializeController() throws IOException {
         if (firstControllerInitialization) {
             logger.debug("First controller initialization, initializing 
controller...");
diff --git 
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardFlowService.java
 
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardFlowService.java
index 1e1a8198e34..f19d6c2d613 100644
--- 
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardFlowService.java
+++ 
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardFlowService.java
@@ -19,13 +19,19 @@ package org.apache.nifi.controller;
 import org.apache.nifi.asset.AssetSynchronizer;
 import org.apache.nifi.authorization.Authorizer;
 import org.apache.nifi.cluster.coordination.ClusterCoordinator;
+import org.apache.nifi.cluster.coordination.node.NodeConnectionState;
 import org.apache.nifi.cluster.coordination.node.NodeConnectionStatus;
+import org.apache.nifi.cluster.protocol.ComponentRevisionSnapshot;
+import org.apache.nifi.cluster.protocol.ConnectionResponse;
 import org.apache.nifi.cluster.protocol.NodeIdentifier;
+import org.apache.nifi.cluster.protocol.StandardDataFlow;
 import org.apache.nifi.cluster.protocol.impl.NodeProtocolSenderListener;
 import org.apache.nifi.cluster.protocol.message.DisconnectMessage;
 import org.apache.nifi.components.state.Scope;
 import org.apache.nifi.components.state.StateManager;
 import org.apache.nifi.components.state.StateManagerProvider;
+import org.apache.nifi.controller.serialization.FlowSynchronizationException;
+import org.apache.nifi.groups.BundleUpdateStrategy;
 import org.apache.nifi.nar.ExtensionManager;
 import org.apache.nifi.nar.NarManager;
 import org.apache.nifi.state.MockStateMap;
@@ -34,18 +40,27 @@ import org.apache.nifi.web.revision.RevisionManager;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
+import org.mockito.ArgumentCaptor;
+import org.mockito.InOrder;
 
 import java.io.IOException;
 import java.nio.file.Path;
 import java.util.Collections;
+import java.util.List;
 import java.util.Map;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.inOrder;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.reset;
+import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -54,6 +69,7 @@ public class TestStandardFlowService {
     private FlowController controller;
     private ClusterCoordinator clusterCoordinator;
     private StandardFlowService flowService;
+    private NiFiProperties nifiProperties;
 
     @TempDir
     private Path tempDir;
@@ -77,7 +93,7 @@ public class TestStandardFlowService {
         
when(controller.getExtensionManager()).thenReturn(mock(ExtensionManager.class));
 
         final Path flowConfigFile = tempDir.resolve("flow.json.gz");
-        final NiFiProperties nifiProperties = 
NiFiProperties.createBasicNiFiProperties(null, Map.of(
+        nifiProperties = NiFiProperties.createBasicNiFiProperties(null, Map.of(
                 NiFiProperties.FLOW_CONFIGURATION_FILE, 
flowConfigFile.toString(),
                 NiFiProperties.FLOW_CONTROLLER_GRACEFUL_SHUTDOWN_PERIOD, "10 
secs",
                 NiFiProperties.WEB_HTTPS_HOST, "localhost",
@@ -124,12 +140,52 @@ public class TestStandardFlowService {
         verify(clusterCoordinator, never()).setConnected(false);
     }
 
+    @Test
+    public void testNodeConnectionStateConnectingBeforeFlowSynchronization() 
throws Exception {
+        final FlowController clusteredController = mock(FlowController.class);
+        final StateManagerProvider stateManagerProvider = 
controller.getStateManagerProvider();
+        final ExtensionManager extensionManager = 
controller.getExtensionManager();
+        final NarManager narManager = mock(NarManager.class);
+        
when(clusteredController.getStateManagerProvider()).thenReturn(stateManagerProvider);
+        
when(clusteredController.getExtensionManager()).thenReturn(extensionManager);
+
+        final StandardFlowService clusteredFlowService = 
spy(StandardFlowService.createClusteredInstance(clusteredController, 
nifiProperties,
+                mock(NodeProtocolSenderListener.class), clusterCoordinator, 
mock(RevisionManager.class), narManager,
+                mock(AssetSynchronizer.class), mock(AssetSynchronizer.class), 
mock(Authorizer.class)));
+        final NodeIdentifier nodeIdentifier = createNodeIdentifier();
+        final NodeConnectionStatus connectingStatus = new 
NodeConnectionStatus(nodeIdentifier, NodeConnectionState.CONNECTING);
+        final ConnectionResponse response = new 
ConnectionResponse(nodeIdentifier,
+                new StandardDataFlow(new byte[0], new byte[0], new byte[0], 
Collections.emptySet()), "instance-1",
+                List.of(connectingStatus), 
mock(ComponentRevisionSnapshot.class));
+        final ArgumentCaptor<NodeConnectionStatus> statusCaptor = 
ArgumentCaptor.forClass(NodeConnectionStatus.class);
+        doAnswer(invocation -> {
+            verify(clusteredController, never()).setClustered(eq(true), any());
+            throw new FlowSynchronizationException("Stop after observing the 
flow synchronization boundary");
+        }).when(clusteredFlowService).loadFromBytes(any(), eq(true), 
eq(BundleUpdateStrategy.USE_SPECIFIED_OR_COMPATIBLE_OR_GHOST));
+
+        assertThrows(FlowSynchronizationException.class, () -> 
clusteredFlowService.loadFromConnectionResponse(response));
+
+        final InOrder inOrder = inOrder(clusterCoordinator, 
clusteredController, narManager, clusteredFlowService);
+        inOrder.verify(clusterCoordinator).resetNodeStatuses(any(Map.class));
+        inOrder.verify(clusteredController).setNodeId(nodeIdentifier);
+        
inOrder.verify(clusteredController).setConnectionStatus(statusCaptor.capture());
+        inOrder.verify(narManager).syncWithClusterCoordinator();
+        inOrder.verify(clusteredFlowService).loadFromBytes(any(), eq(true), 
eq(BundleUpdateStrategy.USE_SPECIFIED_OR_COMPATIBLE_OR_GHOST));
+        assertEquals(nodeIdentifier, 
statusCaptor.getValue().getNodeIdentifier());
+        assertEquals(NodeConnectionState.CONNECTING, 
statusCaptor.getValue().getState());
+        verify(clusteredController, never()).setClustered(anyBoolean(), any());
+    }
+
     private DisconnectMessage createDisconnectMessage(final String 
explanation) {
-        final NodeIdentifier nodeIdentifier = new NodeIdentifier("node-1", 
"localhost", 8443,
-                "localhost", 9090, "localhost", 6342, 10443, false);
+        final NodeIdentifier nodeIdentifier = createNodeIdentifier();
         final DisconnectMessage message = new DisconnectMessage();
         message.setNodeId(nodeIdentifier);
         message.setExplanation(explanation);
         return message;
     }
+
+    private NodeIdentifier createNodeIdentifier() {
+        return new NodeIdentifier("node-1", "localhost", 8443, "localhost", 
9090, "localhost", 6342, 10443, false);
+    }
+
 }

Reply via email to