pvillard31 commented on code in PR #11570:
URL: https://github.com/apache/nifi/pull/11570#discussion_r4147103206


##########
nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/registry/RegistryClientIT.java:
##########
@@ -58,6 +63,153 @@ public class RegistryClientIT extends NiFiSystemIT {
 
     public static final String FIRST_FLOW_ID = "first-flow";
 
+    @Test
+    public void testChangeVersionDrainsRemovedConnectionBeforeUpdate() throws 
Exception {
+        final Path gateFile = Path.of(System.getProperty("java.io.tmpdir"), 
"nifi-removed-connection-drain-" + System.nanoTime());
+        final RemovedConnectionFixture fixture = 
createRemovedConnectionFixture(gateFile, false, true);
+        final NiFiClientUtil util = getClientUtil();
+
+        final VersionedFlowUpdateRequestEntity initiated = 
util.initiateFlowVersionChange(fixture.groupId(), "2");
+        final String requestId = initiated.getRequest().getRequestId();
+        waitFor(() -> "Draining Removed 
Connections".equals(getNifiClient().getVersionsClient().getUpdateRequest(requestId).getRequest().getState()));
+        if (getNumberOfNodes() > 1) {
+            final List<Integer> queuedByNode = 
getNifiClient().getFlowClient().getConnectionStatus(fixture.connectionId(), 
true)
+                    .getConnectionStatus().getNodeSnapshots().stream()
+                    .map(node -> node.getStatusSnapshot().getFlowFilesQueued())
+                    .sorted()
+                    .toList();
+            assertEquals(List.of(0, 1), queuedByNode);
+        }
+        assertEquals("1", 
getNifiClient().getProcessGroupClient().getProcessGroup(fixture.groupId())
+                .getComponent().getVersionControlInformation().getVersion());
+
+        Files.createFile(gateFile);
+        final VersionedFlowUpdateRequestEntity completed = 
util.waitForVersionFlowUpdateComplete(requestId, true);
+        assertTrue(completed.getRequest().isComplete());
+        assertNull(completed.getRequest().getFailureReason());
+        assertEquals("2", 
completed.getRequest().getVersionControlInformation().getVersion());
+        util.waitForRunningProcessor(fixture.sourceId());
+        util.waitForRunningProcessor(fixture.destinationId());
+        
assertTrue(getConnections(fixture.groupId()).stream().noneMatch(connection -> 
fixture.connectionId().equals(connection.getId())));
+
+        Files.deleteIfExists(gateFile);
+    }
+
+    @Test
+    public void testCancelledRemovedConnectionDrainRestoresOriginalFlow() 
throws Exception {
+        final Path gateFile = Path.of(System.getProperty("java.io.tmpdir"), 
"nifi-removed-connection-cancel-" + System.nanoTime());
+        final RemovedConnectionFixture fixture = 
createRemovedConnectionFixture(gateFile, false, true);
+        final NiFiClientUtil util = getClientUtil();
+
+        final VersionedFlowUpdateRequestEntity initiated = 
util.initiateFlowVersionChange(fixture.groupId(), "2");
+        final String requestId = initiated.getRequest().getRequestId();
+        waitFor(() -> "Draining Removed 
Connections".equals(getNifiClient().getVersionsClient().getUpdateRequest(requestId).getRequest().getState()));
+
+        final VersionedFlowUpdateRequestEntity cancelled = 
getNifiClient().getVersionsClient().deleteUpdateRequest(requestId);
+        assertEquals("Request cancelled by user", 
cancelled.getRequest().getFailureReason());
+        assertOriginalFlowRestored(fixture);
+        assertTrue(getConnectionQueueSize(fixture.connectionId()) >= 1);
+    }
+
+    @Test
+    public void testRemovedConnectionDrainTimeoutRestoresOriginalFlow() throws 
Exception {
+        final Path gateFile = Path.of(System.getProperty("java.io.tmpdir"), 
"nifi-removed-connection-timeout-" + System.nanoTime());
+        final RemovedConnectionFixture fixture = 
createRemovedConnectionFixture(gateFile, false, true);
+        final NiFiClientUtil util = getClientUtil();
+
+        final VersionedFlowUpdateRequestEntity initiated = 
util.initiateFlowVersionChange(fixture.groupId(), "2");

Review Comment:
   I kept the end-to-end timeout/restoration check at the fixed production 
30-second deadline; it passes on standalone and clustered NiFi. I prefer not to 
add a runtime timeout override solely to shorten this test, since that would 
introduce a configuration surface for a fixed safety limit. The deterministic 
coordinator test also checks that producer stop and queue wait share one 
deadline. Would you be comfortable retaining the 30-second end-to-end test?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to