This is an automated email from the ASF dual-hosted git repository.
pvillard31 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 e411fa647e1 NIFI-16305 Fix Connector mock framework behavior (#11635)
e411fa647e1 is described below
commit e411fa647e1460967b3bf6cbbb8147b2e90bd7e3
Author: Mark Payne <[email protected]>
AuthorDate: Mon Sep 7 14:30:09 2026 -0400
NIFI-16305 Fix Connector mock framework behavior (#11635)
---
.../mock/connector/server/ConnectorTestRunner.java | 5 +-
.../server/StandardConnectorMockServer.java | 63 ++++-----
.../secrets/ConnectorTestRunnerSecretProvider.java | 12 +-
.../server/StandardConnectorMockServerTest.java | 82 ++++++++++++
.../ConnectorTestRunnerSecretsManagerTest.java | 89 +++++++++++++
.../mock/connectors/tests/CreateConnectorIT.java | 30 +++--
.../nifi/mock/connectors/tests/FlowSnapshotIT.java | 95 ++++++++++++++
.../mock/connectors/FlowSnapshotConnector.java | 142 +++++++++++++++++++++
.../org.apache.nifi.components.connector.Connector | 3 +-
9 files changed, 478 insertions(+), 43 deletions(-)
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-api/src/main/java/org/apache/nifi/mock/connector/server/ConnectorTestRunner.java
b/nifi-connector-mock-bundle/nifi-connector-mock-api/src/main/java/org/apache/nifi/mock/connector/server/ConnectorTestRunner.java
index 42f4c1821ee..9678330f676 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-api/src/main/java/org/apache/nifi/mock/connector/server/ConnectorTestRunner.java
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-api/src/main/java/org/apache/nifi/mock/connector/server/ConnectorTestRunner.java
@@ -136,7 +136,10 @@ public interface ConnectorTestRunner extends Closeable {
AssetReference addAsset(String assetName, InputStream contents);
/**
- * Starts the Connector, beginning processing of data through its managed
flow.
+ * Starts the Connector, beginning processing of data through its managed
flow. Failures detected while
+ * initiating startup are propagated to the caller; lifecycle startup
continues asynchronously.
+ *
+ * @throws IllegalStateException if the Connector cannot begin startup
*/
void startConnector();
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServer.java
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServer.java
index 15e54d6d971..bbca1347a4d 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServer.java
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServer.java
@@ -61,14 +61,12 @@ import org.apache.nifi.encrypt.PropertyEncryptor;
import org.apache.nifi.engine.FlowEngine;
import org.apache.nifi.events.VolatileBulletinRepository;
import org.apache.nifi.flow.VersionedExternalFlow;
-import org.apache.nifi.flow.VersionedParameter;
import org.apache.nifi.flow.VersionedParameterContext;
import org.apache.nifi.flow.VersionedProcessGroup;
import org.apache.nifi.groups.ProcessGroup;
import
org.apache.nifi.mock.connector.server.secrets.ConnectorTestRunnerSecretProvider;
import
org.apache.nifi.mock.connector.server.secrets.ConnectorTestRunnerSecretsManager;
import org.apache.nifi.nar.ExtensionMapping;
-import org.apache.nifi.parameter.Parameter;
import org.apache.nifi.parameter.ParameterContext;
import org.apache.nifi.processor.Processor;
import org.apache.nifi.registry.flow.mapping.ComponentIdLookup;
@@ -97,11 +95,13 @@ import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
-import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.Future;
import java.util.concurrent.TimeoutException;
import java.util.jar.JarFile;
import java.util.stream.Stream;
@@ -328,7 +328,9 @@ public class StandardConnectorMockServer implements
ConnectorMockServer {
@Override
public SecretReference createSecretReference(final String secretName) {
- return new
SecretReference(ConnectorTestRunnerSecretProvider.SECRET_PROVIDER_ID,
ConnectorTestRunnerSecretProvider.SECRET_PROVIDER_NAME, secretName, secretName);
+ final String fullyQualifiedName =
ConnectorTestRunnerSecretProvider.SECRET_PROVIDER_NAME + "."
+ + ConnectorTestRunnerSecretProvider.GROUP_NAME + "." + secretName;
+ return new
SecretReference(ConnectorTestRunnerSecretProvider.SECRET_PROVIDER_ID,
ConnectorTestRunnerSecretProvider.SECRET_PROVIDER_NAME, secretName,
fullyQualifiedName);
}
@Override
@@ -359,7 +361,33 @@ public class StandardConnectorMockServer implements
ConnectorMockServer {
public void startConnector() {
initialFlowFileTransferCounts =
connectorNode.getFlowFileTransferCounts();
- connectorNode.start(flowEngine);
+ final Future<Void> startFuture = connectorNode.start(flowEngine);
+ // Property resolution and validation failures complete the Future
before ConnectorNode.start() returns.
+ // Lifecycle startup is asynchronous and can remain pending while the
node retries, so do not wait for it here.
+ if (!startFuture.isDone()) {
+ return;
+ }
+
+ final IllegalStateException startupFailure;
+ try {
+ startFuture.get();
+ return;
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ startupFailure = new IllegalStateException("Interrupted while
checking Connector startup", e);
+ } catch (final ExecutionException e) {
+ startupFailure = new IllegalStateException("Failed to start
Connector", e.getCause());
+ } catch (final CancellationException e) {
+ startupFailure = new IllegalStateException("Connector startup was
cancelled", e);
+ }
+
+ try {
+ connectorNode.stop(flowEngine);
+ } catch (final RuntimeException stopFailure) {
+ startupFailure.addSuppressed(stopFailure);
+ }
+
+ throw startupFailure;
}
@Override
@@ -505,7 +533,7 @@ public class StandardConnectorMockServer implements
ConnectorMockServer {
final ParameterContext parameterContext =
processGroup.getParameterContext();
if (parameterContext != null) {
final Map<String, VersionedParameterContext> parameterContexts =
new HashMap<>();
- final VersionedParameterContext versionedParameterContext =
createVersionedParameterContext(parameterContext);
+ final VersionedParameterContext versionedParameterContext =
flowMapper.mapParameterContext(parameterContext);
parameterContexts.put(versionedParameterContext.getName(),
versionedParameterContext);
externalFlow.setParameterContexts(parameterContexts);
}
@@ -513,29 +541,6 @@ public class StandardConnectorMockServer implements
ConnectorMockServer {
return externalFlow;
}
- private VersionedParameterContext createVersionedParameterContext(final
ParameterContext parameterContext) {
- final VersionedParameterContext versionedParameterContext = new
VersionedParameterContext();
- versionedParameterContext.setName(parameterContext.getName());
-
versionedParameterContext.setDescription(parameterContext.getDescription());
-
versionedParameterContext.setIdentifier(parameterContext.getIdentifier());
-
- final Set<VersionedParameter> versionedParameters = new
LinkedHashSet<>();
- for (final Parameter parameter :
parameterContext.getParameters().values()) {
- final VersionedParameter versionedParameter = new
VersionedParameter();
- versionedParameter.setName(parameter.getDescriptor().getName());
-
versionedParameter.setDescription(parameter.getDescriptor().getDescription());
-
versionedParameter.setSensitive(parameter.getDescriptor().isSensitive());
-
versionedParameter.setProvided(parameter.getDescriptor().isSensitive());
- if (!parameter.getDescriptor().isSensitive()) {
- versionedParameter.setValue(parameter.getValue());
- }
- versionedParameters.add(versionedParameter);
- }
- versionedParameterContext.setParameters(versionedParameters);
-
- return versionedParameterContext;
- }
-
@Override
public void mockProcessor(final String processorType, final Class<?
extends Processor> mockProcessorClass) {
mockExtensionMapper.mockProcessor(processorType,
mockProcessorClass.getName());
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretProvider.java
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretProvider.java
index 94fd6b9c114..9a512b344d3 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretProvider.java
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/main/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretProvider.java
@@ -60,7 +60,7 @@ public class ConnectorTestRunnerSecretProvider implements
SecretProvider {
.name(entry.getKey())
.value(entry.getValue())
.authorizable(AUTHORIZABLE)
- .fullyQualifiedName(GROUP_NAME + "." + entry.getKey())
+ .fullyQualifiedName(SECRET_PROVIDER_NAME + "." + GROUP_NAME +
"." + entry.getKey())
.build();
secrets.add(secret);
@@ -72,8 +72,14 @@ public class ConnectorTestRunnerSecretProvider implements
SecretProvider {
@Override
public List<Secret> getSecrets(final List<String>
fullyQualifiedSecretNames) {
final List<Secret> matchingSecrets = new ArrayList<>();
+ final String secretNamePrefix = SECRET_PROVIDER_NAME + "." +
GROUP_NAME + ".";
- for (final String secretName : fullyQualifiedSecretNames) {
+ for (final String fullyQualifiedSecretName :
fullyQualifiedSecretNames) {
+ if (!fullyQualifiedSecretName.startsWith(secretNamePrefix)) {
+ continue;
+ }
+
+ final String secretName =
fullyQualifiedSecretName.substring(secretNamePrefix.length());
final String value = secrets.get(secretName);
if (value != null) {
@@ -82,7 +88,7 @@ public class ConnectorTestRunnerSecretProvider implements
SecretProvider {
.providerName(SECRET_PROVIDER_NAME)
.groupName(GROUP_NAME)
.name(secretName)
- .fullyQualifiedName(GROUP_NAME + "." + secretName)
+ .fullyQualifiedName(fullyQualifiedSecretName)
.value(value)
.authorizable(AUTHORIZABLE)
.build();
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServerTest.java
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServerTest.java
new file mode 100644
index 00000000000..b9ac2efb58a
--- /dev/null
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/StandardConnectorMockServerTest.java
@@ -0,0 +1,82 @@
+/*
+ * 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.mock.connector.server;
+
+import org.apache.nifi.components.connector.ConnectorNode;
+import org.apache.nifi.connectable.FlowFileTransferCounts;
+import org.apache.nifi.engine.FlowEngine;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.Future;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@ExtendWith(MockitoExtension.class)
+class StandardConnectorMockServerTest {
+
+ @Mock
+ private ConnectorNode connectorNode;
+
+ @Mock
+ private FlowEngine flowEngine;
+
+ @Mock
+ private Future<Void> startFuture;
+
+ @InjectMocks
+ private StandardConnectorMockServer server;
+
+ @Test
+ void testStartConnectorPropagatesFailureAndRequestsStop() {
+ final IllegalStateException startFailure = new
IllegalStateException("Connector is not valid");
+ when(connectorNode.getFlowFileTransferCounts()).thenReturn(new
FlowFileTransferCounts(0, 0, 0, 0));
+
when(connectorNode.start(flowEngine)).thenReturn(CompletableFuture.failedFuture(startFailure));
+
+ final IllegalStateException exception =
assertThrows(IllegalStateException.class, server::startConnector);
+
+ assertEquals("Failed to start Connector", exception.getMessage());
+ assertSame(startFailure, exception.getCause());
+ verify(connectorNode).stop(flowEngine);
+ }
+
+ @Test
+ void testStartConnectorDoesNotWaitForPendingStart() throws Exception {
+ when(connectorNode.getFlowFileTransferCounts()).thenReturn(new
FlowFileTransferCounts(0, 0, 0, 0));
+ when(connectorNode.start(flowEngine)).thenReturn(startFuture);
+ when(startFuture.isDone()).thenReturn(false);
+
+ server.startConnector();
+
+ verify(startFuture).isDone();
+ verify(startFuture, never()).get();
+ verify(startFuture, never()).get(anyLong(), any());
+ verify(connectorNode, never()).stop(flowEngine);
+ }
+}
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretsManagerTest.java
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretsManagerTest.java
new file mode 100644
index 00000000000..9071c913f49
--- /dev/null
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-server/src/test/java/org/apache/nifi/mock/connector/server/secrets/ConnectorTestRunnerSecretsManagerTest.java
@@ -0,0 +1,89 @@
+/*
+ * 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.mock.connector.server.secrets;
+
+import org.apache.nifi.components.connector.Secret;
+import org.apache.nifi.components.connector.SecretReference;
+import org.apache.nifi.mock.connector.server.StandardConnectorMockServer;
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class ConnectorTestRunnerSecretsManagerTest {
+ private static final String SECRET_NAME = "password";
+ private static final String SECRET_VALUE = "secret-value";
+ private static final String FULLY_QUALIFIED_SECRET_NAME =
ConnectorTestRunnerSecretProvider.SECRET_PROVIDER_NAME + "."
+ + ConnectorTestRunnerSecretProvider.GROUP_NAME + "." + SECRET_NAME;
+
+ @Test
+ void testGetSecretUsingDiscoveredFullyQualifiedName() {
+ final ConnectorTestRunnerSecretsManager secretsManager = new
ConnectorTestRunnerSecretsManager();
+ secretsManager.addSecret(SECRET_NAME, SECRET_VALUE);
+
+ final Secret discoveredSecret =
secretsManager.getAllSecrets().getFirst();
+ assertEquals(FULLY_QUALIFIED_SECRET_NAME,
discoveredSecret.getFullyQualifiedName());
+
+ final SecretReference secretReference = new SecretReference(
+ discoveredSecret.getProviderId(),
+ discoveredSecret.getProviderName(),
+ discoveredSecret.getName(),
+ discoveredSecret.getFullyQualifiedName()
+ );
+
+ final Optional<Secret> resolvedSecret =
secretsManager.getSecret(secretReference);
+
+ assertTrue(resolvedSecret.isPresent(), "Discovered secret should be
resolved using its fully qualified name");
+ final Secret secret = resolvedSecret.orElseThrow();
+ assertEquals(SECRET_NAME, secret.getName());
+ assertEquals(FULLY_QUALIFIED_SECRET_NAME,
secret.getFullyQualifiedName());
+ assertEquals(SECRET_VALUE, secret.getValue());
+ }
+
+ @Test
+ void testCreateSecretReferenceUsesFullyQualifiedName() {
+ final ConnectorTestRunnerSecretsManager secretsManager = new
ConnectorTestRunnerSecretsManager();
+ secretsManager.addSecret(SECRET_NAME, SECRET_VALUE);
+ final StandardConnectorMockServer server = new
StandardConnectorMockServer();
+
+ final SecretReference secretReference =
server.createSecretReference(SECRET_NAME);
+
+ assertEquals(FULLY_QUALIFIED_SECRET_NAME,
secretReference.getFullyQualifiedName());
+ assertEquals(SECRET_VALUE,
secretsManager.getSecret(secretReference).orElseThrow().getValue());
+ }
+
+ @Test
+ void testGetSecretWithDotInName() {
+ final String secretName = "database.password";
+ final ConnectorTestRunnerSecretsManager secretsManager = new
ConnectorTestRunnerSecretsManager();
+ secretsManager.addSecret(secretName, SECRET_VALUE);
+
+ final Secret discoveredSecret =
secretsManager.getAllSecrets().getFirst();
+ final SecretReference secretReference = new SecretReference(
+ discoveredSecret.getProviderId(),
+ discoveredSecret.getProviderName(),
+ discoveredSecret.getName(),
+ discoveredSecret.getFullyQualifiedName()
+ );
+
+ assertEquals("TestRunnerSecretsManager.Default.database.password",
discoveredSecret.getFullyQualifiedName());
+ assertEquals(SECRET_VALUE,
secretsManager.getSecret(secretReference).orElseThrow().getValue());
+ }
+}
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/CreateConnectorIT.java
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/CreateConnectorIT.java
index 1b8c49e49f4..f3d529c6a5e 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/CreateConnectorIT.java
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/CreateConnectorIT.java
@@ -32,9 +32,11 @@ import java.time.Duration;
import java.util.List;
import java.util.Optional;
import java.util.Set;
+import java.util.concurrent.TimeoutException;
-import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class CreateConnectorIT {
@@ -68,7 +70,7 @@ public class CreateConnectorIT {
}
@Test
- public void testStopConnectorWithTimeoutStopsRunningConnector() throws
IOException {
+ public void testStopConnectorWithTimeoutStopsRunningConnector() throws
IOException, TimeoutException {
try (final ConnectorTestRunner testRunner = new
StandardConnectorTestRunner.Builder()
.connectorClassName("org.apache.nifi.mock.connectors.GenerateAndLog")
.narLibraryDirectory(new File("target/libDir"))
@@ -76,13 +78,7 @@ public class CreateConnectorIT {
testRunner.startConnector();
- // Exercises the timeout-aware overload: it initiates the
asynchronous stop and actively polls the
- // Connector's state until it reaches STOPPED, returning as soon
as it does rather than after a single
- // fixed blocking wait. Because start is asynchronous, a stop
issued immediately afterwards may have to
- // ride through the node's internal stop retries (every 10
seconds) before the state settles, so the
- // budget is generous enough to stay deterministic on a slow CI
runner; the poll still returns the
- // instant the Connector reports STOPPED.
- assertDoesNotThrow(() ->
testRunner.stopConnector(Duration.ofSeconds(120)));
+ testRunner.stopConnector(Duration.ofSeconds(120));
}
}
@@ -102,6 +98,22 @@ public class CreateConnectorIT {
}
}
+ @Test
+ public void testConnectorWithMissingBundleFailsStart() throws IOException {
+ try (final ConnectorTestRunner testRunner = new
StandardConnectorTestRunner.Builder()
+
.connectorClassName("org.apache.nifi.mock.connectors.MissingBundleConnector")
+ .narLibraryDirectory(new File("target/libDir"))
+ .build()) {
+
+ final IllegalStateException exception =
assertThrows(IllegalStateException.class, testRunner::startConnector);
+ assertEquals("Failed to start Connector", exception.getMessage());
+
+ final IllegalStateException cause =
assertInstanceOf(IllegalStateException.class, exception.getCause());
+
assertTrue(cause.getMessage().contains("com.example.nonexistent:missing-nar:1.0.0"));
+
assertTrue(cause.getMessage().contains("com.example.nonexistent.MissingProcessor"));
+ }
+ }
+
private Optional<VersionedProcessor> findProcessorByType(final
Set<VersionedProcessor> processors, final String type) {
for (final VersionedProcessor processor : processors) {
if (type.equals(processor.getType())) {
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/FlowSnapshotIT.java
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/FlowSnapshotIT.java
new file mode 100644
index 00000000000..ce0c4465879
--- /dev/null
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-integration-tests/src/test/java/org/apache/nifi/mock/connectors/tests/FlowSnapshotIT.java
@@ -0,0 +1,95 @@
+/*
+ * 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.mock.connectors.tests;
+
+import org.apache.nifi.components.connector.AssetReference;
+import org.apache.nifi.components.connector.FlowUpdateException;
+import org.apache.nifi.flow.VersionedAsset;
+import org.apache.nifi.flow.VersionedExternalFlow;
+import org.apache.nifi.flow.VersionedParameter;
+import org.apache.nifi.flow.VersionedParameterContext;
+import org.apache.nifi.mock.connector.StandardConnectorTestRunner;
+import org.apache.nifi.mock.connector.server.ConnectorTestRunner;
+import org.junit.jupiter.api.Test;
+
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class FlowSnapshotIT {
+ private static final String CONNECTOR_CLASS =
"org.apache.nifi.mock.connectors.FlowSnapshotConnector";
+ private static final String CONFIGURATION_STEP = "Snapshot Configuration";
+ private static final String ASSET_IDENTIFIER_PROPERTY = "Asset Identifier";
+ private static final String ASSET_NAME = "certificate.pem";
+ private static final String SENSITIVE_PARAMETER_NAME =
"sensitive_parameter";
+ private static final String ASSET_PARAMETER_NAME = "asset_parameter";
+
+ @Test
+ public void testWorkingAndActiveFlowSnapshotsPreserveParameterMetadata()
throws IOException, FlowUpdateException {
+ try (final ConnectorTestRunner runner = new
StandardConnectorTestRunner.Builder()
+ .connectorClassName(CONNECTOR_CLASS)
+ .narLibraryDirectory(new File("target/libDir"))
+ .build()) {
+
+ final ByteArrayInputStream assetContents = new
ByteArrayInputStream("certificate contents".getBytes(StandardCharsets.UTF_8));
+ final AssetReference assetReference = runner.addAsset(ASSET_NAME,
assetContents);
+ final String assetIdentifier =
assetReference.getAssetIdentifiers().iterator().next();
+
+ runner.configure(CONFIGURATION_STEP,
Map.of(ASSET_IDENTIFIER_PROPERTY, assetIdentifier));
+ assertSnapshotParameters(runner.getWorkingFlowSnapshot(),
assetIdentifier);
+
+ runner.applyUpdate();
+ assertSnapshotParameters(runner.getActiveFlowSnapshot(),
assetIdentifier);
+ }
+ }
+
+ private void assertSnapshotParameters(final VersionedExternalFlow
snapshot, final String assetIdentifier) {
+ assertEquals(1, snapshot.getParameterContexts().size());
+ final VersionedParameterContext parameterContext =
snapshot.getParameterContexts().values().iterator().next();
+
+ final VersionedParameter sensitiveParameter =
getParameter(parameterContext, SENSITIVE_PARAMETER_NAME);
+ final VersionedParameter assetParameter =
getParameter(parameterContext, ASSET_PARAMETER_NAME);
+ assertTrue(sensitiveParameter.isSensitive());
+ assertFalse(sensitiveParameter.isProvided());
+ assertNull(sensitiveParameter.getValue());
+ assertFalse(assetParameter.isSensitive());
+ assertFalse(assetParameter.isProvided());
+ assertNull(assetParameter.getValue());
+ assertNotNull(assetParameter.getReferencedAssets());
+ assertEquals(1, assetParameter.getReferencedAssets().size());
+
+ final VersionedAsset versionedAsset =
assetParameter.getReferencedAssets().getFirst();
+ assertEquals(assetIdentifier, versionedAsset.getIdentifier());
+ assertEquals(ASSET_NAME, versionedAsset.getName());
+ }
+
+ private VersionedParameter getParameter(final VersionedParameterContext
parameterContext, final String name) {
+ return parameterContext.getParameters().stream()
+ .filter(parameter -> name.equals(parameter.getName()))
+ .findFirst()
+ .orElseThrow();
+ }
+}
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/java/org/apache/nifi/mock/connectors/FlowSnapshotConnector.java
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/java/org/apache/nifi/mock/connectors/FlowSnapshotConnector.java
new file mode 100644
index 00000000000..fb67e70aa1d
--- /dev/null
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/java/org/apache/nifi/mock/connectors/FlowSnapshotConnector.java
@@ -0,0 +1,142 @@
+/*
+ * 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.mock.connectors;
+
+import org.apache.nifi.components.ConfigVerificationResult;
+import org.apache.nifi.components.connector.AbstractConnector;
+import org.apache.nifi.components.connector.BundleCompatibility;
+import org.apache.nifi.components.connector.ConfigurationStep;
+import org.apache.nifi.components.connector.ConnectorPropertyDescriptor;
+import org.apache.nifi.components.connector.ConnectorPropertyGroup;
+import org.apache.nifi.components.connector.FlowUpdateException;
+import org.apache.nifi.components.connector.components.FlowContext;
+import org.apache.nifi.components.connector.util.VersionedFlowUtils;
+import org.apache.nifi.flow.VersionedAsset;
+import org.apache.nifi.flow.VersionedExternalFlow;
+import org.apache.nifi.flow.VersionedParameter;
+import org.apache.nifi.flow.VersionedParameterContext;
+import org.apache.nifi.flow.VersionedProcessor;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * Test Connector that installs sensitive and asset-backed parameters for
exercising flow snapshots.
+ */
+public class FlowSnapshotConnector extends AbstractConnector {
+ private static final String FLOW_RESOURCE =
"flows/Generate_and_Update.json";
+ private static final String STEP_NAME = "Snapshot Configuration";
+ private static final String ASSET_IDENTIFIER_PROPERTY_NAME = "Asset
Identifier";
+ private static final String PARAMETER_CONTEXT_NAME = "Snapshot Parameter
Context";
+ private static final String SENSITIVE_PARAMETER_NAME =
"sensitive_parameter";
+ private static final String ASSET_PARAMETER_NAME = "asset_parameter";
+ private static final String ASSET_NAME = "certificate.pem";
+
+ private static final ConnectorPropertyDescriptor ASSET_IDENTIFIER = new
ConnectorPropertyDescriptor.Builder()
+ .name(ASSET_IDENTIFIER_PROPERTY_NAME)
+ .description("Identifier of the asset referenced by the snapshot test
flow")
+ .required(true)
+ .build();
+
+ private static final ConnectorPropertyGroup PROPERTY_GROUP = new
ConnectorPropertyGroup.Builder()
+ .name("Snapshot Properties")
+ .addProperty(ASSET_IDENTIFIER)
+ .build();
+
+ private static final ConfigurationStep CONFIGURATION_STEP = new
ConfigurationStep.Builder()
+ .name(STEP_NAME)
+ .propertyGroups(List.of(PROPERTY_GROUP))
+ .build();
+
+ @Override
+ public VersionedExternalFlow getInitialFlow() {
+ return VersionedFlowUtils.loadFlowFromResource(FLOW_RESOURCE);
+ }
+
+ @Override
+ public VersionedExternalFlow getActiveFlow(final FlowContext
activeFlowContext) {
+ final String assetIdentifier = getAssetIdentifier(activeFlowContext);
+ return assetIdentifier == null ? getInitialFlow() :
createConfiguredFlow(assetIdentifier);
+ }
+
+ @Override
+ public List<ConfigurationStep> getConfigurationSteps() {
+ return List.of(CONFIGURATION_STEP);
+ }
+
+ @Override
+ protected void onStepConfigured(final String stepName, final FlowContext
workingContext) throws FlowUpdateException {
+ if (STEP_NAME.equals(stepName)) {
+ final String assetIdentifier = getAssetIdentifier(workingContext);
+ if (assetIdentifier != null) {
+ getInitializationContext().updateFlow(workingContext,
createConfiguredFlow(assetIdentifier), BundleCompatibility.RESOLVE_BUNDLE);
+ }
+ }
+ }
+
+ @Override
+ public void applyUpdate(final FlowContext workingContext, final
FlowContext activeContext) throws FlowUpdateException {
+ final String assetIdentifier = getAssetIdentifier(workingContext);
+ if (assetIdentifier != null) {
+ getInitializationContext().updateFlow(activeContext,
createConfiguredFlow(assetIdentifier), BundleCompatibility.RESOLVE_BUNDLE);
+ }
+ }
+
+ @Override
+ public List<ConfigVerificationResult> verifyConfigurationStep(final String
stepName, final Map<String, String> overrides, final FlowContext
workingContext) {
+ return List.of();
+ }
+
+ private String getAssetIdentifier(final FlowContext flowContext) {
+ return
flowContext.getConfigurationContext().getProperty(CONFIGURATION_STEP,
ASSET_IDENTIFIER).getValue();
+ }
+
+ private VersionedExternalFlow createConfiguredFlow(final String
assetIdentifier) {
+ final VersionedExternalFlow flow = getInitialFlow();
+
+ final VersionedParameter sensitiveParameter = new VersionedParameter();
+ sensitiveParameter.setName(SENSITIVE_PARAMETER_NAME);
+ sensitiveParameter.setValue("sensitive-value");
+ sensitiveParameter.setSensitive(true);
+ sensitiveParameter.setProvided(false);
+
+ final VersionedAsset asset = new VersionedAsset();
+ asset.setIdentifier(assetIdentifier);
+ asset.setName(ASSET_NAME);
+
+ final VersionedParameter assetParameter = new VersionedParameter();
+ assetParameter.setName(ASSET_PARAMETER_NAME);
+ assetParameter.setSensitive(false);
+ assetParameter.setProvided(false);
+ assetParameter.setReferencedAssets(List.of(asset));
+
+ final VersionedParameterContext parameterContext = new
VersionedParameterContext();
+ parameterContext.setName(PARAMETER_CONTEXT_NAME);
+ parameterContext.setParameters(Set.of(sensitiveParameter,
assetParameter));
+ flow.setParameterContexts(Map.of(PARAMETER_CONTEXT_NAME,
parameterContext));
+ flow.getFlowContents().setParameterContextName(PARAMETER_CONTEXT_NAME);
+
+ final VersionedProcessor generateFlowFile =
VersionedFlowUtils.findProcessor(flow.getFlowContents(),
+ processor ->
processor.getType().equals("org.apache.nifi.processors.standard.GenerateFlowFile"))
+ .orElseThrow();
+ generateFlowFile.getProperties().put("Custom Text", "#{" +
ASSET_PARAMETER_NAME + "}");
+
+ return flow;
+ }
+}
diff --git
a/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/resources/META-INF/services/org.apache.nifi.components.connector.Connector
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/resources/META-INF/services/org.apache.nifi.components.connector.Connector
index 379d9026452..391ca0e8759 100644
---
a/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/resources/META-INF/services/org.apache.nifi.components.connector.Connector
+++
b/nifi-connector-mock-bundle/nifi-connector-mock-test-bundle/nifi-connector-mock-test-connectors/src/main/resources/META-INF/services/org.apache.nifi.components.connector.Connector
@@ -16,4 +16,5 @@
org.apache.nifi.mock.connectors.AllowableValuesConnector
org.apache.nifi.mock.connectors.GenerateAndLog
org.apache.nifi.mock.connectors.MissingBundleConnector
-org.apache.nifi.mock.connectors.CronScheduleConnector
\ No newline at end of file
+org.apache.nifi.mock.connectors.CronScheduleConnector
+org.apache.nifi.mock.connectors.FlowSnapshotConnector