diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/groups/StandardProcessGroup.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/groups/StandardProcessGroup.java index 600178a68661..c2e3193b915c 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/groups/StandardProcessGroup.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/groups/StandardProcessGroup.java @@ -28,6 +28,7 @@ import org.apache.nifi.authorization.resource.ResourceFactory; import org.apache.nifi.authorization.resource.ResourceType; import org.apache.nifi.components.PropertyDescriptor; +import org.apache.nifi.components.ValidationResult; import org.apache.nifi.components.connector.ConnectorNode; import org.apache.nifi.components.state.StateManager; import org.apache.nifi.components.state.StateManagerProvider; @@ -3880,16 +3881,25 @@ public void synchronizeWithFlowRegistry(final FlowManager flowManager) { return; } + final ValidationStatus validationStatus = flowRegistry.getValidationStatus(10, TimeUnit.SECONDS); + + if (validationStatus == ValidationStatus.VALIDATING) { + return; + } + + if (validationStatus == ValidationStatus.INVALID) { + final String message = buildValidationFailureExplanation(flowRegistry); + versionControlFields.setSyncFailureExplanation(message); + + LOG.error("{} for {}", message, this); + return; + } + final VersionedProcessGroup snapshot = vci.getFlowSnapshot(); if (snapshot == null && vci.getVersion() != null) { // We have not yet obtained the snapshot from the Flow Registry, so we need to request the snapshot of our local version of the flow from the Flow Registry. // This allows us to know whether or not the flow has been modified since it was last synced with the Flow Registry. try { - final ValidationStatus validationStatus = flowRegistry.getValidationStatus(10, TimeUnit.SECONDS); - if (validationStatus == ValidationStatus.VALIDATING) { - throw new FlowRegistryException(flowRegistry + " cannot currently be used to synchronize with Flow Registry because it is currently validating"); - } - final FlowVersionLocation flowVersionLocation = new FlowVersionLocation(vci.getBranch(), vci.getBucketIdentifier(), vci.getFlowIdentifier(), vci.getVersion()); final FlowSnapshotContainer registrySnapshotContainer = flowRegistry.getFlowContents( FlowRegistryClientContextFactory.getAnonymousContext(), flowVersionLocation, false); @@ -3897,8 +3907,9 @@ public void synchronizeWithFlowRegistry(final FlowManager flowManager) { final VersionedProcessGroup registryFlow = registrySnapshot.getFlowContents(); vci.setFlowSnapshot(registryFlow); } catch (final IOException | FlowRegistryException e) { - final String message = String.format("Failed to synchronize Process Group with Flow Registry because could not retrieve version %s of flow with identifier %s in bucket %s", - vci.getVersion(), vci.getFlowIdentifier(), vci.getBucketIdentifier()); + final String message = appendExceptionMessage(String.format( + "Failed to synchronize Process Group with Flow Registry because could not retrieve version %s of flow with identifier %s in bucket %s", + vci.getVersion(), vci.getFlowIdentifier(), vci.getBucketIdentifier()), e); versionControlFields.setSyncFailureExplanation(message); final String logErrorMessage = "Failed to synchronize {} with Flow Registry because could not retrieve version {} of flow with identifier {} in bucket {}"; @@ -3945,6 +3956,30 @@ public void synchronizeWithFlowRegistry(final FlowManager flowManager) { } } + private String buildValidationFailureExplanation(final FlowRegistryClientNode flowRegistry) { + final Collection validationResults = flowRegistry.getValidationErrors(); + final String validationErrors = validationResults == null ? null : validationResults.stream() + .map(validationResult -> StringUtils.isNotBlank(validationResult.getExplanation()) ? validationResult.getExplanation() : validationResult.toString()) + .filter(StringUtils::isNotBlank) + .collect(Collectors.joining("; ")); + + final String message = "Failed to synchronize Process Group with Flow Registry because the Flow Registry failed validation"; + if (StringUtils.isBlank(validationErrors)) { + return message; + } + + return message + ": " + validationErrors; + } + + private String appendExceptionMessage(final String message, final Exception exception) { + final String exceptionMessage = exception.getMessage(); + if (StringUtils.isBlank(exceptionMessage)) { + return message; + } + + return message + ": " + exceptionMessage; + } + @Override public ComponentAdditions addVersionedComponents(final VersionedComponentAdditions additions, final String componentIdSeed) { final ComponentIdGenerator idGenerator = (proposedId, instanceId, destinationGroupId) -> generateUuid(proposedId, destinationGroupId, componentIdSeed); diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java index 2e2659009673..dee057b126c6 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/groups/StandardProcessGroupTest.java @@ -17,10 +17,12 @@ package org.apache.nifi.groups; import org.apache.nifi.asset.AssetManager; +import org.apache.nifi.components.ValidationResult; import org.apache.nifi.components.state.Scope; import org.apache.nifi.components.state.StateManager; import org.apache.nifi.components.state.StateManagerProvider; import org.apache.nifi.components.state.StateMap; +import org.apache.nifi.components.validation.ValidationStatus; import org.apache.nifi.connectable.Connectable; import org.apache.nifi.connectable.ConnectableType; import org.apache.nifi.connectable.Connection; @@ -36,8 +38,16 @@ import org.apache.nifi.controller.queue.QueueSize; import org.apache.nifi.controller.service.ControllerServiceProvider; import org.apache.nifi.flow.ExecutionEngine; +import org.apache.nifi.flow.VersionedProcessGroup; import org.apache.nifi.nar.ExtensionManager; +import org.apache.nifi.registry.flow.FlowRegistryClientNode; +import org.apache.nifi.registry.flow.FlowRegistryException; +import org.apache.nifi.registry.flow.FlowSnapshotContainer; +import org.apache.nifi.registry.flow.RegisteredFlow; +import org.apache.nifi.registry.flow.RegisteredFlowSnapshot; +import org.apache.nifi.registry.flow.StandardVersionControlInformation; import org.apache.nifi.registry.flow.VersionControlInformation; +import org.apache.nifi.registry.flow.VersionedFlowState; import org.apache.nifi.registry.flow.VersionedFlowStatus; import org.apache.nifi.security.encryption.PropertyEncryptionProvider; import org.apache.nifi.util.NiFiProperties; @@ -48,16 +58,21 @@ import org.mockito.junit.jupiter.MockitoExtension; import java.io.IOException; +import java.lang.reflect.Field; +import java.util.List; import java.util.Map; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; 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; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; @@ -83,6 +98,9 @@ class StandardProcessGroupTest { private static final String REGISTERED_FLOW_IDENTIFIER = "87654321-4321"; private static final String REGISTERED_FLOW_VERSION = "1.0.0"; private static final String REGISTERED_FLOW_IDENTIFIER_PATH = "/87654321-4321:1.0.0"; + private static final String REGISTRY_ID = "registry-id"; + private static final String BUCKET_ID = "bucket-id"; + private static final String FLOW_ID = "flow-id"; @Mock private ControllerServiceProvider controllerServiceProvider; @@ -476,10 +494,230 @@ void testDropAllFlowFilesCompletionWaitsForEveryConnection() { assertTrue(completionFuture.isDone()); } + @Test + void testSynchronizeWithFlowRegistryDoesNotLatchRetrievalFailureWhileRegistryIsValidating() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.VALIDATING); + + processGroup.setVersionControlInformation(createVersionControlInformation(null, VersionedFlowState.UP_TO_DATE, null), Map.of()); + + assertDoesNotThrow(() -> processGroup.synchronizeWithFlowRegistry(flowManager)); + + assertNull(getSyncFailureExplanation(processGroup)); + verify(flowRegistry, never()).getFlowContents(any(), any(), eq(false)); + verify(flowRegistry, never()).getFlow(any(), any()); + verify(flowRegistry, never()).getLatestVersion(any(), any()); + } + + @Test + void testSynchronizeWithFlowRegistryPreservesPreviousFailureWhileRegistryIsValidating() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + final String previousFailure = "Previous registry failure"; + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.VALIDATING); + + processGroup.setVersionControlInformation( + createVersionControlInformation(null, VersionedFlowState.SYNC_FAILURE, previousFailure), + Map.of() + ); + + processGroup.synchronizeWithFlowRegistry(flowManager); + + assertEquals(previousFailure, getSyncFailureExplanation(processGroup)); + verify(flowRegistry, never()).getFlowContents(any(), any(), eq(false)); + verify(flowRegistry, never()).getFlow(any(), any()); + verify(flowRegistry, never()).getLatestVersion(any(), any()); + } + + @Test + void testSynchronizeWithFlowRegistryDoesNotReadRegistryMetadataWhileRegistryIsValidatingWithExistingSnapshot() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.VALIDATING); + + processGroup.setVersionControlInformation( + createVersionControlInformation(mock(VersionedProcessGroup.class), VersionedFlowState.UP_TO_DATE, null), + Map.of() + ); + + assertDoesNotThrow(() -> processGroup.synchronizeWithFlowRegistry(flowManager)); + + assertNull(getSyncFailureExplanation(processGroup)); + verify(flowRegistry, never()).getFlow(any(), any()); + verify(flowRegistry, never()).getLatestVersion(any(), any()); + } + + @Test + void testSynchronizeWithFlowRegistryReadsRegistryMetadataWhenRegistryIsValidWithExistingSnapshot() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + final RegisteredFlow versionedFlow = mock(RegisteredFlow.class); + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.VALID); + when(flowRegistry.getFlow(any(), any())).thenReturn(versionedFlow); + when(flowRegistry.getLatestVersion(any(), any())).thenReturn(Optional.of(REGISTERED_FLOW_VERSION)); + when(versionedFlow.getBucketName()).thenReturn("Bucket"); + when(versionedFlow.getName()).thenReturn("Flow"); + when(versionedFlow.getDescription()).thenReturn("Description"); + when(flowRegistry.getName()).thenReturn("Registry"); + + processGroup.setVersionControlInformation( + createVersionControlInformation(mock(VersionedProcessGroup.class), VersionedFlowState.UP_TO_DATE, null), + Map.of() + ); + + processGroup.synchronizeWithFlowRegistry(flowManager); + + verify(flowRegistry).getFlow(any(), any()); + verify(flowRegistry).getLatestVersion(any(), any()); + assertNull(getSyncFailureExplanation(processGroup)); + } + + @Test + void testSynchronizeWithFlowRegistryMarksValidationFailureWithoutRegistryDataCallsWhenRegistryIsInvalid() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.INVALID); + when(flowRegistry.getValidationErrors()).thenReturn(List.of(new ValidationResult.Builder() + .subject("Registry URL") + .input("") + .valid(false) + .explanation("Registry URL is required") + .build())); + + processGroup.setVersionControlInformation(createVersionControlInformation(null, VersionedFlowState.UP_TO_DATE, null), Map.of()); + + processGroup.synchronizeWithFlowRegistry(flowManager); + + final String syncFailureExplanation = getSyncFailureExplanation(processGroup); + assertNotNull(syncFailureExplanation); + assertTrue(syncFailureExplanation.toLowerCase().contains("validation")); + assertTrue(syncFailureExplanation.contains("Registry URL is required")); + verify(flowRegistry, never()).getFlowContents(any(), any(), eq(false)); + verify(flowRegistry, never()).getFlow(any(), any()); + verify(flowRegistry, never()).getLatestVersion(any(), any()); + } + + @Test + void testSynchronizeWithFlowRegistryIncludesUnderlyingRegistryErrorWhenSnapshotFetchFails() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.VALID); + when(flowRegistry.getFlowContents(any(), any(), eq(false))).thenThrow(new FlowRegistryException("Registry offline")); + + processGroup.setVersionControlInformation(createVersionControlInformation(null, VersionedFlowState.UP_TO_DATE, null), Map.of()); + + processGroup.synchronizeWithFlowRegistry(flowManager); + + final String syncFailureExplanation = getSyncFailureExplanation(processGroup); + assertNotNull(syncFailureExplanation); + assertTrue(syncFailureExplanation.contains("could not retrieve version")); + assertTrue(syncFailureExplanation.contains("Registry offline")); + } + + @Test + void testSynchronizeWithFlowRegistryClearsPreviousFailureAfterSuccessfulRetry() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + final FlowSnapshotContainer snapshotContainer = mock(FlowSnapshotContainer.class); + final RegisteredFlowSnapshot registeredFlowSnapshot = mock(RegisteredFlowSnapshot.class); + final RegisteredFlow versionedFlow = mock(RegisteredFlow.class); + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))).thenReturn(ValidationStatus.VALID); + when(flowRegistry.getFlowContents(any(), any(), eq(false))) + .thenThrow(new FlowRegistryException("Registry offline")) + .thenReturn(snapshotContainer); + when(snapshotContainer.getFlowSnapshot()).thenReturn(registeredFlowSnapshot); + when(registeredFlowSnapshot.getFlowContents()).thenReturn(mock(VersionedProcessGroup.class)); + when(flowRegistry.getFlow(any(), any())).thenReturn(versionedFlow); + when(flowRegistry.getLatestVersion(any(), any())).thenReturn(Optional.of(REGISTERED_FLOW_VERSION)); + when(versionedFlow.getBucketName()).thenReturn("Bucket"); + when(versionedFlow.getName()).thenReturn("Flow"); + when(versionedFlow.getDescription()).thenReturn("Description"); + when(flowRegistry.getName()).thenReturn("Registry"); + + processGroup.setVersionControlInformation(createVersionControlInformation(null, VersionedFlowState.UP_TO_DATE, null), Map.of()); + + processGroup.synchronizeWithFlowRegistry(flowManager); + assertNotNull(getSyncFailureExplanation(processGroup)); + + processGroup.synchronizeWithFlowRegistry(flowManager); + + assertNull(getSyncFailureExplanation(processGroup)); + } + + @Test + void testSynchronizeWithFlowRegistryClearsValidationFailureAfterSuccessfulRetry() throws Exception { + final FlowRegistryClientNode flowRegistry = mock(FlowRegistryClientNode.class); + final RegisteredFlow versionedFlow = mock(RegisteredFlow.class); + + when(flowManager.getFlowRegistryClient(REGISTRY_ID)).thenReturn(flowRegistry); + when(flowRegistry.getValidationStatus(anyLong(), eq(TimeUnit.SECONDS))) + .thenReturn(ValidationStatus.INVALID) + .thenReturn(ValidationStatus.VALID); + when(flowRegistry.getValidationErrors()).thenReturn(List.of(new ValidationResult.Builder() + .subject("Registry URL") + .input("") + .valid(false) + .explanation("Registry URL is required") + .build())); + when(flowRegistry.getFlow(any(), any())).thenReturn(versionedFlow); + when(flowRegistry.getLatestVersion(any(), any())).thenReturn(Optional.of(REGISTERED_FLOW_VERSION)); + when(versionedFlow.getBucketName()).thenReturn("Bucket"); + when(versionedFlow.getName()).thenReturn("Flow"); + when(versionedFlow.getDescription()).thenReturn("Description"); + when(flowRegistry.getName()).thenReturn("Registry"); + + processGroup.setVersionControlInformation( + createVersionControlInformation(mock(VersionedProcessGroup.class), VersionedFlowState.UP_TO_DATE, null), + Map.of() + ); + + processGroup.synchronizeWithFlowRegistry(flowManager); + assertNotNull(getSyncFailureExplanation(processGroup)); + + processGroup.synchronizeWithFlowRegistry(flowManager); + + assertNull(getSyncFailureExplanation(processGroup)); + } + private StandardProcessGroup createStandardProcessGroup(final String id) { return createStandardProcessGroup(id, null); } + private StandardVersionControlInformation createVersionControlInformation( + final VersionedProcessGroup flowSnapshot, + final VersionedFlowState state, + final String explanation) { + return new StandardVersionControlInformation( + REGISTRY_ID, + "Registry", + "main", + BUCKET_ID, + FLOW_ID, + REGISTERED_FLOW_VERSION, + null, + flowSnapshot, + new StandardVersionedFlowStatus(state, explanation) + ); + } + + private String getSyncFailureExplanation(final StandardProcessGroup group) { + try { + final Field versionControlFieldsField = StandardProcessGroup.class.getDeclaredField("versionControlFields"); + versionControlFieldsField.setAccessible(true); + final VersionControlFields versionControlFields = (VersionControlFields) versionControlFieldsField.get(group); + return versionControlFields.getSyncFailureExplanation(); + } catch (final ReflectiveOperationException e) { + throw new AssertionError(e); + } + } + private StandardProcessGroup createStandardProcessGroup(final String id, final String connectorId) { return new StandardProcessGroup( id, 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 32cfc3b7ae53..d30a4896c266 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 @@ -268,6 +268,7 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.function.BooleanSupplier; import java.util.function.Supplier; import java.util.stream.Collectors; import javax.management.NotificationEmitter; @@ -363,6 +364,7 @@ public class FlowController implements ReportingTaskProvider, FlowAnalysisRulePr private final RepositoryContextFactory repositoryContextFactory; private final RingBufferGarbageCollectionLog gcLog; private final Optional longRunningTaskMonitorThreadPool; + private volatile RegistryFlowSynchronizationTask registrySynchronizationTask; /** @@ -1394,8 +1396,12 @@ public void initializeFlow(final QueueProvider queueProvider) throws IOException LOG.info("Scheduled Flow Registry with Sync Interval [{} s] Check Interval [{} s]", registrySyncInterval, registrySyncTickSeconds); - final RegistryFlowSynchronizationTask registrySynchronizationTask = new RegistryFlowSynchronizationTask(flowManager, defaultRegistrySyncIntervalSeconds); - timerDrivenEngineRef.get().scheduleWithFixedDelay(registrySynchronizationTask, 300, registrySyncTickSeconds, TimeUnit.SECONDS); + registrySynchronizationTask = scheduleRegistrySynchronizationTask( + timerDrivenEngineRef.get(), + flowManager, + defaultRegistrySyncIntervalSeconds, + registrySyncTickSeconds + ); initialized.set(true); } finally { @@ -1457,6 +1463,7 @@ private void notifyComponentsConfigurationRestored() { * @param startDelayedComponents true if start */ public void onFlowInitialized(final boolean startDelayedComponents) { + Runnable postInitializationRegistrySynchronizationTask = null; writeLock.lock(); try { // Perform validation of all components before attempting to start them. @@ -1621,9 +1628,31 @@ public void trigger(final ComponentNode component) { timerDrivenEngineRef.get().scheduleWithFixedDelay(discoverPythonExtensions, 1, 1, TimeUnit.MINUTES); ComponentAccessPolicyDeprecationLogger.logComponentPolicies(authorizer, flowManager.getRootGroupId()); + postInitializationRegistrySynchronizationTask = registrySynchronizationTask; } finally { writeLock.unlock("onFlowInitialized"); } + + submitPostInitializationRegistrySynchronizationTask(processScheduler, postInitializationRegistrySynchronizationTask, rwLock::isWriteLockedByCurrentThread); + } + + static RegistryFlowSynchronizationTask scheduleRegistrySynchronizationTask(final ScheduledExecutorService timerDrivenEngine, final FlowManager flowManager, + final long defaultRegistrySyncIntervalSeconds, final long registrySyncTickSeconds) { + final RegistryFlowSynchronizationTask registrySynchronizationTask = new RegistryFlowSynchronizationTask(flowManager, defaultRegistrySyncIntervalSeconds); + timerDrivenEngine.scheduleWithFixedDelay(registrySynchronizationTask, 300, registrySyncTickSeconds, TimeUnit.SECONDS); + return registrySynchronizationTask; + } + + static void submitPostInitializationRegistrySynchronizationTask(final ProcessScheduler processScheduler, final Runnable registrySynchronizationTask, + final BooleanSupplier writeLockHeldSupplier) { + if (registrySynchronizationTask == null) { + return; + } + if (writeLockHeldSupplier.getAsBoolean()) { + throw new IllegalStateException("Cannot submit Flow Registry Synchronization Task while write lock is held"); + } + + processScheduler.submitFrameworkTask(registrySynchronizationTask); } private void scheduleBackgroundFlowAnalysis(Supplier rootProcessGroupSupplier) { diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/FlowControllerRegistrySynchronizationLifecycleTest.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/FlowControllerRegistrySynchronizationLifecycleTest.java new file mode 100644 index 000000000000..1b3c00c25a66 --- /dev/null +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/FlowControllerRegistrySynchronizationLifecycleTest.java @@ -0,0 +1,69 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.nifi.controller; + +import org.apache.nifi.controller.flow.FlowManager; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.same; +import static org.mockito.Mockito.verify; + +class FlowControllerRegistrySynchronizationLifecycleTest { + + @Test + void testScheduleAndSubmitUseSameRegistrySynchronizationTaskInstance() { + final FlowManager flowManager = mock(FlowManager.class); + final ScheduledExecutorService timerDrivenEngine = mock(ScheduledExecutorService.class); + final ProcessScheduler processScheduler = mock(ProcessScheduler.class); + + final RegistryFlowSynchronizationTask registrySynchronizationTask = FlowController.scheduleRegistrySynchronizationTask(timerDrivenEngine, flowManager, 1800L, 30L); + + assertDoesNotThrow(() -> FlowController.submitPostInitializationRegistrySynchronizationTask(processScheduler, registrySynchronizationTask, () -> false)); + + verify(timerDrivenEngine).scheduleWithFixedDelay(same(registrySynchronizationTask), eq(300L), eq(30L), eq(TimeUnit.SECONDS)); + verify(processScheduler).submitFrameworkTask(same(registrySynchronizationTask)); + } + + @Test + void testSubmitPostInitializationRegistrySynchronizationTaskRequiresReleasedWriteLock() { + final ProcessScheduler processScheduler = mock(ProcessScheduler.class); + final Runnable registrySynchronizationTask = mock(Runnable.class); + + assertThrows(IllegalStateException.class, + () -> FlowController.submitPostInitializationRegistrySynchronizationTask(processScheduler, registrySynchronizationTask, () -> true)); + + verify(processScheduler, never()).submitFrameworkTask(any(Runnable.class)); + } + + @Test + void testSubmitPostInitializationRegistrySynchronizationTaskIgnoresAbsentTask() { + final ProcessScheduler processScheduler = mock(ProcessScheduler.class); + + assertDoesNotThrow(() -> FlowController.submitPostInitializationRegistrySynchronizationTask(processScheduler, null, () -> false)); + + verify(processScheduler, never()).submitFrameworkTask(any(Runnable.class)); + } +} diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/RegistryFlowSynchronizationTaskTest.java b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/RegistryFlowSynchronizationTaskTest.java index 8683d73cad6c..514d1e1e92f5 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/RegistryFlowSynchronizationTaskTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-core/src/test/java/org/apache/nifi/controller/RegistryFlowSynchronizationTaskTest.java @@ -17,12 +17,20 @@ package org.apache.nifi.controller; import org.apache.nifi.controller.flow.FlowManager; +import org.apache.nifi.groups.ProcessGroup; import org.apache.nifi.registry.flow.AbstractFlowRegistryClient; import org.apache.nifi.registry.flow.FlowRegistryClientNode; +import org.apache.nifi.registry.flow.VersionControlInformation; import org.junit.jupiter.api.Test; +import java.util.ArrayList; +import java.util.List; + import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; class RegistryFlowSynchronizationTaskTest { @@ -66,4 +74,135 @@ void testGetEffectiveIntervalSeconds() { when(flowManager.getFlowRegistryClient("missing")).thenReturn(null); assertEquals(DEFAULT_INTERVAL_SECONDS, task.getEffectiveIntervalSeconds("missing")); } + + @Test + void testFirstRunSynchronizesAllVersionControlledGroups() { + final FlowManager flowManager = mock(FlowManager.class); + final ProcessGroup rootGroup = mock(ProcessGroup.class); + final ProcessGroup childGroup = mock(ProcessGroup.class); + final VersionControlInformation rootVersionControlInformation = versionControlInformation("root-registry"); + final VersionControlInformation childVersionControlInformation = versionControlInformation("child-registry"); + + when(flowManager.getRootGroup()).thenReturn(rootGroup); + when(rootGroup.findAllProcessGroups()).thenReturn(new ArrayList<>(List.of(childGroup))); + when(rootGroup.getVersionControlInformation()).thenReturn(rootVersionControlInformation); + when(childGroup.getVersionControlInformation()).thenReturn(childVersionControlInformation); + + final RegistryFlowSynchronizationTask task = new RegistryFlowSynchronizationTask(flowManager, DEFAULT_INTERVAL_SECONDS); + + task.run(); + + verify(rootGroup).synchronizeWithFlowRegistry(flowManager); + verify(childGroup).synchronizeWithFlowRegistry(flowManager); + } + + @Test + void testGroupsSharingRegistryClientAreProcessedInSingleClientBatch() { + final FlowManager flowManager = mock(FlowManager.class); + final ProcessGroup rootGroup = mock(ProcessGroup.class); + final ProcessGroup firstChild = mock(ProcessGroup.class); + final ProcessGroup secondChild = mock(ProcessGroup.class); + final FlowRegistryClientNode clientNode = mock(FlowRegistryClientNode.class); + final VersionControlInformation firstVersionControlInformation = versionControlInformation("shared-registry"); + final VersionControlInformation secondVersionControlInformation = versionControlInformation("shared-registry"); + + when(flowManager.getRootGroup()).thenReturn(rootGroup); + when(rootGroup.findAllProcessGroups()).thenReturn(new ArrayList<>(List.of(firstChild, secondChild))); + when(rootGroup.getVersionControlInformation()).thenReturn(null); + when(firstChild.getVersionControlInformation()).thenReturn(firstVersionControlInformation); + when(secondChild.getVersionControlInformation()).thenReturn(secondVersionControlInformation); + when(flowManager.getFlowRegistryClient("shared-registry")).thenReturn(clientNode); + when(clientNode.getEffectivePropertyValue(AbstractFlowRegistryClient.SYNCHRONIZATION_INTERVAL)).thenReturn("10 min"); + + final RegistryFlowSynchronizationTask task = new RegistryFlowSynchronizationTask(flowManager, DEFAULT_INTERVAL_SECONDS); + + task.run(); + + verify(flowManager, times(1)).getFlowRegistryClient("shared-registry"); + verify(firstChild).synchronizeWithFlowRegistry(flowManager); + verify(secondChild).synchronizeWithFlowRegistry(flowManager); + } + + @Test + void testSecondRunBeforeIntervalDoesNotRepeatSynchronization() { + final FlowManager flowManager = mock(FlowManager.class); + final ProcessGroup rootGroup = mock(ProcessGroup.class); + final ProcessGroup childGroup = mock(ProcessGroup.class); + final FlowRegistryClientNode clientNode = mock(FlowRegistryClientNode.class); + final VersionControlInformation childVersionControlInformation = versionControlInformation("shared-registry"); + + when(flowManager.getRootGroup()).thenReturn(rootGroup); + when(rootGroup.findAllProcessGroups()).thenReturn(new ArrayList<>(List.of(childGroup))); + when(rootGroup.getVersionControlInformation()).thenReturn(null); + when(childGroup.getVersionControlInformation()).thenReturn(childVersionControlInformation); + when(flowManager.getFlowRegistryClient("shared-registry")).thenReturn(clientNode); + when(clientNode.getEffectivePropertyValue(AbstractFlowRegistryClient.SYNCHRONIZATION_INTERVAL)).thenReturn("10 min"); + + final RegistryFlowSynchronizationTask task = new RegistryFlowSynchronizationTask(flowManager, DEFAULT_INTERVAL_SECONDS); + + task.run(); + task.run(); + + verify(childGroup, times(1)).synchronizeWithFlowRegistry(flowManager); + } + + @Test + void testFailureInOneGroupDoesNotPreventSiblingSynchronization() { + final FlowManager flowManager = mock(FlowManager.class); + final ProcessGroup rootGroup = mock(ProcessGroup.class); + final ProcessGroup failingGroup = mock(ProcessGroup.class); + final ProcessGroup siblingGroup = mock(ProcessGroup.class); + final FlowRegistryClientNode clientNode = mock(FlowRegistryClientNode.class); + final VersionControlInformation failingVersionControlInformation = versionControlInformation("shared-registry"); + final VersionControlInformation siblingVersionControlInformation = versionControlInformation("shared-registry"); + + when(flowManager.getRootGroup()).thenReturn(rootGroup); + when(rootGroup.findAllProcessGroups()).thenReturn(new ArrayList<>(List.of(failingGroup, siblingGroup))); + when(rootGroup.getVersionControlInformation()).thenReturn(null); + when(failingGroup.getVersionControlInformation()).thenReturn(failingVersionControlInformation); + when(siblingGroup.getVersionControlInformation()).thenReturn(siblingVersionControlInformation); + when(flowManager.getFlowRegistryClient("shared-registry")).thenReturn(clientNode); + when(clientNode.getEffectivePropertyValue(AbstractFlowRegistryClient.SYNCHRONIZATION_INTERVAL)).thenReturn("10 min"); + doThrow(new RuntimeException("boom")).when(failingGroup).synchronizeWithFlowRegistry(flowManager); + + final RegistryFlowSynchronizationTask task = new RegistryFlowSynchronizationTask(flowManager, DEFAULT_INTERVAL_SECONDS); + + task.run(); + + verify(failingGroup).synchronizeWithFlowRegistry(flowManager); + verify(siblingGroup).synchronizeWithFlowRegistry(flowManager); + } + + @Test + void testRemovedClientIsForgottenAndSynchronizesImmediatelyWhenReintroduced() { + final FlowManager flowManager = mock(FlowManager.class); + final ProcessGroup rootGroup = mock(ProcessGroup.class); + final ProcessGroup childGroup = mock(ProcessGroup.class); + final FlowRegistryClientNode clientNode = mock(FlowRegistryClientNode.class); + final VersionControlInformation childVersionControlInformation = versionControlInformation("shared-registry"); + + when(flowManager.getRootGroup()).thenReturn(rootGroup); + when(rootGroup.getVersionControlInformation()).thenReturn(null); + when(childGroup.getVersionControlInformation()).thenReturn(childVersionControlInformation); + when(flowManager.getFlowRegistryClient("shared-registry")).thenReturn(clientNode); + when(clientNode.getEffectivePropertyValue(AbstractFlowRegistryClient.SYNCHRONIZATION_INTERVAL)).thenReturn("10 min"); + when(rootGroup.findAllProcessGroups()) + .thenReturn(new ArrayList<>(List.of(childGroup))) + .thenReturn(new ArrayList<>()) + .thenReturn(new ArrayList<>(List.of(childGroup))); + + final RegistryFlowSynchronizationTask task = new RegistryFlowSynchronizationTask(flowManager, DEFAULT_INTERVAL_SECONDS); + + task.run(); + task.run(); + task.run(); + + verify(childGroup, times(2)).synchronizeWithFlowRegistry(flowManager); + } + + private VersionControlInformation versionControlInformation(final String registryIdentifier) { + final VersionControlInformation versionControlInformation = mock(VersionControlInformation.class); + when(versionControlInformation.getRegistryIdentifier()).thenReturn(registryIdentifier); + return versionControlInformation; + } }