From 751a14012abbe2bd1a83203a19a1639d1a62d87a Mon Sep 17 00:00:00 2001 From: Mark Payne Date: Mon, 14 Sep 2026 16:00:22 -0400 Subject: [PATCH 1/3] NIFI-16283 Support changing Connector versions Co-authored-by: Cursor --- .../src/main/asciidoc/toolkit-guide.adoc | 3 +- .../components/connector/ConnectorNode.java | 8 + .../connector/FrameworkFlowContext.java | 6 + .../MutableConnectorConfigurationContext.java | 7 + .../nifi/controller/ReloadComponent.java | 12 + .../ConnectorInstantiationException.java | 30 +++ ...StandardConnectorConfigurationContext.java | 11 + .../connector/StandardConnectorNode.java | 147 +++++++++++- .../connector/StandardFlowContext.java | 12 +- .../controller/StandardReloadComponent.java | 69 ++++++ .../connector/TestStandardConnectorNode.java | 56 +++++ .../TestStandardReloadComponent.java | 172 ++++++++++++++ .../apache/nifi/audit/ConnectorAuditor.java | 80 +++++++ .../nifi/web/StandardNiFiServiceFacade.java | 10 +- .../nifi/web/api/ConnectorResource.java | 2 +- .../apache/nifi/web/api/dto/DtoFactory.java | 14 +- .../web/dao/impl/StandardConnectorDAO.java | 32 +++ .../nifi/audit/TestConnectorAuditor.java | 144 ++++++++++++ .../web/StandardNiFiServiceFacadeTest.java | 70 ++++++ .../nifi/web/api/dto/DtoFactoryTest.java | 50 +++++ .../dao/impl/StandardConnectorDAOTest.java | 62 ++++++ .../connectors/service/connector.service.ts | 11 +- .../connectors-listing.actions.ts | 22 ++ .../connectors-listing.effects.spec.ts | 172 +++++++++++++- .../connectors-listing.effects.ts | 87 +++++++- .../connectors-listing.reducer.ts | 13 +- .../src/app/pages/connectors/state/index.ts | 12 + .../connector-table.component.html | 9 + .../connector-table.component.spec.ts | 34 ++- .../connector-table.component.ts | 20 ++ .../connectors-listing.component.html | 1 + .../connectors-listing.component.ts | 15 ++ .../app/service/extension-types.service.ts | 9 + .../frontend/libs/shared/src/types/index.ts | 3 +- .../utils/connector-permissions.utils.spec.ts | 4 + .../src/utils/connector-permissions.utils.ts | 1 + .../engine/StatelessReloadComponent.java | 7 + .../pom.xml | 66 ++++++ .../nifi-connector-change-version-v1/pom.xml | 46 ++++ .../ChangeVersionTestConnector.java | 101 +++++++++ ...apache.nifi.components.connector.Connector | 15 ++ .../pom.xml | 66 ++++++ .../nifi-connector-change-version-v2/pom.xml | 46 ++++ .../ChangeVersionTestConnector.java | 111 +++++++++ ...apache.nifi.components.connector.Connector | 15 ++ .../pom.xml | 4 + .../nifi-system-test-suite/pom.xml | 44 ++++ .../connectors/ConnectorChangeVersionIT.java | 151 +++++++++++++ .../impl/command/nifi/NiFiCommandGroup.java | 2 + .../connectors/ChangeVersionConnector.java | 199 +++++++++++++++++ .../impl/result/nifi/ConnectorsResult.java | 80 +++++++ .../TestChangeVersionConnector.java | 210 ++++++++++++++++++ 52 files changed, 2545 insertions(+), 28 deletions(-) create mode 100644 nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/exception/ConnectorInstantiationException.java create mode 100644 nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardReloadComponent.java create mode 100644 nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/audit/TestConnectorAuditor.java create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v1-nar/pom.xml create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v1/pom.xml create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v1/src/main/java/org/apache/nifi/connectors/tests/system/changeversion/ChangeVersionTestConnector.java create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v1/src/main/resources/META-INF/services/org.apache.nifi.components.connector.Connector create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v2-nar/pom.xml create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v2/pom.xml create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v2/src/main/java/org/apache/nifi/connectors/tests/system/changeversion/ChangeVersionTestConnector.java create mode 100644 nifi-system-tests/nifi-system-test-nar-provider-bundles/nifi-connector-change-version-v2/src/main/resources/META-INF/services/org.apache.nifi.components.connector.Connector create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/connectors/ConnectorChangeVersionIT.java create mode 100644 nifi-toolkit/nifi-toolkit-cli/src/main/java/org/apache/nifi/toolkit/cli/impl/command/nifi/connectors/ChangeVersionConnector.java create mode 100644 nifi-toolkit/nifi-toolkit-cli/src/main/java/org/apache/nifi/toolkit/cli/impl/result/nifi/ConnectorsResult.java create mode 100644 nifi-toolkit/nifi-toolkit-cli/src/test/java/org/apache/nifi/toolkit/cli/impl/command/nifi/connectors/TestChangeVersionConnector.java diff --git a/nifi-docs/src/main/asciidoc/toolkit-guide.adoc b/nifi-docs/src/main/asciidoc/toolkit-guide.adoc index 23095494171b..581591119e66 100644 --- a/nifi-docs/src/main/asciidoc/toolkit-guide.adoc +++ b/nifi-docs/src/main/asciidoc/toolkit-guide.adoc @@ -142,6 +142,7 @@ The following are available commands: nifi get-controller-configuration nifi update-controller-configuration nifi change-version-processor + nifi change-version-connector nifi upload-nar nifi download-nar nifi delete-nar @@ -313,7 +314,7 @@ For example, typing tab at an empty prompt should display possible commands for Typing "nifi " and then a tab will show the sub-commands for NiFi: #> nifi - change-version-processor delete-flow-analysis-rule export-reporting-task get-policy list-user-groups pg-export pg-delete + change-version-connector change-version-processor delete-flow-analysis-rule export-reporting-task get-policy list-user-groups pg-export pg-delete cluster-summary delete-node export-reporting-tasks get-reg-client-id list-users pg-get-all-versions pg-stop-version-control connect-node delete-param fetch-params get-reporting-task logout-access-token pg-get-param-context set-inherited-param-contexts create-flow-analysis-rule delete-param-context get-access-token get-reporting-tasks merge-param-context pg-get-services remove-inherited-param-contexts diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorNode.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorNode.java index 46c10627b7e7..11a5e6759d2b 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorNode.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorNode.java @@ -58,6 +58,12 @@ public interface ConnectorNode extends ComponentAuthorizable, VersionedComponent void verifyCanStart(); + void verifyCanReload(); + + void verifyCanUpdateBundle(BundleCoordinate bundleCoordinate); + + void replaceConnector(Connector connector, BundleCoordinate bundleCoordinate, ComponentLog componentLog) throws FlowUpdateException; + Connector getConnector(); /** @@ -99,6 +105,8 @@ default void setCustomLoggingAttributes(final Map attributes) { */ boolean isExtensionMissing(); + void resetValidationState(); + List fetchAllowableValues(String stepName, String propertyName); List fetchAllowableValues(String stepName, String propertyName, String filter); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/FrameworkFlowContext.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/FrameworkFlowContext.java index 0b1862cdffd8..545e478a1b69 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/FrameworkFlowContext.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/FrameworkFlowContext.java @@ -19,9 +19,11 @@ import org.apache.nifi.asset.AssetManager; import org.apache.nifi.components.connector.components.FlowContext; +import org.apache.nifi.flow.Bundle; import org.apache.nifi.flow.VersionedExternalFlow; import org.apache.nifi.flow.VersionedProcessGroup; import org.apache.nifi.groups.ProcessGroup; +import org.apache.nifi.logging.ComponentLog; public interface FrameworkFlowContext extends FlowContext { ProcessGroup getManagedProcessGroup(); @@ -53,4 +55,8 @@ public interface FrameworkFlowContext extends FlowContext { * persisted while the Connector was in Troubleshooting mode */ void restoreTroubleshootingFlow(VersionedProcessGroup troubleshootingProcessGroup); + + default void reload(final Bundle bundle, final ComponentLog connectorLog) { + throw new UnsupportedOperationException(); + } } diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/MutableConnectorConfigurationContext.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/MutableConnectorConfigurationContext.java index ed7512e1c830..b902de7de840 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/MutableConnectorConfigurationContext.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/MutableConnectorConfigurationContext.java @@ -49,6 +49,13 @@ public interface MutableConnectorConfigurationContext extends ConnectorConfigura */ ConfigurationUpdateResult replaceProperties(String stepName, StepConfiguration configuration); + /** + * Removes the named configuration step and its resolved values. If the step is not present, this is a no-op. + * + * @param stepName the name of the configuration step to remove + */ + void removeStep(String stepName); + /** * Resolves all existing property values. */ diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/ReloadComponent.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/ReloadComponent.java index 9d68c5c022c2..1cdc347742d8 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/ReloadComponent.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/ReloadComponent.java @@ -17,6 +17,8 @@ package org.apache.nifi.controller; import org.apache.nifi.bundle.BundleCoordinate; +import org.apache.nifi.components.connector.ConnectorNode; +import org.apache.nifi.controller.exception.ConnectorInstantiationException; import org.apache.nifi.controller.exception.ControllerServiceInstantiationException; import org.apache.nifi.controller.exception.ProcessorInstantiationException; import org.apache.nifi.controller.flowanalysis.FlowAnalysisRuleInstantiationException; @@ -104,4 +106,14 @@ void reload(ParameterProviderNode existingNode, String newType, BundleCoordinate */ void reload(FlowRegistryClientNode existingNode, String newType, BundleCoordinate bundleCoordinate, Set additionalUrls) throws FlowRepositoryClientInstantiationException; + + /** + * Changes the underlying Connector held by the node to an instance of the new type. + * + * @param existingNode the ConnectorNode being updated + * @param newType the fully qualified class name of the new type + * @param bundleCoordinate the bundle coordinate of the new type + * @throws ConnectorInstantiationException if unable to create an instance of the new type + */ + void reload(ConnectorNode existingNode, String newType, BundleCoordinate bundleCoordinate) throws ConnectorInstantiationException; } diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/exception/ConnectorInstantiationException.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/exception/ConnectorInstantiationException.java new file mode 100644 index 000000000000..af17acf3499b --- /dev/null +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/controller/exception/ConnectorInstantiationException.java @@ -0,0 +1,30 @@ +/* + * 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.controller.exception; + +public class ConnectorInstantiationException extends Exception { + + private static final long serialVersionUID = 1L; + + public ConnectorInstantiationException(final String message) { + super(message); + } + + public ConnectorInstantiationException(final String message, final Throwable cause) { + super(message, cause); + } +} diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorConfigurationContext.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorConfigurationContext.java index 59b5bfc7f142..34453ca8ba7d 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorConfigurationContext.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorConfigurationContext.java @@ -258,6 +258,17 @@ public ConfigurationUpdateResult replaceProperties(final String stepName, final } } + @Override + public void removeStep(final String stepName) { + writeLock.lock(); + try { + propertyConfigurations.remove(stepName); + resolvedPropertyConfigurations.remove(stepName); + } finally { + writeLock.unlock(); + } + } + @Override public void resolvePropertyValues() { writeLock.lock(); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java index 5daa9ae0fdd4..71acd6b95595 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java @@ -105,10 +105,10 @@ public class StandardConnectorNode implements ConnectorNode, GroupedComponent { private final ExtensionManager extensionManager; private final StateManagerProvider stateManagerProvider; private final Authorizable parentAuthorizable; - private final ConnectorDetails connectorDetails; - private final String componentType; + private volatile ConnectorDetails connectorDetails; + private volatile String componentType; private final String componentCanonicalClass; - private final BundleCoordinate bundleCoordinate; + private volatile BundleCoordinate bundleCoordinate; private final ConnectorStateTransition stateTransition; private final AtomicReference versionedComponentId = new AtomicReference<>(); private final FlowContextFactory flowContextFactory; @@ -116,7 +116,7 @@ public class StandardConnectorNode implements ConnectorNode, GroupedComponent { private final AtomicReference validationState = new AtomicReference<>(new ValidationState(ValidationStatus.VALIDATING, Collections.emptyList())); private final ConnectorValidationTrigger validationTrigger; - private final boolean extensionMissing; + private volatile boolean extensionMissing; private volatile boolean triggerValidation = true; private final AtomicReference> drainFutureRef = new AtomicReference<>(); private volatile ValidationResult unresolvedBundleValidationResult = null; @@ -485,11 +485,15 @@ private Map migrateProperties(final List migrateProperties(final Connector connector, final Map initial) { final Set persistedStepNames = new LinkedHashSet<>(initial.keySet()); final StandardConnectorPropertyConfiguration propertyConfiguration = new StandardConnectorPropertyConfiguration(initial, this.toString()); - try (final NarCloseable ignored = NarCloseable.withComponentNarLoader(extensionManager, getConnector().getClass(), getIdentifier())) { - getConnector().migrateProperties(propertyConfiguration); - return applyMissingRequiredPropertyDefaults(propertyConfiguration.getMutatedProperties(), persistedStepNames, getConnector().getConfigurationSteps()); + try (final NarCloseable ignored = NarCloseable.withComponentNarLoader(connector.getClass().getClassLoader())) { + connector.migrateProperties(propertyConfiguration); + return applyMissingRequiredPropertyDefaults(propertyConfiguration.getMutatedProperties(), persistedStepNames, connector.getConfigurationSteps()); } } @@ -733,7 +737,11 @@ public void replaceWorkingConfiguration(final String stepName, final StepConfigu private void notifyStepConfigured(final String stepName, final FrameworkFlowContext workingContext) throws FlowUpdateException { final Connector connector = connectorDetails.getConnector(); - try (final NarCloseable ignored = NarCloseable.withComponentNarLoader(extensionManager, connector.getClass(), getIdentifier())) { + notifyStepConfigured(connector, stepName, workingContext); + } + + private void notifyStepConfigured(final Connector connector, final String stepName, final FrameworkFlowContext workingContext) throws FlowUpdateException { + try (final NarCloseable ignored = NarCloseable.withComponentNarLoader(connector.getClass().getClassLoader())) { logger.debug("Notifying {} of configuration change for configuration step {}", this, stepName); connector.onConfigurationStepConfigured(stepName, workingContext); @@ -1310,6 +1318,126 @@ public Connector getConnector() { return connectorDetails.getConnector(); } + @Override + public void verifyCanReload() { + if (!isStopped()) { + throw new IllegalStateException("Cannot reload " + this + " because its state is " + getCurrentState()); + } + } + + @Override + public void verifyCanUpdateBundle(final BundleCoordinate incomingCoordinate) { + final BundleCoordinate currentCoordinate = getBundleCoordinate(); + if (!currentCoordinate.getGroup().equals(incomingCoordinate.getGroup()) || !currentCoordinate.getId().equals(incomingCoordinate.getId())) { + throw new IllegalArgumentException("Cannot update " + this + " from " + currentCoordinate.getCoordinate() + " to " + + incomingCoordinate.getCoordinate() + " because the bundle group and artifact must be unchanged"); + } + } + + @Override + public void replaceConnector(final Connector replacement, final BundleCoordinate replacementCoordinate, final ComponentLog replacementLog) throws FlowUpdateException { + verifyCanReload(); + + final FrameworkConnectorInitializationContext replacementInitializationContext = new StandardConnectorInitializationContext.Builder() + .identifier(identifier) + .name(name) + .componentLog(replacementLog) + .secretsManager(initializationContext.getSecretsManager()) + .assetManager(initializationContext.getAssetManager()) + .componentBundleLookup(initializationContext.getComponentBundleLookup()) + .build(); + + try (final NarCloseable ignored = NarCloseable.withComponentNarLoader(replacement.getClass().getClassLoader())) { + replacement.initialize(replacementInitializationContext); + } + + final WorkingFlowContextState workingContextState = acquireWorkingFlowContext(); + final FrameworkFlowContext workingContext = workingContextState.getContext(); + boolean workingContextReleased = false; + try { + final Map originalActiveConfiguration = toConfigurationMap(activeFlowContext.getConfigurationContext().toConnectorConfiguration()); + final Map originalWorkingConfiguration = workingContext == null ? originalActiveConfiguration + : toConfigurationMap(workingContext.getConfigurationContext().toConnectorConfiguration()); + + final Map activeConfiguration = migrateProperties(replacement, originalActiveConfiguration); + final Map workingConfiguration = migrateProperties(replacement, originalWorkingConfiguration); + final Bundle bundle = new Bundle(replacementCoordinate.getGroup(), replacementCoordinate.getId(), replacementCoordinate.getVersion()); + + replaceConfiguration(activeFlowContext.getConfigurationContext(), activeConfiguration); + if (workingContext != null) { + replaceConfiguration(workingContext.getConfigurationContext(), workingConfiguration); + workingContext.reload(bundle, replacementLog); + + try { + for (final String stepName : workingConfiguration.keySet()) { + notifyStepConfigured(replacement, stepName, workingContext); + } + } catch (final FlowUpdateException | RuntimeException e) { + replaceConfiguration(activeFlowContext.getConfigurationContext(), originalActiveConfiguration); + final Bundle originalBundle = new Bundle(bundleCoordinate.getGroup(), bundleCoordinate.getId(), bundleCoordinate.getVersion()); + releaseWorkingFlowContext(workingContextState); + workingContextReleased = true; + + final MutableConnectorConfigurationContext restoredConfiguration = createConfigurationContext(originalWorkingConfiguration); + final WorkingFlowContextState restoredContextState = installReplacementWorkingFlowContext(restoredConfiguration, originalBundle, true); + final FrameworkFlowContext restoredContext = restoredContextState.getContext(); + + try { + for (final String stepName : originalWorkingConfiguration.keySet()) { + try { + notifyStepConfigured(stepName, restoredContext); + } catch (final Exception refreshException) { + logger.warn("Failed to restore configuration for step [{}] of {}", stepName, this, refreshException); + } + } + } finally { + releaseWorkingFlowContext(restoredContextState); + } + throw e; + } + } + + connectorDetails = new ConnectorDetails(replacement, replacementCoordinate, replacementLog); + bundleCoordinate = replacementCoordinate; + extensionMissing = replacement instanceof GhostConnector; + componentType = extensionMissing ? "(Missing) " + getSimpleClassName(componentCanonicalClass) : replacement.getClass().getSimpleName(); + initializationContext = replacementInitializationContext; + + activeFlowContext.reload(bundle, replacementLog); + } finally { + if (!workingContextReleased) { + releaseWorkingFlowContext(workingContextState); + } + } + + rebuildLoggingAttributes(); + } + + private Map toConfigurationMap(final ConnectorConfiguration configuration) { + final Map configurations = new LinkedHashMap<>(); + for (final NamedStepConfiguration namedConfiguration : configuration.getNamedStepConfigurations()) { + configurations.put(namedConfiguration.stepName(), namedConfiguration.configuration()); + } + return configurations; + } + + private void replaceConfiguration(final MutableConnectorConfigurationContext context, final Map replacement) { + final Set existingStepNames = new HashSet<>(); + for (final NamedStepConfiguration namedConfiguration : context.toConnectorConfiguration().getNamedStepConfigurations()) { + existingStepNames.add(namedConfiguration.stepName()); + } + + for (final Map.Entry entry : replacement.entrySet()) { + context.replaceProperties(entry.getKey(), entry.getValue()); + } + + for (final String existingStepName : existingStepNames) { + if (!replacement.containsKey(existingStepName)) { + context.removeStep(existingStepName); + } + } + } + @Override public String getComponentType() { return componentType; @@ -2480,7 +2608,8 @@ public void setVersionedComponentId(final String versionedComponentId) { } } - private void resetValidationState() { + @Override + public void resetValidationState() { validationState.set(new ValidationState(ValidationStatus.VALIDATING, Collections.emptyList())); validationTrigger.triggerAsync(this); logger.debug("Validation state has been reset for {}", this); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardFlowContext.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardFlowContext.java index 667a0b3b8665..820588a1a0c3 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardFlowContext.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardFlowContext.java @@ -36,14 +36,13 @@ public class StandardFlowContext implements FrameworkFlowContext { private final MutableConnectorConfigurationContext configurationContext; private final ProcessGroupFacadeFactory groupFacadeFactory; private final ParameterContextFacadeFactory parameterContextFacadeFactory; - private final ComponentLog connectorLog; private final FlowContextType flowContextType; - private final Bundle bundle; + private volatile ComponentLog connectorLog; + private volatile Bundle bundle; private volatile ProcessGroupFacade rootGroup; private volatile ParameterContextFacade parameterContext; - public StandardFlowContext(final ProcessGroup managedProcessGroup, final MutableConnectorConfigurationContext configurationContext, final ProcessGroupFacadeFactory groupFacadeFactory, final ParameterContextFacadeFactory parameterContextFacadeFactory, final ComponentLog connectorLog, final FlowContextType flowContextType, @@ -140,6 +139,13 @@ public Bundle getBundle() { return bundle; } + @Override + public void reload(final Bundle bundle, final ComponentLog connectorLog) { + this.bundle = bundle; + this.connectorLog = connectorLog; + this.rootGroup = groupFacadeFactory.create(managedProcessGroup, connectorLog); + } + @Override public FlowContextType getType() { return flowContextType; diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardReloadComponent.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardReloadComponent.java index 4ca87b21c05d..01c99a9e8ccc 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardReloadComponent.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/StandardReloadComponent.java @@ -17,9 +17,14 @@ package org.apache.nifi.controller; import org.apache.nifi.annotation.lifecycle.OnRemoved; +import org.apache.nifi.bundle.Bundle; import org.apache.nifi.bundle.BundleCoordinate; import org.apache.nifi.components.AsyncLoadedProcessor; +import org.apache.nifi.components.connector.Connector; +import org.apache.nifi.components.connector.ConnectorNode; +import org.apache.nifi.components.connector.GhostConnector; import org.apache.nifi.components.state.StateManager; +import org.apache.nifi.controller.exception.ConnectorInstantiationException; import org.apache.nifi.controller.exception.ControllerServiceInstantiationException; import org.apache.nifi.controller.flowanalysis.FlowAnalysisRuleInstantiationException; import org.apache.nifi.controller.flowrepository.FlowRepositoryClientInstantiationException; @@ -29,9 +34,11 @@ import org.apache.nifi.controller.service.StandardConfigurationContext; import org.apache.nifi.flowanalysis.FlowAnalysisRule; import org.apache.nifi.logging.ComponentLog; +import org.apache.nifi.logging.GroupedComponent; import org.apache.nifi.logging.LogRepositoryFactory; import org.apache.nifi.logging.StandardLoggingContext; import org.apache.nifi.nar.ExtensionManager; +import org.apache.nifi.nar.InstanceClassLoader; import org.apache.nifi.nar.NarCloseable; import org.apache.nifi.nar.PythonBundle; import org.apache.nifi.parameter.ParameterProvider; @@ -46,6 +53,7 @@ import org.slf4j.LoggerFactory; import java.net.URL; +import java.util.Collections; import java.util.Set; public class StandardReloadComponent implements ReloadComponent { @@ -377,4 +385,65 @@ public void reload( flowController.getValidationTrigger().triggerAsync(existingNode); } + + @Override + public void reload(final ConnectorNode existingNode, final String newType, final BundleCoordinate bundleCoordinate) throws ConnectorInstantiationException { + if (existingNode == null) { + throw new IllegalStateException("Existing ConnectorNode cannot be null"); + } + + existingNode.verifyCanReload(); + existingNode.verifyCanUpdateBundle(bundleCoordinate); + + final String id = existingNode.getIdentifier(); + final ExtensionManager extensionManager = flowController.getExtensionManager(); + final Bundle bundle = extensionManager.getBundle(bundleCoordinate); + if (bundle == null) { + throw new ConnectorInstantiationException("Unable to find bundle " + bundleCoordinate.getCoordinate()); + } + + final InstanceClassLoader candidateClassLoader = extensionManager.createInstanceClassLoader(newType, id, bundle, Collections.emptySet(), false, null); + final Connector connector = createConnector(newType, id, candidateClassLoader, extensionManager); + final boolean extensionMissing = connector instanceof GhostConnector; + + final StandardLoggingContext loggingContext = new StandardLoggingContext(); + if (existingNode instanceof final GroupedComponent groupedComponent) { + loggingContext.setComponent(groupedComponent); + } + + final ComponentLog componentLog = new TerminationAwareLogger(new StandardComponentLog(id, connector, loggingContext)); + try { + existingNode.replaceConnector(connector, bundleCoordinate, componentLog); + } catch (final Exception | LinkageError e) { + if (!extensionMissing) { + extensionManager.closeURLClassLoader(id, candidateClassLoader); + } + + throw new ConnectorInstantiationException("Failed to reload Connector of type " + newType, e); + } + + extensionManager.removeInstanceClassLoader(id); + if (!extensionMissing) { + extensionManager.registerInstanceClassLoader(id, candidateClassLoader); + } + + LogRepositoryFactory.getRepository(id).setLogger(componentLog); + existingNode.resetValidationState(); + logger.info("Reloaded {} using bundle {}", existingNode, bundleCoordinate); + } + + private Connector createConnector(final String type, final String identifier, final InstanceClassLoader instanceClassLoader, final ExtensionManager extensionManager) { + final ClassLoader contextClassLoader = Thread.currentThread().getContextClassLoader(); + try { + final Class connectorClass = Class.forName(type, true, instanceClassLoader); + Thread.currentThread().setContextClassLoader(instanceClassLoader); + return (Connector) connectorClass.getDeclaredConstructor().newInstance(); + } catch (final Exception | LinkageError e) { + extensionManager.closeURLClassLoader(identifier, instanceClassLoader); + final Exception cause = e instanceof final Exception exception ? exception : new ConnectorInstantiationException("Failed to create Connector of type " + type, e); + return new GhostConnector(identifier, type, cause); + } finally { + Thread.currentThread().setContextClassLoader(contextClassLoader); + } + } } diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java index e68ca70c4a85..97b8bd1d37af 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java @@ -77,6 +77,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; @@ -710,6 +712,60 @@ public void testReplaceWorkingConfigurationWaitsForWorkingContextRecreation() th assertEquals(Map.of("propA", new StringLiteralValue("newA")), namedStep.configuration().getPropertyValues()); } + @Test + public void testReplaceConnectorRemovesMigratedConfigurationStep() throws FlowUpdateException { + final StandardConnectorNode connectorNode = createConnectorNode(new TrackingConnector()); + connectorNode.setConfiguration("legacy", createStepConfiguration()); + final FrameworkFlowContext workingFlowContext = connectorNode.getWorkingFlowContext(); + final Connector replacement = new LegacyStepRemovingConnector(); + final BundleCoordinate replacementCoordinate = new BundleCoordinate("org.apache.nifi", "test-standard-connector-node", "2.0.0"); + + connectorNode.replaceConnector(replacement, replacementCoordinate, new MockComponentLog("ReplacementConnector", replacement)); + + assertSame(replacement, connectorNode.getConnector()); + assertSame(workingFlowContext, connectorNode.getWorkingFlowContext()); + assertTrue(connectorNode.getWorkingFlowContext().getConfigurationContext().getPropertyNames("legacy").isEmpty()); + } + + @Test + public void testReplaceConnectorWithGhostMarksExtensionMissing() throws FlowUpdateException { + final StandardConnectorNode connectorNode = createConnectorNode(new TrackingConnector()); + final GhostConnector ghostConnector = new GhostConnector(connectorNode.getIdentifier(), connectorNode.getCanonicalClassName(), new Exception("Instantiation failed")); + final BundleCoordinate replacementCoordinate = new BundleCoordinate("org.apache.nifi", "test-standard-connector-node", "2.0.0"); + + connectorNode.replaceConnector(ghostConnector, replacementCoordinate, new MockComponentLog("GhostConnector", ghostConnector)); + + assertSame(ghostConnector, connectorNode.getConnector()); + assertTrue(connectorNode.isExtensionMissing()); + assertEquals("(Missing) TrackingConnector", connectorNode.getComponentType()); + } + + @Test + public void testReplaceConnectorRestoresConfigurationWhenCallbackFails() throws FlowUpdateException { + final TrackingConnector original = new TrackingConnector(); + final StandardConnectorNode connectorNode = createConnectorNode(original); + seedActiveConfiguration(connectorNode, "active-step", Map.of("active-property", new StringLiteralValue("active-value"))); + connectorNode.setConfiguration("step", createStepConfiguration()); + final FrameworkFlowContext originalWorkingFlowContext = connectorNode.getWorkingFlowContext(); + final FailingStepConnector replacement = new FailingStepConnector("step") { + @Override + public void migrateProperties(final ConnectorPropertyConfiguration configuration) { + configuration.removeStep("active-step"); + } + }; + replacement.setFailOnStep(true); + final BundleCoordinate replacementCoordinate = new BundleCoordinate("org.apache.nifi", "test-standard-connector-node", "2.0.0"); + + assertThrows(FlowUpdateException.class, + () -> connectorNode.replaceConnector(replacement, replacementCoordinate, new MockComponentLog("ReplacementConnector", replacement))); + + assertSame(original, connectorNode.getConnector()); + assertNotSame(originalWorkingFlowContext, connectorNode.getWorkingFlowContext()); + assertEquals("active-value", connectorNode.getActiveFlowContext().getConfigurationContext().getProperty("active-step", "active-property").getValue()); + assertEquals("testValue", connectorNode.getWorkingFlowContext().getConfigurationContext().getProperty("step", "testProperty").getValue()); + assertEquals("1.0.0", connectorNode.getWorkingFlowContext().getBundle().getVersion()); + } + @Test @Timeout(10) public void testRecreationRefreshDoesNotOverwriteConcurrentReplace() throws Exception { diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardReloadComponent.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardReloadComponent.java new file mode 100644 index 000000000000..716c4d3e0e09 --- /dev/null +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/TestStandardReloadComponent.java @@ -0,0 +1,172 @@ +/* + * 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.controller; + +import org.apache.nifi.bundle.Bundle; +import org.apache.nifi.bundle.BundleCoordinate; +import org.apache.nifi.bundle.BundleDetails; +import org.apache.nifi.components.ConfigVerificationResult; +import org.apache.nifi.components.connector.AbstractConnector; +import org.apache.nifi.components.connector.ConfigurationStep; +import org.apache.nifi.components.connector.Connector; +import org.apache.nifi.components.connector.ConnectorNode; +import org.apache.nifi.components.connector.FlowUpdateException; +import org.apache.nifi.components.connector.GhostConnector; +import org.apache.nifi.components.connector.components.FlowContext; +import org.apache.nifi.controller.exception.ConnectorInstantiationException; +import org.apache.nifi.flow.VersionedExternalFlow; +import org.apache.nifi.migration.ConnectorPropertyConfiguration; +import org.apache.nifi.nar.ExtensionManager; +import org.apache.nifi.nar.InstanceClassLoader; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import java.io.File; +import java.util.Collections; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +public class TestStandardReloadComponent { + private static final String CONNECTOR_ID = "connector-id"; + private static final BundleCoordinate COORDINATE = new BundleCoordinate("org.apache.nifi", "test-connector", "2.0.0"); + + @Mock + private FlowController flowController; + @Mock + private ExtensionManager extensionManager; + @Mock + private ConnectorNode connectorNode; + + private StandardReloadComponent reloadComponent; + private InstanceClassLoader candidateClassLoader; + + @BeforeEach + void setUp() { + reloadComponent = new StandardReloadComponent(flowController); + candidateClassLoader = new InstanceClassLoader(CONNECTOR_ID, TestConnector.class.getName(), Collections.emptySet(), Collections.emptySet(), getClass().getClassLoader()); + + final Bundle bundle = new Bundle(new BundleDetails.Builder().coordinate(COORDINATE).workingDir(new File(".")).build(), getClass().getClassLoader()); + when(flowController.getExtensionManager()).thenReturn(extensionManager); + when(extensionManager.getBundle(COORDINATE)).thenReturn(bundle); + when(extensionManager.createInstanceClassLoader(any(), any(), any(), any(), anyBoolean(), any())).thenReturn(candidateClassLoader); + when(connectorNode.getIdentifier()).thenReturn(CONNECTOR_ID); + } + + @Test + void testReloadReplacesConnectorAndClassLoader() throws Exception { + reloadComponent.reload(connectorNode, TestConnector.class.getName(), COORDINATE); + + final ArgumentCaptor connectorCaptor = ArgumentCaptor.forClass(Connector.class); + verify(connectorNode).replaceConnector(connectorCaptor.capture(), eq(COORDINATE), any()); + assertInstanceOf(TestConnector.class, connectorCaptor.getValue()); + verify(extensionManager).removeInstanceClassLoader(CONNECTOR_ID); + verify(extensionManager).registerInstanceClassLoader(CONNECTOR_ID, candidateClassLoader); + verify(connectorNode).resetValidationState(); + } + + @Test + void testReloadClosesCandidateClassLoaderWhenReplacementFails() throws Exception { + doThrow(new FlowUpdateException("replacement failed")).when(connectorNode).replaceConnector(any(), any(), any()); + + assertThrows(ConnectorInstantiationException.class, () -> reloadComponent.reload(connectorNode, TestConnector.class.getName(), COORDINATE)); + + verify(extensionManager).closeURLClassLoader(CONNECTOR_ID, candidateClassLoader); + verify(extensionManager, never()).removeInstanceClassLoader(any()); + verify(extensionManager, never()).registerInstanceClassLoader(any(), any()); + verify(connectorNode, never()).resetValidationState(); + } + + @Test + void testReloadReplacesConnectorWithGhostWhenInstantiationFails() throws Exception { + reloadComponent.reload(connectorNode, "org.apache.nifi.MissingConnector", COORDINATE); + + final ArgumentCaptor connectorCaptor = ArgumentCaptor.forClass(Connector.class); + verify(connectorNode).replaceConnector(connectorCaptor.capture(), eq(COORDINATE), any()); + assertInstanceOf(GhostConnector.class, connectorCaptor.getValue()); + verify(extensionManager).closeURLClassLoader(CONNECTOR_ID, candidateClassLoader); + verify(extensionManager).removeInstanceClassLoader(CONNECTOR_ID); + verify(extensionManager, never()).registerInstanceClassLoader(any(), any()); + verify(connectorNode).resetValidationState(); + } + + @Test + void testReloadRestoresNullContextClassLoader() throws Exception { + final Thread currentThread = Thread.currentThread(); + final ClassLoader originalContextClassLoader = currentThread.getContextClassLoader(); + try { + currentThread.setContextClassLoader(null); + reloadComponent.reload(connectorNode, TestConnector.class.getName(), COORDINATE); + assertNull(currentThread.getContextClassLoader()); + } finally { + currentThread.setContextClassLoader(originalContextClassLoader); + } + } + + public static class TestConnector extends AbstractConnector { + @Override + public VersionedExternalFlow getInitialFlow() { + return null; + } + + @Override + public VersionedExternalFlow getActiveFlow(final FlowContext activeFlowContext) { + return null; + } + + @Override + public void migrateProperties(final ConnectorPropertyConfiguration configuration) { + } + + @Override + public List getConfigurationSteps() { + return List.of(); + } + + @Override + protected void onStepConfigured(final String stepName, final FlowContext workingContext) { + } + + @Override + public void prepareForUpdate(final FlowContext workingContext, final FlowContext activeContext) { + } + + @Override + public void applyUpdate(final FlowContext workingContext, final FlowContext activeContext) { + } + + @Override + public List verifyConfigurationStep(final String stepName, final Map overrides, final FlowContext flowContext) { + return List.of(); + } + } +} diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/audit/ConnectorAuditor.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/audit/ConnectorAuditor.java index eae8d397509f..3723f2f79ff9 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/audit/ConnectorAuditor.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/audit/ConnectorAuditor.java @@ -23,6 +23,7 @@ import org.apache.nifi.action.component.details.FlowChangeExtensionDetails; import org.apache.nifi.action.details.ActionDetails; import org.apache.nifi.action.details.FlowChangeConfigureDetails; +import org.apache.nifi.bundle.BundleCoordinate; import org.apache.nifi.components.connector.AssetReference; import org.apache.nifi.components.connector.ConnectorConfiguration; import org.apache.nifi.components.connector.ConnectorNode; @@ -35,6 +36,7 @@ import org.apache.nifi.components.connector.StringLiteralValue; import org.apache.nifi.web.api.dto.AssetReferenceDTO; import org.apache.nifi.web.api.dto.ConfigurationStepConfigurationDTO; +import org.apache.nifi.web.api.dto.ConnectorDTO; import org.apache.nifi.web.api.dto.ConnectorValueReferenceDTO; import org.apache.nifi.web.api.dto.PropertyGroupConfigurationDTO; import org.apache.nifi.web.dao.ConnectorDAO; @@ -46,6 +48,7 @@ import org.springframework.stereotype.Service; import java.util.ArrayList; +import java.util.Collection; import java.util.Date; import java.util.HashMap; import java.util.List; @@ -63,6 +66,9 @@ public class ConnectorAuditor extends NiFiAuditor { private static final Logger logger = LoggerFactory.getLogger(ConnectorAuditor.class); + private static final String NAME = "Name"; + private static final String EXTENSION_VERSION = "Extension Version"; + /** * Audits the creation of connectors via createConnector(). * @@ -132,6 +138,68 @@ public void startConnectorAdvice(final ProceedingJoinPoint proceedingJoinPoint, } } + /** + * Audits name and extension version changes via updateConnector(). + * + * @param proceedingJoinPoint join point + * @param connectorDTO connector dto + * @param connectorDAO connector dao + * @throws Throwable if an error occurs + */ + @Around("within(org.apache.nifi.web.dao.ConnectorDAO+) && " + + "execution(void updateConnector(org.apache.nifi.web.api.dto.ConnectorDTO)) && " + + "args(connectorDTO) && " + + "target(connectorDAO)") + public void updateConnectorAdvice(final ProceedingJoinPoint proceedingJoinPoint, final ConnectorDTO connectorDTO, final ConnectorDAO connectorDAO) throws Throwable { + ConnectorNode connector = connectorDAO.getConnector(connectorDTO.getId(), ConnectorSyncMode.LOCAL_ONLY); + final Map values = extractConfiguredPropertyValues(connector, connectorDTO); + + proceedingJoinPoint.proceed(); + + connector = connectorDAO.getConnector(connectorDTO.getId(), ConnectorSyncMode.LOCAL_ONLY); + if (!isAuditable()) { + return; + } + + final Map updatedValues = extractConfiguredPropertyValues(connector, connectorDTO); + final FlowChangeExtensionDetails connectorDetails = new FlowChangeExtensionDetails(); + connectorDetails.setType(connector.getComponentType()); + + final Date actionTimestamp = new Date(); + final Collection actions = new ArrayList<>(); + + for (final String property : updatedValues.keySet()) { + final String newValue = updatedValues.get(property); + final String oldValue = values.get(property); + if (oldValue != null && newValue != null && newValue.equals(oldValue)) { + continue; + } + + if (oldValue == null && newValue == null) { + continue; + } + + final FlowChangeConfigureDetails actionDetails = new FlowChangeConfigureDetails(); + actionDetails.setName(property); + actionDetails.setValue(newValue); + actionDetails.setPreviousValue(oldValue); + + final FlowChangeAction configurationAction = createFlowChangeAction(); + configurationAction.setOperation(Operation.Configure); + configurationAction.setTimestamp(actionTimestamp); + configurationAction.setSourceId(connector.getIdentifier()); + configurationAction.setSourceName(connector.getName()); + configurationAction.setComponentDetails(connectorDetails); + configurationAction.setSourceType(Component.Connector); + configurationAction.setActionDetails(actionDetails); + actions.add(configurationAction); + } + + if (!actions.isEmpty()) { + saveActions(actions, logger); + } + } + /** * Audits the stopping of a connector via stopConnector(). * @@ -395,6 +463,18 @@ public void migrateConnectorAdvice(final ProceedingJoinPoint proceedingJoinPoint } } + private Map extractConfiguredPropertyValues(final ConnectorNode connector, final ConnectorDTO connectorDTO) { + final Map values = new HashMap<>(); + if (connectorDTO.getName() != null) { + values.put(NAME, connector.getName()); + } + if (connectorDTO.getBundle() != null) { + final BundleCoordinate bundle = connector.getBundleCoordinate(); + values.put(EXTENSION_VERSION, formatExtensionVersion(connector.getComponentType(), bundle)); + } + return values; + } + /** * Generates an audit record for a connector. * diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java index b178d39d6769..be6e2790ca9f 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java @@ -3696,11 +3696,19 @@ public void verifyUpdateConnector(final ConnectorDTO connectorDTO) { final ConnectorNode connector = connectorDAO.getConnector(connectorDTO.getId()); final ConnectorState currentState = connector.getCurrentState(); - if (connectorDTO.getName() != null && currentState == ConnectorState.TROUBLESHOOTING) { + if ((connectorDTO.getName() != null || connectorDTO.getBundle() != null) && currentState == ConnectorState.TROUBLESHOOTING) { throw new IllegalStateException("Cannot update Connector " + connectorDTO.getId() + " while it is in Troubleshooting mode; exit Troubleshooting mode before modifying the Connector configuration."); } + if (connectorDTO.getBundle() != null) { + final BundleCoordinate incomingCoordinate = BundleUtils.getBundle(controllerFacade.getExtensionManager(), connector.getCanonicalClassName(), connectorDTO.getBundle()); + if (!incomingCoordinate.equals(connector.getBundleCoordinate())) { + connector.verifyCanUpdateBundle(incomingCoordinate); + connector.verifyCanReload(); + } + } + if (connectorDTO.getState() != null) { final ScheduledState desiredState; try { diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java index a0efe4386cd8..9713a20bf8ca 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/ConnectorResource.java @@ -526,7 +526,7 @@ public Response updateConnector( ) @PathParam("id") final String id, @Parameter( - description = "The connector configuration details.", + description = "The connector configuration details. The bundle may be specified to change the NAR version.", required = true ) final ConnectorEntity requestConnectorEntity) { diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java index 46e8ab35d7ad..2d0fae1b3ecd 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/api/dto/DtoFactory.java @@ -5409,7 +5409,19 @@ public ConnectorDTO createConnectorDto(final ConnectorNode connector) { dto.setType(connector.getCanonicalClassName()); dto.setExtensionMissing(connector.isExtensionMissing()); - dto.setBundle(createBundleDto(connector.getBundleCoordinate())); + final BundleCoordinate bundleCoordinate = connector.getBundleCoordinate(); + final List availableBundles = extensionManager.getBundles(connector.getCanonicalClassName()); + int compatibleBundleCount = 0; + for (final Bundle bundle : availableBundles) { + final BundleCoordinate coordinate = bundle.getBundleDetails().getCoordinate(); + if (bundleCoordinate.getGroup().equals(coordinate.getGroup()) && bundleCoordinate.getId().equals(coordinate.getId())) { + compatibleBundleCount++; + } + } + + dto.setMultipleVersionsAvailable(connector.isExtensionMissing() ? compatibleBundleCount > 0 : compatibleBundleCount > 1); + + dto.setBundle(createBundleDto(bundleCoordinate)); dto.setState(connector.getCurrentState().name()); final FrameworkFlowContext activeFlowContext = connector.getActiveFlowContext(); diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectorDAO.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectorDAO.java index 542c4c316cca..8744d6cf695e 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectorDAO.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectorDAO.java @@ -37,12 +37,16 @@ import org.apache.nifi.components.connector.StepConfiguration; import org.apache.nifi.components.connector.StringLiteralValue; import org.apache.nifi.controller.FlowController; +import org.apache.nifi.controller.ReloadComponent; +import org.apache.nifi.controller.exception.ConnectorInstantiationException; import org.apache.nifi.controller.flow.FlowManager; import org.apache.nifi.flow.VersionedExternalFlow; +import org.apache.nifi.util.BundleUtils; import org.apache.nifi.web.NiFiCoreException; import org.apache.nifi.web.NiFiServiceFacade; import org.apache.nifi.web.ResourceNotFoundException; import org.apache.nifi.web.api.dto.AssetReferenceDTO; +import org.apache.nifi.web.api.dto.BundleDTO; import org.apache.nifi.web.api.dto.ConfigurationStepConfigurationDTO; import org.apache.nifi.web.api.dto.ConnectorDTO; import org.apache.nifi.web.api.dto.ConnectorValueReferenceDTO; @@ -74,6 +78,7 @@ public class StandardConnectorDAO implements ConnectorDAO { private FlowController flowController; private NiFiServiceFacade serviceFacade; private ConnectorMigrationManager connectorMigrationManager; + private ReloadComponent reloadComponent; @Autowired public void setFlowController(final FlowController flowController) { @@ -145,11 +150,33 @@ public ConnectorNode createConnector(final String type, final String id, final B @Override public void updateConnector(final ConnectorDTO connectorDTO) { final ConnectorNode connector = requireConnector(connectorDTO.getId(), ConnectorSyncMode.LOCAL_ONLY); + updateBundle(connector, connectorDTO); + if (connectorDTO.getName() != null) { getConnectorRepository().updateConnector(connector, connectorDTO.getName()); } } + private void updateBundle(final ConnectorNode connector, final ConnectorDTO connectorDTO) { + final BundleDTO bundleDTO = connectorDTO.getBundle(); + if (bundleDTO == null) { + return; + } + + final BundleCoordinate incomingCoordinate = BundleUtils.getBundle(flowController.getExtensionManager(), connector.getCanonicalClassName(), bundleDTO); + final BundleCoordinate existingCoordinate = connector.getBundleCoordinate(); + if (existingCoordinate.getCoordinate().equals(incomingCoordinate.getCoordinate())) { + return; + } + + try { + reloadComponent.reload(connector, connector.getCanonicalClassName(), incomingCoordinate); + } catch (final ConnectorInstantiationException | IllegalStateException | IllegalArgumentException e) { + throw new NiFiCoreException(String.format("Unable to update connector %s from %s to %s due to: %s", + connectorDTO.getId(), connector.getBundleCoordinate().getCoordinate(), incomingCoordinate.getCoordinate(), e.getMessage()), e); + } + } + @Override public void verifyDelete(final String id) { final ConnectorNode connector = requireConnector(id, ConnectorSyncMode.LOCAL_ONLY); @@ -386,6 +413,11 @@ public void migrateFromVersionedFlow(final String connectorId, final String proc throw new NiFiCoreException("Failed to migrate Connector from Versioned Flow: " + e.getMessage(), e); } } + + @Autowired + public void setReloadComponent(final ReloadComponent reloadComponent) { + this.reloadComponent = reloadComponent; + } } diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/audit/TestConnectorAuditor.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/audit/TestConnectorAuditor.java new file mode 100644 index 000000000000..8b4469d97bfb --- /dev/null +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/audit/TestConnectorAuditor.java @@ -0,0 +1,144 @@ +/* + * 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.audit; + +import org.apache.nifi.action.Action; +import org.apache.nifi.action.Component; +import org.apache.nifi.action.Operation; +import org.apache.nifi.action.details.ActionDetails; +import org.apache.nifi.action.details.FlowChangeConfigureDetails; +import org.apache.nifi.admin.service.AuditService; +import org.apache.nifi.bundle.BundleCoordinate; +import org.apache.nifi.components.connector.ConnectorNode; +import org.apache.nifi.components.connector.ConnectorSyncMode; +import org.apache.nifi.web.api.dto.BundleDTO; +import org.apache.nifi.web.api.dto.ConnectorDTO; +import org.apache.nifi.web.dao.ConnectorDAO; +import org.aspectj.lang.ProceedingJoinPoint; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; +import org.springframework.security.core.Authentication; +import org.springframework.security.core.context.SecurityContext; +import org.springframework.security.core.context.SecurityContextHolder; + +import java.util.Collection; +import java.util.HashMap; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +public class TestConnectorAuditor { + + private static final String CONNECTOR_ID = "connector-1"; + private static final String COMPONENT_TYPE = "com.example.TestConnector"; + private static final String OLD_NAME = "Old Name"; + private static final String NEW_NAME = "New Name"; + private static final BundleCoordinate OLD_COORDINATE = new BundleCoordinate("com.example", "test-connector", "1.0.0"); + private static final BundleCoordinate NEW_COORDINATE = new BundleCoordinate("com.example", "test-connector", "2.0.0"); + + private ConnectorAuditor auditor; + private AuditService auditService; + private ConnectorDAO connectorDAO; + private ProceedingJoinPoint proceedingJoinPoint; + + @BeforeEach + void setUp() { + auditService = mock(AuditService.class); + connectorDAO = mock(ConnectorDAO.class); + proceedingJoinPoint = mock(ProceedingJoinPoint.class); + + auditor = new ConnectorAuditor(); + auditor.setAuditService(auditService); + + final Authentication authentication = new UsernamePasswordAuthenticationToken("user", "credentials"); + final SecurityContext securityContext = SecurityContextHolder.createEmptyContext(); + securityContext.setAuthentication(authentication); + SecurityContextHolder.setContext(securityContext); + } + + @AfterEach + void clearSecurityContext() { + SecurityContextHolder.clearContext(); + } + + @Test + void testUpdateConnectorAuditsNameAndExtensionVersionChanges() throws Throwable { + final ConnectorNode oldConnector = createConnector(OLD_NAME, OLD_COORDINATE); + final ConnectorNode newConnector = createConnector(NEW_NAME, NEW_COORDINATE); + when(connectorDAO.getConnector(CONNECTOR_ID, ConnectorSyncMode.LOCAL_ONLY)).thenReturn(oldConnector, newConnector); + + final ConnectorDTO connectorDTO = new ConnectorDTO(); + connectorDTO.setId(CONNECTOR_ID); + connectorDTO.setName(NEW_NAME); + connectorDTO.setBundle(new BundleDTO(NEW_COORDINATE.getGroup(), NEW_COORDINATE.getId(), NEW_COORDINATE.getVersion())); + + auditor.updateConnectorAdvice(proceedingJoinPoint, connectorDTO, connectorDAO); + + verify(proceedingJoinPoint).proceed(); + + final Map detailsByProperty = captureConfigureDetails(); + assertEquals(2, detailsByProperty.size()); + + final FlowChangeConfigureDetails nameDetails = detailsByProperty.get("Name"); + assertNotNull(nameDetails); + assertEquals(OLD_NAME, nameDetails.getPreviousValue()); + assertEquals(NEW_NAME, nameDetails.getValue()); + + final FlowChangeConfigureDetails versionDetails = detailsByProperty.get("Extension Version"); + assertNotNull(versionDetails); + assertEquals("com.example.TestConnector 1.0.0 from com.example - test-connector", versionDetails.getPreviousValue()); + assertEquals("com.example.TestConnector 2.0.0 from com.example - test-connector", versionDetails.getValue()); + } + + private ConnectorNode createConnector(final String name, final BundleCoordinate coordinate) { + final ConnectorNode connector = mock(ConnectorNode.class); + when(connector.getIdentifier()).thenReturn(CONNECTOR_ID); + when(connector.getName()).thenReturn(name); + when(connector.getComponentType()).thenReturn(COMPONENT_TYPE); + when(connector.getBundleCoordinate()).thenReturn(coordinate); + return connector; + } + + @SuppressWarnings("unchecked") + private Map captureConfigureDetails() { + final ArgumentCaptor> actionsCaptor = ArgumentCaptor.forClass(Collection.class); + verify(auditService).addActions(actionsCaptor.capture()); + + final Map detailsByProperty = new HashMap<>(); + final Collection actions = actionsCaptor.getValue(); + for (final Action action : actions) { + assertEquals(Operation.Configure, action.getOperation()); + assertEquals(CONNECTOR_ID, action.getSourceId()); + assertEquals(NEW_NAME, action.getSourceName()); + assertEquals(Component.Connector, action.getSourceType()); + + final ActionDetails actionDetails = action.getActionDetails(); + assertNotNull(actionDetails); + final FlowChangeConfigureDetails configureDetails = (FlowChangeConfigureDetails) actionDetails; + detailsByProperty.put(configureDetails.getName(), configureDetails); + } + + return detailsByProperty; + } +} diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java index 079f8cf8629a..69acb39c7777 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/StandardNiFiServiceFacadeTest.java @@ -40,6 +40,9 @@ import org.apache.nifi.authorization.user.NiFiUserDetails; import org.apache.nifi.authorization.user.StandardNiFiUser; import org.apache.nifi.authorization.user.StandardNiFiUser.Builder; +import org.apache.nifi.bundle.Bundle; +import org.apache.nifi.bundle.BundleCoordinate; +import org.apache.nifi.bundle.BundleDetails; import org.apache.nifi.components.Backlog; import org.apache.nifi.components.BacklogReportingException; import org.apache.nifi.components.PropertyDescriptor; @@ -128,6 +131,7 @@ import org.apache.nifi.web.api.dto.BacklogDTO; import org.apache.nifi.web.api.dto.BulletinBoardDTO; import org.apache.nifi.web.api.dto.BulletinQueryDTO; +import org.apache.nifi.web.api.dto.BundleDTO; import org.apache.nifi.web.api.dto.ComponentStateDTO; import org.apache.nifi.web.api.dto.ConnectorDTO; import org.apache.nifi.web.api.dto.CounterDTO; @@ -2366,6 +2370,45 @@ public void testClearConnectorControllerServiceState() { verify(componentStateDAO).clearState(controllerServiceNode, null); } + @Test + public void testVerifyUpdateConnectorSkipsReloadVerificationWhenBundleUnchanged() { + final String connectorId = "connector-id"; + final String type = "org.apache.nifi.connectors.TestConnector"; + final BundleCoordinate coordinate = new BundleCoordinate("org.apache.nifi", "nifi-test-nar", "1.0.0"); + + final ConnectorNode connectorNode = configureConnectorForBundleUpdate(connectorId, type, coordinate, coordinate); + doThrow(new IllegalStateException("Connector cannot be reloaded in its current state")).when(connectorNode).verifyCanReload(); + + final ConnectorDTO connectorDTO = new ConnectorDTO(); + connectorDTO.setId(connectorId); + connectorDTO.setName("Renamed Connector"); + connectorDTO.setBundle(new BundleDTO(coordinate.getGroup(), coordinate.getId(), coordinate.getVersion())); + + serviceFacade.verifyUpdateConnector(connectorDTO); + + verify(connectorNode, never()).verifyCanReload(); + verify(connectorNode, never()).verifyCanUpdateBundle(any()); + } + + @Test + public void testVerifyUpdateConnectorRequiresReloadVerificationWhenBundleChanges() { + final String connectorId = "connector-id"; + final String type = "org.apache.nifi.connectors.TestConnector"; + final BundleCoordinate existingCoordinate = new BundleCoordinate("org.apache.nifi", "nifi-test-nar", "1.0.0"); + final BundleCoordinate incomingCoordinate = new BundleCoordinate("org.apache.nifi", "nifi-test-nar", "2.0.0"); + + final ConnectorNode connectorNode = configureConnectorForBundleUpdate(connectorId, type, existingCoordinate, incomingCoordinate); + + final ConnectorDTO connectorDTO = new ConnectorDTO(); + connectorDTO.setId(connectorId); + connectorDTO.setBundle(new BundleDTO(incomingCoordinate.getGroup(), incomingCoordinate.getId(), incomingCoordinate.getVersion())); + + serviceFacade.verifyUpdateConnector(connectorDTO); + + verify(connectorNode).verifyCanUpdateBundle(incomingCoordinate); + verify(connectorNode).verifyCanReload(); + } + @Test public void testGetConnectorClusterNodeRequest() { final String connectorId = "connector-id"; @@ -3284,4 +3327,31 @@ private AssetManager configureAssets(final Asset asset, final ParameterContext p return assetManager; } + private ConnectorNode configureConnectorForBundleUpdate(final String connectorId, final String type, + final BundleCoordinate existingCoordinate, final BundleCoordinate incomingCoordinate) { + final ConnectorDAO connectorDAO = mock(ConnectorDAO.class); + serviceFacade.setConnectorDAO(connectorDAO); + + final ConnectorNode connectorNode = mock(ConnectorNode.class); + when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode); + when(connectorNode.getCurrentState()).thenReturn(ConnectorState.STOPPED); + when(connectorNode.getCanonicalClassName()).thenReturn(type); + when(connectorNode.getBundleCoordinate()).thenReturn(existingCoordinate); + + final ExtensionManager extensionManager = mock(ExtensionManager.class); + when(flowController.getExtensionManager()).thenReturn(extensionManager); + final Bundle incomingBundle = createBundle(incomingCoordinate); + when(extensionManager.getBundle(incomingCoordinate)).thenReturn(incomingBundle); + when(extensionManager.getBundles(type)).thenReturn(List.of(incomingBundle)); + + return connectorNode; + } + + private Bundle createBundle(final BundleCoordinate coordinate) { + final BundleDetails details = new BundleDetails.Builder() + .workingDir(new File(".")) + .coordinate(coordinate) + .build(); + return new Bundle(details, getClass().getClassLoader()); + } } diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java index a6f4fb1cd812..c0669921142b 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/dto/DtoFactoryTest.java @@ -22,6 +22,10 @@ import org.apache.nifi.bundle.BundleDetails; import org.apache.nifi.components.AllowableValue; import org.apache.nifi.components.PropertyDescriptor; +import org.apache.nifi.components.connector.ConnectorNode; +import org.apache.nifi.components.connector.ConnectorState; +import org.apache.nifi.components.connector.FrameworkFlowContext; +import org.apache.nifi.components.validation.ValidationState; import org.apache.nifi.components.validation.ValidationStatus; import org.apache.nifi.connectable.Connectable; import org.apache.nifi.connectable.ConnectableType; @@ -1043,6 +1047,52 @@ void testCreateParameterDtoSensitiveValueIsMasked() { assertEquals(DtoFactory.SENSITIVE_VALUE_MASK, dto.getValue()); } + @Test + void testConnectorMultipleVersionsAvailable() { + final List oneCompatibleBundle = Collections.singletonList(createBundle("com.example", "test-connector", "1.1.0")); + assertTrue(createConnectorDto(true, oneCompatibleBundle).getMultipleVersionsAvailable()); + assertFalse(createConnectorDto(false, oneCompatibleBundle).getMultipleVersionsAvailable()); + + final List twoCompatibleBundles = Arrays.asList( + createBundle("com.example", "test-connector", "1.0.0"), + createBundle("com.example", "test-connector", "2.0.0")); + assertTrue(createConnectorDto(false, twoCompatibleBundles).getMultipleVersionsAvailable()); + } + + private ConnectorDTO createConnectorDto(final boolean extensionMissing, final List compatibleBundles) { + final String group = "com.example"; + final String id = "test-connector"; + final BundleCoordinate currentCoordinate = new BundleCoordinate(group, id, "1.0.0"); + final String canonicalClassName = "com.example.TestConnector"; + + final ExtensionManager extensionManager = mock(ExtensionManager.class); + when(extensionManager.getBundles(canonicalClassName)).thenReturn(compatibleBundles); + + final ProcessGroup managedProcessGroup = mock(ProcessGroup.class); + when(managedProcessGroup.getIdentifier()).thenReturn("pg-1"); + + final FrameworkFlowContext activeFlowContext = mock(FrameworkFlowContext.class); + when(activeFlowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup); + when(activeFlowContext.getConfigurationContext()).thenReturn(null); + + final ConnectorNode connectorNode = mock(ConnectorNode.class); + when(connectorNode.getIdentifier()).thenReturn("connector-1"); + when(connectorNode.getName()).thenReturn("Connector"); + when(connectorNode.getCanonicalClassName()).thenReturn(canonicalClassName); + when(connectorNode.getBundleCoordinate()).thenReturn(currentCoordinate); + when(connectorNode.isExtensionMissing()).thenReturn(extensionMissing); + when(connectorNode.getCurrentState()).thenReturn(ConnectorState.STOPPED); + when(connectorNode.getValidationState()).thenReturn(new ValidationState(ValidationStatus.VALID, Collections.emptyList())); + when(connectorNode.getActiveFlowContext()).thenReturn(activeFlowContext); + when(connectorNode.getWorkingFlowContext()).thenReturn(null); + when(connectorNode.getConfigurationSteps()).thenReturn(Collections.emptyList()); + when(connectorNode.getAvailableActions()).thenReturn(Collections.emptyList()); + + final DtoFactory dtoFactory = new DtoFactory(); + dtoFactory.setExtensionManager(extensionManager); + return dtoFactory.createConnectorDto(connectorNode); + } + private static DtoFactory newDtoFactoryForParameters() { final DtoFactory dtoFactory = new DtoFactory(); dtoFactory.setEntityFactory(new EntityFactory()); diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectorDAOTest.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectorDAOTest.java index 4b2e9f6dd6d1..7c3b29fbd790 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectorDAOTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectorDAOTest.java @@ -16,6 +16,9 @@ */ package org.apache.nifi.web.dao.impl; +import org.apache.nifi.bundle.Bundle; +import org.apache.nifi.bundle.BundleCoordinate; +import org.apache.nifi.bundle.BundleDetails; import org.apache.nifi.components.AllowableValue; import org.apache.nifi.components.DescribedValue; import org.apache.nifi.components.connector.ConnectorConfiguration; @@ -27,8 +30,11 @@ import org.apache.nifi.components.connector.FrameworkFlowContext; import org.apache.nifi.components.connector.MutableConnectorConfigurationContext; import org.apache.nifi.controller.FlowController; +import org.apache.nifi.controller.ReloadComponent; +import org.apache.nifi.nar.ExtensionManager; import org.apache.nifi.web.NiFiCoreException; import org.apache.nifi.web.ResourceNotFoundException; +import org.apache.nifi.web.api.dto.BundleDTO; import org.apache.nifi.web.api.dto.ConfigurationStepConfigurationDTO; import org.apache.nifi.web.api.dto.ConnectorDTO; import org.junit.jupiter.api.BeforeEach; @@ -40,6 +46,7 @@ import org.mockito.junit.jupiter.MockitoSettings; import org.mockito.quality.Strictness; +import java.io.File; import java.util.List; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -82,6 +89,9 @@ class StandardConnectorDAOTest { private static final String CONNECTOR_ID = "test-connector-id"; private static final String STEP_NAME = "test-step"; private static final String PROPERTY_NAME = "test-property"; + private static final String CANONICAL_CLASS_NAME = "org.apache.nifi.connectors.TestConnector"; + private static final BundleCoordinate EXISTING_COORDINATE = new BundleCoordinate("org.apache.nifi", "nifi-test-nar", "1.0.0"); + private static final BundleCoordinate INCOMING_COORDINATE = new BundleCoordinate("org.apache.nifi", "nifi-test-nar", "2.0.0"); @BeforeEach void setUp() { @@ -330,5 +340,57 @@ void testVerifyCreateWithNullId() { verify(connectorRepository, never()).verifyCreate(any()); } + @Test + void testUpdateConnectorReloadsBundleWhenVersionChanges() throws Exception { + final ReloadComponent reloadComponent = configureBundleUpdate(EXISTING_COORDINATE, INCOMING_COORDINATE); + final ConnectorDTO connectorDTO = createConnectorDto(null, INCOMING_COORDINATE); + + connectorDAO.updateConnector(connectorDTO); + + verify(reloadComponent).reload(connectorNode, CANONICAL_CLASS_NAME, INCOMING_COORDINATE); + verify(connectorRepository, never()).updateConnector(any(), any()); + } + + @Test + void testUpdateConnectorDoesNotReloadWhenBundleUnchanged() throws Exception { + final ReloadComponent reloadComponent = configureBundleUpdate(EXISTING_COORDINATE, EXISTING_COORDINATE); + final ConnectorDTO connectorDTO = createConnectorDto(null, EXISTING_COORDINATE); + + connectorDAO.updateConnector(connectorDTO); + + verify(reloadComponent, never()).reload(any(ConnectorNode.class), any(), any()); + } + + private ReloadComponent configureBundleUpdate(final BundleCoordinate existingCoordinate, final BundleCoordinate incomingCoordinate) { + when(connectorRepository.getConnector(CONNECTOR_ID, ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode); + when(connectorNode.getCanonicalClassName()).thenReturn(CANONICAL_CLASS_NAME); + when(connectorNode.getBundleCoordinate()).thenReturn(existingCoordinate); + + final ExtensionManager extensionManager = mock(ExtensionManager.class); + when(flowController.getExtensionManager()).thenReturn(extensionManager); + when(extensionManager.getBundle(incomingCoordinate)).thenReturn(createBundle(incomingCoordinate)); + when(extensionManager.getBundles(CANONICAL_CLASS_NAME)).thenReturn(List.of(createBundle(incomingCoordinate))); + + final ReloadComponent reloadComponent = mock(ReloadComponent.class); + connectorDAO.setReloadComponent(reloadComponent); + return reloadComponent; + } + + private ConnectorDTO createConnectorDto(final String name, final BundleCoordinate coordinate) { + final ConnectorDTO connectorDTO = new ConnectorDTO(); + connectorDTO.setId(CONNECTOR_ID); + connectorDTO.setName(name); + connectorDTO.setBundle(new BundleDTO(coordinate.getGroup(), coordinate.getId(), coordinate.getVersion())); + return connectorDTO; + } + + private Bundle createBundle(final BundleCoordinate coordinate) { + final BundleDetails details = new BundleDetails.Builder() + .workingDir(new File(".")) + .coordinate(coordinate) + .build(); + return new Bundle(details, getClass().getClassLoader()); + } + } diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/service/connector.service.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/service/connector.service.ts index 2c8dcb56744f..b8df97b8d98f 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/service/connector.service.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/service/connector.service.ts @@ -21,7 +21,7 @@ import { map } from 'rxjs/operators'; import { HttpClient } from '@angular/common/http'; import { Client } from '../../../service/client.service'; import { ClusterConnectionService } from '../../../service/cluster-connection.service'; -import { ConnectorsResponse, CreateConnectorRequest } from '../state'; +import { ChangeConnectorVersionRequest, ConnectorsResponse, CreateConnectorRequest } from '../state'; import { ConnectorEntity } from '@nifi/shared'; import { ControllerServiceEntity, ParameterContextEntity, SearchResultsEntity } from '../../../state/shared'; import { DropRequestEntity } from '../../../state/empty-queue'; @@ -62,6 +62,15 @@ export class ConnectorService { }); } + changeConnectorVersion(request: ChangeConnectorVersionRequest): Observable { + return this.httpClient.put(`${ConnectorService.API}/connectors/${request.id}`, { + revision: this.client.getRevision({ revision: request.payload.revision }), + disconnectedNodeAcknowledged: this.clusterConnectionService.isDisconnectionAcknowledged(), + component: request.payload.component, + id: request.id + }); + } + deleteConnector(connector: ConnectorEntity): Observable { const revision = this.client.getRevision(connector); const params: Record = { diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.actions.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.actions.ts index 24a9fdbdcefd..5c5596671565 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.actions.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.actions.ts @@ -19,6 +19,7 @@ import { createAction, props } from '@ngrx/store'; import { HttpErrorResponse } from '@angular/common/http'; import { CancelDrainConnectorRequest, + ChangeConnectorVersionRequest, ConfigureConnectorRequest, ConnectorActionSuccess, CreateConnectorRequest, @@ -38,6 +39,7 @@ import { ViewConnectorRequest } from '../index'; import { ConnectorEntity } from '@nifi/shared'; +import { FetchComponentVersionsRequest } from '../../../../state/shared'; export const resetConnectorsListingState = createAction('[Connectors Listing] Reset Connectors Listing State'); @@ -150,6 +152,26 @@ export const renameConnectorApiError = createAction( props<{ error: string }>() ); +export const openChangeConnectorVersionDialog = createAction( + '[Connectors Listing] Open Change Connector Version Dialog', + props<{ request: FetchComponentVersionsRequest }>() +); + +export const changeConnectorVersion = createAction( + '[Connectors Listing] Change Connector Version', + props<{ request: ChangeConnectorVersionRequest }>() +); + +export const changeConnectorVersionSuccess = createAction( + '[Connectors Listing] Change Connector Version Success', + props<{ response: ConnectorActionSuccess }>() +); + +export const changeConnectorVersionApiError = createAction( + '[Connectors Listing] Change Connector Version Api Error', + props<{ error: string }>() +); + export const promptDiscardConnectorConfig = createAction( '[Connectors Listing] Prompt Discard Connector Config', props<{ request: DiscardConnectorConfigRequest }>() diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.spec.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.spec.ts index 73b37b5a41bb..0f77ab93a8e1 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.spec.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.spec.ts @@ -23,18 +23,31 @@ import { ConnectorsListingEffects } from './connectors-listing.effects'; import { ConnectorService } from '../../service/connector.service'; import { ErrorHelper } from '../../../../service/error-helper.service'; import { Client } from '../../../../service/client.service'; +import { ExtensionTypesService } from '../../../../service/extension-types.service'; import { MatDialog } from '@angular/material/dialog'; import { Router } from '@angular/router'; import { provideMockStore, MockStore } from '@ngrx/store/testing'; import { HttpErrorResponse } from '@angular/common/http'; import { CreateConnector } from '../../ui/create-connector/create-connector.component'; -import { ConnectorAction, ConnectorActionName, ConnectorEntity, ConnectorStatus, YesNoDialog } from '@nifi/shared'; +import { + Bundle, + ConnectorAction, + ConnectorActionName, + ConnectorEntity, + ConnectorStatus, + YesNoDialog +} from '@nifi/shared'; +import { ChangeConnectorVersionRequest } from '../index'; import * as ErrorActions from '../../../../state/error/error.actions'; import { ErrorContextKey } from '../../../../state/error'; import { selectLoadedTimestamp, selectSaving } from './connectors-listing.selectors'; +import { ChangeComponentVersionDialog } from '../../../../ui/common/change-component-version-dialog/change-component-version-dialog'; import { cancelConnectorDrain, cancelConnectorDrainSuccess, + changeConnectorVersion, + changeConnectorVersionApiError, + changeConnectorVersionSuccess, connectorsListingBannerApiError, createConnector, createConnectorSuccess, @@ -48,6 +61,7 @@ import { loadConnectorsListingError, loadConnectorsListingSuccess, navigateToConfigureConnector, + openChangeConnectorVersionDialog, openNewConnectorDialog, promptConnectorDeletion, promptDiscardConnectorConfig, @@ -110,8 +124,25 @@ describe('ConnectorsListingEffects', () => { }; } + function createChangeVersionRequest( + connector: ConnectorEntity, + bundle: Bundle = connector.component.bundle + ): ChangeConnectorVersionRequest { + return { + id: connector.id, + uri: connector.uri, + payload: { + revision: connector.revision, + component: { + id: connector.id, + bundle + } + } + }; + } + function createMockDialogRef(data: Record | Subject> = {}) { - return { componentInstance: data, afterClosed: () => new Subject() }; + return { componentInstance: data, afterClosed: () => new Subject(), close: vi.fn() }; } async function setup( @@ -127,6 +158,7 @@ describe('ConnectorsListingEffects', () => { createConnector: vi.fn(), deleteConnector: vi.fn(), updateConnector: vi.fn(), + changeConnectorVersion: vi.fn(), updateConnectorRunStatus: vi.fn(), discardConnectorWorkingConfiguration: vi.fn(), drainConnector: vi.fn(), @@ -151,6 +183,10 @@ describe('ConnectorsListingEffects', () => { navigate: vi.fn() }; + const mockExtensionTypesService = { + getConnectorVersionsForType: vi.fn() + }; + await TestBed.configureTestingModule({ providers: [ ConnectorsListingEffects, @@ -166,7 +202,8 @@ describe('ConnectorsListingEffects', () => { { provide: ErrorHelper, useValue: mockErrorHelper }, { provide: Client, useValue: mockClient }, { provide: MatDialog, useValue: mockDialog }, - { provide: Router, useValue: mockRouter } + { provide: Router, useValue: mockRouter }, + { provide: ExtensionTypesService, useValue: mockExtensionTypesService } ] }).compileComponents(); @@ -183,7 +220,8 @@ describe('ConnectorsListingEffects', () => { mockErrorHelper, mockClient, mockDialog, - mockRouter + mockRouter, + mockExtensionTypesService }; } @@ -570,6 +608,132 @@ describe('ConnectorsListingEffects', () => { })); }); + describe('openChangeConnectorVersionDialog$', () => { + it('should open the dialog and dispatch the selected version', () => + new Promise((resolve) => { + setup().then(({ effects, store, actions$, mockDialog, mockExtensionTypesService }) => { + const connector = createMockConnector(); + const request = { + id: connector.id, + uri: connector.uri, + revision: connector.revision, + type: connector.component.type, + bundle: connector.component.bundle + }; + const selectedBundle = { + group: 'org.apache.nifi', + artifact: 'nifi-test-nar', + version: '2.0.0' + }; + const changeVersion = new Subject<{ bundle: typeof selectedBundle }>(); + const dialogRef = createMockDialogRef({ changeVersion }); + (mockExtensionTypesService.getConnectorVersionsForType as Mock).mockReturnValue( + of({ connectorTypes: [] }) + ); + (mockDialog.open as Mock).mockReturnValue(dialogRef); + const dispatchSpy = vi.spyOn(store, 'dispatch'); + actions$(of(openChangeConnectorVersionDialog({ request }))); + + effects.openChangeConnectorVersionDialog$.subscribe(() => { + expect(mockExtensionTypesService.getConnectorVersionsForType).toHaveBeenCalledWith( + request.type, + request.bundle + ); + expect(mockDialog.open).toHaveBeenCalledWith( + ChangeComponentVersionDialog, + expect.objectContaining({ + data: { + fetchRequest: request, + componentVersions: [] + } + }) + ); + + changeVersion.next({ bundle: selectedBundle }); + expect(dispatchSpy).toHaveBeenCalledWith( + changeConnectorVersion({ + request: { + id: request.id, + uri: request.uri, + payload: { + component: { + bundle: selectedBundle, + id: request.id + }, + revision: request.revision + } + } + }) + ); + expect(dialogRef.close).toHaveBeenCalled(); + resolve(); + }); + }); + })); + }); + + describe('changeConnectorVersion$', () => { + it('should change connector version successfully', () => + new Promise((resolve) => { + setup().then(({ effects, actions$, mockConnectorService }) => { + const mockConnector = createMockConnector(); + const updatedConnector = { + ...mockConnector, + component: { + ...mockConnector.component, + bundle: { + group: 'org.apache.nifi', + artifact: 'nifi-test-nar', + version: '2.0.0' + } + } + }; + const request = createChangeVersionRequest(mockConnector); + (mockConnectorService.changeConnectorVersion as Mock).mockReturnValue(of(updatedConnector)); + actions$(of(changeConnectorVersion({ request }))); + + effects.changeConnectorVersion$.subscribe((action) => { + expect(mockConnectorService.changeConnectorVersion).toHaveBeenCalledWith(request); + expect(action).toEqual( + changeConnectorVersionSuccess({ response: { connector: updatedConnector } }) + ); + resolve(); + }); + }); + })); + + it('should handle error when changing connector version', () => + new Promise((resolve) => { + setup().then(({ effects, actions$, mockConnectorService }) => { + const mockConnector = createMockConnector(); + const errorResponse = new HttpErrorResponse({ error: 'Error', status: 500, statusText: 'ISE' }); + (mockConnectorService.changeConnectorVersion as Mock).mockReturnValue( + throwError(() => errorResponse) + ); + actions$(of(changeConnectorVersion({ request: createChangeVersionRequest(mockConnector) }))); + + effects.changeConnectorVersion$.subscribe((action) => { + expect(action).toEqual(changeConnectorVersionApiError({ error: 'Error message' })); + resolve(); + }); + }); + })); + }); + + describe('changeConnectorVersionApiError$', () => { + it('should dispatch snackbar error', () => + new Promise((resolve) => { + setup().then(({ effects, actions$ }) => { + actions$(of(changeConnectorVersionApiError({ error: 'Version change failed' }))); + + effects.changeConnectorVersionApiError$.subscribe((action) => { + expect(action).toEqual(ErrorActions.snackBarError({ error: 'Version change failed' })); + resolve(); + }); + }); + })); + }); + describe('renameConnectorSuccess$', () => { it('should close dialog', () => new Promise((resolve) => { diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.ts index e3c8d3b764bf..6f133baa2410 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.effects.ts @@ -21,7 +21,7 @@ import { Store } from '@ngrx/store'; import { MatDialog } from '@angular/material/dialog'; import { Router } from '@angular/router'; import { concatLatestFrom } from '@ngrx/operators'; -import { catchError, from, map, of, switchMap, take, takeUntil, tap } from 'rxjs'; +import { catchError, concatMap, EMPTY, from, map, of, switchMap, take, takeUntil, tap } from 'rxjs'; import { HttpErrorResponse } from '@angular/common/http'; import { LARGE_DIALOG, MEDIUM_DIALOG, SMALL_DIALOG, YesNoDialog } from '@nifi/shared'; import { NiFiState } from '../../../../state'; @@ -32,12 +32,15 @@ import { CreateConnector } from '../../ui/create-connector/create-connector.comp import { RenameConnectorDialog } from '../../ui/rename-connector-dialog/rename-connector-dialog.component'; import { selectLoadedTimestamp, selectSaving } from './connectors-listing.selectors'; import { initialState } from './connectors-listing.reducer'; -import { DocumentedType } from '../../../../state/shared'; +import { DocumentedType, OpenChangeComponentVersionDialogRequest } from '../../../../state/shared'; import * as ErrorActions from '../../../../state/error/error.actions'; import { ErrorContextKey } from '../../../../state/error'; import { cancelConnectorDrain, cancelConnectorDrainSuccess, + changeConnectorVersion, + changeConnectorVersionApiError, + changeConnectorVersionSuccess, connectorsListingBannerApiError, createConnector, createConnectorSuccess, @@ -54,6 +57,7 @@ import { navigateToManageAccessPolicies, navigateToViewConnector, navigateToViewConnectorDetails, + openChangeConnectorVersionDialog, openNewConnectorDialog, openRenameConnectorDialog, promptConnectorDeletion, @@ -70,6 +74,8 @@ import { } from './connectors-listing.actions'; import { RenameConnectorRequest } from '../index'; import { BackNavigation } from '../../../../state/navigation'; +import { ExtensionTypesService } from '../../../../service/extension-types.service'; +import { ChangeComponentVersionDialog } from '../../../../ui/common/change-component-version-dialog/change-component-version-dialog'; @Injectable() export class ConnectorsListingEffects { @@ -80,6 +86,7 @@ export class ConnectorsListingEffects { private client = inject(Client); private dialog = inject(MatDialog); private router = inject(Router); + private extensionTypesService = inject(ExtensionTypesService); loadConnectorsListing$ = createEffect(() => this.actions$.pipe( @@ -433,6 +440,82 @@ export class ConnectorsListingEffects { ) ); + openChangeConnectorVersionDialog$ = createEffect( + () => + this.actions$.pipe( + ofType(openChangeConnectorVersionDialog), + map((action) => action.request), + switchMap((request) => + from(this.extensionTypesService.getConnectorVersionsForType(request.type, request.bundle)).pipe( + map( + (response) => + ({ + fetchRequest: request, + componentVersions: response.connectorTypes + }) as OpenChangeComponentVersionDialogRequest + ), + catchError((errorResponse: HttpErrorResponse) => { + this.store.dispatch( + ErrorActions.snackBarError({ + error: this.errorHelper.getErrorString(errorResponse) + }) + ); + return EMPTY; + }) + ) + ), + tap((request) => { + const dialogRequest = this.dialog.open(ChangeComponentVersionDialog, { + ...LARGE_DIALOG, + data: request, + autoFocus: false + }); + + dialogRequest.componentInstance.changeVersion.pipe(take(1)).subscribe((newVersion) => { + this.store.dispatch( + changeConnectorVersion({ + request: { + id: request.fetchRequest.id, + uri: request.fetchRequest.uri, + payload: { + component: { + bundle: newVersion.bundle, + id: request.fetchRequest.id + }, + revision: request.fetchRequest.revision + } + } + }) + ); + dialogRequest.close(); + }); + }) + ), + { dispatch: false } + ); + + changeConnectorVersion$ = createEffect(() => + this.actions$.pipe( + ofType(changeConnectorVersion), + map((action) => action.request), + concatMap((request) => + from(this.connectorService.changeConnectorVersion(request)).pipe( + map((response) => changeConnectorVersionSuccess({ response: { connector: response } })), + catchError((errorResponse: HttpErrorResponse) => + of(changeConnectorVersionApiError({ error: this.errorHelper.getErrorString(errorResponse) })) + ) + ) + ) + ) + ); + + changeConnectorVersionApiError$ = createEffect(() => + this.actions$.pipe( + ofType(changeConnectorVersionApiError), + map((action) => ErrorActions.snackBarError({ error: action.error })) + ) + ); + navigateToViewConnector$ = createEffect( () => this.actions$.pipe( diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.reducer.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.reducer.ts index e5c8c6f0bb41..9829ad400f69 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.reducer.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/connectors-listing/connectors-listing.reducer.ts @@ -20,6 +20,9 @@ import { ConnectorsListingState } from '../index'; import { cancelConnectorDrain, cancelConnectorDrainSuccess, + changeConnectorVersion, + changeConnectorVersionApiError, + changeConnectorVersionSuccess, connectorsListingBannerApiError, createConnector, createConnectorSuccess, @@ -112,7 +115,7 @@ export const connectorsListingReducer = createReducer( saving: false })), - on(renameConnector, (state) => ({ + on(renameConnector, changeConnectorVersion, (state) => ({ ...state, saving: true })), @@ -123,7 +126,13 @@ export const connectorsListingReducer = createReducer( saving: false })), - on(renameConnectorApiError, (state) => ({ + on(changeConnectorVersionSuccess, (state, { response }) => ({ + ...state, + connectors: state.connectors.map((c) => (c.id === response.connector.id ? response.connector : c)), + saving: false + })), + + on(renameConnectorApiError, changeConnectorVersionApiError, (state) => ({ ...state, saving: false })), diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/index.ts b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/index.ts index d2a63b574864..3efe057d15e0 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/index.ts +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/state/index.ts @@ -111,6 +111,18 @@ export interface ConnectorActionSuccess { connector: ConnectorEntity; } +export interface ChangeConnectorVersionRequest { + id: string; + uri: string; + payload: { + revision: Revision; + component: { + id: string; + bundle: Bundle; + }; + }; +} + export interface ViewConnectorRequest { connectorId: string; processGroupId: string; diff --git a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/ui/connector-table/connector-table.component.html b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/ui/connector-table/connector-table.component.html index bb9943817f21..0be399c3a27c 100644 --- a/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/ui/connector-table/connector-table.component.html +++ b/nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/ui/connector-table/connector-table.component.html @@ -128,6 +128,15 @@ Rename } + @if (canChangeVersion(item)) { + + } @if (canOperate(item)) { @if (canStart(item)) {