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]