This is an automated email from the ASF dual-hosted git repository.
kevdoran pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new fc08aa89719 NIFI-15919: Ensure that we only sync Connectors with
Connector Configuration Provider when necessary (#11231)
fc08aa89719 is described below
commit fc08aa897197327baca920640393f07a3b5eed5a
Author: Mark Payne <[email protected]>
AuthorDate: Mon May 11 20:04:44 2026 -0400
NIFI-15919: Ensure that we only sync Connectors with Connector
Configuration Provider when necessary (#11231)
Signed-off-by: Kevin Doran <[email protected]>
---
.../server/MockNiFiConnectorWebContext.java | 4 +-
.../server/MockNiFiConnectorWebContextTest.java | 5 +-
.../components/connector/ConnectorRepository.java | 16 +++++-
.../components/connector/ConnectorSyncMode.java | 50 +++++++++++++++++
.../connector/StandardConnectorRepository.java | 18 ++++--
.../org/apache/nifi/controller/FlowController.java | 9 +--
.../nifi/controller/flow/StandardFlowManager.java | 5 +-
.../serialization/VersionedDataflowMapper.java | 3 +-
.../serialization/VersionedFlowSynchronizer.java | 3 +-
.../bootstrap/tasks/ConnectionDiagnosticTask.java | 3 +-
.../connector/StandardConnectorNodeIT.java | 4 +-
.../connector/TestStandardConnectorRepository.java | 63 ++++++++++++++++-----
.../VersionedFlowSynchronizerTest.java | 9 +--
.../org/apache/nifi/audit/ConnectorAuditor.java | 12 ++--
.../authorization/StandardAuthorizableLookup.java | 3 +-
.../apache/nifi/web/StandardNiFiServiceFacade.java | 45 ++++++++-------
.../java/org/apache/nifi/web/dao/ConnectorDAO.java | 3 +
.../nifi/web/dao/impl/StandardConnectionDAO.java | 3 +-
.../nifi/web/dao/impl/StandardConnectorDAO.java | 65 +++++++++++++---------
.../nifi/web/StandardNiFiServiceFacadeTest.java | 27 +++++----
.../web/dao/impl/StandardConnectionDAOTest.java | 5 +-
.../web/dao/impl/StandardConnectorDAOTest.java | 63 ++++++++++++++-------
22 files changed, 284 insertions(+), 134 deletions(-)
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContext.java
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContext.java
index 9e8891ee527..ba585633edf 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContext.java
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContext.java
@@ -19,6 +19,7 @@ package org.apache.nifi.mock.connector.server;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.components.FlowContext;
import org.apache.nifi.web.NiFiConnectorWebContext;
@@ -38,7 +39,8 @@ public class MockNiFiConnectorWebContext implements
NiFiConnectorWebContext {
@Override
@SuppressWarnings("unchecked")
public <T> ConnectorWebContext<T> getConnectorWebContext(final String
connectorId) throws IllegalArgumentException {
- final ConnectorNode connectorNode =
connectorRepository.getConnector(connectorId);
+ // Test runner: there is no real ConnectorConfigurationProvider in
this context, so syncing would be wasted work.
+ final ConnectorNode connectorNode =
connectorRepository.getConnector(connectorId, ConnectorSyncMode.LOCAL_ONLY);
if (connectorNode == null) {
throw new IllegalArgumentException("Unable to find connector with
id: " + connectorId);
}
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContextTest.java
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContextTest.java
index 659fe2e8fb7..115e306188a 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContextTest.java
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/MockNiFiConnectorWebContextTest.java
@@ -20,6 +20,7 @@ package org.apache.nifi.mock.connector.server;
import org.apache.nifi.components.connector.Connector;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.FrameworkFlowContext;
import org.apache.nifi.web.NiFiConnectorWebContext;
import org.apache.nifi.web.NiFiConnectorWebContext.ConnectorWebContext;
@@ -55,7 +56,7 @@ class MockNiFiConnectorWebContextTest {
@Test
void testGetConnectorWebContextReturnsConnectorAndFlowContexts() {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getConnector()).thenReturn(connector);
when(connectorNode.getWorkingFlowContext()).thenReturn(workingFlowContext);
when(connectorNode.getActiveFlowContext()).thenReturn(activeFlowContext);
@@ -71,7 +72,7 @@ class MockNiFiConnectorWebContextTest {
@Test
void testGetConnectorWebContextThrowsForUnknownConnector() {
- when(connectorRepository.getConnector("unknown-id")).thenReturn(null);
+ when(connectorRepository.getConnector("unknown-id",
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(null);
final NiFiConnectorWebContext context = new
MockNiFiConnectorWebContext(connectorRepository);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorRepository.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorRepository.java
index 6c0a3ba43ea..d8d06bc79ef 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorRepository.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorRepository.java
@@ -89,16 +89,26 @@ public interface ConnectorRepository {
void removeConnector(String connectorId);
/**
- * Gets the Connector with the given identifier
+ * Gets the Connector with the given identifier.
+ *
* @param identifier the identifier of the Connector to get
+ * @param syncMode whether to consult the {@link
ConnectorConfigurationProvider} (if configured)
+ * to refresh local state from the external store before
returning, or whether
+ * to return the in-memory copy as-is. See {@link
ConnectorSyncMode}.
* @return the Connector with the given identifier, or null if no such
Connector exists
*/
- ConnectorNode getConnector(String identifier);
+ ConnectorNode getConnector(String identifier, ConnectorSyncMode syncMode);
/**
+ * Returns all Connectors in the Repository.
+ *
+ * @param syncMode whether to consult the {@link
ConnectorConfigurationProvider} (if configured)
+ * to refresh each connector's local state from the
external store before
+ * returning, or whether to return the in-memory copies
as-is. See
+ * {@link ConnectorSyncMode}.
* @return all Connectors in the Repository
*/
- List<ConnectorNode> getConnectors();
+ List<ConnectorNode> getConnectors(ConnectorSyncMode syncMode);
/**
* Starts the given Connector, managing any appropriate lifecycle events.
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorSyncMode.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorSyncMode.java
new file mode 100644
index 00000000000..9d970eee044
--- /dev/null
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core-api/src/main/java/org/apache/nifi/components/connector/ConnectorSyncMode.java
@@ -0,0 +1,50 @@
+/*
+ * 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.components.connector;
+
+/**
+ * Indicates whether retrieval of {@link ConnectorNode} instances from the
+ * {@link ConnectorRepository} should consult the configured
+ * {@link ConnectorConfigurationProvider} to refresh the local state from the
+ * external store, or whether the in-memory copy should be returned as-is.
+ *
+ * <p>Choosing {@link #SYNC_WITH_PROVIDER} causes the repository to call into
+ * the configuration provider (which may perform network I/O, secret lookups,
+ * and other potentially expensive or failure-prone operations) before
returning
+ * the connector. This is appropriate for user-facing reads that must reflect
+ * the latest external state, such as REST API responses. Callers operating on
+ * critical paths that must not block on or fail because of the external store
+ * (for example, periodic flow serialization) should use {@link #LOCAL_ONLY}
+ * instead.</p>
+ */
+public enum ConnectorSyncMode {
+
+ /**
+ * Consult the {@link ConnectorConfigurationProvider} (when configured) to
+ * refresh the local connector state from the external store before
+ * returning. This may perform network I/O and may throw if the external
+ * store is unavailable.
+ */
+ SYNC_WITH_PROVIDER,
+
+ /**
+ * Return the in-memory connector state without consulting the
+ * {@link ConnectorConfigurationProvider}. No external I/O is performed.
+ */
+ LOCAL_ONLY
+}
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorRepository.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorRepository.java
index ffd8f2e7a76..3e2d8ede2ea 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorRepository.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorRepository.java
@@ -459,19 +459,23 @@ public class StandardConnectorRepository implements
ConnectorRepository {
}
@Override
- public ConnectorNode getConnector(final String identifier) {
+ public ConnectorNode getConnector(final String identifier, final
ConnectorSyncMode syncMode) {
+ Objects.requireNonNull(syncMode, "syncMode is required");
final ConnectorNode connector = connectors.get(identifier);
- if (connector != null) {
+ if (connector != null && syncMode ==
ConnectorSyncMode.SYNC_WITH_PROVIDER) {
syncFromProvider(connector);
}
return connector;
}
@Override
- public List<ConnectorNode> getConnectors() {
+ public List<ConnectorNode> getConnectors(final ConnectorSyncMode syncMode)
{
+ Objects.requireNonNull(syncMode, "syncMode is required");
final List<ConnectorNode> connectorList =
List.copyOf(connectors.values());
- for (final ConnectorNode connector : connectorList) {
- syncFromProvider(connector);
+ if (syncMode == ConnectorSyncMode.SYNC_WITH_PROVIDER) {
+ for (final ConnectorNode connector : connectorList) {
+ syncFromProvider(connector);
+ }
}
return connectorList;
}
@@ -665,7 +669,9 @@ public class StandardConnectorRepository implements
ConnectorRepository {
@Override
public void updateConnector(final ConnectorNode connector, final String
name) {
if (configurationProvider != null) {
- final ConnectorWorkingConfiguration workingConfiguration =
buildWorkingConfiguration(connector);
+ // Load the latest provider state so that other in-flight working
changes are not overwritten by a rename.
+ final Optional<ConnectorWorkingConfiguration> externalConfig =
configurationProvider.load(connector.getIdentifier());
+ final ConnectorWorkingConfiguration workingConfiguration =
externalConfig.orElseGet(() -> buildWorkingConfiguration(connector));
workingConfiguration.setName(name);
configurationProvider.save(connector.getIdentifier(),
workingConfiguration);
}
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
index aaa8f90109a..259aed26464 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/FlowController.java
@@ -55,6 +55,7 @@ import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
import
org.apache.nifi.components.connector.ConnectorRepositoryInitializationContext;
import org.apache.nifi.components.connector.ConnectorRequestReplicator;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorValidationTrigger;
import org.apache.nifi.components.connector.FrameworkFlowContext;
import
org.apache.nifi.components.connector.StandardConnectorConfigurationProviderInitializationContext;
@@ -1051,7 +1052,7 @@ public class FlowController implements
ReportingTaskProvider, FlowAnalysisRulePr
return connection;
}
- for (final ConnectorNode connector :
connectorRepository.getConnectors()) {
+ for (final ConnectorNode connector :
connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)) {
final FrameworkFlowContext flowContext =
connector.getActiveFlowContext();
if (flowContext != null) {
final ProcessGroup managedGroup =
flowContext.getManagedProcessGroup();
@@ -1078,7 +1079,7 @@ public class FlowController implements
ReportingTaskProvider, FlowAnalysisRulePr
return remoteGroupPort;
}
- for (final ConnectorNode connector :
connectorRepository.getConnectors()) {
+ for (final ConnectorNode connector :
connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)) {
final FrameworkFlowContext flowContext =
connector.getActiveFlowContext();
if (flowContext != null) {
final ProcessGroup managedGroup =
flowContext.getManagedProcessGroup();
@@ -1508,7 +1509,7 @@ public class FlowController implements
ReportingTaskProvider, FlowAnalysisRulePr
LOG.info("Starting {} Connectors",
startConnectorsAfterInitialization.size());
for (final ConnectorNode connectorNode :
startConnectorsAfterInitialization) {
try {
- final ConnectorNode existingConnector =
connectorRepository.getConnector(connectorNode.getIdentifier());
+ final ConnectorNode existingConnector =
connectorRepository.getConnector(connectorNode.getIdentifier(),
ConnectorSyncMode.LOCAL_ONLY);
if (existingConnector == null) {
LOG.debug("Will not start {} because it no longer
exists", connectorNode);
continue;
@@ -1540,7 +1541,7 @@ public class FlowController implements
ReportingTaskProvider, FlowAnalysisRulePr
// Explicitly stop Connectors so that their state is properly
transitioned from UPDATED to STOPPED.
for (final ConnectorNode connectorNode :
startConnectorsAfterInitialization) {
try {
- final ConnectorNode existingConnector =
connectorRepository.getConnector(connectorNode.getIdentifier());
+ final ConnectorNode existingConnector =
connectorRepository.getConnector(connectorNode.getIdentifier(),
ConnectorSyncMode.LOCAL_ONLY);
if (existingConnector != null) {
connectorRepository.stopConnector(connectorNode);
}
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/flow/StandardFlowManager.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/flow/StandardFlowManager.java
index cf74d624310..59763f81a9b 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/flow/StandardFlowManager.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/flow/StandardFlowManager.java
@@ -32,6 +32,7 @@ import org.apache.nifi.components.connector.Connector;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
import org.apache.nifi.components.connector.ConnectorStateTransition;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.FlowContextFactory;
import org.apache.nifi.components.connector.ProcessGroupFactory;
import org.apache.nifi.components.connector.StandardComponentBundleLookup;
@@ -853,12 +854,12 @@ public class StandardFlowManager extends
AbstractFlowManager implements FlowMana
@Override
public List<ConnectorNode> getAllConnectors() {
- return flowController.getConnectorRepository().getConnectors();
+ return
flowController.getConnectorRepository().getConnectors(ConnectorSyncMode.LOCAL_ONLY);
}
@Override
public ConnectorNode getConnector(final String id) {
- return flowController.getConnectorRepository().getConnector(id);
+ return flowController.getConnectorRepository().getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
}
@Override
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedDataflowMapper.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedDataflowMapper.java
index ccf5391707e..7e74a7f2389 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedDataflowMapper.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedDataflowMapper.java
@@ -18,6 +18,7 @@
package org.apache.nifi.controller.serialization;
import org.apache.nifi.components.connector.ConnectorNode;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.connectable.Port;
import org.apache.nifi.controller.FlowAnalysisRuleNode;
import org.apache.nifi.controller.FlowController;
@@ -97,7 +98,7 @@ public class VersionedDataflowMapper {
private List<VersionedConnector> mapConnectors() {
final List<VersionedConnector> connectors = new ArrayList<>();
- for (final ConnectorNode connectorNode :
flowController.getConnectorRepository().getConnectors()) {
+ for (final ConnectorNode connectorNode :
flowController.getConnectorRepository().getConnectors(ConnectorSyncMode.LOCAL_ONLY))
{
final VersionedConnector versionedConnector =
flowMapper.mapConnector(connectorNode);
if (flowController.isStartAfterInitialization(connectorNode)) {
versionedConnector.setScheduledState(ScheduledState.RUNNING);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizer.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizer.java
index f2893fca8b1..52fc732be7a 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizer.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizer.java
@@ -28,6 +28,7 @@ import org.apache.nifi.cluster.protocol.DataFlow;
import org.apache.nifi.cluster.protocol.StandardDataFlow;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorSyncResult;
import org.apache.nifi.components.validation.ValidationStatus;
import org.apache.nifi.connectable.Connectable;
@@ -1058,7 +1059,7 @@ public class VersionedFlowSynchronizer implements
FlowSynchronizer {
}
}
- for (final ConnectorNode existingConnector :
connectorRepository.getConnectors()) {
+ for (final ConnectorNode existingConnector :
connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)) {
if
(!proposedConnectorIds.contains(existingConnector.getIdentifier())) {
logger.info("Connector [{}] (state={}) is no longer part of
the proposed flow. Stopping and removing.",
existingConnector.getIdentifier(),
existingConnector.getCurrentState());
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/diagnostics/bootstrap/tasks/ConnectionDiagnosticTask.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/diagnostics/bootstrap/tasks/ConnectionDiagnosticTask.java
index a7294307ac4..eff5c724a57 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/diagnostics/bootstrap/tasks/ConnectionDiagnosticTask.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/diagnostics/bootstrap/tasks/ConnectionDiagnosticTask.java
@@ -18,6 +18,7 @@ package org.apache.nifi.diagnostics.bootstrap.tasks;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.connectable.Connection;
import org.apache.nifi.controller.FlowController;
import org.apache.nifi.controller.queue.FlowFileQueue;
@@ -60,7 +61,7 @@ public class ConnectionDiagnosticTask implements
DiagnosticTask {
details.add("");
final ConnectorRepository connectorRepository =
flowController.getConnectorRepository();
- final List<ConnectorNode> connectors = connectorRepository != null ?
connectorRepository.getConnectors() : List.of();
+ final List<ConnectorNode> connectors = connectorRepository != null ?
connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY) : List.of();
if (connectors.isEmpty()) {
details.add("This instance has no Connectors.");
} else {
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/StandardConnectorNodeIT.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/StandardConnectorNodeIT.java
index 6b321eca8e8..0ae7e272a89 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/StandardConnectorNodeIT.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/StandardConnectorNodeIT.java
@@ -263,7 +263,7 @@ public class StandardConnectorNodeIT {
final DynamicFlowConnector flowConnector = (DynamicFlowConnector)
connector;
assertTrue(flowConnector.isInitialized());
- assertEquals(List.of(connectorNode),
connectorRepository.getConnectors());
+ assertEquals(List.of(connectorNode),
connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY));
final ProcessGroup rootGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
assertEquals(3, rootGroup.getProcessGroups().size());
@@ -282,7 +282,7 @@ public class StandardConnectorNodeIT {
private ConnectorNode initializeParameterConnector() {
final ConnectorNode connectorNode =
flowManager.createConnector(ParameterConnector.class.getName(),
"parameter-connector", SystemBundle.SYSTEM_BUNDLE_COORDINATE, true, true);
assertNotNull(connectorNode);
- assertEquals(List.of(connectorNode),
connectorRepository.getConnectors());
+ assertEquals(List.of(connectorNode),
connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY));
final ProcessGroup rootGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
assertEquals(3, rootGroup.getProcessors().size());
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorRepository.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorRepository.java
index 0bcfd2736b2..8560cb6f8bc 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorRepository.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorRepository.java
@@ -75,7 +75,7 @@ public class TestStandardConnectorRepository {
repository.addConnector(connector1);
repository.addConnector(connector2);
- final List<ConnectorNode> connectors = repository.getConnectors();
+ final List<ConnectorNode> connectors =
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertEquals(2, connectors.size());
assertTrue(connectors.contains(connector1));
assertTrue(connectors.contains(connector2));
@@ -98,8 +98,8 @@ public class TestStandardConnectorRepository {
repository.addConnector(connector);
repository.removeConnector("connector-1");
- assertEquals(0, repository.getConnectors().size());
- assertNull(repository.getConnector("connector-1"));
+ assertEquals(0,
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER).size());
+ assertNull(repository.getConnector("connector-1",
ConnectorSyncMode.SYNC_WITH_PROVIDER));
}
@Test
@@ -110,8 +110,8 @@ public class TestStandardConnectorRepository {
when(connector.getIdentifier()).thenReturn("connector-1");
repository.restoreConnector(connector);
- assertEquals(1, repository.getConnectors().size());
- assertEquals(connector, repository.getConnector("connector-1"));
+ assertEquals(1,
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER).size());
+ assertEquals(connector, repository.getConnector("connector-1",
ConnectorSyncMode.SYNC_WITH_PROVIDER));
}
@Test
@@ -126,8 +126,8 @@ public class TestStandardConnectorRepository {
repository.addConnector(connector1);
repository.addConnector(connector2);
- final List<ConnectorNode> connectors1 = repository.getConnectors();
- final List<ConnectorNode> connectors2 = repository.getConnectors();
+ final List<ConnectorNode> connectors1 =
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER);
+ final List<ConnectorNode> connectors2 =
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertEquals(2, connectors1.size());
assertEquals(2, connectors2.size());
@@ -147,9 +147,9 @@ public class TestStandardConnectorRepository {
repository.addConnector(connector1);
repository.addConnector(connector2);
- final List<ConnectorNode> connectors = repository.getConnectors();
+ final List<ConnectorNode> connectors =
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertEquals(1, connectors.size());
- assertEquals(connector2, repository.getConnector("same-id"));
+ assertEquals(connector2, repository.getConnector("same-id",
ConnectorSyncMode.SYNC_WITH_PROVIDER));
}
@Test
@@ -233,7 +233,7 @@ public class TestStandardConnectorRepository {
externalConfig.setWorkingFlowConfiguration(List.of(externalStep));
when(provider.load("connector-1")).thenReturn(Optional.of(externalConfig));
- final ConnectorNode result = repository.getConnector("connector-1");
+ final ConnectorNode result = repository.getConnector("connector-1",
ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertNotNull(result);
verify(connector).setName("External Name");
@@ -250,7 +250,7 @@ public class TestStandardConnectorRepository {
when(provider.load("connector-1")).thenReturn(Optional.empty());
- final ConnectorNode result = repository.getConnector("connector-1");
+ final ConnectorNode result = repository.getConnector("connector-1",
ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertNotNull(result);
verify(connector, never()).setName(anyString());
@@ -266,7 +266,7 @@ public class TestStandardConnectorRepository {
when(provider.load("connector-1")).thenThrow(new
ConnectorConfigurationProviderException("Provider failure"));
- assertThrows(ConnectorConfigurationProviderException.class, () ->
repository.getConnector("connector-1"));
+ assertThrows(ConnectorConfigurationProviderException.class, () ->
repository.getConnector("connector-1", ConnectorSyncMode.SYNC_WITH_PROVIDER));
verify(connector, never()).setName(anyString());
}
@@ -277,12 +277,45 @@ public class TestStandardConnectorRepository {
final ConnectorNode connector =
createSimpleConnectorNode("connector-1", "Original Name");
repository.addConnector(connector);
- final ConnectorNode result = repository.getConnector("connector-1");
+ final ConnectorNode result = repository.getConnector("connector-1",
ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertNotNull(result);
verify(connector, never()).setName(anyString());
}
+ @Test
+ public void testGetConnectorLocalOnlyDoesNotCallProvider() {
+ final ConnectorConfigurationProvider provider =
mock(ConnectorConfigurationProvider.class);
+ final StandardConnectorRepository repository =
createRepositoryWithProvider(provider);
+
+ final ConnectorNode connector =
createSimpleConnectorNode("connector-1", "Local Name");
+ repository.restoreConnector(connector);
+
+ final ConnectorNode result = repository.getConnector("connector-1",
ConnectorSyncMode.LOCAL_ONLY);
+
+ assertNotNull(result);
+ verifyNoInteractions(provider);
+ verify(connector, never()).setName(anyString());
+ }
+
+ @Test
+ public void testGetConnectorsLocalOnlyDoesNotCallProvider() {
+ final ConnectorConfigurationProvider provider =
mock(ConnectorConfigurationProvider.class);
+ final StandardConnectorRepository repository =
createRepositoryWithProvider(provider);
+
+ final ConnectorNode connector1 =
createSimpleConnectorNode("connector-1", "Local Name 1");
+ final ConnectorNode connector2 =
createSimpleConnectorNode("connector-2", "Local Name 2");
+ repository.restoreConnector(connector1);
+ repository.restoreConnector(connector2);
+
+ final List<ConnectorNode> results =
repository.getConnectors(ConnectorSyncMode.LOCAL_ONLY);
+
+ assertEquals(2, results.size());
+ verifyNoInteractions(provider);
+ verify(connector1, never()).setName(anyString());
+ verify(connector2, never()).setName(anyString());
+ }
+
@Test
public void testGetConnectorsWithProviderOverrides() {
final ConnectorConfigurationProvider provider =
mock(ConnectorConfigurationProvider.class);
@@ -301,7 +334,7 @@ public class TestStandardConnectorRepository {
when(provider.load("connector-1")).thenReturn(Optional.of(externalConfig1));
when(provider.load("connector-2")).thenReturn(Optional.empty());
- final List<ConnectorNode> results = repository.getConnectors();
+ final List<ConnectorNode> results =
repository.getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER);
assertEquals(2, results.size());
verify(connector1).setName("External Name 1");
@@ -675,7 +708,7 @@ public class TestStandardConnectorRepository {
config.setWorkingFlowConfiguration(List.of(step));
when(provider.load("connector-1")).thenReturn(Optional.of(config));
- repository.getConnector("connector-1");
+ repository.getConnector("connector-1",
ConnectorSyncMode.SYNC_WITH_PROVIDER);
// Working config is updated with NiFi UUIDs as-is -- no translation
in the repository
verify(workingConfigContext).replaceProperties(eq("step1"),
any(StepConfiguration.class));
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizerTest.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizerTest.java
index 47b5c2c02c1..2684627afc3 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizerTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/serialization/VersionedFlowSynchronizerTest.java
@@ -20,6 +20,7 @@ import org.apache.nifi.cluster.protocol.DataFlow;
import org.apache.nifi.components.PropertyDescriptor;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorSyncResult;
import org.apache.nifi.controller.FlowController;
import org.apache.nifi.controller.ReportingTaskNode;
@@ -289,7 +290,7 @@ class VersionedFlowSynchronizerTest {
when(flowController.createVersionedComponentStateLookup(any())).thenReturn(stateLookup);
when(flowController.getControllerServiceProvider()).thenReturn(controllerServiceProvider);
-
when(connectorRepository.getConnectors()).thenReturn(Collections.emptyList());
+
when(connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)).thenReturn(Collections.emptyList());
when(flowController.getConnectorRepository()).thenReturn(connectorRepository);
}
@@ -376,7 +377,7 @@ class VersionedFlowSynchronizerTest {
.thenReturn(java.util.concurrent.CompletableFuture.completedFuture(null));
setFlowController(connectorRepository);
-
when(connectorRepository.getConnectors()).thenReturn(List.of(orphanConnector));
+
when(connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)).thenReturn(List.of(orphanConnector));
when(versionedDataflow.getConnectors()).thenReturn(List.of(proposedConnector));
versionedFlowSynchronizer.sync(flowController, dataFlow, flowService,
BundleUpdateStrategy.USE_SPECIFIED_OR_GHOST);
@@ -410,7 +411,7 @@ class VersionedFlowSynchronizerTest {
.thenReturn(java.util.concurrent.CompletableFuture.completedFuture(null));
setFlowController(connectorRepository);
-
when(connectorRepository.getConnectors()).thenReturn(List.of(existingConnector));
+
when(connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)).thenReturn(List.of(existingConnector));
when(versionedDataflow.getConnectors()).thenReturn(List.of(versionedConnector));
versionedFlowSynchronizer.sync(flowController, dataFlow, flowService,
BundleUpdateStrategy.USE_SPECIFIED_OR_GHOST);
@@ -447,7 +448,7 @@ class VersionedFlowSynchronizerTest {
.thenReturn(java.util.concurrent.CompletableFuture.completedFuture(null));
setFlowController(connectorRepository);
-
when(connectorRepository.getConnectors()).thenReturn(List.of(orphanConnector));
+
when(connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)).thenReturn(List.of(orphanConnector));
when(versionedDataflow.getConnectors()).thenReturn(List.of(proposedConnector));
versionedFlowSynchronizer.sync(flowController, dataFlow, flowService,
BundleUpdateStrategy.USE_SPECIFIED_OR_GHOST);
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 57c762451fe..9108e3a3b34 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
@@ -27,6 +27,7 @@ import org.apache.nifi.components.connector.AssetReference;
import org.apache.nifi.components.connector.ConnectorConfiguration;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorState;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorValueReference;
import org.apache.nifi.components.connector.NamedStepConfiguration;
import org.apache.nifi.components.connector.SecretReference;
@@ -95,7 +96,7 @@ public class ConnectorAuditor extends NiFiAuditor {
+ "args(connectorId) && "
+ "target(connectorDAO)")
public void removeConnectorAdvice(final ProceedingJoinPoint
proceedingJoinPoint, final String connectorId, final ConnectorDAO connectorDAO)
throws Throwable {
- final ConnectorNode connector = connectorDAO.getConnector(connectorId);
+ final ConnectorNode connector = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
proceedingJoinPoint.proceed();
@@ -118,7 +119,7 @@ public class ConnectorAuditor extends NiFiAuditor {
+ "args(connectorId) && "
+ "target(connectorDAO)")
public void startConnectorAdvice(final ProceedingJoinPoint
proceedingJoinPoint, final String connectorId, final ConnectorDAO connectorDAO)
throws Throwable {
- final ConnectorNode connector = connectorDAO.getConnector(connectorId);
+ final ConnectorNode connector = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorState previousState = connector.getCurrentState();
proceedingJoinPoint.proceed();
@@ -144,7 +145,7 @@ public class ConnectorAuditor extends NiFiAuditor {
+ "args(connectorId) && "
+ "target(connectorDAO)")
public void stopConnectorAdvice(final ProceedingJoinPoint
proceedingJoinPoint, final String connectorId, final ConnectorDAO connectorDAO)
throws Throwable {
- final ConnectorNode connector = connectorDAO.getConnector(connectorId);
+ final ConnectorNode connector = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorState previousState = connector.getCurrentState();
proceedingJoinPoint.proceed();
@@ -173,9 +174,8 @@ public class ConnectorAuditor extends NiFiAuditor {
+ "target(connectorDAO)")
public void updateConfigurationStepAdvice(final ProceedingJoinPoint
proceedingJoinPoint, final String connectorId, final String
configurationStepName,
final
ConfigurationStepConfigurationDTO configurationStepConfiguration, final
ConnectorDAO connectorDAO) throws Throwable {
- final ConnectorNode connector = connectorDAO.getConnector(connectorId);
+ final ConnectorNode connector = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
- // Capture the current property values before the update (flat map:
property name -> value)
final Map<String, String> previousValues =
extractCurrentPropertyValues(connector, configurationStepName);
proceedingJoinPoint.proceed();
@@ -346,7 +346,7 @@ public class ConnectorAuditor extends NiFiAuditor {
+ "args(connectorId) && "
+ "target(connectorDAO)")
public void applyConnectorUpdateAdvice(final ProceedingJoinPoint
proceedingJoinPoint, final String connectorId, final ConnectorDAO connectorDAO)
throws Throwable {
- final ConnectorNode connector = connectorDAO.getConnector(connectorId);
+ final ConnectorNode connector = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
proceedingJoinPoint.proceed();
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/authorization/StandardAuthorizableLookup.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/authorization/StandardAuthorizableLookup.java
index 9e9339493cc..be703d121f4 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/authorization/StandardAuthorizableLookup.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/authorization/StandardAuthorizableLookup.java
@@ -31,6 +31,7 @@ import org.apache.nifi.authorization.user.NiFiUser;
import org.apache.nifi.bundle.BundleCoordinate;
import org.apache.nifi.components.ConfigurableComponent;
import org.apache.nifi.components.PropertyDescriptor;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.connectable.Connectable;
import org.apache.nifi.connectable.Connection;
import org.apache.nifi.connectable.Port;
@@ -660,7 +661,7 @@ public class StandardAuthorizableLookup implements
AuthorizableLookup {
@Override
public Authorizable getConnector(final String connectorId) {
- return connectorDAO.getConnector(connectorId);
+ return connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
}
private Authorizable handleResourceTypeContainingOtherResourceType(final
String resource, final ResourceType resourceType) {
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 784710c3a34..c6fcb7fb93e 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
@@ -81,6 +81,7 @@ import org.apache.nifi.components.Validator;
import org.apache.nifi.components.connector.Connector;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorState;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorUpdateContext;
import org.apache.nifi.components.connector.Secret;
import org.apache.nifi.components.connector.secrets.AuthorizableSecret;
@@ -2139,7 +2140,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
}
private ProcessorNode locateConnectorProcessor(final String connectorId,
final String processorId) {
- final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId);
+ final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId, ConnectorSyncMode.LOCAL_ONLY);
final ProcessGroup managedGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
final ProcessorNode processor =
managedGroup.findProcessor(processorId);
if (processor == null) {
@@ -2149,7 +2150,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
}
private ControllerServiceNode locateConnectorControllerService(final
String connectorId, final String controllerServiceId) {
- final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId);
+ final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId, ConnectorSyncMode.LOCAL_ONLY);
final ProcessGroup managedGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
final ControllerServiceNode controllerService =
managedGroup.findControllerService(controllerServiceId, false, true);
if (controllerService == null) {
@@ -3562,7 +3563,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode connector =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode connector =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(connector);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(connector));
final ConnectorStatusDTO status = createConnectorStatusDto(connector);
@@ -3631,13 +3632,13 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
controllerFacade.save();
- final ConnectorNode node =
connectorDAO.getConnector(connectorDTO.getId());
+ final ConnectorNode node =
connectorDAO.getConnector(connectorDTO.getId(), ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(node);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(node));
final ConnectorStatusDTO statusDto = createConnectorStatusDto(node);
@@ -3687,13 +3688,13 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
}
controllerFacade.save();
- final ConnectorNode node = connectorDAO.getConnector(id);
+ final ConnectorNode node = connectorDAO.getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(node);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(node));
final ConnectorStatusDTO statusDto = createConnectorStatusDto(node);
@@ -3702,7 +3703,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
@Override
public void verifyDrainConnector(final String id) {
- final ConnectorNode connector = connectorDAO.getConnector(id);
+ final ConnectorNode connector = connectorDAO.getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorState currentState = connector.getCurrentState();
if (currentState != ConnectorState.STOPPED) {
throw new IllegalStateException("Cannot drain FlowFiles for
Connector " + id + " because it is not currently stopped. Current state: " +
currentState);
@@ -3718,13 +3719,13 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
connectorDAO.drainFlowFiles(id);
controllerFacade.save();
- final ConnectorNode node = connectorDAO.getConnector(id);
+ final ConnectorNode node = connectorDAO.getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(node);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(node));
final ConnectorStatusDTO statusDto = createConnectorStatusDto(node);
@@ -3745,13 +3746,13 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
connectorDAO.cancelDrainFlowFiles(id);
controllerFacade.save();
- final ConnectorNode node = connectorDAO.getConnector(id);
+ final ConnectorNode node = connectorDAO.getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(node);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(node));
final ConnectorStatusDTO statusDto = createConnectorStatusDto(node);
@@ -3782,12 +3783,10 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
final RevisionClaim claim = new StandardRevisionClaim(revision);
final RevisionUpdate<ConnectorDTO> snapshot =
revisionManager.updateRevision(claim, user, () -> {
- // Update the configuration step
connectorDAO.updateConnectorConfigurationStep(id,
configurationStepName, configurationStepConfiguration);
controllerFacade.save();
- // Return updated connector DTO
- final ConnectorNode node = connectorDAO.getConnector(id);
+ final ConnectorNode node = connectorDAO.getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
@@ -3807,13 +3806,13 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
final RevisionUpdate<ConnectorDTO> snapshot =
revisionManager.updateRevision(claim, user, () -> {
connectorDAO.applyConnectorUpdate(connectorId, updateContext);
- final ConnectorNode node = connectorDAO.getConnector(connectorId);
+ final ConnectorNode node = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(node);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(node));
final ConnectorStatusDTO statusDto = createConnectorStatusDto(node);
@@ -3828,13 +3827,13 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
final RevisionUpdate<ConnectorDTO> snapshot =
revisionManager.updateRevision(claim, user, () -> {
connectorDAO.discardWorkingConfiguration(connectorId);
- final ConnectorNode node = connectorDAO.getConnector(connectorId);
+ final ConnectorNode node = connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
final ConnectorDTO dto = dtoFactory.createConnectorDto(node);
final FlowModification lastMod = new
FlowModification(revision.incrementRevision(revision.getClientId()),
user.getIdentity());
return new StandardRevisionUpdate<>(dto, lastMod);
});
- final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId());
+ final ConnectorNode node =
connectorDAO.getConnector(snapshot.getComponent().getId(),
ConnectorSyncMode.LOCAL_ONLY);
final PermissionsDTO permissions =
dtoFactory.createPermissionsDto(node);
final PermissionsDTO operatePermissions =
dtoFactory.createPermissionsDto(new OperationAuthorizable(node));
final ConnectorStatusDTO statusDto = createConnectorStatusDto(node);
@@ -3843,7 +3842,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
@Override
public ProcessGroupFlowEntity getConnectorFlow(final String connectorId,
final String processGroupId, final boolean uiOnly) {
- final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId);
+ final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId, ConnectorSyncMode.LOCAL_ONLY);
final ProcessGroup managedProcessGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
final ProcessGroup targetProcessGroup =
managedProcessGroup.findProcessGroup(processGroupId);
if (targetProcessGroup == null) {
@@ -3854,7 +3853,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
@Override
public ProcessGroupStatusEntity getConnectorProcessGroupStatus(final
String id, final Boolean recursive) {
- final ConnectorNode connectorNode = connectorDAO.getConnector(id);
+ final ConnectorNode connectorNode = connectorDAO.getConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
final ProcessGroup managedProcessGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
final String processGroupId = managedProcessGroup.getIdentifier();
@@ -3877,7 +3876,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
@Override
public Set<ControllerServiceEntity> getConnectorControllerServices(final
String connectorId, final String processGroupId,
final boolean includeAncestorGroups, final boolean
includeDescendantGroups, final boolean includeReferencingComponents) {
- final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId);
+ final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId, ConnectorSyncMode.LOCAL_ONLY);
final ProcessGroup managedProcessGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
final ProcessGroup targetProcessGroup =
managedProcessGroup.findProcessGroup(processGroupId);
if (targetProcessGroup == null) {
@@ -3912,7 +3911,7 @@ public class StandardNiFiServiceFacade implements
NiFiServiceFacade {
@Override
public SearchResultsDTO searchConnector(final String connectorId, final
String query) {
- final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId);
+ final ConnectorNode connectorNode =
connectorDAO.getConnector(connectorId, ConnectorSyncMode.LOCAL_ONLY);
final ProcessGroup managedProcessGroup =
connectorNode.getActiveFlowContext().getManagedProcessGroup();
return controllerFacade.searchConnector(query, managedProcessGroup);
}
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/ConnectorDAO.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/ConnectorDAO.java
index 70801f1dcdc..acf2e7f5390 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/ConnectorDAO.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/ConnectorDAO.java
@@ -21,6 +21,7 @@ import org.apache.nifi.bundle.BundleCoordinate;
import org.apache.nifi.components.ConfigVerificationResult;
import org.apache.nifi.components.DescribedValue;
import org.apache.nifi.components.connector.ConnectorNode;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorUpdateContext;
import org.apache.nifi.web.api.dto.ConfigurationStepConfigurationDTO;
import org.apache.nifi.web.api.dto.ConnectorDTO;
@@ -38,6 +39,8 @@ public interface ConnectorDAO {
ConnectorNode getConnector(String id);
+ ConnectorNode getConnector(String id, ConnectorSyncMode syncMode);
+
List<ConnectorNode> getConnectors();
ConnectorNode createConnector(String type, String id, BundleCoordinate
bundleCoordinate, boolean firstTimeAdded, boolean registerLogObserver);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectionDAO.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectionDAO.java
index 47aa7a14dc8..5b4ab5c1420 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectionDAO.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/dao/impl/StandardConnectionDAO.java
@@ -24,6 +24,7 @@ import
org.apache.nifi.authorization.resource.DataAuthorizable;
import org.apache.nifi.authorization.user.NiFiUser;
import org.apache.nifi.authorization.user.NiFiUserUtils;
import org.apache.nifi.components.connector.ConnectorNode;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.FrameworkFlowContext;
import org.apache.nifi.connectable.Connectable;
import org.apache.nifi.connectable.ConnectableType;
@@ -92,7 +93,7 @@ public class StandardConnectionDAO extends ComponentDAO
implements ConnectionDAO
// Optionally search Connector-managed ProcessGroups
if (includeConnectorManaged) {
- for (final ConnectorNode connector :
flowController.getConnectorRepository().getConnectors()) {
+ for (final ConnectorNode connector :
flowController.getConnectorRepository().getConnectors(ConnectorSyncMode.LOCAL_ONLY))
{
final FrameworkFlowContext flowContext =
connector.getActiveFlowContext();
if (flowContext != null) {
final ProcessGroup managedGroup =
flowContext.getManagedProcessGroup();
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 697e8c5de6c..9a507e16e5c 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
@@ -23,6 +23,7 @@ import org.apache.nifi.components.DescribedValue;
import org.apache.nifi.components.connector.AssetReference;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorUpdateContext;
import org.apache.nifi.components.connector.ConnectorValueReference;
import org.apache.nifi.components.connector.ConnectorValueType;
@@ -85,21 +86,24 @@ public class StandardConnectorDAO implements ConnectorDAO {
@Override
public boolean hasConnector(final String id) {
- return getConnectorRepository().getConnector(id) != null;
+ return getConnectorRepository().getConnector(id,
ConnectorSyncMode.LOCAL_ONLY) != null;
}
@Override
public ConnectorNode getConnector(final String id) {
- final ConnectorNode connector =
getConnectorRepository().getConnector(id);
- if (connector == null) {
- throw new ResourceNotFoundException("Could not find Connector with
ID " + id);
- }
- return connector;
+ // Public read returned to clients; must reflect the latest
configuration from the provider.
+ return requireConnector(id, ConnectorSyncMode.SYNC_WITH_PROVIDER);
+ }
+
+ @Override
+ public ConnectorNode getConnector(final String id, final ConnectorSyncMode
syncMode) {
+ return requireConnector(id, syncMode);
}
@Override
public List<ConnectorNode> getConnectors() {
- return getConnectorRepository().getConnectors();
+ // Public read returned to clients; must reflect the latest
configuration from the provider.
+ return
getConnectorRepository().getConnectors(ConnectorSyncMode.SYNC_WITH_PROVIDER);
}
@Override
@@ -111,7 +115,7 @@ public class StandardConnectorDAO implements ConnectorDAO {
@Override
public void updateConnector(final ConnectorDTO connectorDTO) {
- final ConnectorNode connector = getConnector(connectorDTO.getId());
+ final ConnectorNode connector = requireConnector(connectorDTO.getId(),
ConnectorSyncMode.LOCAL_ONLY);
if (connectorDTO.getName() != null) {
getConnectorRepository().updateConnector(connector,
connectorDTO.getName());
}
@@ -125,43 +129,43 @@ public class StandardConnectorDAO implements ConnectorDAO
{
@Override
public void startConnector(final String id) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
getConnectorRepository().startConnector(connector);
}
@Override
public void stopConnector(final String id) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
getConnectorRepository().stopConnector(connector);
}
@Override
public void drainFlowFiles(final String id) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
connector.drainFlowFiles();
}
@Override
public void cancelDrainFlowFiles(final String id) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
connector.cancelDrainFlowFiles();
}
@Override
public void verifyCancelDrainFlowFile(final String id) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
connector.verifyCancelDrainFlowFiles();
}
@Override
public void verifyPurgeFlowFiles(final String id) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
connector.verifyCanPurgeFlowFiles();
}
@Override
public void purgeFlowFiles(final String id, final String requestor) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
try {
connector.purgeFlowFiles(requestor).get();
} catch (final InterruptedException e) {
@@ -174,12 +178,12 @@ public class StandardConnectorDAO implements ConnectorDAO
{
@Override
public void updateConnectorConfigurationStep(final String id, final String
configurationStepName, final ConfigurationStepConfigurationDTO
configurationStepDto) {
- final ConnectorNode connector = getConnector(id);
+ // ConnectorRepository.configureConnector consults the provider
directly when merging the step,
+ // so a sync at the lookup would be redundant.
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
- // Convert DTO to domain object - flatten all property groups into a
single StepConfiguration
final StepConfiguration stepConfiguration =
convertToStepConfiguration(configurationStepDto);
- // Update the connector configuration through the repository
try {
getConnectorRepository().configureConnector(connector,
configurationStepName, stepConfiguration);
} catch (final Exception e) {
@@ -187,6 +191,14 @@ public class StandardConnectorDAO implements ConnectorDAO {
}
}
+ private ConnectorNode requireConnector(final String id, final
ConnectorSyncMode syncMode) {
+ final ConnectorNode connector =
getConnectorRepository().getConnector(id, syncMode);
+ if (connector == null) {
+ throw new ResourceNotFoundException("Could not find Connector with
ID " + id);
+ }
+ return connector;
+ }
+
private StepConfiguration convertToStepConfiguration(final
ConfigurationStepConfigurationDTO dto) {
final Map<String, ConnectorValueReference> propertyValues = new
HashMap<>();
if (dto.getPropertyGroupConfigurations() != null) {
@@ -222,7 +234,9 @@ public class StandardConnectorDAO implements ConnectorDAO {
@Override
public void applyConnectorUpdate(final String id, final
ConnectorUpdateContext updateContext) {
- final ConnectorNode connector = getConnector(id);
+ // ConnectorRepository.applyUpdate calls syncAssetsFromProvider
internally, which refreshes
+ // the connector from the provider; a sync at the lookup would be
redundant.
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
try {
getConnectorRepository().applyUpdate(connector, updateContext);
} catch (final Exception e) {
@@ -232,19 +246,20 @@ public class StandardConnectorDAO implements ConnectorDAO
{
@Override
public void discardWorkingConfiguration(final String id) {
- final ConnectorNode connector = getConnector(id);
+ // The working configuration is being thrown away; reading the latest
from the provider first serves no purpose.
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
getConnectorRepository().discardWorkingConfiguration(connector);
}
@Override
public void verifyCanVerifyConfigurationStep(final String id, final String
configurationStepName) {
- // Verify that the connector exists
- getConnector(id);
+ requireConnector(id, ConnectorSyncMode.LOCAL_ONLY);
}
@Override
public List<ConfigVerificationResult> verifyConfigurationStep(final String
id, final String configurationStepName, final ConfigurationStepConfigurationDTO
configurationStepDto) {
- final ConnectorNode connector = getConnector(id);
+ // syncAssetsFromProvider below performs the provider sync, so the
lookup itself does not need to.
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
getConnectorRepository().syncAssetsFromProvider(connector);
final StepConfiguration stepConfiguration =
convertToStepConfiguration(configurationStepDto);
return connector.verifyConfigurationStep(configurationStepName,
stepConfiguration);
@@ -252,7 +267,7 @@ public class StandardConnectorDAO implements ConnectorDAO {
@Override
public List<DescribedValue> fetchAllowableValues(final String id, final
String stepName, final String propertyName, final String filter) {
- final ConnectorNode connector = getConnector(id);
+ final ConnectorNode connector = requireConnector(id,
ConnectorSyncMode.LOCAL_ONLY);
if (filter == null || filter.isEmpty()) {
return connector.fetchAllowableValues(stepName, propertyName);
} else {
@@ -262,7 +277,7 @@ public class StandardConnectorDAO implements ConnectorDAO {
@Override
public void verifyCreateAsset(final String id) {
- getConnector(id);
+ requireConnector(id, ConnectorSyncMode.LOCAL_ONLY);
}
@Override
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 3de939f586e..aba36685007 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
@@ -37,6 +37,7 @@ import org.apache.nifi.authorization.user.NiFiUserDetails;
import org.apache.nifi.authorization.user.StandardNiFiUser.Builder;
import org.apache.nifi.components.PropertyDescriptor;
import org.apache.nifi.components.connector.ConnectorNode;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.FrameworkFlowContext;
import org.apache.nifi.components.connector.Secret;
import org.apache.nifi.components.connector.secrets.AuthorizableSecret;
@@ -1576,7 +1577,7 @@ public class StandardNiFiServiceFacadeTest {
final FrameworkFlowContext flowContext =
mock(FrameworkFlowContext.class);
final ProcessGroup managedProcessGroup = mock(ProcessGroup.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.getIdentifier()).thenReturn(managedGroupId);
@@ -1589,7 +1590,7 @@ public class StandardNiFiServiceFacadeTest {
final SearchResultsDTO results =
serviceFacade.searchConnector(connectorId, searchQuery);
assertNotNull(results);
- verify(connectorDAO).getConnector(connectorId);
+ verify(connectorDAO).getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
verify(connectorNode).getActiveFlowContext();
verify(flowContext).getManagedProcessGroup();
verify(controllerFacade).searchConnector(searchQuery,
managedProcessGroup);
@@ -1809,7 +1810,7 @@ public class StandardNiFiServiceFacadeTest {
final Processor processor = mock(Processor.class);
final StateMap localStateMap = mock(StateMap.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findProcessor(processorId)).thenReturn(processorNode);
@@ -1824,7 +1825,7 @@ public class StandardNiFiServiceFacadeTest {
assertNotNull(result);
assertEquals(processorId, result.getComponentId());
- verify(connectorDAO).getConnector(connectorId);
+ verify(connectorDAO).getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
verify(managedProcessGroup).findProcessor(processorId);
verify(componentStateDAO).getState(processorNode, Scope.LOCAL);
}
@@ -1841,7 +1842,7 @@ public class StandardNiFiServiceFacadeTest {
final FrameworkFlowContext flowContext =
mock(FrameworkFlowContext.class);
final ProcessGroup managedProcessGroup = mock(ProcessGroup.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findProcessor(processorId)).thenReturn(null);
@@ -1862,12 +1863,11 @@ public class StandardNiFiServiceFacadeTest {
final ProcessGroup managedProcessGroup = mock(ProcessGroup.class);
final ProcessorNode processorNode = mock(ProcessorNode.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findProcessor(processorId)).thenReturn(processorNode);
- // Should not throw
serviceFacade.verifyCanClearConnectorProcessorState(connectorId,
processorId);
verify(processorNode).verifyCanClearState();
@@ -1892,7 +1892,7 @@ public class StandardNiFiServiceFacadeTest {
final Processor processor = mock(Processor.class);
final StateMap localStateMap = mock(StateMap.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findProcessor(processorId)).thenReturn(processorNode);
@@ -1929,7 +1929,7 @@ public class StandardNiFiServiceFacadeTest {
final ControllerService controllerService =
mock(ControllerService.class);
final StateMap localStateMap = mock(StateMap.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findControllerService(controllerServiceId,
false, true)).thenReturn(controllerServiceNode);
@@ -1944,7 +1944,7 @@ public class StandardNiFiServiceFacadeTest {
assertNotNull(result);
assertEquals(controllerServiceId, result.getComponentId());
- verify(connectorDAO).getConnector(connectorId);
+ verify(connectorDAO).getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY);
verify(managedProcessGroup).findControllerService(controllerServiceId,
false, true);
verify(componentStateDAO).getState(controllerServiceNode, Scope.LOCAL);
}
@@ -1961,7 +1961,7 @@ public class StandardNiFiServiceFacadeTest {
final FrameworkFlowContext flowContext =
mock(FrameworkFlowContext.class);
final ProcessGroup managedProcessGroup = mock(ProcessGroup.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findControllerService(controllerServiceId,
false, true)).thenReturn(null);
@@ -1982,12 +1982,11 @@ public class StandardNiFiServiceFacadeTest {
final ProcessGroup managedProcessGroup = mock(ProcessGroup.class);
final ControllerServiceNode controllerServiceNode =
mock(ControllerServiceNode.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findControllerService(controllerServiceId,
false, true)).thenReturn(controllerServiceNode);
- // Should not throw
serviceFacade.verifyCanClearConnectorControllerServiceState(connectorId,
controllerServiceId);
verify(controllerServiceNode).verifyCanClearState();
@@ -2012,7 +2011,7 @@ public class StandardNiFiServiceFacadeTest {
final ControllerService controllerService =
mock(ControllerService.class);
final StateMap localStateMap = mock(StateMap.class);
- when(connectorDAO.getConnector(connectorId)).thenReturn(connectorNode);
+ when(connectorDAO.getConnector(connectorId,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.getActiveFlowContext()).thenReturn(flowContext);
when(flowContext.getManagedProcessGroup()).thenReturn(managedProcessGroup);
when(managedProcessGroup.findControllerService(controllerServiceId,
false, true)).thenReturn(controllerServiceNode);
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectionDAOTest.java
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectionDAOTest.java
index 7aa6dc19d22..db8a22f34e9 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectionDAOTest.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/dao/impl/StandardConnectionDAOTest.java
@@ -18,6 +18,7 @@ package org.apache.nifi.web.dao.impl;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.FrameworkFlowContext;
import org.apache.nifi.connectable.Connection;
import org.apache.nifi.controller.FlowController;
@@ -92,7 +93,7 @@ class StandardConnectionDAOTest {
when(rootGroup.findConnection(NON_EXISTENT_ID)).thenReturn(null);
// Setup connector managed group
-
when(connectorRepository.getConnectors()).thenReturn(List.of(connectorNode));
+
when(connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)).thenReturn(List.of(connectorNode));
when(connectorNode.getActiveFlowContext()).thenReturn(frameworkFlowContext);
when(frameworkFlowContext.getManagedProcessGroup()).thenReturn(connectorManagedGroup);
when(connectorManagedGroup.findConnection(CONNECTOR_CONNECTION_ID)).thenReturn(connectorConnection);
@@ -184,7 +185,7 @@ class StandardConnectionDAOTest {
final Connection connectionInSecondConnector =
org.mockito.Mockito.mock(Connection.class);
final String secondConnectorConnectionId =
"second-connector-connection-id";
-
when(connectorRepository.getConnectors()).thenReturn(List.of(connectorNode,
connectorNode2));
+
when(connectorRepository.getConnectors(ConnectorSyncMode.LOCAL_ONLY)).thenReturn(List.of(connectorNode,
connectorNode2));
when(connectorNode2.getActiveFlowContext()).thenReturn(flowContext2);
when(flowContext2.getManagedProcessGroup()).thenReturn(managedGroup2);
when(managedGroup2.findConnection(secondConnectorConnectionId)).thenReturn(connectionInSecondConnector);
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 d1c589fd229..4b2e9f6dd6d 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
@@ -21,6 +21,7 @@ import org.apache.nifi.components.DescribedValue;
import org.apache.nifi.components.connector.ConnectorConfiguration;
import org.apache.nifi.components.connector.ConnectorNode;
import org.apache.nifi.components.connector.ConnectorRepository;
+import org.apache.nifi.components.connector.ConnectorSyncMode;
import org.apache.nifi.components.connector.ConnectorUpdateContext;
import org.apache.nifi.components.connector.FlowUpdateException;
import org.apache.nifi.components.connector.FrameworkFlowContext;
@@ -98,30 +99,30 @@ class StandardConnectorDAOTest {
@Test
void testApplyConnectorUpdate() throws Exception {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
connectorDAO.applyConnectorUpdate(CONNECTOR_ID,
connectorUpdateContext);
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
verify(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
}
@Test
void testApplyConnectorUpdateWithNonExistentConnector() throws Exception {
- when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(null);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(null);
final ResourceNotFoundException exception =
assertThrows(ResourceNotFoundException.class, () ->
connectorDAO.applyConnectorUpdate(CONNECTOR_ID,
connectorUpdateContext)
);
assertEquals("Could not find Connector with ID " + CONNECTOR_ID,
exception.getMessage());
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
verify(connectorRepository,
never()).applyUpdate(any(ConnectorNode.class),
any(ConnectorUpdateContext.class));
}
@Test
void testApplyConnectorUpdateWithFlowUpdateException() throws Exception {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
doThrow(new FlowUpdateException("Flow update
failed")).when(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
final NiFiCoreException exception =
assertThrows(NiFiCoreException.class, () ->
@@ -129,13 +130,13 @@ class StandardConnectorDAOTest {
);
assertEquals("Failed to apply connector update:
org.apache.nifi.components.connector.FlowUpdateException: Flow update failed",
exception.getMessage());
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
verify(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
}
@Test
void testApplyConnectorUpdateWithRuntimeException() throws Exception {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
doThrow(new RuntimeException("Test
exception")).when(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
final NiFiCoreException exception =
assertThrows(NiFiCoreException.class, () ->
@@ -143,13 +144,13 @@ class StandardConnectorDAOTest {
);
assertEquals("Failed to apply connector update:
java.lang.RuntimeException: Test exception", exception.getMessage());
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
verify(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
}
@Test
void testApplyConnectorUpdateWithNullException() throws Exception {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
doThrow(new
RuntimeException()).when(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
final NiFiCoreException exception =
assertThrows(NiFiCoreException.class, () ->
@@ -157,29 +158,51 @@ class StandardConnectorDAOTest {
);
assertEquals("Failed to apply connector update:
java.lang.RuntimeException", exception.getMessage());
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
verify(connectorRepository).applyUpdate(connectorNode,
connectorUpdateContext);
}
@Test
void testGetConnectorWithNonExistentId() {
- when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(null);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.SYNC_WITH_PROVIDER)).thenReturn(null);
assertThrows(ResourceNotFoundException.class, () ->
connectorDAO.getConnector(CONNECTOR_ID)
);
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.SYNC_WITH_PROVIDER);
}
@Test
void testGetConnectorSuccess() {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.SYNC_WITH_PROVIDER)).thenReturn(connectorNode);
final ConnectorNode result = connectorDAO.getConnector(CONNECTOR_ID);
assertEquals(connectorNode, result);
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.SYNC_WITH_PROVIDER);
+ }
+
+ @Test
+ void testGetConnectorLocalOnly() {
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
+
+ final ConnectorNode result = connectorDAO.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
+
+ assertEquals(connectorNode, result);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
+ verify(connectorRepository, never()).getConnector(CONNECTOR_ID,
ConnectorSyncMode.SYNC_WITH_PROVIDER);
+ }
+
+ @Test
+ void testGetConnectorLocalOnlyWithNonExistentId() {
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(null);
+
+ assertThrows(ResourceNotFoundException.class, () ->
+ connectorDAO.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)
+ );
+
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
}
@Test
@@ -188,7 +211,7 @@ class StandardConnectorDAOTest {
new AllowableValue("value1", "Value 1", "First value"),
new AllowableValue("value2", "Value 2", "Second value")
);
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.fetchAllowableValues(STEP_NAME,
PROPERTY_NAME)).thenReturn(expectedValues);
final List<DescribedValue> result =
connectorDAO.fetchAllowableValues(CONNECTOR_ID, STEP_NAME, PROPERTY_NAME, null);
@@ -203,7 +226,7 @@ class StandardConnectorDAOTest {
final List<DescribedValue> expectedValues = List.of(
new AllowableValue("value1", "Value 1", "First value")
);
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.fetchAllowableValues(STEP_NAME,
PROPERTY_NAME)).thenReturn(expectedValues);
final List<DescribedValue> result =
connectorDAO.fetchAllowableValues(CONNECTOR_ID, STEP_NAME, PROPERTY_NAME, "");
@@ -219,7 +242,7 @@ class StandardConnectorDAOTest {
final List<DescribedValue> expectedValues = List.of(
new AllowableValue("filtered-value", "Filtered Value", "Filtered
result")
);
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
when(connectorNode.fetchAllowableValues(STEP_NAME, PROPERTY_NAME,
filter)).thenReturn(expectedValues);
final List<DescribedValue> result =
connectorDAO.fetchAllowableValues(CONNECTOR_ID, STEP_NAME, PROPERTY_NAME,
filter);
@@ -231,18 +254,18 @@ class StandardConnectorDAOTest {
@Test
void testFetchAllowableValuesWithNonExistentConnector() {
- when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(null);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(null);
assertThrows(ResourceNotFoundException.class, () ->
connectorDAO.fetchAllowableValues(CONNECTOR_ID, STEP_NAME,
PROPERTY_NAME, null)
);
- verify(connectorRepository).getConnector(CONNECTOR_ID);
+ verify(connectorRepository).getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY);
}
@Test
void testVerifyConfigurationStepSyncsAssetsBeforeVerification() {
-
when(connectorRepository.getConnector(CONNECTOR_ID)).thenReturn(connectorNode);
+ when(connectorRepository.getConnector(CONNECTOR_ID,
ConnectorSyncMode.LOCAL_ONLY)).thenReturn(connectorNode);
final ConfigurationStepConfigurationDTO stepConfigDto = new
ConfigurationStepConfigurationDTO();
connectorDAO.verifyConfigurationStep(CONNECTOR_ID, STEP_NAME,
stepConfigDto);