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);
+ }
+
}