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


##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/util/ClusterReplicationComponentLifecycle.java:
##########
@@ -680,6 +739,183 @@ private boolean waitForControllerServiceValidation(final 
NiFiUser user, final UR
         return false;
     }
 
+    private ListingRequestResult createFlowFileListingRequest(final NiFiUser 
user, final URI originalUri, final String connectionId,
+                                                              final 
Set<NodeIdentifier> expectedNodes) throws LifecycleManagementException {
+        final URI createListingRequestUri;
+        try {
+            createListingRequestUri = new URI(originalUri.getScheme(), 
originalUri.getUserInfo(), originalUri.getHost(), originalUri.getPort(),
+                "/nifi-api/flowfile-queues/" + connectionId + 
"/listing-requests", null, originalUri.getFragment());
+        } catch (final URISyntaxException e) {
+            throw new RuntimeException(e);
+        }
+
+        try {
+            final AsyncClusterResponse clusterResponse = 
replicateFlowFileListingRequest(expectedNodes, user, HttpMethod.POST, 
createListingRequestUri);

Review Comment:
   [GPT-5.6 Sol] **Cluster updates now require an unrelated data permission.** 
This replicated FlowFile-listing request runs as the user who requested the 
flow update. The listing endpoint requires Read Source Data permission, while 
flow updates require component read and write permission. As a result, a user 
who is allowed to update the flow but not inspect queued data can update 
successfully on standalone NiFi but receives a 403 and fails here in a cluster. 
Please use an internal queue-count operation that does not expose FlowFile 
listings, and cover this with a user that lacks data-read permission.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/util/ClusterReplicationComponentLifecycle.java:
##########
@@ -680,6 +739,183 @@ private boolean waitForControllerServiceValidation(final 
NiFiUser user, final UR
         return false;
     }
 
+    private ListingRequestResult createFlowFileListingRequest(final NiFiUser 
user, final URI originalUri, final String connectionId,
+                                                              final 
Set<NodeIdentifier> expectedNodes) throws LifecycleManagementException {
+        final URI createListingRequestUri;
+        try {
+            createListingRequestUri = new URI(originalUri.getScheme(), 
originalUri.getUserInfo(), originalUri.getHost(), originalUri.getPort(),
+                "/nifi-api/flowfile-queues/" + connectionId + 
"/listing-requests", null, originalUri.getFragment());
+        } catch (final URISyntaxException e) {
+            throw new RuntimeException(e);
+        }
+
+        try {
+            final AsyncClusterResponse clusterResponse = 
replicateFlowFileListingRequest(expectedNodes, user, HttpMethod.POST, 
createListingRequestUri);
+

Review Comment:
   [GPT-5.6 Sol] **The 30-second drain limit does not cover this wait.** 
`awaitMergedResponse()` has no deadline from the drain pause, and cancellation 
cannot wake it. If one node does not finish the listing request, producers can 
remain stopped beyond the promised 30 seconds. Please bound the cluster wait 
using the remaining drain time and make cancellation able to end it.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/RemovedConnectionDrainCoordinator.java:
##########
@@ -0,0 +1,567 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.web;
+
+import org.apache.nifi.connectable.ConnectableType;
+import org.apache.nifi.controller.ScheduledState;
+import org.apache.nifi.web.api.dto.AffectedComponentDTO;
+import org.apache.nifi.web.api.entity.AffectedComponentEntity;
+import org.apache.nifi.web.util.CancellableTimedPause;
+import org.apache.nifi.web.util.ComponentLifecycle;
+import org.apache.nifi.web.util.InvalidComponentAction;
+import org.apache.nifi.web.util.LifecycleManagementException;
+import org.apache.nifi.web.util.Pause;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.net.URI;
+import java.time.Duration;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.function.LongSupplier;
+import java.util.stream.Collectors;
+
+public final class RemovedConnectionDrainCoordinator {
+    private static final Logger logger = 
LoggerFactory.getLogger(RemovedConnectionDrainCoordinator.class);
+    static final Duration DEFAULT_DRAIN_TIMEOUT = Duration.ofSeconds(30);
+    private static final Duration DEFAULT_POLL_INTERVAL = 
Duration.ofMillis(250);
+
+    private final RemovedConnectionDrainClassifier classifier;
+    private final PauseFactory pauseFactory;
+    private final Duration drainTimeout;
+
+    public RemovedConnectionDrainCoordinator() {
+        this(new RemovedConnectionDrainClassifier(), new 
MonotonicPauseFactory(DEFAULT_POLL_INTERVAL, System::nanoTime), 
DEFAULT_DRAIN_TIMEOUT);
+    }
+
+    RemovedConnectionDrainCoordinator(final RemovedConnectionDrainClassifier 
classifier, final PauseFactory pauseFactory, final Duration drainTimeout) {
+        this.classifier = Objects.requireNonNull(classifier, "Removed 
Connection Drain Classifier required");
+        this.pauseFactory = Objects.requireNonNull(pauseFactory, "Pause 
Factory required");
+        this.drainTimeout = Objects.requireNonNull(drainTimeout, "Drain 
Timeout required");
+    }
+
+    public DrainResult coordinateDrain(final FlowUpdateImpact 
flowUpdateImpact, final RemovedConnectionDrainClassifier.Context context,
+                                       final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                       final CancellationHandle 
cancellationHandle) throws LifecycleManagementException {
+        Objects.requireNonNull(flowUpdateImpact, "Flow Update Impact 
required");
+        Objects.requireNonNull(context, "Removed Connection Drain Context 
required");
+        Objects.requireNonNull(componentLifecycle, "Component Lifecycle 
required");
+        Objects.requireNonNull(requestUri, "Request URI required");
+        Objects.requireNonNull(groupId, "Group ID required");
+        Objects.requireNonNull(cancellationHandle, "Cancellation Handle 
required");
+
+        final RemovedConnectionDrainClassifier.Context queueAwareContext = 
createQueueAwareContext(flowUpdateImpact, context, componentLifecycle, 
requestUri);
+        final RemovedConnectionDrainClassifier.BatchResult batchResult = 
classifier.classify(flowUpdateImpact, queueAwareContext);
+        if (!batchResult.isSupported()) {
+            throw new 
LifecycleManagementException(buildClassificationFailureMessage(batchResult));
+        }
+
+        final Set<String> candidateConnectionIds = 
batchResult.connectionResults().stream()
+                .filter(result -> result.classification() == 
RemovedConnectionDrainClassifier.Classification.CANDIDATE)
+                .map(result -> result.connection().getConnectionInstanceId())
+                .collect(Collectors.toCollection(LinkedHashSet::new));
+
+        if (candidateConnectionIds.isEmpty()) {
+            return DrainResult.success(Collections.emptySet(), 
Collections.emptySet());
+        }
+
+        final Map<String, AffectedComponentEntity> affectedComponentsById = 
flowUpdateImpact.getAffectedComponents().stream()
+                .collect(Collectors.toMap(AffectedComponentEntity::getId, 
entity -> entity, (left, right) -> left, LinkedHashMap::new));
+        final Set<AffectedComponentEntity> componentsToStop = new 
LinkedHashSet<>();
+        for (final String producerBarrierComponentId : 
batchResult.producerBarrierComponentIds()) {
+            final AffectedComponentEntity entity = 
getProducerBarrierEntity(affectedComponentsById, queueAwareContext, 
producerBarrierComponentId);
+            if (entity == null || entity.getComponent() == null) {
+                continue;
+            }
+
+            if (isActive(entity.getComponent())) {
+                componentsToStop.add(entity);
+            }
+        }
+
+        final List<String> orderedCandidateConnectionIds = 
candidateConnectionIds.stream().sorted().toList();
+        final List<String> orderedProducerBarrierIds = 
componentsToStop.stream().map(AffectedComponentEntity::getId).sorted().toList();
+        logger.info("Starting drain of removed connections {} with producer 
barriers {}", orderedCandidateConnectionIds, orderedProducerBarrierIds);
+
+        final DeadlinePause drainPause = 
pauseFactory.createDrainPause(drainTimeout);
+        cancellationHandle.setCancelCallback(drainPause::cancel);
+
+        final Set<AffectedComponentEntity> drainStoppedComponents = new 
LinkedHashSet<>();
+        try {
+            if (!componentsToStop.isEmpty()) {
+                final Set<AffectedComponentEntity> updatedStoppedComponents = 
componentLifecycle.scheduleComponents(
+                        requestUri, groupId, componentsToStop, 
ScheduledState.STOPPED, drainPause, InvalidComponentAction.SKIP);
+                
drainStoppedComponents.addAll(getStoppedComponents(componentsToStop, 
updatedStoppedComponents));
+
+                if (!allComponentsStopped(componentsToStop, 
updatedStoppedComponents)) {
+                    if (cancellationHandle.isCancelled()) {
+                        return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                    }
+
+                    final Set<String> producerBarrierIds = 
componentsToStop.stream()
+                            .map(AffectedComponentEntity::getId)
+                            
.collect(Collectors.toCollection(LinkedHashSet::new));
+                    throw new 
LifecycleManagementException(buildStopTimeoutMessage(producerBarrierIds));
+                }
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            final boolean queuesDrained = 
componentLifecycle.waitForConnectionQueuesEmpty(requestUri, 
candidateConnectionIds, drainPause);
+            if (queuesDrained) {
+                if (cancellationHandle.isCancelled()) {
+                    return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                }
+
+                logger.info("Completed draining removed connections {}", 
orderedCandidateConnectionIds);
+                return DrainResult.success(candidateConnectionIds, 
drainStoppedComponents);

Review Comment:
   [GPT-5.6 Sol] **There is still a cancellation gap after a successful queue 
wait.** Cancellation can happen after the check on line 135 and before this 
success result reaches `FlowUpdateResource`. The caller then adds the 
drain-stopped producers to `runningComponents`, begins stopping all affected 
components, and returns on its cancellation check without restoring them. The 
handoff of the stopped components to the caller must account for cancellation 
that happens between this check and the next lifecycle step.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/RemovedConnectionDrainCoordinator.java:
##########
@@ -0,0 +1,567 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.web;
+
+import org.apache.nifi.connectable.ConnectableType;
+import org.apache.nifi.controller.ScheduledState;
+import org.apache.nifi.web.api.dto.AffectedComponentDTO;
+import org.apache.nifi.web.api.entity.AffectedComponentEntity;
+import org.apache.nifi.web.util.CancellableTimedPause;
+import org.apache.nifi.web.util.ComponentLifecycle;
+import org.apache.nifi.web.util.InvalidComponentAction;
+import org.apache.nifi.web.util.LifecycleManagementException;
+import org.apache.nifi.web.util.Pause;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.net.URI;
+import java.time.Duration;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.function.LongSupplier;
+import java.util.stream.Collectors;
+
+public final class RemovedConnectionDrainCoordinator {
+    private static final Logger logger = 
LoggerFactory.getLogger(RemovedConnectionDrainCoordinator.class);
+    static final Duration DEFAULT_DRAIN_TIMEOUT = Duration.ofSeconds(30);
+    private static final Duration DEFAULT_POLL_INTERVAL = 
Duration.ofMillis(250);
+
+    private final RemovedConnectionDrainClassifier classifier;
+    private final PauseFactory pauseFactory;
+    private final Duration drainTimeout;
+
+    public RemovedConnectionDrainCoordinator() {
+        this(new RemovedConnectionDrainClassifier(), new 
MonotonicPauseFactory(DEFAULT_POLL_INTERVAL, System::nanoTime), 
DEFAULT_DRAIN_TIMEOUT);
+    }
+
+    RemovedConnectionDrainCoordinator(final RemovedConnectionDrainClassifier 
classifier, final PauseFactory pauseFactory, final Duration drainTimeout) {
+        this.classifier = Objects.requireNonNull(classifier, "Removed 
Connection Drain Classifier required");
+        this.pauseFactory = Objects.requireNonNull(pauseFactory, "Pause 
Factory required");
+        this.drainTimeout = Objects.requireNonNull(drainTimeout, "Drain 
Timeout required");
+    }
+
+    public DrainResult coordinateDrain(final FlowUpdateImpact 
flowUpdateImpact, final RemovedConnectionDrainClassifier.Context context,
+                                       final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                       final CancellationHandle 
cancellationHandle) throws LifecycleManagementException {
+        Objects.requireNonNull(flowUpdateImpact, "Flow Update Impact 
required");
+        Objects.requireNonNull(context, "Removed Connection Drain Context 
required");
+        Objects.requireNonNull(componentLifecycle, "Component Lifecycle 
required");
+        Objects.requireNonNull(requestUri, "Request URI required");
+        Objects.requireNonNull(groupId, "Group ID required");
+        Objects.requireNonNull(cancellationHandle, "Cancellation Handle 
required");
+
+        final RemovedConnectionDrainClassifier.Context queueAwareContext = 
createQueueAwareContext(flowUpdateImpact, context, componentLifecycle, 
requestUri);
+        final RemovedConnectionDrainClassifier.BatchResult batchResult = 
classifier.classify(flowUpdateImpact, queueAwareContext);
+        if (!batchResult.isSupported()) {
+            throw new 
LifecycleManagementException(buildClassificationFailureMessage(batchResult));
+        }
+
+        final Set<String> candidateConnectionIds = 
batchResult.connectionResults().stream()
+                .filter(result -> result.classification() == 
RemovedConnectionDrainClassifier.Classification.CANDIDATE)
+                .map(result -> result.connection().getConnectionInstanceId())
+                .collect(Collectors.toCollection(LinkedHashSet::new));
+
+        if (candidateConnectionIds.isEmpty()) {
+            return DrainResult.success(Collections.emptySet(), 
Collections.emptySet());
+        }
+
+        final Map<String, AffectedComponentEntity> affectedComponentsById = 
flowUpdateImpact.getAffectedComponents().stream()
+                .collect(Collectors.toMap(AffectedComponentEntity::getId, 
entity -> entity, (left, right) -> left, LinkedHashMap::new));
+        final Set<AffectedComponentEntity> componentsToStop = new 
LinkedHashSet<>();
+        for (final String producerBarrierComponentId : 
batchResult.producerBarrierComponentIds()) {
+            final AffectedComponentEntity entity = 
getProducerBarrierEntity(affectedComponentsById, queueAwareContext, 
producerBarrierComponentId);
+            if (entity == null || entity.getComponent() == null) {
+                continue;
+            }
+
+            if (isActive(entity.getComponent())) {
+                componentsToStop.add(entity);
+            }
+        }
+
+        final List<String> orderedCandidateConnectionIds = 
candidateConnectionIds.stream().sorted().toList();
+        final List<String> orderedProducerBarrierIds = 
componentsToStop.stream().map(AffectedComponentEntity::getId).sorted().toList();
+        logger.info("Starting drain of removed connections {} with producer 
barriers {}", orderedCandidateConnectionIds, orderedProducerBarrierIds);
+
+        final DeadlinePause drainPause = 
pauseFactory.createDrainPause(drainTimeout);
+        cancellationHandle.setCancelCallback(drainPause::cancel);
+
+        final Set<AffectedComponentEntity> drainStoppedComponents = new 
LinkedHashSet<>();
+        try {
+            if (!componentsToStop.isEmpty()) {
+                final Set<AffectedComponentEntity> updatedStoppedComponents = 
componentLifecycle.scheduleComponents(
+                        requestUri, groupId, componentsToStop, 
ScheduledState.STOPPED, drainPause, InvalidComponentAction.SKIP);
+                
drainStoppedComponents.addAll(getStoppedComponents(componentsToStop, 
updatedStoppedComponents));
+
+                if (!allComponentsStopped(componentsToStop, 
updatedStoppedComponents)) {
+                    if (cancellationHandle.isCancelled()) {
+                        return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                    }
+
+                    final Set<String> producerBarrierIds = 
componentsToStop.stream()
+                            .map(AffectedComponentEntity::getId)
+                            
.collect(Collectors.toCollection(LinkedHashSet::new));
+                    throw new 
LifecycleManagementException(buildStopTimeoutMessage(producerBarrierIds));
+                }
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            final boolean queuesDrained = 
componentLifecycle.waitForConnectionQueuesEmpty(requestUri, 
candidateConnectionIds, drainPause);
+            if (queuesDrained) {
+                if (cancellationHandle.isCancelled()) {
+                    return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                }
+
+                logger.info("Completed draining removed connections {}", 
orderedCandidateConnectionIds);
+                return DrainResult.success(candidateConnectionIds, 
drainStoppedComponents);
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            throw new 
LifecycleManagementException(buildQueueTimeoutMessage(candidateConnectionIds));
+        } catch (final LifecycleManagementException e) {
+            final Set<AffectedComponentEntity> stoppedComponents = 
getStoppedComponentsToRestore(queueAwareContext, componentsToStop, 
drainStoppedComponents);
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, stoppedComponents);
+            }
+
+            final LifecycleManagementException failure = decorateFailure(e, 
componentsToStop, candidateConnectionIds);
+            logger.warn("Removed connection drain failed for connections {}", 
orderedCandidateConnectionIds, failure);
+            restoreOrSuppress(componentLifecycle, requestUri, groupId, 
stoppedComponents, failure);
+            throw failure;
+        } catch (final RuntimeException e) {
+            final Set<AffectedComponentEntity> stoppedComponents = 
getStoppedComponentsToRestore(queueAwareContext, componentsToStop, 
drainStoppedComponents);
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, stoppedComponents);
+            }
+
+            final LifecycleManagementException failure = new 
LifecycleManagementException(
+                    "Removed connection drain failed for connections " + 
candidateConnectionIds.stream().sorted().toList(), e);
+            logger.warn("Removed connection drain failed for connections {}", 
orderedCandidateConnectionIds, failure);
+            restoreOrSuppress(componentLifecycle, requestUri, groupId, 
stoppedComponents, failure);
+            throw failure;
+        } finally {
+            cancellationHandle.setCancelCallback(null);
+        }
+    }
+
+    private DrainResult restoreAfterCancellation(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                                 final Set<String> 
candidateConnectionIds,
+                                                 final 
Set<AffectedComponentEntity> drainStoppedComponents) {
+        LifecycleManagementException restorationFailure = null;
+        try {
+            restoreStoppedComponents(componentLifecycle, requestUri, groupId, 
drainStoppedComponents);
+        } catch (final LifecycleManagementException e) {

Review Comment:
   [GPT-5.6 Sol] **An unchecked exception during restoration is still lost.** 
Cluster scheduling can throw unchecked cluster-state exceptions. This catch, 
and the corresponding catch in `restoreOrSuppress`, handle only 
`LifecycleManagementException`. An unchecked exception therefore escapes 
without being recorded in the cancellation result or attached to the original 
failure, and the producers can remain stopped. Please handle unchecked 
restoration failures, log them, and preserve them with the request failure.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/RemovedConnectionDrainCoordinator.java:
##########
@@ -0,0 +1,567 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.web;
+
+import org.apache.nifi.connectable.ConnectableType;
+import org.apache.nifi.controller.ScheduledState;
+import org.apache.nifi.web.api.dto.AffectedComponentDTO;
+import org.apache.nifi.web.api.entity.AffectedComponentEntity;
+import org.apache.nifi.web.util.CancellableTimedPause;
+import org.apache.nifi.web.util.ComponentLifecycle;
+import org.apache.nifi.web.util.InvalidComponentAction;
+import org.apache.nifi.web.util.LifecycleManagementException;
+import org.apache.nifi.web.util.Pause;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.net.URI;
+import java.time.Duration;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.function.LongSupplier;
+import java.util.stream.Collectors;
+
+public final class RemovedConnectionDrainCoordinator {
+    private static final Logger logger = 
LoggerFactory.getLogger(RemovedConnectionDrainCoordinator.class);
+    static final Duration DEFAULT_DRAIN_TIMEOUT = Duration.ofSeconds(30);
+    private static final Duration DEFAULT_POLL_INTERVAL = 
Duration.ofMillis(250);
+
+    private final RemovedConnectionDrainClassifier classifier;
+    private final PauseFactory pauseFactory;
+    private final Duration drainTimeout;
+
+    public RemovedConnectionDrainCoordinator() {
+        this(new RemovedConnectionDrainClassifier(), new 
MonotonicPauseFactory(DEFAULT_POLL_INTERVAL, System::nanoTime), 
DEFAULT_DRAIN_TIMEOUT);
+    }
+
+    RemovedConnectionDrainCoordinator(final RemovedConnectionDrainClassifier 
classifier, final PauseFactory pauseFactory, final Duration drainTimeout) {
+        this.classifier = Objects.requireNonNull(classifier, "Removed 
Connection Drain Classifier required");
+        this.pauseFactory = Objects.requireNonNull(pauseFactory, "Pause 
Factory required");
+        this.drainTimeout = Objects.requireNonNull(drainTimeout, "Drain 
Timeout required");
+    }
+
+    public DrainResult coordinateDrain(final FlowUpdateImpact 
flowUpdateImpact, final RemovedConnectionDrainClassifier.Context context,
+                                       final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                       final CancellationHandle 
cancellationHandle) throws LifecycleManagementException {
+        Objects.requireNonNull(flowUpdateImpact, "Flow Update Impact 
required");
+        Objects.requireNonNull(context, "Removed Connection Drain Context 
required");
+        Objects.requireNonNull(componentLifecycle, "Component Lifecycle 
required");
+        Objects.requireNonNull(requestUri, "Request URI required");
+        Objects.requireNonNull(groupId, "Group ID required");
+        Objects.requireNonNull(cancellationHandle, "Cancellation Handle 
required");
+
+        final RemovedConnectionDrainClassifier.Context queueAwareContext = 
createQueueAwareContext(flowUpdateImpact, context, componentLifecycle, 
requestUri);
+        final RemovedConnectionDrainClassifier.BatchResult batchResult = 
classifier.classify(flowUpdateImpact, queueAwareContext);
+        if (!batchResult.isSupported()) {
+            throw new 
LifecycleManagementException(buildClassificationFailureMessage(batchResult));
+        }
+
+        final Set<String> candidateConnectionIds = 
batchResult.connectionResults().stream()
+                .filter(result -> result.classification() == 
RemovedConnectionDrainClassifier.Classification.CANDIDATE)
+                .map(result -> result.connection().getConnectionInstanceId())
+                .collect(Collectors.toCollection(LinkedHashSet::new));
+
+        if (candidateConnectionIds.isEmpty()) {
+            return DrainResult.success(Collections.emptySet(), 
Collections.emptySet());
+        }
+
+        final Map<String, AffectedComponentEntity> affectedComponentsById = 
flowUpdateImpact.getAffectedComponents().stream()
+                .collect(Collectors.toMap(AffectedComponentEntity::getId, 
entity -> entity, (left, right) -> left, LinkedHashMap::new));
+        final Set<AffectedComponentEntity> componentsToStop = new 
LinkedHashSet<>();
+        for (final String producerBarrierComponentId : 
batchResult.producerBarrierComponentIds()) {
+            final AffectedComponentEntity entity = 
getProducerBarrierEntity(affectedComponentsById, queueAwareContext, 
producerBarrierComponentId);
+            if (entity == null || entity.getComponent() == null) {
+                continue;
+            }
+
+            if (isActive(entity.getComponent())) {
+                componentsToStop.add(entity);
+            }
+        }
+
+        final List<String> orderedCandidateConnectionIds = 
candidateConnectionIds.stream().sorted().toList();
+        final List<String> orderedProducerBarrierIds = 
componentsToStop.stream().map(AffectedComponentEntity::getId).sorted().toList();
+        logger.info("Starting drain of removed connections {} with producer 
barriers {}", orderedCandidateConnectionIds, orderedProducerBarrierIds);
+
+        final DeadlinePause drainPause = 
pauseFactory.createDrainPause(drainTimeout);
+        cancellationHandle.setCancelCallback(drainPause::cancel);
+
+        final Set<AffectedComponentEntity> drainStoppedComponents = new 
LinkedHashSet<>();
+        try {
+            if (!componentsToStop.isEmpty()) {
+                final Set<AffectedComponentEntity> updatedStoppedComponents = 
componentLifecycle.scheduleComponents(
+                        requestUri, groupId, componentsToStop, 
ScheduledState.STOPPED, drainPause, InvalidComponentAction.SKIP);
+                
drainStoppedComponents.addAll(getStoppedComponents(componentsToStop, 
updatedStoppedComponents));
+
+                if (!allComponentsStopped(componentsToStop, 
updatedStoppedComponents)) {
+                    if (cancellationHandle.isCancelled()) {
+                        return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                    }
+
+                    final Set<String> producerBarrierIds = 
componentsToStop.stream()
+                            .map(AffectedComponentEntity::getId)
+                            
.collect(Collectors.toCollection(LinkedHashSet::new));
+                    throw new 
LifecycleManagementException(buildStopTimeoutMessage(producerBarrierIds));
+                }
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            final boolean queuesDrained = 
componentLifecycle.waitForConnectionQueuesEmpty(requestUri, 
candidateConnectionIds, drainPause);
+            if (queuesDrained) {
+                if (cancellationHandle.isCancelled()) {
+                    return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                }
+
+                logger.info("Completed draining removed connections {}", 
orderedCandidateConnectionIds);
+                return DrainResult.success(candidateConnectionIds, 
drainStoppedComponents);
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            throw new 
LifecycleManagementException(buildQueueTimeoutMessage(candidateConnectionIds));
+        } catch (final LifecycleManagementException e) {
+            final Set<AffectedComponentEntity> stoppedComponents = 
getStoppedComponentsToRestore(queueAwareContext, componentsToStop, 
drainStoppedComponents);
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, stoppedComponents);
+            }
+
+            final LifecycleManagementException failure = decorateFailure(e, 
componentsToStop, candidateConnectionIds);
+            logger.warn("Removed connection drain failed for connections {}", 
orderedCandidateConnectionIds, failure);
+            restoreOrSuppress(componentLifecycle, requestUri, groupId, 
stoppedComponents, failure);
+            throw failure;
+        } catch (final RuntimeException e) {
+            final Set<AffectedComponentEntity> stoppedComponents = 
getStoppedComponentsToRestore(queueAwareContext, componentsToStop, 
drainStoppedComponents);
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, stoppedComponents);
+            }
+
+            final LifecycleManagementException failure = new 
LifecycleManagementException(
+                    "Removed connection drain failed for connections " + 
candidateConnectionIds.stream().sorted().toList(), e);
+            logger.warn("Removed connection drain failed for connections {}", 
orderedCandidateConnectionIds, failure);
+            restoreOrSuppress(componentLifecycle, requestUri, groupId, 
stoppedComponents, failure);
+            throw failure;
+        } finally {
+            cancellationHandle.setCancelCallback(null);
+        }
+    }
+
+    private DrainResult restoreAfterCancellation(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                                 final Set<String> 
candidateConnectionIds,
+                                                 final 
Set<AffectedComponentEntity> drainStoppedComponents) {
+        LifecycleManagementException restorationFailure = null;
+        try {
+            restoreStoppedComponents(componentLifecycle, requestUri, groupId, 
drainStoppedComponents);
+        } catch (final LifecycleManagementException e) {
+            restorationFailure = e;
+            logger.warn("Failed to restore producer barriers {} after removed 
connection drain cancellation",
+                    
drainStoppedComponents.stream().map(AffectedComponentEntity::getId).sorted().toList(),
 e);
+        }
+
+        return DrainResult.cancelled(candidateConnectionIds, 
drainStoppedComponents, restorationFailure);
+    }
+
+    private String buildClassificationFailureMessage(final 
RemovedConnectionDrainClassifier.BatchResult batchResult) {
+        final List<String> unsupportedConnections = 
batchResult.connectionResults().stream()
+                .filter(result -> result.classification() == 
RemovedConnectionDrainClassifier.Classification.UNSUPPORTED)
+                .map(result -> result.connection().getConnectionInstanceId() + 
"[reason=" + result.unsupportedReason().name() + "]")
+                .toList();
+        return "Removed connection drain preflight failed: " + 
unsupportedConnections;
+    }
+
+    private String buildStopTimeoutMessage(final Set<String> 
producerBarrierIds) {
+        return "Removed connection drain timed out [timedOut=true, 
cancelled=false, phase=stopping-producer-barriers, timeout="
+                + drainTimeout.toSeconds() + "s, componentIds=" + 
producerBarrierIds.stream().sorted().toList() + "]";
+    }
+
+    private String buildQueueTimeoutMessage(final Set<String> connectionIds) {
+        return "Removed connection drain timed out [timedOut=true, 
cancelled=false, phase=waiting-for-queues, timeout="
+                + drainTimeout.toSeconds() + "s, connectionIds=" + 
connectionIds.stream().sorted().toList() + "]";
+    }
+
+    private AffectedComponentEntity getProducerBarrierEntity(final Map<String, 
AffectedComponentEntity> affectedComponentsById,
+                                                             final 
RemovedConnectionDrainClassifier.Context context,
+                                                             final String 
producerBarrierComponentId) {
+        final AffectedComponentEntity affectedComponentEntity = 
affectedComponentsById.get(producerBarrierComponentId);
+        if (affectedComponentEntity != null && 
affectedComponentEntity.getComponent() != null) {
+            return affectedComponentEntity;
+        }
+
+        final RemovedConnectionDrainClassifier.LiveConnectable liveConnectable 
= context.getConnectable(producerBarrierComponentId);
+        if (liveConnectable == null) {
+            return null;
+        }
+
+        final String referenceType = getReferenceType(liveConnectable.type());
+        if (referenceType == null) {
+            return null;
+        }
+
+        final AffectedComponentDTO componentDto = new AffectedComponentDTO();
+        componentDto.setId(liveConnectable.id());
+        componentDto.setName(liveConnectable.id());
+        componentDto.setProcessGroupId(liveConnectable.processGroupId());
+        componentDto.setReferenceType(referenceType);
+        componentDto.setState(getState(liveConnectable));
+
+        final AffectedComponentEntity componentEntity = new 
AffectedComponentEntity();
+        componentEntity.setId(liveConnectable.id());
+        componentEntity.setReferenceType(referenceType);
+        componentEntity.setComponent(componentDto);
+        return componentEntity;
+    }
+
+    private String getReferenceType(final ConnectableType connectableType) {
+        if (connectableType == ConnectableType.PROCESSOR) {
+            return AffectedComponentDTO.COMPONENT_TYPE_PROCESSOR;
+        }
+        if (connectableType == ConnectableType.INPUT_PORT) {
+            return AffectedComponentDTO.COMPONENT_TYPE_INPUT_PORT;
+        }
+        if (connectableType == ConnectableType.OUTPUT_PORT) {
+            return AffectedComponentDTO.COMPONENT_TYPE_OUTPUT_PORT;
+        }
+
+        return null;
+    }
+
+    private String getState(final 
RemovedConnectionDrainClassifier.LiveConnectable liveConnectable) {
+        if (liveConnectable.type() == ConnectableType.PROCESSOR) {
+            return liveConnectable.physicalScheduledState() == null ? null : 
liveConnectable.physicalScheduledState().name();
+        }
+
+        return liveConnectable.running() ? "RUNNING" : "STOPPED";
+    }
+
+    private LifecycleManagementException decorateFailure(final 
LifecycleManagementException failure,
+                                                         final 
Set<AffectedComponentEntity> componentsToStop,
+                                                         final Set<String> 
candidateConnectionIds) {
+        final String message = failure.getMessage();
+        if (message != null && message.startsWith("Removed connection drain")) 
{
+            return failure;
+        }
+
+        final String decoratedMessage;
+        if (message != null && message.contains("waiting for connection 
queues")) {
+            decoratedMessage = "Removed connection drain failed while waiting 
for connections "
+                    + candidateConnectionIds.stream().sorted().toList() + ": " 
+ message;
+        } else {
+            final List<String> componentIds = 
componentsToStop.stream().map(AffectedComponentEntity::getId).sorted().toList();
+            decoratedMessage = "Removed connection drain failed while stopping 
producer barriers " + componentIds + ": " + message;
+        }
+
+        return new LifecycleManagementException(decoratedMessage, failure);
+    }
+
+    private void restoreOrSuppress(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                   final Set<AffectedComponentEntity> 
stoppedComponents, final LifecycleManagementException failure) {
+        try {
+            restoreStoppedComponents(componentLifecycle, requestUri, groupId, 
stoppedComponents);
+        } catch (final LifecycleManagementException restorationFailure) {
+            logger.warn("Failed to restore producer barriers {} after removed 
connection drain failure",
+                    
stoppedComponents.stream().map(AffectedComponentEntity::getId).sorted().toList(),
 restorationFailure);
+            failure.addSuppressed(restorationFailure);
+        }
+    }
+
+    private void restoreStoppedComponents(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                          final Set<AffectedComponentEntity> 
stoppedComponents) throws LifecycleManagementException {
+        if (stoppedComponents.isEmpty()) {
+            return;
+        }
+
+        componentLifecycle.scheduleComponents(requestUri, groupId, 
stoppedComponents, ScheduledState.RUNNING,
+                pauseFactory.createRestorationPause(), 
InvalidComponentAction.SKIP);
+    }
+
+    private RemovedConnectionDrainClassifier.Context 
createQueueAwareContext(final FlowUpdateImpact flowUpdateImpact,
+                                                                             
final RemovedConnectionDrainClassifier.Context context,
+                                                                             
final ComponentLifecycle componentLifecycle,
+                                                                             
final URI requestUri) throws LifecycleManagementException {
+        final Map<String, Boolean> knownQueueEmptyByConnectionId = new 
LinkedHashMap<>();
+        final Pause noWaitPause = NoWaitPause.INSTANCE;
+        for (final RemovedConnectionDescriptor removedConnection : 
flowUpdateImpact.getRemovedConnections()) {
+            final String connectionId = 
removedConnection.getConnectionInstanceId();
+            if (connectionId == null) {
+                continue;
+            }
+
+            final boolean knownQueueEmpty = 
componentLifecycle.waitForConnectionQueuesEmpty(requestUri, 
Set.of(connectionId), noWaitPause);

Review Comment:
   [GPT-5.6 Sol] This no-wait probe intentionally returns `false` when a 
removed connection has queued data. Both lifecycle implementations log the end 
of a non-empty queue wait at WARN, so every normal successful drain begins with 
a warning that looks like a failure. Please keep the preflight probe from 
producing the timeout warning, or move timeout logging to the coordinator where 
the caller knows whether a real wait expired.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/RemovedConnectionDrainCoordinator.java:
##########
@@ -0,0 +1,567 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.web;
+
+import org.apache.nifi.connectable.ConnectableType;
+import org.apache.nifi.controller.ScheduledState;
+import org.apache.nifi.web.api.dto.AffectedComponentDTO;
+import org.apache.nifi.web.api.entity.AffectedComponentEntity;
+import org.apache.nifi.web.util.CancellableTimedPause;
+import org.apache.nifi.web.util.ComponentLifecycle;
+import org.apache.nifi.web.util.InvalidComponentAction;
+import org.apache.nifi.web.util.LifecycleManagementException;
+import org.apache.nifi.web.util.Pause;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.net.URI;
+import java.time.Duration;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.function.LongSupplier;
+import java.util.stream.Collectors;
+
+public final class RemovedConnectionDrainCoordinator {
+    private static final Logger logger = 
LoggerFactory.getLogger(RemovedConnectionDrainCoordinator.class);
+    static final Duration DEFAULT_DRAIN_TIMEOUT = Duration.ofSeconds(30);
+    private static final Duration DEFAULT_POLL_INTERVAL = 
Duration.ofMillis(250);
+
+    private final RemovedConnectionDrainClassifier classifier;
+    private final PauseFactory pauseFactory;
+    private final Duration drainTimeout;
+
+    public RemovedConnectionDrainCoordinator() {
+        this(new RemovedConnectionDrainClassifier(), new 
MonotonicPauseFactory(DEFAULT_POLL_INTERVAL, System::nanoTime), 
DEFAULT_DRAIN_TIMEOUT);
+    }
+
+    RemovedConnectionDrainCoordinator(final RemovedConnectionDrainClassifier 
classifier, final PauseFactory pauseFactory, final Duration drainTimeout) {
+        this.classifier = Objects.requireNonNull(classifier, "Removed 
Connection Drain Classifier required");
+        this.pauseFactory = Objects.requireNonNull(pauseFactory, "Pause 
Factory required");
+        this.drainTimeout = Objects.requireNonNull(drainTimeout, "Drain 
Timeout required");
+    }
+
+    public DrainResult coordinateDrain(final FlowUpdateImpact 
flowUpdateImpact, final RemovedConnectionDrainClassifier.Context context,
+                                       final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                       final CancellationHandle 
cancellationHandle) throws LifecycleManagementException {
+        Objects.requireNonNull(flowUpdateImpact, "Flow Update Impact 
required");
+        Objects.requireNonNull(context, "Removed Connection Drain Context 
required");
+        Objects.requireNonNull(componentLifecycle, "Component Lifecycle 
required");
+        Objects.requireNonNull(requestUri, "Request URI required");
+        Objects.requireNonNull(groupId, "Group ID required");
+        Objects.requireNonNull(cancellationHandle, "Cancellation Handle 
required");
+
+        final RemovedConnectionDrainClassifier.Context queueAwareContext = 
createQueueAwareContext(flowUpdateImpact, context, componentLifecycle, 
requestUri);
+        final RemovedConnectionDrainClassifier.BatchResult batchResult = 
classifier.classify(flowUpdateImpact, queueAwareContext);
+        if (!batchResult.isSupported()) {
+            throw new 
LifecycleManagementException(buildClassificationFailureMessage(batchResult));
+        }
+
+        final Set<String> candidateConnectionIds = 
batchResult.connectionResults().stream()
+                .filter(result -> result.classification() == 
RemovedConnectionDrainClassifier.Classification.CANDIDATE)
+                .map(result -> result.connection().getConnectionInstanceId())
+                .collect(Collectors.toCollection(LinkedHashSet::new));
+
+        if (candidateConnectionIds.isEmpty()) {
+            return DrainResult.success(Collections.emptySet(), 
Collections.emptySet());
+        }
+
+        final Map<String, AffectedComponentEntity> affectedComponentsById = 
flowUpdateImpact.getAffectedComponents().stream()
+                .collect(Collectors.toMap(AffectedComponentEntity::getId, 
entity -> entity, (left, right) -> left, LinkedHashMap::new));
+        final Set<AffectedComponentEntity> componentsToStop = new 
LinkedHashSet<>();
+        for (final String producerBarrierComponentId : 
batchResult.producerBarrierComponentIds()) {
+            final AffectedComponentEntity entity = 
getProducerBarrierEntity(affectedComponentsById, queueAwareContext, 
producerBarrierComponentId);
+            if (entity == null || entity.getComponent() == null) {
+                continue;
+            }
+
+            if (isActive(entity.getComponent())) {
+                componentsToStop.add(entity);
+            }
+        }
+
+        final List<String> orderedCandidateConnectionIds = 
candidateConnectionIds.stream().sorted().toList();
+        final List<String> orderedProducerBarrierIds = 
componentsToStop.stream().map(AffectedComponentEntity::getId).sorted().toList();
+        logger.info("Starting drain of removed connections {} with producer 
barriers {}", orderedCandidateConnectionIds, orderedProducerBarrierIds);
+
+        final DeadlinePause drainPause = 
pauseFactory.createDrainPause(drainTimeout);
+        cancellationHandle.setCancelCallback(drainPause::cancel);
+
+        final Set<AffectedComponentEntity> drainStoppedComponents = new 
LinkedHashSet<>();
+        try {
+            if (!componentsToStop.isEmpty()) {
+                final Set<AffectedComponentEntity> updatedStoppedComponents = 
componentLifecycle.scheduleComponents(
+                        requestUri, groupId, componentsToStop, 
ScheduledState.STOPPED, drainPause, InvalidComponentAction.SKIP);
+                
drainStoppedComponents.addAll(getStoppedComponents(componentsToStop, 
updatedStoppedComponents));
+
+                if (!allComponentsStopped(componentsToStop, 
updatedStoppedComponents)) {
+                    if (cancellationHandle.isCancelled()) {
+                        return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                    }
+
+                    final Set<String> producerBarrierIds = 
componentsToStop.stream()
+                            .map(AffectedComponentEntity::getId)
+                            
.collect(Collectors.toCollection(LinkedHashSet::new));
+                    throw new 
LifecycleManagementException(buildStopTimeoutMessage(producerBarrierIds));
+                }
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            final boolean queuesDrained = 
componentLifecycle.waitForConnectionQueuesEmpty(requestUri, 
candidateConnectionIds, drainPause);
+            if (queuesDrained) {
+                if (cancellationHandle.isCancelled()) {
+                    return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+                }
+
+                logger.info("Completed draining removed connections {}", 
orderedCandidateConnectionIds);
+                return DrainResult.success(candidateConnectionIds, 
drainStoppedComponents);
+            }
+
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, drainStoppedComponents);
+            }
+
+            throw new 
LifecycleManagementException(buildQueueTimeoutMessage(candidateConnectionIds));
+        } catch (final LifecycleManagementException e) {
+            final Set<AffectedComponentEntity> stoppedComponents = 
getStoppedComponentsToRestore(queueAwareContext, componentsToStop, 
drainStoppedComponents);
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, stoppedComponents);
+            }
+
+            final LifecycleManagementException failure = decorateFailure(e, 
componentsToStop, candidateConnectionIds);
+            logger.warn("Removed connection drain failed for connections {}", 
orderedCandidateConnectionIds, failure);
+            restoreOrSuppress(componentLifecycle, requestUri, groupId, 
stoppedComponents, failure);
+            throw failure;
+        } catch (final RuntimeException e) {
+            final Set<AffectedComponentEntity> stoppedComponents = 
getStoppedComponentsToRestore(queueAwareContext, componentsToStop, 
drainStoppedComponents);
+            if (cancellationHandle.isCancelled()) {
+                return restoreAfterCancellation(componentLifecycle, 
requestUri, groupId, candidateConnectionIds, stoppedComponents);
+            }
+
+            final LifecycleManagementException failure = new 
LifecycleManagementException(
+                    "Removed connection drain failed for connections " + 
candidateConnectionIds.stream().sorted().toList(), e);
+            logger.warn("Removed connection drain failed for connections {}", 
orderedCandidateConnectionIds, failure);
+            restoreOrSuppress(componentLifecycle, requestUri, groupId, 
stoppedComponents, failure);
+            throw failure;
+        } finally {
+            cancellationHandle.setCancelCallback(null);
+        }
+    }
+
+    private DrainResult restoreAfterCancellation(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                                 final Set<String> 
candidateConnectionIds,
+                                                 final 
Set<AffectedComponentEntity> drainStoppedComponents) {
+        LifecycleManagementException restorationFailure = null;
+        try {
+            restoreStoppedComponents(componentLifecycle, requestUri, groupId, 
drainStoppedComponents);
+        } catch (final LifecycleManagementException e) {
+            restorationFailure = e;
+            logger.warn("Failed to restore producer barriers {} after removed 
connection drain cancellation",
+                    
drainStoppedComponents.stream().map(AffectedComponentEntity::getId).sorted().toList(),
 e);
+        }
+
+        return DrainResult.cancelled(candidateConnectionIds, 
drainStoppedComponents, restorationFailure);
+    }
+
+    private String buildClassificationFailureMessage(final 
RemovedConnectionDrainClassifier.BatchResult batchResult) {
+        final List<String> unsupportedConnections = 
batchResult.connectionResults().stream()
+                .filter(result -> result.classification() == 
RemovedConnectionDrainClassifier.Classification.UNSUPPORTED)
+                .map(result -> result.connection().getConnectionInstanceId() + 
"[reason=" + result.unsupportedReason().name() + "]")
+                .toList();
+        return "Removed connection drain preflight failed: " + 
unsupportedConnections;
+    }
+
+    private String buildStopTimeoutMessage(final Set<String> 
producerBarrierIds) {
+        return "Removed connection drain timed out [timedOut=true, 
cancelled=false, phase=stopping-producer-barriers, timeout="
+                + drainTimeout.toSeconds() + "s, componentIds=" + 
producerBarrierIds.stream().sorted().toList() + "]";
+    }
+
+    private String buildQueueTimeoutMessage(final Set<String> connectionIds) {
+        return "Removed connection drain timed out [timedOut=true, 
cancelled=false, phase=waiting-for-queues, timeout="
+                + drainTimeout.toSeconds() + "s, connectionIds=" + 
connectionIds.stream().sorted().toList() + "]";
+    }
+
+    private AffectedComponentEntity getProducerBarrierEntity(final Map<String, 
AffectedComponentEntity> affectedComponentsById,
+                                                             final 
RemovedConnectionDrainClassifier.Context context,
+                                                             final String 
producerBarrierComponentId) {
+        final AffectedComponentEntity affectedComponentEntity = 
affectedComponentsById.get(producerBarrierComponentId);
+        if (affectedComponentEntity != null && 
affectedComponentEntity.getComponent() != null) {
+            return affectedComponentEntity;
+        }
+
+        final RemovedConnectionDrainClassifier.LiveConnectable liveConnectable 
= context.getConnectable(producerBarrierComponentId);
+        if (liveConnectable == null) {
+            return null;
+        }
+
+        final String referenceType = getReferenceType(liveConnectable.type());
+        if (referenceType == null) {
+            return null;
+        }
+
+        final AffectedComponentDTO componentDto = new AffectedComponentDTO();
+        componentDto.setId(liveConnectable.id());
+        componentDto.setName(liveConnectable.id());
+        componentDto.setProcessGroupId(liveConnectable.processGroupId());
+        componentDto.setReferenceType(referenceType);
+        componentDto.setState(getState(liveConnectable));
+
+        final AffectedComponentEntity componentEntity = new 
AffectedComponentEntity();
+        componentEntity.setId(liveConnectable.id());
+        componentEntity.setReferenceType(referenceType);
+        componentEntity.setComponent(componentDto);
+        return componentEntity;
+    }
+
+    private String getReferenceType(final ConnectableType connectableType) {
+        if (connectableType == ConnectableType.PROCESSOR) {
+            return AffectedComponentDTO.COMPONENT_TYPE_PROCESSOR;
+        }
+        if (connectableType == ConnectableType.INPUT_PORT) {
+            return AffectedComponentDTO.COMPONENT_TYPE_INPUT_PORT;
+        }
+        if (connectableType == ConnectableType.OUTPUT_PORT) {
+            return AffectedComponentDTO.COMPONENT_TYPE_OUTPUT_PORT;
+        }
+
+        return null;
+    }
+
+    private String getState(final 
RemovedConnectionDrainClassifier.LiveConnectable liveConnectable) {
+        if (liveConnectable.type() == ConnectableType.PROCESSOR) {
+            return liveConnectable.physicalScheduledState() == null ? null : 
liveConnectable.physicalScheduledState().name();
+        }
+
+        return liveConnectable.running() ? "RUNNING" : "STOPPED";
+    }
+
+    private LifecycleManagementException decorateFailure(final 
LifecycleManagementException failure,
+                                                         final 
Set<AffectedComponentEntity> componentsToStop,
+                                                         final Set<String> 
candidateConnectionIds) {
+        final String message = failure.getMessage();
+        if (message != null && message.startsWith("Removed connection drain")) 
{
+            return failure;
+        }
+
+        final String decoratedMessage;
+        if (message != null && message.contains("waiting for connection 
queues")) {
+            decoratedMessage = "Removed connection drain failed while waiting 
for connections "
+                    + candidateConnectionIds.stream().sorted().toList() + ": " 
+ message;
+        } else {
+            final List<String> componentIds = 
componentsToStop.stream().map(AffectedComponentEntity::getId).sorted().toList();
+            decoratedMessage = "Removed connection drain failed while stopping 
producer barriers " + componentIds + ": " + message;
+        }
+
+        return new LifecycleManagementException(decoratedMessage, failure);
+    }
+
+    private void restoreOrSuppress(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                   final Set<AffectedComponentEntity> 
stoppedComponents, final LifecycleManagementException failure) {
+        try {
+            restoreStoppedComponents(componentLifecycle, requestUri, groupId, 
stoppedComponents);
+        } catch (final LifecycleManagementException restorationFailure) {
+            logger.warn("Failed to restore producer barriers {} after removed 
connection drain failure",
+                    
stoppedComponents.stream().map(AffectedComponentEntity::getId).sorted().toList(),
 restorationFailure);
+            failure.addSuppressed(restorationFailure);
+        }
+    }
+
+    private void restoreStoppedComponents(final ComponentLifecycle 
componentLifecycle, final URI requestUri, final String groupId,
+                                          final Set<AffectedComponentEntity> 
stoppedComponents) throws LifecycleManagementException {
+        if (stoppedComponents.isEmpty()) {
+            return;
+        }
+
+        componentLifecycle.scheduleComponents(requestUri, groupId, 
stoppedComponents, ScheduledState.RUNNING,
+                pauseFactory.createRestorationPause(), 
InvalidComponentAction.SKIP);
+    }
+
+    private RemovedConnectionDrainClassifier.Context 
createQueueAwareContext(final FlowUpdateImpact flowUpdateImpact,
+                                                                             
final RemovedConnectionDrainClassifier.Context context,
+                                                                             
final ComponentLifecycle componentLifecycle,
+                                                                             
final URI requestUri) throws LifecycleManagementException {
+        final Map<String, Boolean> knownQueueEmptyByConnectionId = new 
LinkedHashMap<>();
+        final Pause noWaitPause = NoWaitPause.INSTANCE;
+        for (final RemovedConnectionDescriptor removedConnection : 
flowUpdateImpact.getRemovedConnections()) {
+            final String connectionId = 
removedConnection.getConnectionInstanceId();
+            if (connectionId == null) {
+                continue;
+            }
+
+            final boolean knownQueueEmpty = 
componentLifecycle.waitForConnectionQueuesEmpty(requestUri, 
Set.of(connectionId), noWaitPause);
+            knownQueueEmptyByConnectionId.put(connectionId, knownQueueEmpty);
+        }
+
+        return new QueueAwareContext(context, knownQueueEmptyByConnectionId);
+    }
+
+    private boolean allComponentsStopped(final Set<AffectedComponentEntity> 
componentsToStop, final Set<AffectedComponentEntity> updatedStoppedComponents) {
+        if (componentsToStop.isEmpty()) {
+            return true;
+        }
+
+        final Map<String, AffectedComponentEntity> updatedComponentsById = 
toOrderedMap(updatedStoppedComponents);
+        for (final AffectedComponentEntity componentToStop : componentsToStop) 
{
+            final AffectedComponentEntity updatedComponent = 
updatedComponentsById.getOrDefault(componentToStop.getId(), componentToStop);
+            if (updatedComponent.getComponent() == null || 
isActive(updatedComponent.getComponent())) {
+                return false;
+            }
+        }
+
+        return true;
+    }
+
+    private Set<AffectedComponentEntity> getStoppedComponents(final 
Collection<AffectedComponentEntity> componentsToStop,
+                                                              final 
Collection<AffectedComponentEntity> updatedComponents) {
+        if (componentsToStop == null || componentsToStop.isEmpty()) {
+            return Collections.emptySet();
+        }
+
+        final Map<String, AffectedComponentEntity> updatedComponentsById = 
toOrderedMap(updatedComponents);
+        final Set<AffectedComponentEntity> stoppedComponents = new 
LinkedHashSet<>();
+        for (final AffectedComponentEntity componentToStop : componentsToStop) 
{
+            final AffectedComponentEntity updatedComponent = 
updatedComponentsById.getOrDefault(componentToStop.getId(), componentToStop);
+            if (updatedComponent.getComponent() != null && 
!isActive(updatedComponent.getComponent())) {
+                stoppedComponents.add(updatedComponent);
+            }
+        }
+
+        return stoppedComponents;
+    }
+
+    private Set<AffectedComponentEntity> getStoppedComponentsToRestore(final 
RemovedConnectionDrainClassifier.Context context,
+                                                                       final 
Collection<AffectedComponentEntity> componentsToStop,
+                                                                       final 
Collection<AffectedComponentEntity> updatedComponents) {
+        if (componentsToStop == null || componentsToStop.isEmpty()) {
+            return Collections.emptySet();
+        }
+
+        final Map<String, AffectedComponentEntity> updatedComponentsById = 
toOrderedMap(updatedComponents);
+        final Set<AffectedComponentEntity> stoppedComponents = new 
LinkedHashSet<>();
+        for (final AffectedComponentEntity componentToStop : componentsToStop) 
{
+            final AffectedComponentEntity updatedComponent = 
updatedComponentsById.get(componentToStop.getId());
+            if (updatedComponent != null) {
+                if (updatedComponent.getComponent() != null && 
!isActive(updatedComponent.getComponent())) {
+                    stoppedComponents.add(updatedComponent);
+                }
+                continue;
+            }
+
+            final RemovedConnectionDrainClassifier.LiveConnectable 
liveConnectable = context.getConnectable(componentToStop.getId());

Review Comment:
   [GPT-5.6 Sol] **This cannot identify a partial stop across cluster nodes.** 
If the replicated stop succeeds on node A but fails on the coordinator node B, 
`scheduleComponents` throws before returning updated entities. This lookup sees 
only node B, finds the producer still running, and omits it from restoration, 
leaving it stopped on node A. Every component in `componentsToStop` was active 
before this operation, so failure recovery needs to restore conservatively 
across the cluster instead of relying on one node's current state.



##########
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()));

Review Comment:
   [GPT-5.6 Sol] The main system test should verify the feature's key lifecycle 
behavior while this state is visible: the source must be stopped and the 
destination must remain running, on each node in a cluster. `GenerateFlowFile` 
is limited to one FlowFile, so the queue can drain successfully after opening 
the gate even if producer stopping is broken. The current final-state 
assertions would not catch that.



##########
nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/registry/RegistryClientIT.java:
##########
@@ -739,4 +891,82 @@ public void 
testStopVersionControlThenSetVersionControlInfo() throws NiFiClientE
         assertEquals(vci.getVersionControlInformation().getVersion(), 
groupAfterSetVersionInfo.getComponent().getVersionControlInformation().getVersion());
         assertEquals("UP_TO_DATE", 
groupAfterSetVersionInfo.getComponent().getVersionControlInformation().getState());
     }
+
+    private RemovedConnectionFixture createRemovedConnectionFixture(final Path 
gateFile, final boolean removeDestination,
+                                                                     final 
boolean queueFlowFile) throws Exception {
+        Files.deleteIfExists(gateFile);
+        final FlowRegistryClientEntity clientEntity = registerClient();
+        final NiFiClientUtil util = getClientUtil();
+        final ProcessGroupEntity group = util.createProcessGroup("Removed 
Connection Drain", "root");
+        final ProcessorEntity generate = 
util.createProcessor("GenerateFlowFile", group.getId());
+        util.updateProcessorProperties(generate, Map.of("Max FlowFiles", "1", 
"State Scope", "CLUSTER"));
+        if (getNumberOfNodes() > 1) {
+            util.updateProcessorExecutionNode(generate, ExecutionNode.PRIMARY);
+        }
+
+        final ProcessorEntity gated = util.createProcessor("GatedPassThrough", 
group.getId());

Review Comment:
   [GPT-5.6 Sol] The existing `TerminateFlowFile` test Processor already has an 
optional `Gate File` property and waits without consuming until that file 
exists. It appears this fixture can use that Processor directly, which would 
remove `GatedPassThrough` and its service registration entirely. Please 
simplify this unless pass-through behavior is required for an assertion that I 
missed.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/NiFiServiceFacade.java:
##########
@@ -2079,6 +2079,18 @@ VersionControlInformationEntity 
setVersionControlInformation(Revision processGro
      */
     String getFlowRegistryName(String flowRegistryId);
 
+    /**
+     * Determines which components currently exist in the Process Group with 
the given identifier and calculates which of those components
+     * would be impacted by updating the Process Group to the provided snapshot
+     *
+     * @param processGroupId the ID of the Process Group to update
+     * @param updatedSnapshot the snapshot to update the Process Group to
+     * @return the impact of updating the Process Group
+     */
+    FlowUpdateImpact getFlowUpdateImpact(String processGroupId, 
RegisteredFlowSnapshot updatedSnapshot);
+
+    RemovedConnectionDrainClassifier.Context 
getRemovedConnectionDrainContext();

Review Comment:
   [GPT-5.6 Sol] Please add JavaDoc for this new service-facade method, 
including what runtime view the returned context provides and whether callers 
can expect it to reflect later component-state changes. The restoration code 
depends on that behavior.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/util/ComponentLifecycle.java:
##########
@@ -59,4 +59,6 @@ Set<AffectedComponentEntity> scheduleComponents(URI 
exampleUri, String groupId,
      */
     Set<AffectedComponentEntity> activateControllerServices(URI exampleUri, 
String groupId, Set<AffectedComponentEntity> servicesToUpdate, 
Set<AffectedComponentEntity> servicesRequiringDesiredState,
         ControllerServiceState desiredState, Pause pause, 
InvalidComponentAction invalidComponentAction) throws 
LifecycleManagementException;
+
+    boolean waitForConnectionQueuesEmpty(URI exampleUri, Set<String> 
connectionIds, Pause pause) throws LifecycleManagementException;

Review Comment:
   [GPT-5.6 Sol] Please document this new interface method. Its contract needs 
to state what `true` means in a cluster, how the `Pause` controls waiting and 
cancellation, and which failures are reported through 
`LifecycleManagementException`. Those details are central to safe connection 
removal.



##########
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:
   [GPT-5.6 Sol] This test waits for the full 30-second production timeout and 
is inherited by both standalone and clustered suites. Please provide a 
test-only way to use a short drain timeout while retaining the end-to-end 
restoration check.



##########
nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java:
##########
@@ -462,6 +478,143 @@ public void 
testGetComponentsAffectedByFlowUpdate_WithNewStatelessProcessGroup_R
         assertTrue(affected.isEmpty(), "No local components should be affected 
for added Stateless group");
     }
 
+    @Test
+    public void 
testGetFlowUpdateImpactExtractsRemovedConnectionMetadataAndAffectedProjection() 
{
+        final String rootGroupId = "root-group-instance";
+        final String removedGroupId = "removed-group-instance";
+
+        final ProcessGroup rootGroup = mock(ProcessGroup.class);
+        
when(processGroupDAO.getProcessGroup(rootGroupId)).thenReturn(rootGroup);
+
+        final ProcessGroup rootGroupMetadata = mock(ProcessGroup.class);
+        when(rootGroupMetadata.getIdentifier()).thenReturn(rootGroupId);
+        when(rootGroupMetadata.getName()).thenReturn("Root Group");
+
+        final ProcessGroup removedGroup = mock(ProcessGroup.class);
+        when(removedGroup.getIdentifier()).thenReturn(removedGroupId);
+        
when(processGroupDAO.getProcessGroup(removedGroupId)).thenReturn(removedGroup);
+        when(removedGroup.findAllProcessors()).thenReturn(List.of());
+        when(removedGroup.findAllFunnels()).thenReturn(List.of());
+        when(removedGroup.findAllInputPorts()).thenReturn(List.of());
+        when(removedGroup.findAllOutputPorts()).thenReturn(List.of());
+        when(removedGroup.findAllRemoteProcessGroups()).thenReturn(List.of());
+        when(removedGroup.findAllControllerServices()).thenReturn(Set.of());
+
+        final ProcessorNode removedSource = createLocalConnectableProcessor(
+                "removed-source-instance", "removed-source-versioned", 
ConnectableType.PROCESSOR, rootGroupMetadata);
+        final Port removedDestination = createLocalConnectablePort(
+                "removed-destination-instance", 
"removed-destination-versioned", ConnectableType.OUTPUT_PORT, 
rootGroupMetadata);
+        final Port changedSource = createLocalConnectablePort(
+                "changed-source-instance", "changed-source-versioned", 
ConnectableType.INPUT_PORT, rootGroupMetadata);
+        final ProcessorNode changedDestination = 
createLocalConnectableProcessor(
+                "changed-destination-instance", 
"changed-destination-versioned", ConnectableType.PROCESSOR, rootGroupMetadata);
+
+        when(rootGroup.findAllProcessors()).thenReturn(List.of(removedSource, 
changedDestination));
+        when(rootGroup.findAllFunnels()).thenReturn(List.of());
+        when(rootGroup.findAllInputPorts()).thenReturn(List.of(changedSource));
+        
when(rootGroup.findAllOutputPorts()).thenReturn(List.of(removedDestination));
+        when(rootGroup.findAllRemoteProcessGroups()).thenReturn(List.of());
+
+        stubLocalConnectableAuthorizable(removedSource.getIdentifier());
+        stubLocalConnectableAuthorizable(removedDestination.getIdentifier());
+        stubLocalConnectableAuthorizable(changedSource.getIdentifier());
+        stubLocalConnectableAuthorizable(changedDestination.getIdentifier());
+        stubConnectionAuthorizable("removed-connection-instance");
+        stubConnectionAuthorizable("changed-connection-instance");
+        stubProcessGroupAuthorizable(removedGroupId, removedGroup);
+        stubInputPortAuthorizable("removed-endpoint-instance");
+
+        final FlowComparison comparison = mock(FlowComparison.class);
+        when(comparison.getDifferences()).thenReturn(Set.of(
+                new StandardFlowDifference(DifferenceType.COMPONENT_REMOVED, 
createRemovedConnection(), null, null, null, "Removed connection"),
+                new StandardFlowDifference(DifferenceType.SOURCE_CHANGED, 
createSourceChangedConnection(), createReplacementConnection(), null, null, 
"Source changed connection"),
+                new StandardFlowDifference(DifferenceType.COMPONENT_REMOVED, 
createRemovedGroup(removedGroupId), null, null, null, "Removed group"),
+                new StandardFlowDifference(DifferenceType.COMPONENT_REMOVED, 
createRemovedEndpoint(removedGroupId), null, null, null, "Removed endpoint")
+        ));
+
+        final StandardNiFiServiceFacade serviceFacadeSpy = spy(serviceFacade);
+        
doReturn(comparison).when(serviceFacadeSpy).compareFlowUpdate(eq(rootGroup), 
any(RegisteredFlowSnapshot.class));
+
+        final RegisteredFlowSnapshot updatedSnapshot = new 
RegisteredFlowSnapshot();
+        final FlowUpdateImpact impact = 
serviceFacadeSpy.getFlowUpdateImpact(rootGroupId, updatedSnapshot);
+
+        assertNotNull(impact);
+        assertEquals(Set.of(removedGroupId), 
impact.getRemovedProcessGroupIds());
+        assertEquals(Set.of("removed-endpoint-instance"), 
impact.getRemovedEndpointIds());
+
+        final Map<String, RemovedConnectionDescriptor> 
removedConnectionByInstanceId = impact.getRemovedConnections().stream()
+                
.collect(Collectors.toMap(RemovedConnectionDescriptor::getConnectionInstanceId, 
Function.identity()));
+        assertEquals(Set.of("removed-connection-instance", 
"changed-connection-instance"), removedConnectionByInstanceId.keySet());
+
+        final RemovedConnectionDescriptor pureRemoval = 
removedConnectionByInstanceId.get("removed-connection-instance");
+        assertEquals(RemovalReason.COMPONENT_REMOVED, 
pureRemoval.getRemovalReason());
+        assertEquals("removed-connection-versioned", 
pureRemoval.getConnectionVersionedId());
+        assertEquals("root-group-instance", 
pureRemoval.getContainingProcessGroupId());
+        assertEquals("removed-source-instance", 
pureRemoval.getSourceInstanceId());
+        assertEquals("removed-source-versioned", 
pureRemoval.getSourceVersionedId());
+        assertEquals("removed-source-group-instance", 
pureRemoval.getSourceProcessGroupId());
+        assertEquals(ConnectableType.PROCESSOR, pureRemoval.getSourceType());
+        assertEquals("removed-destination-instance", 
pureRemoval.getDestinationInstanceId());
+        assertEquals("removed-destination-versioned", 
pureRemoval.getDestinationVersionedId());
+        assertEquals("removed-destination-group-instance", 
pureRemoval.getDestinationProcessGroupId());
+        assertEquals(ConnectableType.OUTPUT_PORT, 
pureRemoval.getDestinationType());
+
+        final RemovedConnectionDescriptor sourceChanged = 
removedConnectionByInstanceId.get("changed-connection-instance");
+        assertEquals(RemovalReason.SOURCE_CHANGED, 
sourceChanged.getRemovalReason());
+        assertEquals("changed-source-instance", 
sourceChanged.getSourceInstanceId());
+        assertEquals("changed-destination-instance", 
sourceChanged.getDestinationInstanceId());
+
+        final Set<String> affectedIdsFromImpact = 
impact.getAffectedComponents().stream()
+                .map(AffectedComponentEntity::getId)
+                .collect(Collectors.toCollection(LinkedHashSet::new));
+        assertTrue(affectedIdsFromImpact.contains("removed-source-instance"));

Review Comment:
   [GPT-5.6 Sol] These `contains` checks allow unexpected affected components, 
and the later projection comparison calls a method that delegates to the same 
implementation. Please assert the exact expected identifier set and remove the 
redundant self-comparison so the test catches accidental additions.



-- 
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