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 3766b74a50d NIFI-16320 Prevent connector GET from racing
working-context recreation. (#11647)
3766b74a50d is described below
commit 3766b74a50d507ff8ee44f5d3127d40f852f183c
Author: Mark Payne <[email protected]>
AuthorDate: Wed Sep 9 08:35:51 2026 -0400
NIFI-16320 Prevent connector GET from racing working-context recreation.
(#11647)
Connector GET can sync from the provider while the working flow context is
being recreated. Keep a live working context visible to callers by swapping a
per-context holder under the monitor. Destroy the previous working process
group before creating the replacement so clustered load-balanced connections
are not registered twice, then destroy leftover holders after unlock once their
use count reaches zero.
---
.../connector/StandardConnectorNode.java | 496 +++++++++++++++------
.../connector/TestStandardConnectorNode.java | 245 +++++++---
2 files changed, 540 insertions(+), 201 deletions(-)
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java
index 6a883dd710d..5d3f7c582dc 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/components/connector/StandardConnectorNode.java
@@ -120,7 +120,9 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
private final AtomicReference<CompletableFuture<Void>> drainFutureRef =
new AtomicReference<>();
private volatile ValidationResult unresolvedBundleValidationResult = null;
- private volatile FrameworkFlowContext workingFlowContext;
+ private final Object workingFlowContextLock = new Object();
+ private volatile WorkingFlowContextState workingFlowContextState = new
WorkingFlowContextState(null);
+ private boolean workingContextReplacementInProgress;
private volatile String name;
private volatile FrameworkConnectorInitializationContext
initializationContext;
@@ -297,20 +299,26 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
}
logger.debug("Preparing {} for update", this);
- try (final NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
- getConnector().prepareForUpdate(workingFlowContext,
activeFlowContext);
- stateTransition.setCurrentState(ConnectorState.UPDATING);
- logger.debug("Successfully prepared {} for update", this);
- } catch (final Throwable t) {
- logger.error("Failed to prepare update for {}", this, t);
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext workingContext =
workingContextState.getContext();
+ try {
+ try (final NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
+ getConnector().prepareForUpdate(workingContext,
activeFlowContext);
+ stateTransition.setCurrentState(ConnectorState.UPDATING);
+ logger.debug("Successfully prepared {} for update", this);
+ } catch (final Throwable t) {
+ logger.error("Failed to prepare update for {}", this, t);
- try {
- abortUpdate(t);
- } catch (final Throwable abortFailure) {
- logger.error("Failed to abort update preparation for {}",
this, abortFailure);
- }
+ try {
+ abortUpdate(t);
+ } catch (final Throwable abortFailure) {
+ logger.error("Failed to abort update preparation for {}",
this, abortFailure);
+ }
- throw t;
+ throw t;
+ }
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
}
@@ -395,22 +403,77 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
// two lists actually diverge.
applyUpdate(inheritContext);
- // Tear down the working context that applyUpdate created aliased to
active, and rebuild it around an
- // independent configuration seeded from migratedWorkingProperties.
Then fire onConfigurationStepConfigured
- // for every step so renamed steps trigger the flow-builder callback
under their new name and any
- // value-derived flow state (resolved asset paths, secret values,
etc.) is populated against the fresh
- // working context.
- destroyWorkingContext();
+ // Replace the working context that applyUpdate created aliased to
active with an independent context
+ // seeded from migratedWorkingProperties. Then fire
onConfigurationStepConfigured for every step so
+ // renamed steps trigger the flow-builder callback under their new
name and any value-derived flow
+ // state (resolved asset paths, secret values, etc.) is populated
against the fresh working context.
final MutableConnectorConfigurationContext workingConfigContext =
createConfigurationContext(migratedWorkingProperties);
- workingFlowContext =
flowContextFactory.createWorkingFlowContext(identifier,
connectorDetails.getComponentLog(), workingConfigContext, flowContextBundle);
+ final WorkingFlowContextState independentWorkingContextState =
installReplacementWorkingFlowContext(workingConfigContext, flowContextBundle,
true);
+ final FrameworkFlowContext independentWorkingContext =
independentWorkingContextState.getContext();
+
getComponentLog().info("Working Flow Context has been rebuilt with
independent configuration");
- for (final String stepName : migratedWorkingProperties.keySet()) {
- notifyStepConfigured(stepName);
+
+ try {
+ for (final String stepName : migratedWorkingProperties.keySet()) {
+ notifyStepConfigured(stepName, independentWorkingContext);
+ }
+ } finally {
+ releaseWorkingFlowContext(independentWorkingContextState);
}
logger.debug("Successfully inherited configuration for {}", this);
}
+ /**
+ * Removes the current working process group before creating its
replacement. Working-context copies reuse the same
+ * connection identifiers as the active flow, and the cluster load-balance
client registry allows only one
+ * registration per connection ID, so the previous group must be gone
before the factory copies the active group.
+ * The published working context is never set to null: callers that arrive
while the previous group is being
+ * destroyed still see that context, and callers that arrive while the
replacement is created wait on the monitor.
+ */
+ private WorkingFlowContextState installReplacementWorkingFlowContext(final
MutableConnectorConfigurationContext configurationContext, final Bundle bundle,
final boolean incrementUseCount) {
+ final WorkingFlowContextState previousWorkingFlowContextState;
+ final boolean destroyPrevious;
+ synchronized (workingFlowContextLock) {
+ while (workingContextReplacementInProgress) {
+ try {
+ workingFlowContextLock.wait();
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Interrupted while waiting
to replace the working flow context of " + this, e);
+ }
+ }
+
+ workingContextReplacementInProgress = true;
+ previousWorkingFlowContextState = workingFlowContextState;
+ previousWorkingFlowContextState.retire();
+ destroyPrevious =
previousWorkingFlowContextState.claimDestruction();
+ }
+
+ if (destroyPrevious) {
+
destroyWorkingContext(previousWorkingFlowContextState.getContext());
+ }
+
+ final WorkingFlowContextState replacementWorkingFlowContextState;
+ synchronized (workingFlowContextLock) {
+ try {
+ final FrameworkFlowContext replacementWorkingFlowContext =
flowContextFactory.createWorkingFlowContext(identifier,
+ connectorDetails.getComponentLog(), configurationContext,
bundle);
+ replacementWorkingFlowContextState = new
WorkingFlowContextState(replacementWorkingFlowContext);
+ if (incrementUseCount) {
+ replacementWorkingFlowContextState.incrementUseCount();
+ }
+ workingFlowContextState = replacementWorkingFlowContextState;
+ } finally {
+ workingContextReplacementInProgress = false;
+ workingFlowContextLock.notifyAll();
+ }
+ }
+
+ getComponentLog().info("Working Flow Context has been set");
+ return replacementWorkingFlowContextState;
+ }
+
private Map<String, StepConfiguration> migrateProperties(final
List<VersionedConfigurationStep> flowConfiguration) {
// Preserve persisted step order so the notifyStepConfigured loop in
inheritConfiguration fires in a
// deterministic order matching the flow definition.
@@ -504,8 +567,10 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
@Override
public void applyUpdate() throws FlowUpdateException {
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext contextToInherit =
workingContextState.getContext();
try {
- applyUpdate(workingFlowContext);
+ applyUpdate(contextToInherit);
} catch (final FlowUpdateException e) {
// Since we failed to update, make sure that we stop the
Connector. Note that we do not do this for all
// throwables because IllegalStateException for example indicates
that we did not even attempt to perform the update.
@@ -516,6 +581,8 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
}
throw e;
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
}
@@ -539,7 +606,7 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
// The update has been completed. Tear down and recreate the
working flow context to ensure it is in a clean state.
resetValidationState();
- recreateWorkingFlowContext();
+ replaceWorkingFlowContextFromActive();
} catch (final Throwable t) {
logger.error("Failed to finish update for {}", this, t);
stateTransition.setCurrentState(ConnectorState.UPDATE_FAILED);
@@ -553,20 +620,36 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
logger.info("Successfully applied update for {}", this);
}
- private void destroyWorkingContext() {
- if (this.workingFlowContext == null) {
+ private void destroyWorkingContext(final FrameworkFlowContext context) {
+ if (context == null) {
return;
}
try {
- workingFlowContext.getManagedProcessGroup().purge().get(1,
TimeUnit.MINUTES);
+ context.getManagedProcessGroup().purge().get(1, TimeUnit.MINUTES);
} catch (final Exception e) {
logger.warn("Failed to purge working flow context for {}", this,
e);
}
-
flowManager.onProcessGroupRemoved(workingFlowContext.getManagedProcessGroup());
+ flowManager.onProcessGroupRemoved(context.getManagedProcessGroup());
+ }
- this.workingFlowContext = null;
+ private WorkingFlowContextState acquireWorkingFlowContext() {
+ synchronized (workingFlowContextLock) {
+ workingFlowContextState.incrementUseCount();
+ return workingFlowContextState;
+ }
+ }
+
+ private void releaseWorkingFlowContext(final WorkingFlowContextState
workingContextState) {
+ final boolean destroyNow;
+ synchronized (workingFlowContextLock) {
+ destroyNow = workingContextState.decrementUseCount();
+ }
+
+ if (destroyNow) {
+ destroyWorkingContext(workingContextState.getContext());
+ }
}
@Override
@@ -574,9 +657,14 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
stateTransition.setCurrentState(ConnectorState.UPDATE_FAILED);
stateTransition.setDesiredState(ConnectorState.UPDATE_FAILED);
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext workingContext =
workingContextState.getContext();
try (final NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
- getConnector().abortUpdate(workingFlowContext, cause);
+ getConnector().abortUpdate(workingContext, cause);
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
+
logger.debug("Aborted update for {}", this);
}
@@ -596,33 +684,56 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
+ " while it is in Troubleshooting mode; exit Troubleshooting
mode before modifying Connector configuration.");
}
- setConfiguration(stepName, configuration, false);
- }
+ final FrameworkFlowContext workingContext;
+ final WorkingFlowContextState workingContextState;
+ synchronized (workingFlowContextLock) {
+ workingContextState = this.workingFlowContextState;
+ workingContext = workingContextState.getContext();
+ final ConfigurationUpdateResult updateResult =
workingContext.getConfigurationContext().setProperties(stepName, configuration);
+ if (updateResult == ConfigurationUpdateResult.NO_CHANGES) {
+ return;
+ }
- private void setConfiguration(final String stepName, final
StepConfiguration configuration, final boolean
forceOnConfigurationStepConfigured) throws FlowUpdateException {
- final ConfigurationUpdateResult updateResult =
workingFlowContext.getConfigurationContext().setProperties(stepName,
configuration);
- if (updateResult == ConfigurationUpdateResult.NO_CHANGES &&
!forceOnConfigurationStepConfigured) {
- return;
+ workingContextState.incrementUseCount();
+ }
+
+ try {
+ notifyStepConfigured(stepName, workingContext);
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
- notifyStepConfigured(stepName);
}
@Override
public void replaceWorkingConfiguration(final String stepName, final
StepConfiguration configuration) throws FlowUpdateException {
// The configuration provider's view is authoritative: any property
absent from the provided
// configuration is removed from the step.
- final ConfigurationUpdateResult updateResult =
workingFlowContext.getConfigurationContext().replaceProperties(stepName,
configuration);
- if (updateResult == ConfigurationUpdateResult.NO_CHANGES) {
- return;
+ final FrameworkFlowContext workingContext;
+ final WorkingFlowContextState workingContextState;
+ synchronized (workingFlowContextLock) {
+ workingContextState = this.workingFlowContextState;
+ workingContext = workingContextState.getContext();
+ final ConfigurationUpdateResult updateResult =
workingContext.getConfigurationContext().replaceProperties(stepName,
configuration);
+ if (updateResult == ConfigurationUpdateResult.NO_CHANGES) {
+ return;
+ }
+
+ workingContextState.incrementUseCount();
+ }
+
+ try {
+ notifyStepConfigured(stepName, workingContext);
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
- notifyStepConfigured(stepName);
}
- private void notifyStepConfigured(final String stepName) throws
FlowUpdateException {
+ private void notifyStepConfigured(final String stepName, final
FrameworkFlowContext workingContext) throws FlowUpdateException {
final Connector connector = connectorDetails.getConnector();
try (final NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager, connector.getClass(),
getIdentifier())) {
logger.debug("Notifying {} of configuration change for
configuration step {}", this, stepName);
- connector.onConfigurationStepConfigured(stepName,
workingFlowContext);
+ connector.onConfigurationStepConfigured(stepName, workingContext);
+
logger.debug("Successfully notified {} of configuration change for
step {}", this, stepName);
} catch (final FlowUpdateException e) {
throw e;
@@ -1165,29 +1276,41 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
@Override
public List<DescribedValue> fetchAllowableValues(final String stepName,
final String propertyName) {
- if (workingFlowContext == null) {
- throw new IllegalStateException("Cannot fetch Allowable Values for
%s.%s because %s is not being updated.".formatted(
- stepName, propertyName, this));
- }
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext workingContext =
workingContextState.getContext();
+ try {
+ if (workingContext == null) {
+ throw new IllegalStateException("Cannot fetch Allowable Values
for %s.%s because %s is not being updated.".formatted(
+ stepName, propertyName, this));
+ }
- workingFlowContext.getConfigurationContext().resolvePropertyValues();
+ workingContext.getConfigurationContext().resolvePropertyValues();
- try (NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
- return getConnector().fetchAllowableValues(stepName, propertyName,
workingFlowContext);
+ try (NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
+ return getConnector().fetchAllowableValues(stepName,
propertyName, workingContext);
+ }
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
}
@Override
public List<DescribedValue> fetchAllowableValues(final String stepName,
final String propertyName, final String filter) {
- if (workingFlowContext == null) {
- throw new IllegalStateException("Cannot fetch Allowable Values for
%s.%s because %s is not being updated.".formatted(
- stepName, propertyName, this));
- }
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext workingContext =
workingContextState.getContext();
+ try {
+ if (workingContext == null) {
+ throw new IllegalStateException("Cannot fetch Allowable Values
for %s.%s because %s is not being updated.".formatted(
+ stepName, propertyName, this));
+ }
- workingFlowContext.getConfigurationContext().resolvePropertyValues();
+ workingContext.getConfigurationContext().resolvePropertyValues();
- try (NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
- return getConnector().fetchAllowableValues(stepName, propertyName,
workingFlowContext, filter);
+ try (NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
+ return getConnector().fetchAllowableValues(stepName,
propertyName, workingContext, filter);
+ }
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
}
@@ -1270,12 +1393,14 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
public boolean isModified() {
final Map<String, ConfigurationStep> stepsByName =
indexConfigurationSteps();
- if (configurationDiffersFromDefaults(activeFlowContext, stepsByName)
- || configurationDiffersFromDefaults(workingFlowContext,
stepsByName)) {
- return true;
- }
+ synchronized (workingFlowContextLock) {
+ final FrameworkFlowContext workingFlowContext =
workingFlowContextState.getContext();
+ if (configurationDiffersFromDefaults(activeFlowContext,
stepsByName) || configurationDiffersFromDefaults(workingFlowContext,
stepsByName)) {
+ return true;
+ }
- return hasComponentState(activeFlowContext) ||
hasComponentState(workingFlowContext);
+ return hasComponentState(activeFlowContext) ||
hasComponentState(workingFlowContext);
+ }
}
/**
@@ -1541,30 +1666,39 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
@Override
public void recreateWorkingFlowContext() {
- destroyWorkingContext();
- workingFlowContext =
flowContextFactory.createWorkingFlowContext(identifier,
- connectorDetails.getComponentLog(),
activeFlowContext.getConfigurationContext(), activeFlowContext.getBundle());
+ replaceWorkingFlowContextFromActive();
+ }
- getComponentLog().info("Working Flow Context has been set");
+ private void replaceWorkingFlowContextFromActive() {
+ final boolean notifyReplacementSteps = initializationContext != null;
+ final WorkingFlowContextState replacementWorkingFlowContextState =
installReplacementWorkingFlowContext(
+ activeFlowContext.getConfigurationContext(),
activeFlowContext.getBundle(), notifyReplacementSteps);
+ final FrameworkFlowContext replacementWorkingFlowContext =
replacementWorkingFlowContextState.getContext();
- // Re-fire onConfigurationStepConfigured for every step so flow
parameters derived from the
- // configuration (e.g., resolved asset paths, secret values) are
refreshed against the new
- // working context. Step failures are logged so the remaining steps
can still be refreshed.
- // Skipped before the connector has been initialized because there is
no flow to update yet.
- if (initializationContext == null) {
- return;
- }
- final ConnectorConfiguration config =
workingFlowContext.getConfigurationContext().toConnectorConfiguration();
- for (final NamedStepConfiguration stepConfig :
config.getNamedStepConfigurations()) {
- try {
- setConfiguration(stepConfig.stepName(),
stepConfig.configuration(), true);
- } catch (final Exception e) {
- logger.warn("Failed to refresh resolved configuration for step
[{}] of {}",
- stepConfig.stepName(), this, e);
+ try {
+ // Re-fire onConfigurationStepConfigured for every step so flow
parameters derived from the
+ // configuration (e.g., resolved asset paths, secret values) are
refreshed against the new
+ // working context. Step failures are logged so the remaining
steps can still be refreshed.
+ // Skipped before the connector has been initialized because there
is no flow to update yet.
+ // Configuration values are already present on the replacement;
this loop only notifies.
+ if (notifyReplacementSteps) {
+ final ConnectorConfiguration config =
replacementWorkingFlowContext.getConfigurationContext().toConnectorConfiguration();
+ for (final NamedStepConfiguration stepConfig :
config.getNamedStepConfigurations()) {
+ try {
+ notifyStepConfigured(stepConfig.stepName(),
replacementWorkingFlowContext);
+ } catch (final Exception e) {
+ logger.warn("Failed to refresh resolved configuration
for step [{}] of {}",
+ stepConfig.stepName(), this, e);
+ }
+ }
+
+ getComponentLog().info("Working Flow Context configuration has
been refreshed");
+ }
+ } finally {
+ if (notifyReplacementSteps) {
+ releaseWorkingFlowContext(replacementWorkingFlowContextState);
}
}
-
- getComponentLog().info("Working Flow Context configuration has been
refreshed");
}
@Override
@@ -1592,66 +1726,73 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
}
logger.debug("Verifying configuration step {} for {}", stepName, this);
- final List<ConfigVerificationResult> results = new ArrayList<>();
- try (final NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext workingContext =
workingContextState.getContext();
+ try {
+ final List<ConfigVerificationResult> results = new ArrayList<>();
+ try (final NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
+
+ final Optional<ConfigurationStep> optionalStep =
getConfigurationStep(stepName);
+ if (optionalStep.isEmpty()) {
+ results.add(new ConfigVerificationResult.Builder()
+ .verificationStepName("Property Validation")
+ .outcome(Outcome.FAILED)
+ .explanation("Configuration step with name '" +
stepName + "' does not exist.")
+ .build());
+ return results;
+ }
- final Optional<ConfigurationStep> optionalStep =
getConfigurationStep(stepName);
- if (optionalStep.isEmpty()) {
- results.add(new ConfigVerificationResult.Builder()
- .verificationStepName("Property Validation")
- .outcome(Outcome.FAILED)
- .explanation("Configuration step with name '" + stepName +
"' does not exist.")
- .build());
- return results;
- }
+ final ConfigurationStep configurationStep = optionalStep.get();
+ final List<SecretReference> invalidSecretRefs = new
ArrayList<>();
+ final List<AssetReference> invalidAssetRefs = new
ArrayList<>();
- final ConfigurationStep configurationStep = optionalStep.get();
- final List<SecretReference> invalidSecretRefs = new ArrayList<>();
- final List<AssetReference> invalidAssetRefs = new ArrayList<>();
- // Bypass the Secret value cache during verification so the user
sees results based on the current
- // Secret values rather than potentially stale cached values
awaiting TTL expiration.
- final Map<String, String> resolvedPropertyOverrides =
resolvePropertyReferences(configurationStep, configurationOverrides,
invalidSecretRefs, invalidAssetRefs, false);
+ // Bypass the Secret value cache during verification so the
user sees results based on the current
+ // Secret values rather than potentially stale cached values
awaiting TTL expiration.
+ final Map<String, String> resolvedPropertyOverrides =
resolvePropertyReferences(workingContext, configurationStep,
configurationOverrides, invalidSecretRefs, invalidAssetRefs, false);
- final DescribedValueProvider allowableValueProvider = (step,
propertyName) -> fetchAllowableValues(step, propertyName, workingFlowContext);
+ final DescribedValueProvider allowableValueProvider = (step,
propertyName) -> fetchAllowableValues(step, propertyName, workingContext);
- final MutableConnectorConfigurationContext configContext =
workingFlowContext.getConfigurationContext().createWithOverrides(stepName,
resolvedPropertyOverrides);
- final ConnectorConfiguration connectorConfig =
configContext.toConnectorConfiguration();
- final ParameterContextFacade paramContext =
workingFlowContext.getParameterContext();
- final ConnectorValidationContext validationContext = new
StandardConnectorValidationContext(connectorConfig, allowableValueProvider,
paramContext);
+ final MutableConnectorConfigurationContext configContext =
workingContext.getConfigurationContext().createWithOverrides(stepName,
resolvedPropertyOverrides);
+ final ConnectorConfiguration connectorConfig =
configContext.toConnectorConfiguration();
+ final ParameterContextFacade paramContext =
workingContext.getParameterContext();
+ final ConnectorValidationContext validationContext = new
StandardConnectorValidationContext(connectorConfig, allowableValueProvider,
paramContext);
- final List<ValidationResult> validationResults = new ArrayList<>();
- validatePropertyReferences(configurationStep,
configurationOverrides, validationResults);
+ final List<ValidationResult> validationResults = new
ArrayList<>();
+ validatePropertyReferences(configurationStep,
configurationOverrides, validationResults);
- // If there are any invalid secrets or assets referenced, add
Validation Results for them.
- addInvalidReferenceResults(validationResults, invalidSecretRefs,
invalidAssetRefs);
+ // If there are any invalid secrets or assets referenced, add
Validation Results for them.
+ addInvalidReferenceResults(validationResults,
invalidSecretRefs, invalidAssetRefs);
- // If there are any framework-level validation failures, we do not
run the Connector-specific validation because
- // doing so would mean that we must provide weak guarantees about
the state of the configuration when the Connector's
- // validation is invoked. But if there are no framework-level
validation failures, we can proceed to invoke the
- // Connector's validation logic.
- if (validationResults.isEmpty()) {
- final List<ValidationResult> implValidationResults =
getConnector().validateConfigurationStep(configurationStep, configContext,
validationContext);
- validationResults.addAll(implValidationResults);
- }
+ // If there are any framework-level validation failures, we do
not run the Connector-specific validation because
+ // doing so would mean that we must provide weak guarantees
about the state of the configuration when the Connector's
+ // validation is invoked. But if there are no framework-level
validation failures, we can proceed to invoke the
+ // Connector's validation logic.
+ if (validationResults.isEmpty()) {
+ final List<ValidationResult> implValidationResults =
getConnector().validateConfigurationStep(configurationStep, configContext,
validationContext);
+ validationResults.addAll(implValidationResults);
+ }
- final List<ConfigVerificationResult> invalidConfigResults =
validationResults.stream()
- .filter(result -> !result.isValid())
- .map(this::createConfigVerificationResult)
- .toList();
+ final List<ConfigVerificationResult> invalidConfigResults =
validationResults.stream()
+ .filter(result -> !result.isValid())
+ .map(this::createConfigVerificationResult)
+ .toList();
- if (invalidConfigResults.isEmpty()) {
- results.add(new ConfigVerificationResult.Builder()
- .verificationStepName("Property Validation")
- .outcome(Outcome.SUCCESSFUL)
- .build());
+ if (invalidConfigResults.isEmpty()) {
+ results.add(new ConfigVerificationResult.Builder()
+ .verificationStepName("Property Validation")
+ .outcome(Outcome.SUCCESSFUL)
+ .build());
-
results.addAll(getConnector().verifyConfigurationStep(stepName,
resolvedPropertyOverrides, workingFlowContext));
- } else {
- results.addAll(invalidConfigResults);
- }
+
results.addAll(getConnector().verifyConfigurationStep(stepName,
resolvedPropertyOverrides, workingContext));
+ } else {
+ results.addAll(invalidConfigResults);
+ }
- logger.debug("Completed verification of configuration step {} for
{}", stepName, this);
- return results;
+ logger.debug("Completed verification of configuration step {}
for {}", stepName, this);
+ return results;
+ }
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
}
@@ -1664,12 +1805,12 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
.build();
}
- private Map<String, String> resolvePropertyReferences(final
ConfigurationStep configurationStep, final StepConfiguration
configurationOverrides,
+ private Map<String, String> resolvePropertyReferences(final
FrameworkFlowContext workingContext, final ConfigurationStep configurationStep,
final StepConfiguration configurationOverrides,
final
List<SecretReference> invalidSecretRefs, final List<AssetReference>
invalidAssetRefs, final boolean useCache) {
final Map<String, String> resolvedProperties = new HashMap<>();
final Map<String, ConnectorPropertyDescriptor> descriptorLookup =
buildPropertyDescriptorLookup(configurationStep);
- final StepConfiguration effectiveConfiguration =
createEffectiveStepConfiguration(configurationStep.getName(),
configurationOverrides);
+ final StepConfiguration effectiveConfiguration =
createEffectiveStepConfiguration(workingContext, configurationStep.getName(),
configurationOverrides);
try {
// Secret References can be expensive to lookup so we don't want
to call getSecret() for each one. Instead, we
@@ -1749,9 +1890,9 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
return resolvedProperties;
}
- private StepConfiguration createEffectiveStepConfiguration(final String
stepName, final StepConfiguration configurationOverrides) {
+ private StepConfiguration createEffectiveStepConfiguration(final
FrameworkFlowContext workingContext, final String stepName, final
StepConfiguration configurationOverrides) {
final Map<String, ConnectorValueReference> effectiveProperties = new
HashMap<>();
- final NamedStepConfiguration workingStepConfiguration =
workingFlowContext.getConfigurationContext()
+ final NamedStepConfiguration workingStepConfiguration =
workingContext.getConfigurationContext()
.toConnectorConfiguration()
.getNamedStepConfiguration(stepName);
if (workingStepConfiguration != null) {
@@ -1932,10 +2073,16 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
return results;
}
- workingFlowContext.getConfigurationContext().resolvePropertyValues();
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ final FrameworkFlowContext workingContext =
workingContextState.getContext();
+ try {
+ workingContext.getConfigurationContext().resolvePropertyValues();
- try (NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
- results.addAll(getConnector().verify(workingFlowContext));
+ try (NarCloseable ignored =
NarCloseable.withComponentNarLoader(extensionManager,
getConnector().getClass(), getIdentifier())) {
+ results.addAll(getConnector().verify(workingContext));
+ }
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
}
logger.debug("Completed verification for {}", this);
@@ -1972,7 +2119,9 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
@Override
public FrameworkFlowContext getWorkingFlowContext() {
- return workingFlowContext;
+ synchronized (workingFlowContextLock) {
+ return workingFlowContextState.getContext();
+ }
}
@Override
@@ -2230,14 +2379,16 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
}
private boolean hasWorkingConfigurationChanges() {
- final FrameworkFlowContext workingContext = this.workingFlowContext;
- if (workingContext == null) {
- return false;
- }
+ synchronized (workingFlowContextLock) {
+ final FrameworkFlowContext workingContext =
workingFlowContextState.getContext();
+ if (workingContext == null) {
+ return false;
+ }
- final ConnectorConfiguration activeConfig =
activeFlowContext.getConfigurationContext().toConnectorConfiguration();
- final ConnectorConfiguration workingConfig =
workingContext.getConfigurationContext().toConnectorConfiguration();
- return !Objects.equals(activeConfig, workingConfig);
+ final ConnectorConfiguration activeConfig =
activeFlowContext.getConfigurationContext().toConnectorConfiguration();
+ final ConnectorConfiguration workingConfig =
workingContext.getConfigurationContext().toConnectorConfiguration();
+ return !Objects.equals(activeConfig, workingConfig);
+ }
}
@Override
@@ -2429,7 +2580,12 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
final List<AssetReference> invalidAssets = new ArrayList<>();
// Regular validation may run frequently, so cached Secret values
are used here to avoid
// repeatedly fetching from the underlying Secret Providers on
every validation cycle.
- resolvePropertyReferences(step, stepConfiguration, invalidSecrets,
invalidAssets, true);
+ final WorkingFlowContextState workingContextState =
acquireWorkingFlowContext();
+ try {
+ resolvePropertyReferences(workingContextState.getContext(),
step, stepConfiguration, invalidSecrets, invalidAssets, true);
+ } finally {
+ releaseWorkingFlowContext(workingContextState);
+ }
addInvalidReferenceResults(allResults, invalidSecrets,
invalidAssets);
}
}
@@ -2540,6 +2696,52 @@ public class StandardConnectorNode implements
ConnectorNode, GroupedComponent {
return allowableValues;
}
+ private static final class WorkingFlowContextState {
+ private final FrameworkFlowContext context;
+ private int useCount;
+ private boolean retired;
+ private boolean destroyed;
+
+ private WorkingFlowContextState(final FrameworkFlowContext context) {
+ this.context = context;
+ }
+
+ private FrameworkFlowContext getContext() {
+ return context;
+ }
+
+ private void incrementUseCount() {
+ useCount++;
+ }
+
+ private boolean decrementUseCount() {
+ if (useCount == 0) {
+ throw new IllegalStateException("Cannot release a working flow
context that is not in use");
+ }
+
+ useCount--;
+ return retired && useCount == 0 && claimDestruction();
+ }
+
+ private boolean retire() {
+ if (retired) {
+ return false;
+ }
+
+ retired = true;
+ return useCount == 0;
+ }
+
+ private boolean claimDestruction() {
+ if (destroyed) {
+ return false;
+ }
+
+ destroyed = true;
+ return true;
+ }
+ }
+
@Override
public boolean equals(final Object o) {
if (o == null || getClass() != o.getClass()) {
diff --git
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java
index d729c888e72..05682b556a6 100644
---
a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java
+++
b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/components/connector/TestStandardConnectorNode.java
@@ -70,6 +70,7 @@ import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
@@ -497,19 +498,7 @@ public class TestStandardConnectorNode {
}
@Test
- public void testSetConfigurationCallsOnConfigured() throws
FlowUpdateException {
- final TrackingConnector trackingConnector = new TrackingConnector();
- final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
- assertEquals(ConnectorState.STOPPED, connectorNode.getCurrentState());
-
- connectorNode.transitionStateForUpdating();
- connectorNode.prepareForUpdate();
- connectorNode.setConfiguration("testGroup", createStepConfiguration());
- connectorNode.applyUpdate();
- }
-
- @Test
- public void
testSetConfigurationCallsOnPropertyGroupConfiguredForChangedConfigurationSteps()
throws FlowUpdateException, ExecutionException, InterruptedException,
TimeoutException {
+ public void testSetConfigurationCallsOnStepConfiguredWhenChanged() throws
FlowUpdateException {
final TrackingConnector trackingConnector = new TrackingConnector();
final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
assertEquals(ConnectorState.STOPPED, connectorNode.getCurrentState());
@@ -523,26 +512,8 @@ public class TestStandardConnectorNode {
connectorNode.transitionStateForUpdating();
connectorNode.prepareForUpdate();
connectorNode.setConfiguration("configurationStep1",
createStepConfiguration(Map.of("prop1", "value2")));
- connectorNode.applyUpdate();
-
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("configurationStep1"));
- }
-
- @Test
- public void testDiscardWorkingConfigurationCallsOnStepConfigured() throws
FlowUpdateException {
- final TrackingConnector trackingConnector = new TrackingConnector();
- final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
-
- connectorNode.transitionStateForUpdating();
- connectorNode.prepareForUpdate();
- connectorNode.setConfiguration("step1",
createStepConfiguration(Map.of("prop1", "value1")));
- connectorNode.applyUpdate();
-
- trackingConnector.reset();
- connectorNode.discardWorkingConfiguration();
-
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("configurationStep1"));
}
@Test
@@ -563,12 +534,12 @@ public class TestStandardConnectorNode {
connectorNode.setConfiguration("step1",
createStepConfiguration(Map.of("prop1", "value2")));
connectorNode.applyUpdate();
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step2"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("step1"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("step2"));
}
@Test
- public void
testDiscardWorkingConfigurationCallsOnStepConfiguredForMultipleSteps() throws
FlowUpdateException {
+ public void
testDiscardWorkingConfigurationCallsOnStepConfiguredForEveryStep() throws
FlowUpdateException {
final TrackingConnector trackingConnector = new TrackingConnector();
final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
@@ -582,8 +553,8 @@ public class TestStandardConnectorNode {
connectorNode.discardWorkingConfiguration();
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step2"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("step1"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("step2"));
}
@Test
@@ -602,7 +573,7 @@ public class TestStandardConnectorNode {
connectorNode.prepareForUpdate();
connectorNode.applyUpdate();
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("step1"));
}
@Test
@@ -619,7 +590,7 @@ public class TestStandardConnectorNode {
connectorNode.replaceWorkingConfiguration("step1",
createStepConfiguration(Map.of("propA", "newA")));
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
+
assertTrue(trackingConnector.wasOnConfigurationStepConfiguredCalled("step1"));
final ConnectorConfiguration workingConfig =
connectorNode.getWorkingFlowContext().getConfigurationContext().toConnectorConfiguration();
final NamedStepConfiguration namedStep =
workingConfig.getNamedStepConfigurations().iterator().next();
assertEquals("step1", namedStep.stepName());
@@ -640,29 +611,145 @@ public class TestStandardConnectorNode {
connectorNode.replaceWorkingConfiguration("step1",
createStepConfiguration(Map.of("propA", "valueA")));
-
assertFalse(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
+
assertFalse(trackingConnector.wasOnConfigurationStepConfiguredCalled("step1"));
}
@Test
- public void
testDiscardWorkingConfigurationFiresOnConfiguredForEveryWorkingStep() throws
FlowUpdateException {
- final TrackingConnector trackingConnector = new TrackingConnector();
- final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
+ @Timeout(10)
+ public void
testReplaceWorkingConfigurationWaitsForWorkingContextRecreation() throws
Exception {
+ final BlockingWorkingFlowContextFactory blockingFlowContextFactory =
new BlockingWorkingFlowContextFactory(flowContextFactory);
+ flowContextFactory = blockingFlowContextFactory;
+ final StandardConnectorNode connectorNode = createConnectorNode(new
TrackingConnector());
connectorNode.transitionStateForUpdating();
connectorNode.prepareForUpdate();
- connectorNode.setConfiguration("step1",
createStepConfiguration(Map.of("propA", "valueA")));
- connectorNode.setConfiguration("step2",
createStepConfiguration(Map.of("propB", "valueB")));
+ connectorNode.setConfiguration("step1",
createStepConfiguration(Map.of("propA", "oldA")));
connectorNode.applyUpdate();
+ blockingFlowContextFactory.blockNextWorkingContextCreation();
- trackingConnector.reset();
+ final ExecutorService executor = Executors.newFixedThreadPool(2);
+ try {
+ final Future<?> recreationFuture =
executor.submit(connectorNode::recreateWorkingFlowContext);
+
assertTrue(blockingFlowContextFactory.awaitWorkingContextCreation(5,
TimeUnit.SECONDS));
+
+ final CountDownLatch replaceStarted = new CountDownLatch(1);
+ final Future<?> replacementFuture = executor.submit(() -> {
+ replaceStarted.countDown();
+ connectorNode.replaceWorkingConfiguration("step1",
createStepConfiguration(Map.of("propA", "newA")));
+ return null;
+ });
+ assertTrue(replaceStarted.await(5, TimeUnit.SECONDS));
- // Recreating the working flow context from the active flow must fire
onConfigurationStepConfigured
- // for every working configuration step so that flow parameters
derived from the configuration
- // (resolved asset paths, secrets, etc.) are refreshed.
- connectorNode.discardWorkingConfiguration();
+ try {
+ assertThrows(TimeoutException.class, () ->
replacementFuture.get(STOP_NOT_EXPECTED_MILLIS, TimeUnit.MILLISECONDS));
+ } finally {
+ blockingFlowContextFactory.releaseWorkingContextCreation();
+ }
+
+ recreationFuture.get(5, TimeUnit.SECONDS);
+ replacementFuture.get(5, TimeUnit.SECONDS);
+ } finally {
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+
+ final ConnectorConfiguration workingConfiguration =
connectorNode.getWorkingFlowContext().getConfigurationContext().toConnectorConfiguration();
+ final NamedStepConfiguration namedStep =
workingConfiguration.getNamedStepConfigurations().iterator().next();
+ assertEquals(Map.of("propA", new StringLiteralValue("newA")),
namedStep.configuration().getPropertyValues());
+ }
+
+ @Test
+ @Timeout(10)
+ public void testRecreationRefreshDoesNotOverwriteConcurrentReplace()
throws Exception {
+ final CountDownLatch refreshStarted = new CountDownLatch(1);
+ final CountDownLatch permitRefresh = new CountDownLatch(1);
+ final AtomicBoolean blockNextRefresh = new AtomicBoolean();
+ final AtomicReference<String> refreshingStepName = new
AtomicReference<>();
+ final TrackingConnector trackingConnector = new TrackingConnector() {
+ @Override
+ protected void onStepConfigured(final String stepName, final
FlowContext workingContext) throws FlowUpdateException {
+ if (!blockNextRefresh.compareAndSet(true, false)) {
+ return;
+ }
+
+ refreshingStepName.set(stepName);
+ refreshStarted.countDown();
+ try {
+ permitRefresh.await();
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new FlowUpdateException("Interrupted while waiting
to refresh the working flow context", e);
+ }
+ }
+ };
+
+ final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
+ connectorNode.transitionStateForUpdating();
+ connectorNode.prepareForUpdate();
+ connectorNode.setConfiguration("step1",
createStepConfiguration(Map.of("propA", "oldA")));
+ connectorNode.setConfiguration("step2",
createStepConfiguration(Map.of("propA", "oldB")));
+ connectorNode.applyUpdate();
+ blockNextRefresh.set(true);
+
+ final ExecutorService executor = Executors.newSingleThreadExecutor();
+ try {
+ final Future<?> recreationFuture =
executor.submit(connectorNode::recreateWorkingFlowContext);
+ final String replacedStepName;
+ try {
+ assertTrue(refreshStarted.await(5, TimeUnit.SECONDS));
+ replacedStepName = "step1".equals(refreshingStepName.get()) ?
"step2" : "step1";
+ connectorNode.replaceWorkingConfiguration(replacedStepName,
createStepConfiguration(Map.of("propA", "newA")));
+ } finally {
+ permitRefresh.countDown();
+ }
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step1"));
-
assertTrue(trackingConnector.wasOnPropertyGroupConfiguredCalled("step2"));
+ recreationFuture.get(5, TimeUnit.SECONDS);
+
+ final ConnectorConfiguration workingConfiguration =
connectorNode.getWorkingFlowContext().getConfigurationContext().toConnectorConfiguration();
+ final NamedStepConfiguration namedStep =
workingConfiguration.getNamedStepConfiguration(replacedStepName);
+ assertEquals(Map.of("propA", new StringLiteralValue("newA")),
namedStep.configuration().getPropertyValues());
+ } finally {
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+ }
+
+ @Test
+ @Timeout(10)
+ public void
testOnConfigurationStepConfiguredCanWaitForWorkingContextRecreation() throws
Exception {
+ final AtomicReference<StandardConnectorNode> nodeReference = new
AtomicReference<>();
+ final AtomicBoolean waitForWorkingContextRecreation = new
AtomicBoolean();
+ final ExecutorService executor = Executors.newSingleThreadExecutor();
+ try {
+ final TrackingConnector trackingConnector = new
TrackingConnector() {
+ @Override
+ protected void onStepConfigured(final String stepName, final
FlowContext workingContext) throws FlowUpdateException {
+ if (!waitForWorkingContextRecreation.compareAndSet(true,
false)) {
+ return;
+ }
+
+ final StandardConnectorNode connectorNode =
nodeReference.get();
+ try {
+
executor.submit(connectorNode::recreateWorkingFlowContext).get(5,
TimeUnit.SECONDS);
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new FlowUpdateException("Interrupted while
waiting to recreate the working flow context", e);
+ } catch (final ExecutionException | TimeoutException e) {
+ throw new FlowUpdateException("Failed to recreate the
working flow context from onConfigurationStepConfigured", e);
+ }
+ }
+ };
+
+ final StandardConnectorNode connectorNode =
createConnectorNode(trackingConnector);
+ nodeReference.set(connectorNode);
+ connectorNode.transitionStateForUpdating();
+ connectorNode.prepareForUpdate();
+ waitForWorkingContextRecreation.set(true);
+ connectorNode.setConfiguration("step1",
createStepConfiguration(Map.of("propA", "valueA")));
+ } finally {
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
}
@Test
@@ -750,7 +837,7 @@ public class TestStandardConnectorNode {
connectorNode.discardWorkingConfiguration();
-
assertTrue(failingStepConnector.wasOnPropertyGroupConfiguredCalled("successStep"));
+
assertTrue(failingStepConnector.wasOnConfigurationStepConfiguredCalled("successStep"));
}
@Test
@@ -1549,6 +1636,56 @@ public class TestStandardConnectorNode {
}
}
+ /**
+ * Blocks the calling thread inside {@link #createWorkingFlowContext} while
+ * {@link StandardConnectorNode#recreateWorkingFlowContext()} is in
progress, until
+ * {@link #releaseWorkingContextCreation()} is invoked.
+ */
+ private static class BlockingWorkingFlowContextFactory implements
FlowContextFactory {
+ private final FlowContextFactory delegate;
+ private final CountDownLatch workingContextCreationStarted = new
CountDownLatch(1);
+ private final CountDownLatch permitWorkingContextCreation = new
CountDownLatch(1);
+ private volatile boolean blockWorkingContextCreation;
+
+ private BlockingWorkingFlowContextFactory(final FlowContextFactory
delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public FrameworkFlowContext createActiveFlowContext(final String
connectorId, final ComponentLog connectorLogger, final Bundle bundle) {
+ return delegate.createActiveFlowContext(connectorId,
connectorLogger, bundle);
+ }
+
+ @Override
+ public FrameworkFlowContext createWorkingFlowContext(final String
connectorId, final ComponentLog connectorLogger,
+ final MutableConnectorConfigurationContext
currentConfiguration, final Bundle bundle) {
+
+ if (blockWorkingContextCreation) {
+ workingContextCreationStarted.countDown();
+ try {
+ permitWorkingContextCreation.await();
+ } catch (final InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Interrupted while waiting
to create the working flow context", e);
+ }
+ }
+
+ return delegate.createWorkingFlowContext(connectorId,
connectorLogger, currentConfiguration, bundle);
+ }
+
+ private void blockNextWorkingContextCreation() {
+ blockWorkingContextCreation = true;
+ }
+
+ private boolean awaitWorkingContextCreation(final long timeout, final
TimeUnit timeUnit) throws InterruptedException {
+ return workingContextCreationStarted.await(timeout, timeUnit);
+ }
+
+ private void releaseWorkingContextCreation() {
+ permitWorkingContextCreation.countDown();
+ }
+ }
+
private ConnectorDetails createConnectorDetails(final Connector connector)
{
final ComponentLog componentLog = new
MockComponentLog("TestConnector", connector);
final BundleCoordinate bundleCoordinate = new
BundleCoordinate("org.apache.nifi", "test-standard-connector-node", "1.0.0");
@@ -1626,7 +1763,7 @@ public class TestStandardConnectorNode {
return List.of();
}
- public boolean wasOnPropertyGroupConfiguredCalled(final String
stepName) {
+ public boolean wasOnConfigurationStepConfiguredCalled(final String
stepName) {
return onConfigurationStepConfiguredCalls.contains(stepName);
}