From b0b479bf0747bdc694dd634da62b32463611e024 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alaksiej=20=C5=A0=C4=8Darbaty?= Date: Fri, 11 Sep 2026 16:12:17 +0200 Subject: [PATCH 1/3] NIFI-16310 Preserve Controller Services created by migrateProperties during versioned flow synchronization --- ...tandardVersionedComponentSynchronizer.java | 246 +++++++++++- ...ardVersionedComponentSynchronizerTest.java | 291 ++++++++++++++ .../tests/system/StateBackedStoreService.java | 97 +++++ .../nifi/cs/tests/system/StoreService.java | 29 ++ .../system/MigrateToControllerService.java | 90 +++++ ...g.apache.nifi.controller.ControllerService | 1 + .../org.apache.nifi.processor.Processor | 1 + .../nifi-system-test-extensions/pom.xml | 4 - .../system/MigrateToControllerService.java | 68 ++++ .../org.apache.nifi.processor.Processor | 1 + .../pom.xml | 46 +++ .../nifi-system-test-flow-registry/pom.xml | 42 ++ .../FileSystemFlowRegistryClient.java | 0 ...ache.nifi.registry.flow.FlowRegistryClient | 0 .../pom.xml | 31 ++ .../nifi-system-test-suite/pom.xml | 6 + .../src/test/assembly/dependencies.xml | 1 + .../nifi/tests/system/NiFiClientUtil.java | 2 +- .../nifi/tests/system/NiFiSystemIT.java | 1 + ...nCreatedControllerServiceVersioningIT.java | 365 ++++++++++++++++++ .../1/snapshot.json | 82 ++++ .../2/snapshot.json | 112 ++++++ .../1/snapshot.json | 82 ++++ .../2/snapshot.json | 118 ++++++ nifi-system-tests/pom.xml | 1 + .../client/ControllerServicesClient.java | 3 + .../impl/JerseyControllerServicesClient.java | 14 + 27 files changed, 1716 insertions(+), 18 deletions(-) create mode 100644 nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java create mode 100644 nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StoreService.java create mode 100644 nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java create mode 100644 nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java create mode 100644 nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml create mode 100644 nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml rename nifi-system-tests/{nifi-system-test-extensions-bundle/nifi-system-test-extensions => nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry}/src/main/java/org/apache/nifi/flow/registry/FileSystemFlowRegistryClient.java (100%) rename nifi-system-tests/{nifi-system-test-extensions-bundle/nifi-system-test-extensions => nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry}/src/main/resources/META-INF/services/org.apache.nifi.registry.flow.FlowRegistryClient (100%) create mode 100644 nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/1/snapshot.json create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/2/snapshot.json create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/1/snapshot.json create mode 100644 nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/2/snapshot.json diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java index e7930d84c329..274940379700 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java @@ -51,6 +51,7 @@ import org.apache.nifi.controller.reporting.ReportingTaskInstantiationException; import org.apache.nifi.controller.service.ControllerServiceNode; import org.apache.nifi.controller.service.ControllerServiceProvider; +import org.apache.nifi.controller.service.ControllerServiceReference; import org.apache.nifi.controller.service.ControllerServiceState; import org.apache.nifi.flow.BatchSize; import org.apache.nifi.flow.Bundle; @@ -145,6 +146,7 @@ import java.net.URL; import java.time.Duration; import java.util.ArrayList; +import java.util.BitSet; import java.util.Collection; import java.util.Collections; import java.util.HashMap; @@ -552,18 +554,19 @@ private void synchronize(final ProcessGroup group, final VersionedProcessGroup p // // The sequence of steps / order of operations are as follows: // - // 1. Remove any Controller Services that do not exist in the proposed group - // 2. Add any Controller Services that are in the proposed group that are not in the current flow - // 3. Update Controller Services to match those in the proposed group - // 4. Remove any connections that do not exist in the proposed group - // 5. For any connection that does exist, if the proposed group has a different destination for the connection, update the destination. + // 1. Assign proposed versioned ids to Controller Services that property migration created + // 2. Remove any Controller Services that do not exist in the proposed group + // 3. Add any Controller Services that are in the proposed group that are not in the current flow + // 4. Update Controller Services to match those in the proposed group + // 5. Remove any connections that do not exist in the proposed group + // 6. For any connection that does exist, if the proposed group has a different destination for the connection, update the destination. // If the new destination does not yet exist in the flow, set the destination as some temporary component. - // 6. Remove any other components that do not exist in the proposed group. - // 7. Add any components, other than Connections, that exist in the proposed group but not in the current flow - // 8. Update components, other than Connections, to match those in the proposed group - // 9. Add connections that exist in the proposed group that are not in the current flow - // 10. Update connections to match those in the proposed group - // 11. Delete the temporary destination that was created above + // 7. Remove any other components that do not exist in the proposed group. + // 8. Add any components, other than Connections, that exist in the proposed group but not in the current flow + // 9. Update components, other than Connections, to match those in the proposed group + // 10. Add connections that exist in the proposed group that are not in the current flow + // 11. Update connections to match those in the proposed group + // 12. Delete the temporary destination that was created above // During the flow update, we will use temporary names for process group ports. This is because port names must be // unique within a process group, but during an update we might temporarily be in a state where two ports have the same name. @@ -574,6 +577,7 @@ private void synchronize(final ProcessGroup group, final VersionedProcessGroup p final Map proposedPortFinalNames = new HashMap<>(); // Controller Services + assignVersionedIdsToMigrationCreatedControllerServices(group, proposed); final Map controllerServicesByVersionedId = componentsById(group, grp -> grp.getControllerServices(false), ControllerServiceNode::getIdentifier, ControllerServiceNode::getVersionedComponentId); removeMissingControllerServices(group, proposed, controllerServicesByVersionedId); @@ -1169,9 +1173,225 @@ private void removeMissingRpg(final ProcessGroup group, final VersionedProcessGr removeMissingComponents(group, proposed, rpgsByVersionedId, VersionedProcessGroup::getRemoteProcessGroups, ProcessGroup::removeRemoteProcessGroup); } + /** + * Assigns a proposed versioned id to a Controller Service created by property migration. + * The service must have exactly one referencer. + * The proposed counterpart of that referencer must point at a Controller Service of the same type. + * If those do not hold, the service stays unversioned. + * A proposed id is assigned to at most one local service. + */ + private void assignVersionedIdsToMigrationCreatedControllerServices(final ProcessGroup group, final VersionedProcessGroup proposed) { + final Collection groupServices = group.getControllerServices(false); + if (groupServices == null || groupServices.isEmpty()) { + return; + } + + final List migrationCreatedServices = new ArrayList<>(); + for (final ControllerServiceNode localService : groupServices) { + if (localService.getVersionedComponentId().isEmpty() && isMigrationCreated(localService)) { + migrationCreatedServices.add(localService); + } + } + + if (migrationCreatedServices.isEmpty()) { + return; + } + + final Set claimedVersionedIds = HashSet.newHashSet(groupServices.size()); + for (final ControllerServiceNode localService : groupServices) { + localService.getVersionedComponentId().ifPresent(claimedVersionedIds::add); + } + + final Map proposedComponentsByVersionedId = indexByVersionedId(proposed.getControllerServices(), proposed.getProcessors()); + + for (final ControllerServiceNode localService : orderByReferencerChain(migrationCreatedServices)) { + final ComponentNode referencer = getSoleReferencer(localService); + if (referencer == null) { + LOG.debug("Leaving {} in {} unversioned because it is not referenced by exactly one component", localService, group); + continue; + } + + final VersionedConfigurableExtension proposedReferencer = getProposedReferencer(referencer, proposedComponentsByVersionedId); + if (proposedReferencer == null) { + LOG.debug("Leaving {} in {} unversioned because its referencer {} has no counterpart in the proposed flow", localService, group, referencer); + continue; + } + + final ProposedControllerServiceMatch match = findMatchingProposedControllerService(localService, referencer, proposedReferencer, proposedComponentsByVersionedId, group); + if (match == null) { + LOG.debug("Leaving {} in {} unversioned because no proposed Controller Service matches the referencing property of {}", localService, group, referencer); + continue; + } + + if (!claimedVersionedIds.add(match.versionedId())) { + LOG.debug("Leaving {} in {} unversioned because versioned id {} is already used by another Controller Service", localService, group, match.versionedId()); + continue; + } + + localService.setVersionedComponentId(match.versionedId()); + updatedVersionedComponentIds.add(match.versionedId()); + LOG.info("Matched {} in {} to the Controller Service with versioned id {} that the proposed flow declares, based on the {} property of {}", + localService, group, match.versionedId(), match.propertyName(), referencer); + } + } + + private List orderByReferencerChain(final List migrationCreatedServices) { + final BitSet visited = new BitSet(migrationCreatedServices.size()); + final List ordered = new ArrayList<>(migrationCreatedServices.size()); + + for (int i = 0; i < migrationCreatedServices.size(); i++) { + appendReferencerChain(i, migrationCreatedServices, visited, ordered); + } + + return ordered; + } + + private void appendReferencerChain( + final int index, + final List migrationCreatedServices, + final BitSet visited, + final List ordered + ) { + if (visited.get(index)) { + return; + } + visited.set(index); + + final ControllerServiceNode localService = migrationCreatedServices.get(index); + final ComponentNode referencer = getSoleReferencer(localService); + if (referencer instanceof ControllerServiceNode referencingService && isMigrationCreated(referencingService)) { + final int referencerIndex = migrationCreatedServices.indexOf(referencingService); + if (referencerIndex >= 0) { + // Visit the creator first so we can assign the versioned id to the creator first. + appendReferencerChain(referencerIndex, migrationCreatedServices, visited, ordered); + } + } + + ordered.add(localService); + } + + private boolean isMigrationCreated(final ControllerServiceNode service) { + return StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT.equals(service.getComments()); + } + + private Map indexByVersionedId( + final Collection controllerServices, + final Collection processors + ) { + final int serviceCount = controllerServices == null ? 0 : controllerServices.size(); + final int processorCount = processors == null ? 0 : processors.size(); + final Map byVersionedId = HashMap.newHashMap(serviceCount + processorCount); + addByVersionedId(byVersionedId, controllerServices); + addByVersionedId(byVersionedId, processors); + return byVersionedId; + } + + private void addByVersionedId( + final Map byVersionedId, + final Collection components + ) { + if (components == null) { + return; + } + + for (final VersionedConfigurableExtension component : components) { + byVersionedId.put(component.getIdentifier(), component); + } + } + + private ComponentNode getSoleReferencer(final ControllerServiceNode localService) { + final ControllerServiceReference references = localService.getReferences(); + if (references == null) { + return null; + } + + final Set referencers = references.getReferencingComponents(); + if (referencers == null || referencers.size() != 1) { + return null; + } + + return referencers.iterator().next(); + } + + private VersionedConfigurableExtension getProposedReferencer( + final ComponentNode referencer, + final Map proposedComponentsByVersionedId + ) { + if (!(referencer instanceof org.apache.nifi.components.VersionedComponent versionedReferencer)) { + return null; + } + + return versionedReferencer.getVersionedComponentId() + .map(proposedComponentsByVersionedId::get) + .orElse(null); + } + + private ProposedControllerServiceMatch findMatchingProposedControllerService( + final ControllerServiceNode localService, + final ComponentNode referencer, + final VersionedConfigurableExtension proposedReferencer, + final Map proposedComponentsByVersionedId, + final ProcessGroup group + ) { + final Map rawPropertyValues = referencer.getRawPropertyValues(); + if (rawPropertyValues == null) { + LOG.debug("Leaving {} in {} unversioned because referencer {} has no property values", localService, group, referencer); + return null; + } + + for (final Map.Entry propertyEntry : rawPropertyValues.entrySet()) { + final PropertyDescriptor descriptor = propertyEntry.getKey(); + final String propertyName = descriptor.getName(); + if (descriptor.getControllerServiceDefinition() == null) { + continue; + } + + if (!localService.getIdentifier().equals(propertyEntry.getValue())) { + continue; + } + + final Map proposedProperties = proposedReferencer.getProperties(); + final String proposedServiceId = proposedProperties == null ? null : proposedProperties.get(propertyName); + if (proposedServiceId == null) { + LOG.debug("Leaving {} in {} unversioned because the proposed {} does not set the {} property", + localService, group, proposedReferencer, propertyName); + continue; + } + + // In versioned flow, the service identifier is the versioned component id of the service. + final VersionedConfigurableExtension proposedService = proposedComponentsByVersionedId.get(proposedServiceId); + if (proposedService == null || !proposedService.getType().equals(localService.getCanonicalClassName())) { + LOG.debug("Leaving {} in {} unversioned because proposed Controller Service {} is missing or has a different type than {}", + localService, group, proposedServiceId, localService); + continue; + } + + return new ProposedControllerServiceMatch(proposedServiceId, propertyName); + } + + return null; + } + + private record ProposedControllerServiceMatch(String versionedId, String propertyName) { + } + private void removeMissingControllerServices(final ProcessGroup group, final VersionedProcessGroup proposed, final Map servicesByVersionedId) { - final BiConsumer componentRemoval = (grp, service) -> context.getControllerServiceProvider().removeControllerService(service); - removeMissingComponents(group, proposed, servicesByVersionedId, VersionedProcessGroup::getControllerServices, componentRemoval); + // Do not remove Controller Services created by migrateProperties. + final Map servicesEligibleForRemoval = HashMap.newHashMap(servicesByVersionedId.size()); + for (final Map.Entry entry : servicesByVersionedId.entrySet()) { + final ControllerServiceNode service = entry.getValue(); + if (isMigrationCreated(service)) { + if (service.getVersionedComponentId().isEmpty()) { + LOG.info("Keeping {} in {} because it was created by property migration and is not present in the proposed flow", + service, group); + } + } else { + servicesEligibleForRemoval.put(entry.getKey(), service); + } + } + + final BiConsumer componentRemoval = (procGroup, service) -> context.getControllerServiceProvider().removeControllerService(service); + removeMissingComponents(group, proposed, servicesEligibleForRemoval, VersionedProcessGroup::getControllerServices, componentRemoval); } private void removeMissingChildGroups(final ProcessGroup group, final VersionedProcessGroup proposed, final Map groupsByVersionedId) { diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java index 86cdbf215bd1..173dda001c6f 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java @@ -67,6 +67,7 @@ import org.apache.nifi.groups.ScheduledStateChangeListener; import org.apache.nifi.groups.VersionedComponentAdditions; import org.apache.nifi.logging.LogLevel; +import org.apache.nifi.migration.StandardControllerServiceFactory; import org.apache.nifi.nar.ExtensionManager; import org.apache.nifi.parameter.Parameter; import org.apache.nifi.parameter.ParameterContext; @@ -92,6 +93,7 @@ import org.apache.nifi.security.encryption.InternalPassThroughPropertyEncryptionProvider; import org.apache.nifi.security.encryption.PropertyEncryptionEncoder; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; import org.mockito.stubbing.Answer; @@ -106,6 +108,7 @@ import java.util.HashSet; import java.util.HexFormat; import java.util.Iterator; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Optional; @@ -342,6 +345,7 @@ private ControllerServiceNode createMockControllerService() { when(service.getState()).thenReturn(ControllerServiceState.DISABLED); when(service.getBulletinLevel()).thenReturn(LogLevel.WARN); when(service.getControllerServiceImplementation()).thenReturn(new TestControllerService()); + setReferences(service); return service; } @@ -374,6 +378,15 @@ private ProcessGroup createMockProcessGroup() { when(processGroup.getFlowFileConcurrency()).thenReturn(FlowFileConcurrency.UNBOUNDED); when(processGroup.getFlowFileOutboundPolicy()).thenReturn(FlowFileOutboundPolicy.BATCH_OUTPUT); when(processGroup.getExecutionEngine()).thenReturn(ExecutionEngine.STANDARD); + when(processGroup.getProcessors()).thenReturn(Collections.emptySet()); + when(processGroup.getControllerServices(anyBoolean())).thenReturn(Collections.emptySet()); + when(processGroup.getConnections()).thenReturn(Collections.emptySet()); + when(processGroup.getInputPorts()).thenReturn(Collections.emptySet()); + when(processGroup.getOutputPorts()).thenReturn(Collections.emptySet()); + when(processGroup.getFunnels()).thenReturn(Collections.emptySet()); + when(processGroup.getLabels()).thenReturn(Collections.emptySet()); + when(processGroup.getProcessGroups()).thenReturn(Collections.emptySet()); + when(processGroup.getRemoteProcessGroups()).thenReturn(Collections.emptySet()); return processGroup; } @@ -552,6 +565,284 @@ public void testAddProcessorWithServiceAndMigration() { assertEquals(controllerServiceNode.getIdentifier(), migratedProperties.get("cs")); } + @Test + public void testUserAddedControllerServiceRemovedWhenAbsentFromProposedFlow() { + final ProcessGroup processGroup = createMockProcessGroup(); + + final PropertyDescriptor descriptor = new PropertyDescriptor.Builder().name("abc").build(); + final ControllerServiceNode serviceNode = createMockControllerService(); + when(serviceNode.getComments()).thenReturn("Added by a user"); + when(serviceNode.getName()).thenReturn("name"); + when(serviceNode.getCanonicalClassName()).thenReturn("ControllerServiceImpl"); + when(serviceNode.getProperties()).thenReturn(Map.of(descriptor, new PropertyConfiguration("123", null, null, null))); + when(serviceNode.getRawPropertyValues()).thenReturn(Map.of(descriptor, "123")); + when(serviceNode.getVersionedComponentId()).thenReturn(Optional.empty()); + when(processGroup.getControllerServices(false)).thenReturn(Set.of(serviceNode)); + + final VersionedProcessGroup versionedGroup = new VersionedProcessGroup(); + versionedGroup.setIdentifier("pg-v2"); + versionedGroup.setControllerServices(Collections.emptySet()); + versionedGroup.setProcessors(Collections.emptySet()); + + final VersionedExternalFlow externalFlow = new VersionedExternalFlow(); + externalFlow.setFlowContents(versionedGroup); + + synchronizer.synchronize(processGroup, externalFlow, synchronizationOptions); + + verify(controllerServiceProvider).removeControllerService(serviceNode); + } + + @Nested + class MigrationCreatedControllerService { + + @Test + public void doesNotRemoveWhenProposedFlowDeclaresUnrelatedService() { + final ProcessGroup processGroup = createMockProcessGroup(); + + final PropertyDescriptor descriptor = new PropertyDescriptor.Builder().name("abc").build(); + + final ControllerServiceNode localOnlyService = createMockControllerService(); + when(localOnlyService.getComments()).thenReturn(StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT); + when(localOnlyService.getName()).thenReturn("ServiceName"); + when(localOnlyService.getCanonicalClassName()).thenReturn("ServiceType"); + when(localOnlyService.getBundleCoordinate()).thenReturn(bundleCoordinate); + when(localOnlyService.getProperties()).thenReturn(Map.of(descriptor, new PropertyConfiguration("123", null, null, null))); + when(localOnlyService.getRawPropertyValues()).thenReturn(Map.of(descriptor, "123")); + when(localOnlyService.getState()).thenReturn(ControllerServiceState.DISABLED); + trackVersionedComponentId(localOnlyService); + + when(processGroup.getControllerServices(false)).thenReturn(Set.of(localOnlyService)); + + synchronizeProposedFlow(processGroup, Set.of(createMinimalVersionedControllerService()), Collections.emptySet()); + + verify(controllerServiceProvider, never()).removeControllerService(localOnlyService); + assertUnversioned(localOnlyService); + verifyControllerServiceAdded(processGroup); + } + + @Test + public void assignsProposedVersionedIdWithoutCreatingDuplicate() { + final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); + + synchronizeMigrationMatch(setup); + + assertVersionedId(setup.localService(), setup.proposedService().getIdentifier()); + verify(setup.localService()).setComments(setup.proposedService().getComments()); + verify(setup.localService()).setName(setup.proposedService().getName()); + verify(controllerServiceProvider, never()).removeControllerService(setup.localService()); + verifyNoControllerServiceAdded(setup.processGroup()); + } + + @Test + public void doesNotAssignWhenReferencedByMultipleComponents() { + final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); + + final ProcessorNode additionalReferencer = createMappableProcessor(setup.processGroup()); + setReferences(setup.localService(), setup.processor(), additionalReferencer); + + synchronizeMigrationMatch(setup); + + assertUnversioned(setup.localService()); + verify(controllerServiceProvider, never()).removeControllerService(setup.localService()); + verifyControllerServiceAdded(setup.processGroup()); + } + + @Test + public void doesNotAssignWhenProposedServiceTypeDiffers() { + final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); + when(setup.localService().getCanonicalClassName()).thenReturn("org.apache.nifi.cs.DifferentService"); + + synchronizeMigrationMatch(setup); + + assertUnversioned(setup.localService()); + verify(controllerServiceProvider, never()).removeControllerService(setup.localService()); + verifyControllerServiceAdded(setup.processGroup()); + } + + @Test + public void assignsNestedServicesWhenInnerListedBeforeOuter() { + final MigrationCreatedMatchSetup outer = newMigrationCreatedMatchSetup(); + + final VersionedControllerService proposedInnerService = createMinimalVersionedControllerService(); + proposedInnerService.setIdentifier("inner-service-versioned-id"); + + final ControllerServiceNode innerService = createMockControllerService(); + stubMigrationCreatedService(innerService, proposedInnerService); + stubControllerServiceReferenceProperty(outer.localService(), storeServicePropertyDescriptor(), innerService.getIdentifier()); + setReferences(innerService, outer.localService()); + when(outer.processGroup().findControllerService(eq(innerService.getIdentifier()), anyBoolean(), anyBoolean())).thenReturn(innerService); + + outer.proposedService().setProperties(Map.of("Store Service", proposedInnerService.getIdentifier())); + outer.proposedService().setPropertyDescriptors(Map.of("Store Service", storeServiceVersionedDescriptor())); + + final Set localServices = new LinkedHashSet<>(); + localServices.add(innerService); + localServices.add(outer.localService()); + + synchronizeMigrationMatch(outer, localServices, Set.of(outer.proposedService(), proposedInnerService)); + + assertVersionedId(outer.localService(), outer.proposedService().getIdentifier()); + assertVersionedId(innerService, proposedInnerService.getIdentifier()); + verify(controllerServiceProvider, never()).removeControllerService(innerService); + verify(controllerServiceProvider, never()).removeControllerService(outer.localService()); + verifyNoControllerServiceAdded(outer.processGroup()); + } + + @Test + public void doesNotAssignWhenProposedVersionedIdAlreadyUsed() { + final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); + + final ControllerServiceNode existingService = createMockControllerService(); + when(existingService.getVersionedComponentId()).thenReturn(Optional.of(setup.proposedService().getIdentifier())); + when(existingService.getComments()).thenReturn(""); + stubControllerServiceListedInFlow(existingService, setup.proposedService()); + + synchronizeMigrationMatch(setup, Set.of(existingService, setup.localService())); + + assertUnversioned(setup.localService()); + verify(controllerServiceProvider, never()).removeControllerService(setup.localService()); + verifyNoControllerServiceAdded(setup.processGroup()); + } + + private MigrationCreatedMatchSetup newMigrationCreatedMatchSetup() { + final ProcessGroup processGroup = createMockProcessGroup(); + + final VersionedControllerService proposedService = createMinimalVersionedControllerService(); + proposedService.setIdentifier("declared-service-versioned-id"); + + final PropertyDescriptor storeServiceDescriptor = storeServicePropertyDescriptor(); + + final ControllerServiceNode localService = createMockControllerService(); + final String localServiceId = localService.getIdentifier(); + stubMigrationCreatedService(localService, proposedService); + + final ProcessorNode processor = createMappableProcessor(processGroup); + when(processor.getVersionedComponentId()).thenReturn(Optional.of("processor-versioned-id")); + stubControllerServiceReferenceProperty(processor, storeServiceDescriptor, localServiceId); + setReferences(localService, processor); + + when(processGroup.getControllerServices(false)).thenReturn(Set.of(localService)); + when(processGroup.getProcessors()).thenReturn(Set.of(processor)); + when(processGroup.findControllerService(eq(localServiceId), anyBoolean(), anyBoolean())).thenReturn(localService); + + final VersionedProcessor proposedProcessor = createMinimalVersionedProcessor(); + proposedProcessor.setIdentifier("processor-versioned-id"); + proposedProcessor.setProperties(Map.of("Store Service", proposedService.getIdentifier())); + proposedProcessor.setPropertyDescriptors(Map.of("Store Service", storeServiceVersionedDescriptor())); + + return new MigrationCreatedMatchSetup(processGroup, localService, processor, proposedService, proposedProcessor); + } + + private void synchronizeMigrationMatch(final MigrationCreatedMatchSetup setup) { + synchronizeMigrationMatch(setup, Set.of(setup.localService())); + } + + private void synchronizeMigrationMatch(final MigrationCreatedMatchSetup setup, final Set localServices) { + synchronizeMigrationMatch(setup, localServices, Set.of(setup.proposedService())); + } + + private void synchronizeMigrationMatch( + final MigrationCreatedMatchSetup setup, + final Set localServices, + final Set proposedServices + ) { + when(setup.processGroup().getControllerServices(false)).thenReturn(localServices); + synchronizeProposedFlow(setup.processGroup(), proposedServices, Set.of(setup.proposedProcessor())); + } + + private record MigrationCreatedMatchSetup( + ProcessGroup processGroup, + ControllerServiceNode localService, + ProcessorNode processor, + VersionedControllerService proposedService, + VersionedProcessor proposedProcessor + ) { + } + + private PropertyDescriptor storeServicePropertyDescriptor() { + return new PropertyDescriptor.Builder() + .name("Store Service") + .identifiesControllerService(ControllerService.class) + .build(); + } + + private VersionedPropertyDescriptor storeServiceVersionedDescriptor() { + final VersionedPropertyDescriptor proposedDescriptor = new VersionedPropertyDescriptor(); + proposedDescriptor.setName("Store Service"); + proposedDescriptor.setIdentifiesControllerService(true); + return proposedDescriptor; + } + + private void stubControllerServiceListedInFlow(final ControllerServiceNode localService, final VersionedControllerService proposed) { + when(localService.getCanonicalClassName()).thenReturn(proposed.getType()); + when(localService.getName()).thenReturn(proposed.getName()); + when(localService.getProperties()).thenReturn(Map.of(new PropertyDescriptor.Builder().name("abc").build(), + new PropertyConfiguration("123", null, null, null))); + when(localService.getRawPropertyValues()).thenReturn(Map.of(new PropertyDescriptor.Builder().name("abc").build(), "123")); + } + + private void stubMigrationCreatedService(final ControllerServiceNode localService, final VersionedControllerService proposed) { + stubControllerServiceListedInFlow(localService, proposed); + when(localService.getComments()).thenReturn(StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT); + trackVersionedComponentId(localService); + } + + private void stubControllerServiceReferenceProperty(final ComponentNode component, final PropertyDescriptor descriptor, final String serviceId) { + final PropertyConfiguration configuration = new PropertyConfiguration(serviceId, null, null, null); + when(component.getRawPropertyValues()).thenReturn(Map.of(descriptor, serviceId)); + when(component.getProperties()).thenReturn(Map.of(descriptor, configuration)); + when(component.getProperty(eq(descriptor))).thenReturn(configuration); + when(component.getPropertyDescriptor(eq(descriptor.getName()))).thenReturn(descriptor); + } + + /** + * Makes the mock report the versioned component id that synchronization assigns to it, so that a test can + * assert on the resulting state rather than on the setter alone. + */ + private void trackVersionedComponentId(final ControllerServiceNode service) { + final AtomicReference assignedVersionedComponentId = new AtomicReference<>(); + when(service.getVersionedComponentId()).thenAnswer(invocation -> Optional.ofNullable(assignedVersionedComponentId.get())); + doAnswer(invocation -> { + assignedVersionedComponentId.set(invocation.getArgument(0)); + return null; + }).when(service).setVersionedComponentId(any()); + } + + private void assertUnversioned(final ControllerServiceNode service) { + assertTrue(service.getVersionedComponentId().isEmpty(), "Expected no versioned component id but found " + service.getVersionedComponentId()); + } + + private void assertVersionedId(final ControllerServiceNode service, final String expectedVersionedId) { + assertEquals(Optional.of(expectedVersionedId), service.getVersionedComponentId()); + } + + private void verifyControllerServiceAdded(final ProcessGroup processGroup) { + verify(flowManager).createControllerService(any(), any(), any(), anySet(), anyBoolean(), anyBoolean(), nullable(String.class)); + verify(processGroup).addControllerService(any(ControllerServiceNode.class)); + } + + private void verifyNoControllerServiceAdded(final ProcessGroup processGroup) { + verify(flowManager, never()).createControllerService(any(), any(), any(), anySet(), anyBoolean(), anyBoolean(), nullable(String.class)); + verify(processGroup, never()).addControllerService(any(ControllerServiceNode.class)); + } + + private void synchronizeProposedFlow( + final ProcessGroup processGroup, + final Set proposedServices, + final Set proposedProcessors + ) { + final VersionedProcessGroup versionedGroup = new VersionedProcessGroup(); + versionedGroup.setIdentifier("pg-v2"); + versionedGroup.setControllerServices(proposedServices); + versionedGroup.setProcessors(proposedProcessors); + + final VersionedExternalFlow externalFlow = new VersionedExternalFlow(); + externalFlow.setFlowContents(versionedGroup); + + synchronizer.synchronize(processGroup, externalFlow, synchronizationOptions); + } + } + @Test public void testSynchronizeProcessorSensitiveDynamicProperties() throws FlowSynchronizationException, InterruptedException, TimeoutException { final Map versionedProperties = Collections.singletonMap(SENSITIVE_PROPERTY_NAME, ENCRYPTED_PROPERTY_VALUE); diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java new file mode 100644 index 000000000000..e0ce89410f1b --- /dev/null +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java @@ -0,0 +1,97 @@ +/* + * 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.cs.tests.system; + +import org.apache.nifi.annotation.behavior.Stateful; +import org.apache.nifi.annotation.lifecycle.OnEnabled; +import org.apache.nifi.components.PropertyDescriptor; +import org.apache.nifi.components.state.Scope; +import org.apache.nifi.components.state.StateManager; +import org.apache.nifi.components.state.StateMap; +import org.apache.nifi.controller.AbstractControllerService; +import org.apache.nifi.controller.ConfigurationContext; +import org.apache.nifi.processor.util.StandardValidators; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +/** + * Store whose contents live in the component's local state, so that they are subject to the framework's own state + * lifecycle. The store records when it was first created and how many rows it holds. + * + * Removing a Controller Service clears its state, so a service that was torn down and recreated comes back with a + * later creation timestamp and no rows. Comparing the creation timestamp across an operation therefore distinguishes + * a preserved service from one that was replaced, even when the replacement carries the same identifier. + */ +@Stateful(scopes = Scope.LOCAL, description = "Holds the creation timestamp of the store and the number of rows written to it.") +public class StateBackedStoreService extends AbstractControllerService implements StoreService { + + static final PropertyDescriptor STORE_NAME = new PropertyDescriptor.Builder() + .name("Store Name") + .required(true) + .defaultValue("store") + .addValidator(StandardValidators.NON_EMPTY_VALIDATOR) + .build(); + + static final String CREATED_KEY = "created"; + static final String ROW_COUNT_KEY = "rowCount"; + static final String LAST_ROW_KEY = "lastRow"; + + private static final List PROPERTIES = List.of(STORE_NAME); + + @Override + protected List getSupportedPropertyDescriptors() { + return PROPERTIES; + } + + /** + * Establishes the store on first enable and leaves it alone afterwards, so that the creation timestamp survives + * the disable and re-enable cycle that a version change performs. + */ + @OnEnabled + public void onEnabled(final ConfigurationContext context) throws IOException { + final StateManager stateManager = getStateManager(); + final StateMap stateMap = stateManager.getState(Scope.LOCAL); + if (stateMap.get(CREATED_KEY) != null) { + return; + } + + final Map state = new HashMap<>(); + state.put(CREATED_KEY, String.valueOf(System.currentTimeMillis())); + state.put(ROW_COUNT_KEY, "0"); + stateManager.setState(state, Scope.LOCAL); + } + + @Override + public synchronized void append(final String row) { + try { + final StateManager stateManager = getStateManager(); + final Map state = new HashMap<>(stateManager.getState(Scope.LOCAL).toMap()); + final long rowCount = Long.parseLong(state.getOrDefault(ROW_COUNT_KEY, "0")); + state.put(ROW_COUNT_KEY, String.valueOf(rowCount + 1)); + state.put(LAST_ROW_KEY, row); + stateManager.setState(state, Scope.LOCAL); + } catch (final IOException e) { + throw new UncheckedIOException(e); + } + } + +} diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StoreService.java b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StoreService.java new file mode 100644 index 000000000000..121a4fe9e098 --- /dev/null +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StoreService.java @@ -0,0 +1,29 @@ +/* + * 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.cs.tests.system; + +import org.apache.nifi.controller.ControllerService; + +/** + * A store that accumulates rows and whose contents must survive flow and runtime upgrades. + */ +public interface StoreService extends ControllerService { + + void append(String row); + +} diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java new file mode 100644 index 000000000000..64a6582a6c6c --- /dev/null +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java @@ -0,0 +1,90 @@ +/* + * 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.processors.tests.system; + +import org.apache.nifi.annotation.behavior.InputRequirement; +import org.apache.nifi.annotation.behavior.InputRequirement.Requirement; +import org.apache.nifi.components.PropertyDescriptor; +import org.apache.nifi.cs.tests.system.StateBackedStoreService; +import org.apache.nifi.cs.tests.system.StoreService; +import org.apache.nifi.migration.PropertyConfiguration; +import org.apache.nifi.processor.AbstractProcessor; +import org.apache.nifi.processor.ProcessContext; +import org.apache.nifi.processor.ProcessSession; +import org.apache.nifi.processor.Relationship; +import org.apache.nifi.processor.exception.ProcessException; + +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.atomic.AtomicLong; + +/** + * Post-upgrade shape of a processor whose property migration creates the store Controller Service + * that the pre-upgrade shape did not have. Each execution appends a row to the store so that tests + * can observe whether store contents survive flow and runtime upgrades. + */ +@InputRequirement(Requirement.INPUT_FORBIDDEN) +public class MigrateToControllerService extends AbstractProcessor { + + static final PropertyDescriptor STORE_SERVICE = new PropertyDescriptor.Builder() + .name("Store Service") + .required(false) + .identifiesControllerService(StoreService.class) + .build(); + + static final Relationship REL_SUCCESS = new Relationship.Builder().name("success").build(); + + private static final List PROPERTIES = List.of(STORE_SERVICE); + private static final Set RELATIONSHIPS = Set.of(REL_SUCCESS); + + private final AtomicLong rowCounter = new AtomicLong(0L); + + @Override + protected List getSupportedPropertyDescriptors() { + return PROPERTIES; + } + + @Override + public Set getRelationships() { + return RELATIONSHIPS; + } + + @Override + public void migrateProperties(final PropertyConfiguration config) { + final String storeName = config.getPropertyValue("store-name").orElse(null); + config.removeProperty("store-name"); + + if (storeName != null) { + final String serviceId = config.createControllerService(StateBackedStoreService.class.getName(), Map.of("Store Name", storeName)); + config.setProperty(STORE_SERVICE, serviceId); + } + } + + @Override + public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException { + final StoreService storeService = context.getProperty(STORE_SERVICE).asControllerService(StoreService.class); + if (storeService == null) { + context.yield(); + return; + } + + storeService.append("row-" + rowCounter.getAndIncrement()); + } + +} diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.controller.ControllerService b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.controller.ControllerService index 4fe1a28f3574..2fb8f48cbb51 100644 --- a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.controller.ControllerService +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.controller.ControllerService @@ -14,3 +14,4 @@ # limitations under the License. org.apache.nifi.cs.tests.system.MigrationService +org.apache.nifi.cs.tests.system.StateBackedStoreService diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor index c112e05097f0..e19a568035b1 100644 --- a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor @@ -14,3 +14,4 @@ # limitations under the License. org.apache.nifi.processors.tests.system.MigrateProperties +org.apache.nifi.processors.tests.system.MigrateToControllerService diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/pom.xml b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/pom.xml index 875b26ea181d..90e9c6d12550 100644 --- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/pom.xml +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/pom.xml @@ -45,10 +45,6 @@ nifi-bin-manager 2.13.0-SNAPSHOT - - com.fasterxml.jackson.core - jackson-databind - org.apache.nifi nifi-connector-common diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java new file mode 100644 index 000000000000..9b0ce040aa28 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java @@ -0,0 +1,68 @@ +/* + * 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.processors.tests.system; + +import org.apache.nifi.annotation.behavior.InputRequirement; +import org.apache.nifi.annotation.behavior.InputRequirement.Requirement; +import org.apache.nifi.components.PropertyDescriptor; +import org.apache.nifi.processor.AbstractProcessor; +import org.apache.nifi.processor.ProcessContext; +import org.apache.nifi.processor.ProcessSession; +import org.apache.nifi.processor.Relationship; +import org.apache.nifi.processor.exception.ProcessException; +import org.apache.nifi.processor.util.StandardValidators; + +import java.util.List; +import java.util.Set; + +/** + * Pre-upgrade shape of a processor that keeps its store location in a plain property. + * The post-upgrade shape of the same processor, in the alternate-config extensions bundle, + * migrates that property into a Controller Service. + */ +@InputRequirement(Requirement.INPUT_FORBIDDEN) +public class MigrateToControllerService extends AbstractProcessor { + + static final PropertyDescriptor STORE_NAME = new PropertyDescriptor.Builder() + .name("store-name") + .displayName("Store Name") + .required(true) + .addValidator(StandardValidators.NON_EMPTY_VALIDATOR) + .build(); + + static final Relationship REL_SUCCESS = new Relationship.Builder().name("success").build(); + + private static final List PROPERTIES = List.of(STORE_NAME); + private static final Set RELATIONSHIPS = Set.of(REL_SUCCESS); + + @Override + protected List getSupportedPropertyDescriptors() { + return PROPERTIES; + } + + @Override + public Set getRelationships() { + return RELATIONSHIPS; + } + + @Override + public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException { + context.yield(); + } + +} diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor index 7d9389d9be3e..3e267e1906a6 100644 --- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.processor.Processor @@ -38,6 +38,7 @@ org.apache.nifi.processors.tests.system.HoldInput org.apache.nifi.processors.tests.system.IngestFile org.apache.nifi.processors.tests.system.LoopFlowFile org.apache.nifi.processors.tests.system.MigrateProperties +org.apache.nifi.processors.tests.system.MigrateToControllerService org.apache.nifi.processors.tests.system.MultiKeyState org.apache.nifi.processors.tests.system.MultiKeyStateNotDroppable org.apache.nifi.processors.tests.system.PartitionText diff --git a/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml new file mode 100644 index 000000000000..e3e6c968afe7 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml @@ -0,0 +1,46 @@ + + + + 4.0.0 + + + org.apache.nifi + nifi-system-test-flow-registry-bundle + 2.12.0-SNAPSHOT + + + nifi-system-test-flow-registry-nar + nar + + + + org.apache.nifi + nifi-api + provided + + + org.apache.nifi + nifi-framework-api + 2.12.0-SNAPSHOT + provided + + + org.apache.nifi + nifi-system-test-flow-registry + 2.12.0-SNAPSHOT + + + diff --git a/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml new file mode 100644 index 000000000000..8e212da4cfd1 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml @@ -0,0 +1,42 @@ + + + + + nifi-system-test-flow-registry-bundle + org.apache.nifi + 2.12.0-SNAPSHOT + + 4.0.0 + + nifi-system-test-flow-registry + + + + org.apache.nifi + nifi-framework-api + 2.12.0-SNAPSHOT + + + org.apache.nifi + nifi-utils + 2.12.0-SNAPSHOT + + + com.fasterxml.jackson.core + jackson-databind + + + diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/flow/registry/FileSystemFlowRegistryClient.java b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/src/main/java/org/apache/nifi/flow/registry/FileSystemFlowRegistryClient.java similarity index 100% rename from nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/flow/registry/FileSystemFlowRegistryClient.java rename to nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/src/main/java/org/apache/nifi/flow/registry/FileSystemFlowRegistryClient.java diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.registry.flow.FlowRegistryClient b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/src/main/resources/META-INF/services/org.apache.nifi.registry.flow.FlowRegistryClient similarity index 100% rename from nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/resources/META-INF/services/org.apache.nifi.registry.flow.FlowRegistryClient rename to nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/src/main/resources/META-INF/services/org.apache.nifi.registry.flow.FlowRegistryClient diff --git a/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml b/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml new file mode 100644 index 000000000000..ae125c11e997 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml @@ -0,0 +1,31 @@ + + + + + nifi-system-tests + org.apache.nifi + 2.12.0-SNAPSHOT + + 4.0.0 + + nifi-system-test-flow-registry-bundle + pom + + + nifi-system-test-flow-registry + nifi-system-test-flow-registry-nar + + diff --git a/nifi-system-tests/nifi-system-test-suite/pom.xml b/nifi-system-tests/nifi-system-test-suite/pom.xml index 387c9f7ec9c3..699d3deacfc5 100644 --- a/nifi-system-tests/nifi-system-test-suite/pom.xml +++ b/nifi-system-tests/nifi-system-test-suite/pom.xml @@ -379,6 +379,12 @@ 2.13.0-SNAPSHOT nar + + org.apache.nifi + nifi-system-test-flow-registry-nar + 2.12.0-SNAPSHOT + nar + org.apache.nifi nifi-system-test-flow-action-reporter-nar diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/assembly/dependencies.xml b/nifi-system-tests/nifi-system-test-suite/src/test/assembly/dependencies.xml index 6c8cf4c7640d..4ca2ce1ec7db 100644 --- a/nifi-system-tests/nifi-system-test-suite/src/test/assembly/dependencies.xml +++ b/nifi-system-tests/nifi-system-test-suite/src/test/assembly/dependencies.xml @@ -68,6 +68,7 @@ *:nifi-system-test-extensions-services-nar *:nifi-system-test-extensions-services-api-nar *:nifi-system-test-extensions2-nar + *:nifi-system-test-flow-registry-nar *:nifi-system-test-flow-action-reporter-nar *:nifi-system-test-component-metric-reporter-nar *:nifi-jetty-nar diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java index faef5b9a5245..e0e763a60cc6 100644 --- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java +++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiClientUtil.java @@ -2617,7 +2617,7 @@ public ReportingTaskEntity updateReportingTask(final ReportingTaskEntity current } public FlowRegistryClientEntity createFlowRegistryClient(final String name) throws NiFiClientException, IOException { - final BundleDTO bundleDto = new BundleDTO(NiFiSystemIT.NIFI_GROUP_ID, NiFiSystemIT.TEST_EXTENSIONS_ARTIFACT_ID, nifiVersion); + final BundleDTO bundleDto = new BundleDTO(NiFiSystemIT.NIFI_GROUP_ID, NiFiSystemIT.TEST_FLOW_REGISTRY_ARTIFACT_ID, nifiVersion); final FlowRegistryClientDTO clientDto = new FlowRegistryClientDTO(); clientDto.setBundle(bundleDto); clientDto.setName(name); diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiSystemIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiSystemIT.java index fae7e7bd76ee..5763ffcc66e1 100644 --- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiSystemIT.java +++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/NiFiSystemIT.java @@ -85,6 +85,7 @@ public abstract class NiFiSystemIT implements NiFiInstanceProvider { public static final int STANDALONE_CLIENT_API_BASE_PORT = 5670; public static final String NIFI_GROUP_ID = "org.apache.nifi"; public static final String TEST_EXTENSIONS_ARTIFACT_ID = "nifi-system-test-extensions-nar"; + public static final String TEST_FLOW_REGISTRY_ARTIFACT_ID = "nifi-system-test-flow-registry-nar"; public static final String TEST_EXTENSIONS_SERVICES_ARTIFACT_ID = "nifi-system-test-extensions-services-nar"; public static final String TEST_PYTHON_EXTENSIONS_ARTIFACT_ID = "python-extensions"; public static final String TEST_PARAM_PROVIDERS_PACKAGE = "org.apache.nifi.parameter.tests.system"; diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java new file mode 100644 index 000000000000..0d1521903f09 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java @@ -0,0 +1,365 @@ +/* + * 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.tests.system.migration; + +import org.apache.nifi.migration.StandardControllerServiceFactory; +import org.apache.nifi.tests.system.AbstractNarSwapMigrationIT; +import org.apache.nifi.tests.system.ExceptionalBooleanSupplier; +import org.apache.nifi.toolkit.client.NiFiClientException; +import org.apache.nifi.web.api.dto.ComponentStateDTO; +import org.apache.nifi.web.api.dto.StateEntryDTO; +import org.apache.nifi.web.api.entity.ComponentStateEntity; +import org.apache.nifi.web.api.entity.ControllerServiceEntity; +import org.apache.nifi.web.api.entity.FlowRegistryClientEntity; +import org.apache.nifi.web.api.entity.ProcessGroupEntity; +import org.apache.nifi.web.api.entity.ProcessorEntity; +import org.apache.nifi.web.api.entity.VersionedFlowUpdateRequestEntity; +import org.junit.jupiter.api.Test; + +import java.io.File; +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Collection; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.fail; + +/** + * Verifies that a Controller Service created by property migration survives the operations a deployed flow goes + * through: a runtime upgrade that introduces the migration, a plain runtime restart, and version changes of the + * enclosing versioned Process Group. + * + * Two independent flow lineages are pre-seeded in the registry. In the first, a later version declares the store + * Controller Service itself, as a published flow would once the vendor adds it. In the second, no version ever + * declares the service, so the flow relies entirely on property migration to create it, and the later version only + * adds an unrelated processor. + */ +public class MigrationCreatedControllerServiceVersioningIT extends AbstractNarSwapMigrationIT { + private static final String TEST_FLOWS_BUCKET = "test-flows"; + private static final String SERVICE_DECLARED_FLOW_ID = "11111111-2222-3333-4444-555555555555"; + private static final String SERVICE_ABSENT_FLOW_ID = "22222222-3333-4444-5555-666666666666"; + private static final String DECLARED_SERVICE_VERSIONED_ID = "99999999-8888-7777-6666-555555555555"; + private static final String STORE_SERVICE_PROPERTY = "Store Service"; + private static final String STORE_SERVICE_TYPE = "org.apache.nifi.cs.tests.system.StateBackedStoreService"; + private static final String PROCESSOR_TYPE = "org.apache.nifi.processors.tests.system.MigrateToControllerService"; + private static final String MIGRATING_PROCESSOR_NAME = "MigrateToControllerService"; + private static final String ADDED_PROCESSOR_NAME = "Added Processor"; + private static final String VERSIONED_FLOWS_DIRECTORY = "src/test/resources/versioned-flows"; + private static final String CREATED_STATE_KEY = "created"; + private static final String ROW_COUNT_STATE_KEY = "rowCount"; + private static final Duration CONDITION_TIMEOUT = Duration.ofSeconds(30); + private static final Duration CONDITION_POLL = Duration.ofMillis(100); + + /** + * After the runtime is upgraded, the Controller Service that property migration creates must be present, + * enabled and referenced, and the flow that was running before the upgrade must be running again, with no manual action. + */ + @Test + public void testRuntimeUpgradeCreatesEnabledServiceAndKeepsFlowRunning() throws NiFiClientException, IOException, InterruptedException { + final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_DECLARED_FLOW_ID); + final ControllerServiceEntity service = waitForSingleStoreService(flow.groupId()); + final String serviceId = service.getComponent().getId(); + + assertEquals(StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT, service.getComponent().getComments()); + assertBelongsToLocalFlowOnly(service); + getClientUtil().waitForControllerServiceRunStatus(serviceId, "ENABLED"); + getClientUtil().waitForRunningProcessor(flow.processorId()); + assertEquals(serviceId, getStoreServiceId(flow.processorId())); + + final Collection validationErrors = getNifiClient().getProcessorClient().getProcessor(flow.processorId()).getComponent().getValidationErrors(); + final boolean processorValid = validationErrors == null || validationErrors.isEmpty(); + assertTrue(processorValid, "Processor must be valid after the runtime upgrade"); + + waitForCondition(() -> countRows(serviceId) > 0, "store row count > 0"); + + final String versionedFlowState = getClientUtil().getVersionedFlowState(flow.groupId(), "root"); + assertNotEquals("LOCALLY_MODIFIED", versionedFlowState, "The migration-created Controller Service must not make the flow dirty"); + assertNotEquals("LOCALLY_MODIFIED_AND_STALE", versionedFlowState, "The migration-created Controller Service must not make the flow dirty"); + + final boolean serviceReportedAsLocalModification = getNifiClient().getProcessGroupClient().getLocalModifications(flow.groupId()) + .getComponentDifferences().stream() + .anyMatch(diff -> serviceId.equals(diff.getComponentId())); + assertFalse(serviceReportedAsLocalModification); + } + + /** + * Upgrading to a flow version that declares the store Controller Service must keep using the service that + * property migration already created, rather than removing it and substituting the one the published version declares. + */ + @Test + public void testFlowUpgradePreservesMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { + final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_DECLARED_FLOW_ID); + final MigratedStore store = awaitPopulatedStoreService(flow); + assertBelongsToLocalFlowOnly(waitForSingleStoreService(flow.groupId())); + + final VersionedFlowUpdateRequestEntity upgradeRequest = getClientUtil().changeFlowVersion(flow.groupId(), "2", false); + + assertStorePreserved(flow, store, "flow upgrade"); + getClientUtil().assertFlowUpToDate(flow.groupId()); + + final ControllerServiceEntity serviceAfterUpgrade = waitForSingleStoreService(flow.groupId()); + assertEquals(DECLARED_SERVICE_VERSIONED_ID, serviceAfterUpgrade.getComponent().getVersionedComponentId(), + "The service must track the Controller Service that version 2 declares, instead of belonging to the local flow only"); + assertEquals("", serviceAfterUpgrade.getComponent().getComments(), + "Comments must come from version 2 once the service tracks the declared Controller Service"); + assertNull(upgradeRequest.getRequest().getFailureReason(), + "The migration-created Controller Service must not prevent the versioned flow from upgrading"); + } + + /** + * A flow whose definition never declares the store Controller Service relies entirely on property migration to + * create it. Upgrading such a flow to a version that only adds an unrelated processor, leaving the migrating + * processor untouched, must not disturb the service: the version change has nothing to say about it. + */ + @Test + public void testFlowUpgradeAddingUnrelatedProcessorPreservesMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { + final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_ABSENT_FLOW_ID); + final MigratedStore store = awaitPopulatedStoreService(flow); + + final VersionedFlowUpdateRequestEntity upgradeRequest = getClientUtil().changeFlowVersion(flow.groupId(), "2", false); + + assertTrue(hasProcessorNamed(flow.groupId(), ADDED_PROCESSOR_NAME), "The upgrade must have added the unrelated processor"); + + assertStorePreserved(flow, store, "flow upgrade"); + getClientUtil().assertFlowUpToDate(flow.groupId()); + + final ControllerServiceEntity serviceAfterUpgrade = waitForSingleStoreService(flow.groupId()); + assertBelongsToLocalFlowOnly(serviceAfterUpgrade); + assertEquals(StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT, serviceAfterUpgrade.getComponent().getComments()); + assertNull(upgradeRequest.getRequest().getFailureReason(), + "The migration-created Controller Service must not prevent the versioned flow from upgrading"); + } + + /** + * Restarting the runtime with no NAR change must reuse the Controller Service that property migration + * created previously rather than creating a second one. + */ + @Test + public void testRuntimeRestartDoesNotRecreateMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { + final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_DECLARED_FLOW_ID); + final MigratedStore store = awaitPopulatedStoreService(flow); + + getNiFiInstance().stop(); + getNiFiInstance().start(true); + + assertStorePreserved(flow, store, "runtime restart"); + assertBelongsToLocalFlowOnly(waitForSingleStoreService(flow.groupId())); + } + + /** + * Asserts that the given Controller Service has no counterpart in the flow definition. Such a service either + * carries no versioned component id at all, or carries the placeholder that the framework derives from its own + * instance id while mapping the group. A service that tracks a Controller Service declared by the flow definition + * carries the identifier from the definition instead. + */ + private void assertBelongsToLocalFlowOnly(final ControllerServiceEntity service) { + final String versionedComponentId = service.getComponent().getVersionedComponentId(); + if (versionedComponentId == null) { + return; + } + + final String derivedFromInstanceId = UUID.nameUUIDFromBytes(service.getComponent().getId().getBytes(StandardCharsets.UTF_8)).toString(); + assertEquals(derivedFromInstanceId, versionedComponentId, + "The service tracks a Controller Service declared by the flow definition, so it no longer belongs to the local flow only"); + } + + /** + * Imports version 1 of the given pre-seeded flow, waits for its processor to be running, and then simulates a + * runtime upgrade by stopping NiFi, swapping in the alternate-config extensions and starting NiFi again. Nothing + * in the flow is stopped or started by hand, so the assertions that follow observe what the upgrade does on its own. + */ + private MigratedFlow importAndUpgradeRuntime(final String flowId) throws NiFiClientException, IOException, InterruptedException { + final FlowRegistryClientEntity registryClient = registerClient(new File(VERSIONED_FLOWS_DIRECTORY)); + final ProcessGroupEntity group = getClientUtil().importFlowFromRegistry("root", registryClient.getId(), TEST_FLOWS_BUCKET, flowId, "1"); + final ProcessorEntity processor = findProcessor(group.getId(), MIGRATING_PROCESSOR_NAME); + getClientUtil().waitForRunningProcessor(processor.getId()); + + getNiFiInstance().stop(); + switchOutNars(); + getNiFiInstance().start(true); + + return new MigratedFlow(group.getId(), processor.getId()); + } + + /** + * Waits until the migration-created Controller Service is enabled and its store has accumulated at least one row, + * so that the flow is known to be working before the operation under test runs, and returns its identifier. + */ + private MigratedStore awaitPopulatedStoreService(final MigratedFlow flow) throws NiFiClientException, IOException, InterruptedException { + final ControllerServiceEntity service = waitForSingleStoreService(flow.groupId()); + final String serviceId = service.getComponent().getId(); + getClientUtil().waitForControllerServiceRunStatus(serviceId, "ENABLED"); + waitForCondition(() -> countRows(serviceId) > 0, "store row count > 0"); + + return new MigratedStore(serviceId, readCreatedTimestamp(flow, serviceId)); + } + + /** + * Asserts that the given operation left the migration-created Controller Service in place and the flow working. + * + * Two independent properties are checked. An unchanged creation timestamp proves the service was never torn down, + * since removing a Controller Service clears its state and a replacement records a new timestamp. A row count that + * keeps climbing afterwards proves the flow is doing work again, which a RUNNING state alone does not show: a + * processor whose Controller Service reference is broken still reports RUNNING while writing nothing. + * + * The store is read through the service the processor references now rather than the one recorded earlier, so + * that a service substituted by the operation is detected rather than silently passing. + */ + private void assertStorePreserved(final MigratedFlow flow, final MigratedStore store, final String operation) + throws NiFiClientException, IOException, InterruptedException { + + waitForCondition(() -> !findStoreServices(flow.groupId()).isEmpty(), "store Controller Service found"); + + final List servicesAfter = findControllerServicesInGroup(flow.groupId()); + assertEquals(1, servicesAfter.size(), + "The " + operation + " left " + servicesAfter.size() + " Controller Services in the group, so one was created alongside the migration-created one"); + + final ControllerServiceEntity serviceAfter = servicesAfter.getFirst(); + assertEquals(store.serviceId(), serviceAfter.getComponent().getId(), + "The " + operation + " replaced the migration-created Controller Service with a different one"); + assertEquals(store.serviceId(), getStoreServiceId(flow.processorId())); + + getClientUtil().waitForControllerServiceRunStatus(store.serviceId(), "ENABLED"); + getClientUtil().waitForRunningProcessor(flow.processorId()); + + assertEquals(store.created(), readCreatedTimestamp(flow, store.serviceId()), + "The migration-created Controller Service was removed during the " + operation + ": its store was created again, so the state it held is gone"); + + final long rowsAfterOperation = countRows(getStoreServiceId(flow.processorId())); + waitForCondition(() -> countRows(getStoreServiceId(flow.processorId())) > rowsAfterOperation, "store row count increase"); + } + + + private ProcessorEntity findProcessor(final String groupId, final String name) throws NiFiClientException, IOException { + final List matching = getNifiClient().getFlowClient().getProcessGroup(groupId).getProcessGroupFlow().getFlow().getProcessors().stream() + .filter(processor -> PROCESSOR_TYPE.equals(processor.getComponent().getType())) + .filter(processor -> name.equals(processor.getComponent().getName())) + .toList(); + + if (matching.size() != 1) { + throw new AssertionError("Expected exactly one processor named " + name + " in group " + groupId + " but found " + matching.size()); + } + + return matching.getFirst(); + } + + private boolean hasProcessorNamed(final String groupId, final String name) throws NiFiClientException, IOException { + return getNifiClient().getFlowClient().getProcessGroup(groupId).getProcessGroupFlow().getFlow().getProcessors().stream() + .anyMatch(processor -> name.equals(processor.getComponent().getName())); + } + + private List findControllerServicesInGroup(final String groupId) throws NiFiClientException, IOException { + return getNifiClient().getFlowClient().getControllerServices(groupId).getControllerServices().stream() + .filter(service -> groupId.equals(service.getComponent().getParentGroupId())) + .toList(); + } + + private List findStoreServices(final String groupId) throws NiFiClientException, IOException { + return findControllerServicesInGroup(groupId).stream() + .filter(service -> STORE_SERVICE_TYPE.equals(service.getComponent().getType())) + .toList(); + } + + private ControllerServiceEntity waitForSingleStoreService(final String groupId) throws NiFiClientException, IOException, InterruptedException { + waitForCondition(() -> findStoreServices(groupId).size() == 1, "exactly one store Controller Service"); + return findStoreServices(groupId).getFirst(); + } + + private String getStoreServiceId(final String processorId) throws NiFiClientException, IOException { + final Map properties = getNifiClient().getProcessorClient().getProcessor(processorId).getComponent().getConfig().getProperties(); + return properties.get(STORE_SERVICE_PROPERTY); + } + + /** + * Reads the store's local component state. Removing a Controller Service clears its state, so an empty map means + * either that the service has not been enabled yet or that it was torn down. + */ + private Map readStoreState(final String serviceId) throws NiFiClientException, IOException { + final ComponentStateEntity stateEntity = getNifiClient().getControllerServicesClient().getControllerServiceState(serviceId); + final ComponentStateDTO componentState = stateEntity.getComponentState(); + if (componentState == null || componentState.getLocalState() == null || componentState.getLocalState().getState() == null) { + return Map.of(); + } + + final Map state = new HashMap<>(); + for (final StateEntryDTO entry : componentState.getLocalState().getState()) { + state.put(entry.getKey(), entry.getValue()); + } + + return state; + } + + /** + * Returns the timestamp recorded when the store was established, or null if the service is no longer in the group. + * A service that has been removed outright has no state to read, and reporting that as a missing timestamp lets the + * caller fail on the comparison rather than on a lookup error. + */ + private String readCreatedTimestamp(final MigratedFlow flow, final String serviceId) throws NiFiClientException, IOException { + final boolean stillPresent = findStoreServices(flow.groupId()).stream() + .anyMatch(service -> serviceId.equals(service.getComponent().getId())); + if (!stillPresent) { + return null; + } + + return readStoreState(serviceId).get(CREATED_STATE_KEY); + } + + private long countRows(final String serviceId) throws NiFiClientException, IOException { + final String rowCount = readStoreState(serviceId).get(ROW_COUNT_STATE_KEY); + return rowCount == null ? 0 : Long.parseLong(rowCount); + } + + /** + * Polls the given condition until it holds, failing the test once the timeout elapses. Conditions query a NiFi that + * may still be starting up or replacing components, so a failing query is treated as the condition not holding yet. + */ + private void waitForCondition(final ExceptionalBooleanSupplier condition, final String conditionDescription) throws InterruptedException { + final long deadline = System.nanoTime() + CONDITION_TIMEOUT.toNanos(); + + while (System.nanoTime() < deadline) { + try { + if (condition.getAsBoolean()) { + return; + } + } catch (final InterruptedException ie) { + throw ie; + } catch (final Exception ignored) { + } + + Thread.sleep(CONDITION_POLL); + } + + fail("Timed out after " + CONDITION_TIMEOUT + " waiting for " + conditionDescription); + } + + private record MigratedFlow(String groupId, String processorId) { + } + + private record MigratedStore(String serviceId, String created) { + } + +} diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/1/snapshot.json b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/1/snapshot.json new file mode 100644 index 000000000000..916fe7f6f999 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/1/snapshot.json @@ -0,0 +1,82 @@ +{ + "externalControllerServices": {}, + "flowContents": { + "comments": "", + "componentType": "PROCESS_GROUP", + "connections": [], + "controllerServices": [], + "defaultBackPressureDataSizeThreshold": "1 GB", + "defaultBackPressureObjectThreshold": 10000, + "defaultFlowFileExpiration": "0 sec", + "flowFileConcurrency": "UNBOUNDED", + "flowFileOutboundPolicy": "STREAM_WHEN_AVAILABLE", + "funnels": [], + "identifier": "aaaaaaaa-1111-1111-1111-111111111111", + "inputPorts": [], + "labels": [], + "name": "Store Group", + "outputPorts": [], + "position": { + "x": 0.0, + "y": 0.0 + }, + "processGroups": [], + "processors": [ + { + "autoTerminatedRelationships": [ + "success" + ], + "backoffMechanism": "PENALIZE_FLOWFILE", + "bulletinLevel": "WARN", + "bundle": { + "artifact": "nifi-system-test-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "comments": "", + "componentType": "PROCESSOR", + "concurrentlySchedulableTaskCount": 1, + "executionNode": "ALL", + "groupIdentifier": "aaaaaaaa-1111-1111-1111-111111111111", + "identifier": "bbbbbbbb-2222-2222-2222-222222222222", + "maxBackoffPeriod": "10 mins", + "name": "MigrateToControllerService", + "penaltyDuration": "30 sec", + "position": { + "x": 100.0, + "y": 100.0 + }, + "properties": { + "store-name": "store" + }, + "propertyDescriptors": { + "store-name": { + "displayName": "Store Name", + "identifiesControllerService": false, + "name": "store-name", + "sensitive": false + } + }, + "retriedRelationships": [], + "style": {}, + "retryCount": 0, + "runDurationMillis": 0, + "scheduledState": "RUNNING", + "schedulingPeriod": "100 ms", + "schedulingStrategy": "TIMER_DRIVEN", + "type": "org.apache.nifi.processors.tests.system.MigrateToControllerService", + "yieldDuration": "1 sec" + } + ], + "remoteProcessGroups": [], + "variables": {} + }, + "flowEncodingVersion": "1.0", + "parameterContexts": {}, + "snapshotMetadata": { + "author": "system-test", + "comments": "", + "timestamp": 1700000000000, + "version": 1 + } +} diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/2/snapshot.json b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/2/snapshot.json new file mode 100644 index 000000000000..c448840bc950 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/11111111-2222-3333-4444-555555555555/2/snapshot.json @@ -0,0 +1,112 @@ +{ + "externalControllerServices": {}, + "flowContents": { + "comments": "", + "componentType": "PROCESS_GROUP", + "connections": [], + "controllerServices": [ + { + "bulletinLevel": "WARN", + "bundle": { + "artifact": "nifi-alternate-config-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "comments": "", + "componentType": "CONTROLLER_SERVICE", + "controllerServiceApis": [ + { + "bundle": { + "artifact": "nifi-alternate-config-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "type": "org.apache.nifi.cs.tests.system.StoreService" + } + ], + "groupIdentifier": "aaaaaaaa-1111-1111-1111-111111111111", + "identifier": "99999999-8888-7777-6666-555555555555", + "name": "StateBackedStoreService", + "properties": { + "Store Name": "store" + }, + "propertyDescriptors": {}, + "scheduledState": "ENABLED", + "type": "org.apache.nifi.cs.tests.system.StateBackedStoreService" + } + ], + "defaultBackPressureDataSizeThreshold": "1 GB", + "defaultBackPressureObjectThreshold": 10000, + "defaultFlowFileExpiration": "0 sec", + "flowFileConcurrency": "UNBOUNDED", + "flowFileOutboundPolicy": "STREAM_WHEN_AVAILABLE", + "funnels": [], + "identifier": "aaaaaaaa-1111-1111-1111-111111111111", + "inputPorts": [], + "labels": [], + "name": "Store Group", + "outputPorts": [], + "position": { + "x": 0.0, + "y": 0.0 + }, + "processGroups": [], + "processors": [ + { + "autoTerminatedRelationships": [ + "success" + ], + "backoffMechanism": "PENALIZE_FLOWFILE", + "bulletinLevel": "WARN", + "bundle": { + "artifact": "nifi-alternate-config-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "comments": "", + "componentType": "PROCESSOR", + "concurrentlySchedulableTaskCount": 1, + "executionNode": "ALL", + "groupIdentifier": "aaaaaaaa-1111-1111-1111-111111111111", + "identifier": "bbbbbbbb-2222-2222-2222-222222222222", + "maxBackoffPeriod": "10 mins", + "name": "MigrateToControllerService", + "penaltyDuration": "30 sec", + "position": { + "x": 100.0, + "y": 100.0 + }, + "properties": { + "Store Service": "99999999-8888-7777-6666-555555555555" + }, + "propertyDescriptors": { + "Store Service": { + "displayName": "Store Service", + "identifiesControllerService": true, + "name": "Store Service", + "sensitive": false + } + }, + "retriedRelationships": [], + "style": {}, + "retryCount": 0, + "runDurationMillis": 0, + "scheduledState": "RUNNING", + "schedulingPeriod": "100 ms", + "schedulingStrategy": "TIMER_DRIVEN", + "type": "org.apache.nifi.processors.tests.system.MigrateToControllerService", + "yieldDuration": "1 sec" + } + ], + "remoteProcessGroups": [], + "variables": {} + }, + "flowEncodingVersion": "1.0", + "parameterContexts": {}, + "snapshotMetadata": { + "author": "system-test", + "comments": "", + "timestamp": 1700000000000, + "version": 2 + } +} diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/1/snapshot.json b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/1/snapshot.json new file mode 100644 index 000000000000..74017984ed11 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/1/snapshot.json @@ -0,0 +1,82 @@ +{ + "externalControllerServices": {}, + "flowContents": { + "comments": "", + "componentType": "PROCESS_GROUP", + "connections": [], + "controllerServices": [], + "defaultBackPressureDataSizeThreshold": "1 GB", + "defaultBackPressureObjectThreshold": 10000, + "defaultFlowFileExpiration": "0 sec", + "flowFileConcurrency": "UNBOUNDED", + "flowFileOutboundPolicy": "STREAM_WHEN_AVAILABLE", + "funnels": [], + "identifier": "dddddddd-1111-1111-1111-111111111111", + "inputPorts": [], + "labels": [], + "name": "Store Group", + "outputPorts": [], + "position": { + "x": 0.0, + "y": 0.0 + }, + "processGroups": [], + "processors": [ + { + "autoTerminatedRelationships": [ + "success" + ], + "backoffMechanism": "PENALIZE_FLOWFILE", + "bulletinLevel": "WARN", + "bundle": { + "artifact": "nifi-system-test-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "comments": "", + "componentType": "PROCESSOR", + "concurrentlySchedulableTaskCount": 1, + "executionNode": "ALL", + "groupIdentifier": "dddddddd-1111-1111-1111-111111111111", + "identifier": "eeeeeeee-2222-2222-2222-222222222222", + "maxBackoffPeriod": "10 mins", + "name": "MigrateToControllerService", + "penaltyDuration": "30 sec", + "position": { + "x": 100.0, + "y": 100.0 + }, + "properties": { + "store-name": "store" + }, + "propertyDescriptors": { + "store-name": { + "displayName": "Store Name", + "identifiesControllerService": false, + "name": "store-name", + "sensitive": false + } + }, + "retriedRelationships": [], + "style": {}, + "retryCount": 0, + "runDurationMillis": 0, + "scheduledState": "RUNNING", + "schedulingPeriod": "100 ms", + "schedulingStrategy": "TIMER_DRIVEN", + "type": "org.apache.nifi.processors.tests.system.MigrateToControllerService", + "yieldDuration": "1 sec" + } + ], + "remoteProcessGroups": [], + "variables": {} + }, + "flowEncodingVersion": "1.0", + "parameterContexts": {}, + "snapshotMetadata": { + "author": "system-test", + "comments": "", + "timestamp": 1700000000000, + "version": 1 + } +} diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/2/snapshot.json b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/2/snapshot.json new file mode 100644 index 000000000000..86e93efddbc2 --- /dev/null +++ b/nifi-system-tests/nifi-system-test-suite/src/test/resources/versioned-flows/test-flows/22222222-3333-4444-5555-666666666666/2/snapshot.json @@ -0,0 +1,118 @@ +{ + "externalControllerServices": {}, + "flowContents": { + "comments": "", + "componentType": "PROCESS_GROUP", + "connections": [], + "controllerServices": [], + "defaultBackPressureDataSizeThreshold": "1 GB", + "defaultBackPressureObjectThreshold": 10000, + "defaultFlowFileExpiration": "0 sec", + "flowFileConcurrency": "UNBOUNDED", + "flowFileOutboundPolicy": "STREAM_WHEN_AVAILABLE", + "funnels": [], + "identifier": "dddddddd-1111-1111-1111-111111111111", + "inputPorts": [], + "labels": [], + "name": "Store Group", + "outputPorts": [], + "position": { + "x": 0.0, + "y": 0.0 + }, + "processGroups": [], + "processors": [ + { + "autoTerminatedRelationships": [ + "success" + ], + "backoffMechanism": "PENALIZE_FLOWFILE", + "bulletinLevel": "WARN", + "bundle": { + "artifact": "nifi-system-test-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "comments": "", + "componentType": "PROCESSOR", + "concurrentlySchedulableTaskCount": 1, + "executionNode": "ALL", + "groupIdentifier": "dddddddd-1111-1111-1111-111111111111", + "identifier": "eeeeeeee-2222-2222-2222-222222222222", + "maxBackoffPeriod": "10 mins", + "name": "MigrateToControllerService", + "penaltyDuration": "30 sec", + "position": { + "x": 100.0, + "y": 100.0 + }, + "properties": { + "store-name": "store" + }, + "propertyDescriptors": { + "store-name": { + "displayName": "Store Name", + "identifiesControllerService": false, + "name": "store-name", + "sensitive": false + } + }, + "retriedRelationships": [], + "style": {}, + "retryCount": 0, + "runDurationMillis": 0, + "scheduledState": "RUNNING", + "schedulingPeriod": "100 ms", + "schedulingStrategy": "TIMER_DRIVEN", + "type": "org.apache.nifi.processors.tests.system.MigrateToControllerService", + "yieldDuration": "1 sec" + }, + { + "autoTerminatedRelationships": [ + "success" + ], + "backoffMechanism": "PENALIZE_FLOWFILE", + "bulletinLevel": "WARN", + "bundle": { + "artifact": "nifi-system-test-extensions-nar", + "group": "org.apache.nifi", + "version": "2.12.0-SNAPSHOT" + }, + "comments": "", + "componentType": "PROCESSOR", + "concurrentlySchedulableTaskCount": 1, + "executionNode": "ALL", + "groupIdentifier": "dddddddd-1111-1111-1111-111111111111", + "identifier": "ffffffff-3333-3333-3333-333333333333", + "maxBackoffPeriod": "10 mins", + "name": "Added Processor", + "penaltyDuration": "30 sec", + "position": { + "x": 400.0, + "y": 100.0 + }, + "properties": {}, + "propertyDescriptors": {}, + "retriedRelationships": [], + "style": {}, + "retryCount": 0, + "runDurationMillis": 0, + "scheduledState": "ENABLED", + "schedulingPeriod": "100 ms", + "schedulingStrategy": "TIMER_DRIVEN", + "type": "org.apache.nifi.processors.tests.system.MigrateToControllerService", + "yieldDuration": "1 sec" + } + ], + "remoteProcessGroups": [], + "variables": {} + }, + "flowEncodingVersion": "1.0", + "parameterContexts": {}, + "snapshotMetadata": { + "author": "system-test", + "comments": "", + "timestamp": 1700000000000, + "version": 2 + } +} diff --git a/nifi-system-tests/pom.xml b/nifi-system-tests/pom.xml index e6c97cba6a45..ad928867a79e 100644 --- a/nifi-system-tests/pom.xml +++ b/nifi-system-tests/pom.xml @@ -30,6 +30,7 @@ nifi-system-test-extensions-bundle nifi-system-test-extensions2-bundle nifi-system-test-flow-action-reporter-bundle + nifi-system-test-flow-registry-bundle nifi-alternate-config-extensions-bundle nifi-system-test-nar-provider-bundles nifi-python-test-extensions-nar diff --git a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/ControllerServicesClient.java b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/ControllerServicesClient.java index e6a44fbd1d87..5aaaf6141ebc 100644 --- a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/ControllerServicesClient.java +++ b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/ControllerServicesClient.java @@ -16,6 +16,7 @@ */ package org.apache.nifi.toolkit.client; +import org.apache.nifi.web.api.entity.ComponentStateEntity; import org.apache.nifi.web.api.entity.ControllerServiceEntity; import org.apache.nifi.web.api.entity.ControllerServiceReferencingComponentsEntity; import org.apache.nifi.web.api.entity.ControllerServiceRunStatusEntity; @@ -51,4 +52,6 @@ public interface ControllerServicesClient { VerifyConfigRequestEntity deleteConfigVerificationRequest(String serviceId, String verificationRequestId) throws NiFiClientException, IOException; PropertyDescriptorEntity getPropertyDescriptor(String serviceId, String propertyName, Boolean sensitive) throws NiFiClientException, IOException; + + ComponentStateEntity getControllerServiceState(String serviceId) throws NiFiClientException, IOException; } diff --git a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyControllerServicesClient.java b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyControllerServicesClient.java index ba5b6cf33090..a7ee3293e026 100644 --- a/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyControllerServicesClient.java +++ b/nifi-toolkit/nifi-toolkit-client/src/main/java/org/apache/nifi/toolkit/client/impl/JerseyControllerServicesClient.java @@ -24,6 +24,7 @@ import org.apache.nifi.toolkit.client.NiFiClientException; import org.apache.nifi.toolkit.client.RequestConfig; import org.apache.nifi.web.api.dto.RevisionDTO; +import org.apache.nifi.web.api.entity.ComponentStateEntity; import org.apache.nifi.web.api.entity.ControllerServiceEntity; import org.apache.nifi.web.api.entity.ControllerServiceReferencingComponentsEntity; import org.apache.nifi.web.api.entity.ControllerServiceRunStatusEntity; @@ -255,4 +256,17 @@ public PropertyDescriptorEntity getPropertyDescriptor(final String serviceId, fi return getRequestBuilder(target).get(PropertyDescriptorEntity.class); }); } + + @Override + public ComponentStateEntity getControllerServiceState(final String serviceId) throws NiFiClientException, IOException { + Objects.requireNonNull(serviceId, "Service ID required"); + + return executeAction("Error getting state of the Controller Service", () -> { + final WebTarget target = controllerServicesTarget + .path("{id}/state") + .resolveTemplate("id", serviceId); + + return getRequestBuilder(target).get(ComponentStateEntity.class); + }); + } } From 5f66f59dda9f0c00108fc68b1132d738befb49c3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alaksiej=20=C5=A0=C4=8Darbaty?= Date: Tue, 22 Sep 2026 16:09:05 +0200 Subject: [PATCH 2/3] NIFI-16310 Address review comments --- ...tandardVersionedComponentSynchronizer.java | 110 ++++++++-------- ...ardVersionedComponentSynchronizerTest.java | 57 +++++---- .../tests/system/StateBackedStoreService.java | 2 +- .../system/MigrateToControllerService.java | 10 +- .../system/MigrateToControllerService.java | 10 +- ...nCreatedControllerServiceVersioningIT.java | 119 ++++++------------ 6 files changed, 136 insertions(+), 172 deletions(-) diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java index 274940379700..b0c17606088f 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizer.java @@ -51,7 +51,6 @@ import org.apache.nifi.controller.reporting.ReportingTaskInstantiationException; import org.apache.nifi.controller.service.ControllerServiceNode; import org.apache.nifi.controller.service.ControllerServiceProvider; -import org.apache.nifi.controller.service.ControllerServiceReference; import org.apache.nifi.controller.service.ControllerServiceState; import org.apache.nifi.flow.BatchSize; import org.apache.nifi.flow.Bundle; @@ -1182,57 +1181,77 @@ private void removeMissingRpg(final ProcessGroup group, final VersionedProcessGr */ private void assignVersionedIdsToMigrationCreatedControllerServices(final ProcessGroup group, final VersionedProcessGroup proposed) { final Collection groupServices = group.getControllerServices(false); - if (groupServices == null || groupServices.isEmpty()) { + if (groupServices.isEmpty()) { return; } + final List migrationCreatedServices = getUnversionedMigrationCreatedServices(groupServices); + if (migrationCreatedServices.isEmpty()) { + return; + } + + final Set claimedVersionedIds = getClaimedVersionedIds(groupServices); + + final Set proposedControllerServices = Objects.requireNonNullElse(proposed.getControllerServices(), Set.of()); + final Set proposedProcessors = Objects.requireNonNullElse(proposed.getProcessors(), Set.of()); + final Map proposedComponentsByVersionedId = indexByVersionedId(proposedControllerServices, proposedProcessors); + + for (final ControllerServiceNode localService : orderByReferencerChain(migrationCreatedServices)) { + assignMatchingVersionedId(group, localService, claimedVersionedIds, proposedComponentsByVersionedId); + } + } + + private List getUnversionedMigrationCreatedServices(final Collection groupServices) { final List migrationCreatedServices = new ArrayList<>(); for (final ControllerServiceNode localService : groupServices) { if (localService.getVersionedComponentId().isEmpty() && isMigrationCreated(localService)) { migrationCreatedServices.add(localService); } } + return migrationCreatedServices; + } - if (migrationCreatedServices.isEmpty()) { - return; - } - + private Set getClaimedVersionedIds(final Collection groupServices) { final Set claimedVersionedIds = HashSet.newHashSet(groupServices.size()); for (final ControllerServiceNode localService : groupServices) { localService.getVersionedComponentId().ifPresent(claimedVersionedIds::add); } + return claimedVersionedIds; + } - final Map proposedComponentsByVersionedId = indexByVersionedId(proposed.getControllerServices(), proposed.getProcessors()); - - for (final ControllerServiceNode localService : orderByReferencerChain(migrationCreatedServices)) { - final ComponentNode referencer = getSoleReferencer(localService); - if (referencer == null) { - LOG.debug("Leaving {} in {} unversioned because it is not referenced by exactly one component", localService, group); - continue; - } - - final VersionedConfigurableExtension proposedReferencer = getProposedReferencer(referencer, proposedComponentsByVersionedId); - if (proposedReferencer == null) { - LOG.debug("Leaving {} in {} unversioned because its referencer {} has no counterpart in the proposed flow", localService, group, referencer); - continue; - } + private void assignMatchingVersionedId( + final ProcessGroup group, + final ControllerServiceNode localService, + final Set claimedVersionedIds, + final Map proposedComponentsByVersionedId + ) { + final ComponentNode referencer = getSoleReferencer(localService); + if (referencer == null) { + LOG.debug("Leaving {} in {} unversioned because it is not referenced by exactly one component", localService, group); + return; + } - final ProposedControllerServiceMatch match = findMatchingProposedControllerService(localService, referencer, proposedReferencer, proposedComponentsByVersionedId, group); - if (match == null) { - LOG.debug("Leaving {} in {} unversioned because no proposed Controller Service matches the referencing property of {}", localService, group, referencer); - continue; - } + final VersionedConfigurableExtension proposedReferencer = getProposedReferencer(referencer, proposedComponentsByVersionedId); + if (proposedReferencer == null) { + LOG.debug("Leaving {} in {} unversioned because its referencer {} has no counterpart in the proposed flow", localService, group, referencer); + return; + } - if (!claimedVersionedIds.add(match.versionedId())) { - LOG.debug("Leaving {} in {} unversioned because versioned id {} is already used by another Controller Service", localService, group, match.versionedId()); - continue; - } + final ProposedControllerServiceMatch match = findMatchingProposedControllerService(localService, referencer, proposedReferencer, proposedComponentsByVersionedId, group); + if (match == null) { + LOG.debug("Leaving {} in {} unversioned because no proposed Controller Service matches the referencing property of {}", localService, group, referencer); + return; + } - localService.setVersionedComponentId(match.versionedId()); - updatedVersionedComponentIds.add(match.versionedId()); - LOG.info("Matched {} in {} to the Controller Service with versioned id {} that the proposed flow declares, based on the {} property of {}", - localService, group, match.versionedId(), match.propertyName(), referencer); + if (!claimedVersionedIds.add(match.versionedId())) { + LOG.debug("Leaving {} in {} unversioned because versioned id {} is already used by another Controller Service", localService, group, match.versionedId()); + return; } + + localService.setVersionedComponentId(match.versionedId()); + updatedVersionedComponentIds.add(match.versionedId()); + LOG.info("Matched {} in {} to the Controller Service with versioned id {} that the proposed flow declares, based on the {} property of {}", + localService, group, match.versionedId(), match.propertyName(), referencer); } private List orderByReferencerChain(final List migrationCreatedServices) { @@ -1278,9 +1297,7 @@ private Map indexByVersionedId( final Collection controllerServices, final Collection processors ) { - final int serviceCount = controllerServices == null ? 0 : controllerServices.size(); - final int processorCount = processors == null ? 0 : processors.size(); - final Map byVersionedId = HashMap.newHashMap(serviceCount + processorCount); + final Map byVersionedId = HashMap.newHashMap(controllerServices.size() + processors.size()); addByVersionedId(byVersionedId, controllerServices); addByVersionedId(byVersionedId, processors); return byVersionedId; @@ -1290,23 +1307,14 @@ private void addByVersionedId( final Map byVersionedId, final Collection components ) { - if (components == null) { - return; - } - for (final VersionedConfigurableExtension component : components) { byVersionedId.put(component.getIdentifier(), component); } } private ComponentNode getSoleReferencer(final ControllerServiceNode localService) { - final ControllerServiceReference references = localService.getReferences(); - if (references == null) { - return null; - } - - final Set referencers = references.getReferencingComponents(); - if (referencers == null || referencers.size() != 1) { + final Set referencers = localService.getReferences().getReferencingComponents(); + if (referencers.size() != 1) { return null; } @@ -1333,13 +1341,7 @@ private ProposedControllerServiceMatch findMatchingProposedControllerService( final Map proposedComponentsByVersionedId, final ProcessGroup group ) { - final Map rawPropertyValues = referencer.getRawPropertyValues(); - if (rawPropertyValues == null) { - LOG.debug("Leaving {} in {} unversioned because referencer {} has no property values", localService, group, referencer); - return null; - } - - for (final Map.Entry propertyEntry : rawPropertyValues.entrySet()) { + for (final Map.Entry propertyEntry : referencer.getRawPropertyValues().entrySet()) { final PropertyDescriptor descriptor = propertyEntry.getKey(); final String propertyName = descriptor.getName(); if (descriptor.getControllerServiceDefinition() == null) { diff --git a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java index 173dda001c6f..4913d07d35a7 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java +++ b/nifi-framework-bundle/nifi-framework/nifi-framework-components/src/test/java/org/apache/nifi/flow/synchronization/StandardVersionedComponentSynchronizerTest.java @@ -569,13 +569,13 @@ public void testAddProcessorWithServiceAndMigration() { public void testUserAddedControllerServiceRemovedWhenAbsentFromProposedFlow() { final ProcessGroup processGroup = createMockProcessGroup(); - final PropertyDescriptor descriptor = new PropertyDescriptor.Builder().name("abc").build(); + final PropertyDescriptor descriptor = new PropertyDescriptor.Builder().name(PARAM_ABC).build(); final ControllerServiceNode serviceNode = createMockControllerService(); when(serviceNode.getComments()).thenReturn("Added by a user"); when(serviceNode.getName()).thenReturn("name"); when(serviceNode.getCanonicalClassName()).thenReturn("ControllerServiceImpl"); - when(serviceNode.getProperties()).thenReturn(Map.of(descriptor, new PropertyConfiguration("123", null, null, null))); - when(serviceNode.getRawPropertyValues()).thenReturn(Map.of(descriptor, "123")); + when(serviceNode.getProperties()).thenReturn(Map.of(descriptor, new PropertyConfiguration(VALUE_123, null, null, null))); + when(serviceNode.getRawPropertyValues()).thenReturn(Map.of(descriptor, VALUE_123)); when(serviceNode.getVersionedComponentId()).thenReturn(Optional.empty()); when(processGroup.getControllerServices(false)).thenReturn(Set.of(serviceNode)); @@ -595,19 +595,26 @@ public void testUserAddedControllerServiceRemovedWhenAbsentFromProposedFlow() { @Nested class MigrationCreatedControllerService { + private static final String STORE_SERVICE_PROPERTY = "Store Service"; + private static final String STORE_NAME_PROPERTY = "Store Name"; + private static final String STORE_NAME = "store"; + private static final String PROCESSOR_VERSIONED_ID = "processor-versioned-id"; + private static final String DECLARED_SERVICE_VERSIONED_ID = "declared-service-versioned-id"; + private static final String INNER_SERVICE_VERSIONED_ID = "inner-service-versioned-id"; + @Test - public void doesNotRemoveWhenProposedFlowDeclaresUnrelatedService() { + void doesNotRemoveWhenProposedFlowDeclaresUnrelatedService() { final ProcessGroup processGroup = createMockProcessGroup(); - final PropertyDescriptor descriptor = new PropertyDescriptor.Builder().name("abc").build(); + final PropertyDescriptor descriptor = new PropertyDescriptor.Builder().name(STORE_NAME_PROPERTY).build(); final ControllerServiceNode localOnlyService = createMockControllerService(); when(localOnlyService.getComments()).thenReturn(StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT); when(localOnlyService.getName()).thenReturn("ServiceName"); when(localOnlyService.getCanonicalClassName()).thenReturn("ServiceType"); when(localOnlyService.getBundleCoordinate()).thenReturn(bundleCoordinate); - when(localOnlyService.getProperties()).thenReturn(Map.of(descriptor, new PropertyConfiguration("123", null, null, null))); - when(localOnlyService.getRawPropertyValues()).thenReturn(Map.of(descriptor, "123")); + when(localOnlyService.getProperties()).thenReturn(Map.of(descriptor, new PropertyConfiguration(STORE_NAME, null, null, null))); + when(localOnlyService.getRawPropertyValues()).thenReturn(Map.of(descriptor, STORE_NAME)); when(localOnlyService.getState()).thenReturn(ControllerServiceState.DISABLED); trackVersionedComponentId(localOnlyService); @@ -621,7 +628,7 @@ public void doesNotRemoveWhenProposedFlowDeclaresUnrelatedService() { } @Test - public void assignsProposedVersionedIdWithoutCreatingDuplicate() { + void assignsProposedVersionedIdWithoutCreatingDuplicate() { final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); synchronizeMigrationMatch(setup); @@ -634,7 +641,7 @@ public void assignsProposedVersionedIdWithoutCreatingDuplicate() { } @Test - public void doesNotAssignWhenReferencedByMultipleComponents() { + void doesNotAssignWhenReferencedByMultipleComponents() { final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); final ProcessorNode additionalReferencer = createMappableProcessor(setup.processGroup()); @@ -648,7 +655,7 @@ public void doesNotAssignWhenReferencedByMultipleComponents() { } @Test - public void doesNotAssignWhenProposedServiceTypeDiffers() { + void doesNotAssignWhenProposedServiceTypeDiffers() { final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); when(setup.localService().getCanonicalClassName()).thenReturn("org.apache.nifi.cs.DifferentService"); @@ -660,11 +667,11 @@ public void doesNotAssignWhenProposedServiceTypeDiffers() { } @Test - public void assignsNestedServicesWhenInnerListedBeforeOuter() { + void assignsNestedServicesWhenInnerListedBeforeOuter() { final MigrationCreatedMatchSetup outer = newMigrationCreatedMatchSetup(); final VersionedControllerService proposedInnerService = createMinimalVersionedControllerService(); - proposedInnerService.setIdentifier("inner-service-versioned-id"); + proposedInnerService.setIdentifier(INNER_SERVICE_VERSIONED_ID); final ControllerServiceNode innerService = createMockControllerService(); stubMigrationCreatedService(innerService, proposedInnerService); @@ -672,8 +679,8 @@ public void assignsNestedServicesWhenInnerListedBeforeOuter() { setReferences(innerService, outer.localService()); when(outer.processGroup().findControllerService(eq(innerService.getIdentifier()), anyBoolean(), anyBoolean())).thenReturn(innerService); - outer.proposedService().setProperties(Map.of("Store Service", proposedInnerService.getIdentifier())); - outer.proposedService().setPropertyDescriptors(Map.of("Store Service", storeServiceVersionedDescriptor())); + outer.proposedService().setProperties(Map.of(STORE_SERVICE_PROPERTY, proposedInnerService.getIdentifier())); + outer.proposedService().setPropertyDescriptors(Map.of(STORE_SERVICE_PROPERTY, storeServiceVersionedDescriptor())); final Set localServices = new LinkedHashSet<>(); localServices.add(innerService); @@ -689,7 +696,7 @@ public void assignsNestedServicesWhenInnerListedBeforeOuter() { } @Test - public void doesNotAssignWhenProposedVersionedIdAlreadyUsed() { + void doesNotAssignWhenProposedVersionedIdAlreadyUsed() { final MigrationCreatedMatchSetup setup = newMigrationCreatedMatchSetup(); final ControllerServiceNode existingService = createMockControllerService(); @@ -708,7 +715,7 @@ private MigrationCreatedMatchSetup newMigrationCreatedMatchSetup() { final ProcessGroup processGroup = createMockProcessGroup(); final VersionedControllerService proposedService = createMinimalVersionedControllerService(); - proposedService.setIdentifier("declared-service-versioned-id"); + proposedService.setIdentifier(DECLARED_SERVICE_VERSIONED_ID); final PropertyDescriptor storeServiceDescriptor = storeServicePropertyDescriptor(); @@ -717,7 +724,7 @@ private MigrationCreatedMatchSetup newMigrationCreatedMatchSetup() { stubMigrationCreatedService(localService, proposedService); final ProcessorNode processor = createMappableProcessor(processGroup); - when(processor.getVersionedComponentId()).thenReturn(Optional.of("processor-versioned-id")); + when(processor.getVersionedComponentId()).thenReturn(Optional.of(PROCESSOR_VERSIONED_ID)); stubControllerServiceReferenceProperty(processor, storeServiceDescriptor, localServiceId); setReferences(localService, processor); @@ -726,9 +733,9 @@ private MigrationCreatedMatchSetup newMigrationCreatedMatchSetup() { when(processGroup.findControllerService(eq(localServiceId), anyBoolean(), anyBoolean())).thenReturn(localService); final VersionedProcessor proposedProcessor = createMinimalVersionedProcessor(); - proposedProcessor.setIdentifier("processor-versioned-id"); - proposedProcessor.setProperties(Map.of("Store Service", proposedService.getIdentifier())); - proposedProcessor.setPropertyDescriptors(Map.of("Store Service", storeServiceVersionedDescriptor())); + proposedProcessor.setIdentifier(PROCESSOR_VERSIONED_ID); + proposedProcessor.setProperties(Map.of(STORE_SERVICE_PROPERTY, proposedService.getIdentifier())); + proposedProcessor.setPropertyDescriptors(Map.of(STORE_SERVICE_PROPERTY, storeServiceVersionedDescriptor())); return new MigrationCreatedMatchSetup(processGroup, localService, processor, proposedService, proposedProcessor); } @@ -761,14 +768,14 @@ private record MigrationCreatedMatchSetup( private PropertyDescriptor storeServicePropertyDescriptor() { return new PropertyDescriptor.Builder() - .name("Store Service") + .name(STORE_SERVICE_PROPERTY) .identifiesControllerService(ControllerService.class) .build(); } private VersionedPropertyDescriptor storeServiceVersionedDescriptor() { final VersionedPropertyDescriptor proposedDescriptor = new VersionedPropertyDescriptor(); - proposedDescriptor.setName("Store Service"); + proposedDescriptor.setName(STORE_SERVICE_PROPERTY); proposedDescriptor.setIdentifiesControllerService(true); return proposedDescriptor; } @@ -776,9 +783,9 @@ private VersionedPropertyDescriptor storeServiceVersionedDescriptor() { private void stubControllerServiceListedInFlow(final ControllerServiceNode localService, final VersionedControllerService proposed) { when(localService.getCanonicalClassName()).thenReturn(proposed.getType()); when(localService.getName()).thenReturn(proposed.getName()); - when(localService.getProperties()).thenReturn(Map.of(new PropertyDescriptor.Builder().name("abc").build(), - new PropertyConfiguration("123", null, null, null))); - when(localService.getRawPropertyValues()).thenReturn(Map.of(new PropertyDescriptor.Builder().name("abc").build(), "123")); + when(localService.getProperties()).thenReturn(Map.of(new PropertyDescriptor.Builder().name(STORE_NAME_PROPERTY).build(), + new PropertyConfiguration(STORE_NAME, null, null, null))); + when(localService.getRawPropertyValues()).thenReturn(Map.of(new PropertyDescriptor.Builder().name(STORE_NAME_PROPERTY).build(), STORE_NAME)); } private void stubMigrationCreatedService(final ControllerServiceNode localService, final VersionedControllerService proposed) { diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java index e0ce89410f1b..899b2f5ebdeb 100644 --- a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/cs/tests/system/StateBackedStoreService.java @@ -81,7 +81,7 @@ public void onEnabled(final ConfigurationContext context) throws IOException { } @Override - public synchronized void append(final String row) { + public void append(final String row) { try { final StateManager stateManager = getStateManager(); final Map state = new HashMap<>(stateManager.getState(Scope.LOCAL).toMap()); diff --git a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java index 64a6582a6c6c..608b0f4582e4 100644 --- a/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java +++ b/nifi-system-tests/nifi-alternate-config-extensions-bundle/nifi-alternate-config-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java @@ -19,6 +19,7 @@ import org.apache.nifi.annotation.behavior.InputRequirement; import org.apache.nifi.annotation.behavior.InputRequirement.Requirement; +import org.apache.nifi.annotation.documentation.CapabilityDescription; import org.apache.nifi.components.PropertyDescriptor; import org.apache.nifi.cs.tests.system.StateBackedStoreService; import org.apache.nifi.cs.tests.system.StoreService; @@ -34,11 +35,10 @@ import java.util.Set; import java.util.concurrent.atomic.AtomicLong; -/** - * Post-upgrade shape of a processor whose property migration creates the store Controller Service - * that the pre-upgrade shape did not have. Each execution appends a row to the store so that tests - * can observe whether store contents survive flow and runtime upgrades. - */ +@CapabilityDescription(""" + Post-upgrade shape of a processor whose property migration creates the store Controller Service that the pre-upgrade shape did not have. + Each execution appends a row to the store so that tests can observe whether store contents survive flow and runtime upgrades. + """) @InputRequirement(Requirement.INPUT_FORBIDDEN) public class MigrateToControllerService extends AbstractProcessor { diff --git a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java index 9b0ce040aa28..05eb343c665b 100644 --- a/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java +++ b/nifi-system-tests/nifi-system-test-extensions-bundle/nifi-system-test-extensions/src/main/java/org/apache/nifi/processors/tests/system/MigrateToControllerService.java @@ -19,6 +19,7 @@ import org.apache.nifi.annotation.behavior.InputRequirement; import org.apache.nifi.annotation.behavior.InputRequirement.Requirement; +import org.apache.nifi.annotation.documentation.CapabilityDescription; import org.apache.nifi.components.PropertyDescriptor; import org.apache.nifi.processor.AbstractProcessor; import org.apache.nifi.processor.ProcessContext; @@ -30,11 +31,10 @@ import java.util.List; import java.util.Set; -/** - * Pre-upgrade shape of a processor that keeps its store location in a plain property. - * The post-upgrade shape of the same processor, in the alternate-config extensions bundle, - * migrates that property into a Controller Service. - */ +@CapabilityDescription(""" + Pre-upgrade shape of a processor that keeps its store location in a plain property. + The post-upgrade shape of the same processor, in the alternate-config extensions bundle, migrates that property into a Controller Service. + """) @InputRequirement(Requirement.INPUT_FORBIDDEN) public class MigrateToControllerService extends AbstractProcessor { diff --git a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java index 0d1521903f09..ae340eb0bfab 100644 --- a/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java +++ b/nifi-system-tests/nifi-system-test-suite/src/test/java/org/apache/nifi/tests/system/migration/MigrationCreatedControllerServiceVersioningIT.java @@ -19,9 +19,9 @@ import org.apache.nifi.migration.StandardControllerServiceFactory; import org.apache.nifi.tests.system.AbstractNarSwapMigrationIT; -import org.apache.nifi.tests.system.ExceptionalBooleanSupplier; import org.apache.nifi.toolkit.client.NiFiClientException; import org.apache.nifi.web.api.dto.ComponentStateDTO; +import org.apache.nifi.web.api.dto.ProcessorDTO; import org.apache.nifi.web.api.dto.StateEntryDTO; import org.apache.nifi.web.api.entity.ComponentStateEntity; import org.apache.nifi.web.api.entity.ControllerServiceEntity; @@ -34,8 +34,6 @@ import java.io.File; import java.io.IOException; import java.nio.charset.StandardCharsets; -import java.time.Duration; -import java.util.Collection; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -43,10 +41,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; -import static org.junit.jupiter.api.Assertions.fail; /** * Verifies that a Controller Service created by property migration survives the operations a deployed flow goes @@ -58,7 +54,7 @@ * declares the service, so the flow relies entirely on property migration to create it, and the later version only * adds an unrelated processor. */ -public class MigrationCreatedControllerServiceVersioningIT extends AbstractNarSwapMigrationIT { +class MigrationCreatedControllerServiceVersioningIT extends AbstractNarSwapMigrationIT { private static final String TEST_FLOWS_BUCKET = "test-flows"; private static final String SERVICE_DECLARED_FLOW_ID = "11111111-2222-3333-4444-555555555555"; private static final String SERVICE_ABSENT_FLOW_ID = "22222222-3333-4444-5555-666666666666"; @@ -71,15 +67,14 @@ public class MigrationCreatedControllerServiceVersioningIT extends AbstractNarSw private static final String VERSIONED_FLOWS_DIRECTORY = "src/test/resources/versioned-flows"; private static final String CREATED_STATE_KEY = "created"; private static final String ROW_COUNT_STATE_KEY = "rowCount"; - private static final Duration CONDITION_TIMEOUT = Duration.ofSeconds(30); - private static final Duration CONDITION_POLL = Duration.ofMillis(100); + private static final long CONDITION_POLL_MILLIS = 100L; /** * After the runtime is upgraded, the Controller Service that property migration creates must be present, * enabled and referenced, and the flow that was running before the upgrade must be running again, with no manual action. */ @Test - public void testRuntimeUpgradeCreatesEnabledServiceAndKeepsFlowRunning() throws NiFiClientException, IOException, InterruptedException { + void testRuntimeUpgradeCreatesEnabledServiceAndKeepsFlowRunning() throws NiFiClientException, IOException, InterruptedException { final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_DECLARED_FLOW_ID); final ControllerServiceEntity service = waitForSingleStoreService(flow.groupId()); final String serviceId = service.getComponent().getId(); @@ -88,17 +83,16 @@ public void testRuntimeUpgradeCreatesEnabledServiceAndKeepsFlowRunning() throws assertBelongsToLocalFlowOnly(service); getClientUtil().waitForControllerServiceRunStatus(serviceId, "ENABLED"); getClientUtil().waitForRunningProcessor(flow.processorId()); - assertEquals(serviceId, getStoreServiceId(flow.processorId())); - final Collection validationErrors = getNifiClient().getProcessorClient().getProcessor(flow.processorId()).getComponent().getValidationErrors(); - final boolean processorValid = validationErrors == null || validationErrors.isEmpty(); - assertTrue(processorValid, "Processor must be valid after the runtime upgrade"); + final String referencedServiceId = getStoreServiceId(flow.processorId()); + assertEquals(serviceId, referencedServiceId); - waitForCondition(() -> countRows(serviceId) > 0, "store row count > 0"); + final String validationStatus = getNifiClient().getProcessorClient().getProcessor(flow.processorId()).getComponent().getValidationStatus(); + assertEquals(ProcessorDTO.VALID, validationStatus); - final String versionedFlowState = getClientUtil().getVersionedFlowState(flow.groupId(), "root"); - assertNotEquals("LOCALLY_MODIFIED", versionedFlowState, "The migration-created Controller Service must not make the flow dirty"); - assertNotEquals("LOCALLY_MODIFIED_AND_STALE", versionedFlowState, "The migration-created Controller Service must not make the flow dirty"); + waitFor(() -> countRows(serviceId) > 0, CONDITION_POLL_MILLIS, "store row count > 0"); + + getClientUtil().assertFlowStaleAndUnmodified(flow.groupId()); final boolean serviceReportedAsLocalModification = getNifiClient().getProcessGroupClient().getLocalModifications(flow.groupId()) .getComponentDifferences().stream() @@ -111,7 +105,7 @@ public void testRuntimeUpgradeCreatesEnabledServiceAndKeepsFlowRunning() throws * property migration already created, rather than removing it and substituting the one the published version declares. */ @Test - public void testFlowUpgradePreservesMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { + void testFlowUpgradePreservesMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_DECLARED_FLOW_ID); final MigratedStore store = awaitPopulatedStoreService(flow); assertBelongsToLocalFlowOnly(waitForSingleStoreService(flow.groupId())); @@ -122,12 +116,13 @@ public void testFlowUpgradePreservesMigrationCreatedControllerService() throws N getClientUtil().assertFlowUpToDate(flow.groupId()); final ControllerServiceEntity serviceAfterUpgrade = waitForSingleStoreService(flow.groupId()); - assertEquals(DECLARED_SERVICE_VERSIONED_ID, serviceAfterUpgrade.getComponent().getVersionedComponentId(), - "The service must track the Controller Service that version 2 declares, instead of belonging to the local flow only"); - assertEquals("", serviceAfterUpgrade.getComponent().getComments(), - "Comments must come from version 2 once the service tracks the declared Controller Service"); - assertNull(upgradeRequest.getRequest().getFailureReason(), - "The migration-created Controller Service must not prevent the versioned flow from upgrading"); + final String versionedComponentId = serviceAfterUpgrade.getComponent().getVersionedComponentId(); + assertEquals(DECLARED_SERVICE_VERSIONED_ID, versionedComponentId); + + final String comments = serviceAfterUpgrade.getComponent().getComments(); + assertTrue(comments.isEmpty()); + + assertNull(upgradeRequest.getRequest().getFailureReason()); } /** @@ -136,13 +131,14 @@ public void testFlowUpgradePreservesMigrationCreatedControllerService() throws N * processor untouched, must not disturb the service: the version change has nothing to say about it. */ @Test - public void testFlowUpgradeAddingUnrelatedProcessorPreservesMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { + void testFlowUpgradeAddingUnrelatedProcessorPreservesMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_ABSENT_FLOW_ID); final MigratedStore store = awaitPopulatedStoreService(flow); final VersionedFlowUpdateRequestEntity upgradeRequest = getClientUtil().changeFlowVersion(flow.groupId(), "2", false); - assertTrue(hasProcessorNamed(flow.groupId(), ADDED_PROCESSOR_NAME), "The upgrade must have added the unrelated processor"); + final boolean addedProcessorPresent = hasProcessorNamed(flow.groupId(), ADDED_PROCESSOR_NAME); + assertTrue(addedProcessorPresent, "unrelated processor not added by upgrade"); assertStorePreserved(flow, store, "flow upgrade"); getClientUtil().assertFlowUpToDate(flow.groupId()); @@ -150,8 +146,7 @@ public void testFlowUpgradeAddingUnrelatedProcessorPreservesMigrationCreatedCont final ControllerServiceEntity serviceAfterUpgrade = waitForSingleStoreService(flow.groupId()); assertBelongsToLocalFlowOnly(serviceAfterUpgrade); assertEquals(StandardControllerServiceFactory.MIGRATION_CREATED_COMMENT, serviceAfterUpgrade.getComponent().getComments()); - assertNull(upgradeRequest.getRequest().getFailureReason(), - "The migration-created Controller Service must not prevent the versioned flow from upgrading"); + assertNull(upgradeRequest.getRequest().getFailureReason()); } /** @@ -159,7 +154,7 @@ public void testFlowUpgradeAddingUnrelatedProcessorPreservesMigrationCreatedCont * created previously rather than creating a second one. */ @Test - public void testRuntimeRestartDoesNotRecreateMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { + void testRuntimeRestartDoesNotRecreateMigrationCreatedControllerService() throws NiFiClientException, IOException, InterruptedException { final MigratedFlow flow = importAndUpgradeRuntime(SERVICE_DECLARED_FLOW_ID); final MigratedStore store = awaitPopulatedStoreService(flow); @@ -170,12 +165,6 @@ public void testRuntimeRestartDoesNotRecreateMigrationCreatedControllerService() assertBelongsToLocalFlowOnly(waitForSingleStoreService(flow.groupId())); } - /** - * Asserts that the given Controller Service has no counterpart in the flow definition. Such a service either - * carries no versioned component id at all, or carries the placeholder that the framework derives from its own - * instance id while mapping the group. A service that tracks a Controller Service declared by the flow definition - * carries the identifier from the definition instead. - */ private void assertBelongsToLocalFlowOnly(final ControllerServiceEntity service) { final String versionedComponentId = service.getComponent().getVersionedComponentId(); if (versionedComponentId == null) { @@ -183,8 +172,7 @@ private void assertBelongsToLocalFlowOnly(final ControllerServiceEntity service) } final String derivedFromInstanceId = UUID.nameUUIDFromBytes(service.getComponent().getId().getBytes(StandardCharsets.UTF_8)).toString(); - assertEquals(derivedFromInstanceId, versionedComponentId, - "The service tracks a Controller Service declared by the flow definition, so it no longer belongs to the local flow only"); + assertEquals(derivedFromInstanceId, versionedComponentId, "versioned component id not derived from instance id"); } /** @@ -213,44 +201,34 @@ private MigratedStore awaitPopulatedStoreService(final MigratedFlow flow) throws final ControllerServiceEntity service = waitForSingleStoreService(flow.groupId()); final String serviceId = service.getComponent().getId(); getClientUtil().waitForControllerServiceRunStatus(serviceId, "ENABLED"); - waitForCondition(() -> countRows(serviceId) > 0, "store row count > 0"); + waitFor(() -> countRows(serviceId) > 0, CONDITION_POLL_MILLIS, "store row count > 0"); return new MigratedStore(serviceId, readCreatedTimestamp(flow, serviceId)); } - /** - * Asserts that the given operation left the migration-created Controller Service in place and the flow working. - * - * Two independent properties are checked. An unchanged creation timestamp proves the service was never torn down, - * since removing a Controller Service clears its state and a replacement records a new timestamp. A row count that - * keeps climbing afterwards proves the flow is doing work again, which a RUNNING state alone does not show: a - * processor whose Controller Service reference is broken still reports RUNNING while writing nothing. - * - * The store is read through the service the processor references now rather than the one recorded earlier, so - * that a service substituted by the operation is detected rather than silently passing. - */ private void assertStorePreserved(final MigratedFlow flow, final MigratedStore store, final String operation) throws NiFiClientException, IOException, InterruptedException { - waitForCondition(() -> !findStoreServices(flow.groupId()).isEmpty(), "store Controller Service found"); + waitFor(() -> !findStoreServices(flow.groupId()).isEmpty(), CONDITION_POLL_MILLIS, "store Controller Service found"); final List servicesAfter = findControllerServicesInGroup(flow.groupId()); - assertEquals(1, servicesAfter.size(), - "The " + operation + " left " + servicesAfter.size() + " Controller Services in the group, so one was created alongside the migration-created one"); + assertEquals(1, servicesAfter.size(), operation); final ControllerServiceEntity serviceAfter = servicesAfter.getFirst(); - assertEquals(store.serviceId(), serviceAfter.getComponent().getId(), - "The " + operation + " replaced the migration-created Controller Service with a different one"); - assertEquals(store.serviceId(), getStoreServiceId(flow.processorId())); + final String remainingServiceId = serviceAfter.getComponent().getId(); + assertEquals(store.serviceId(), remainingServiceId, operation + " replaced the Controller Service"); + + final String referencedServiceId = getStoreServiceId(flow.processorId()); + assertEquals(store.serviceId(), referencedServiceId); getClientUtil().waitForControllerServiceRunStatus(store.serviceId(), "ENABLED"); getClientUtil().waitForRunningProcessor(flow.processorId()); - assertEquals(store.created(), readCreatedTimestamp(flow, store.serviceId()), - "The migration-created Controller Service was removed during the " + operation + ": its store was created again, so the state it held is gone"); + final String createdTimestamp = readCreatedTimestamp(flow, store.serviceId()); + assertEquals(store.created(), createdTimestamp, "store was recreated during " + operation); final long rowsAfterOperation = countRows(getStoreServiceId(flow.processorId())); - waitForCondition(() -> countRows(getStoreServiceId(flow.processorId())) > rowsAfterOperation, "store row count increase"); + waitFor(() -> countRows(getStoreServiceId(flow.processorId())) > rowsAfterOperation, CONDITION_POLL_MILLIS, "store row count increase"); } @@ -285,7 +263,7 @@ private List findStoreServices(final String groupId) th } private ControllerServiceEntity waitForSingleStoreService(final String groupId) throws NiFiClientException, IOException, InterruptedException { - waitForCondition(() -> findStoreServices(groupId).size() == 1, "exactly one store Controller Service"); + waitFor(() -> findStoreServices(groupId).size() == 1, CONDITION_POLL_MILLIS, "exactly one store Controller Service"); return findStoreServices(groupId).getFirst(); } @@ -333,29 +311,6 @@ private long countRows(final String serviceId) throws NiFiClientException, IOExc return rowCount == null ? 0 : Long.parseLong(rowCount); } - /** - * Polls the given condition until it holds, failing the test once the timeout elapses. Conditions query a NiFi that - * may still be starting up or replacing components, so a failing query is treated as the condition not holding yet. - */ - private void waitForCondition(final ExceptionalBooleanSupplier condition, final String conditionDescription) throws InterruptedException { - final long deadline = System.nanoTime() + CONDITION_TIMEOUT.toNanos(); - - while (System.nanoTime() < deadline) { - try { - if (condition.getAsBoolean()) { - return; - } - } catch (final InterruptedException ie) { - throw ie; - } catch (final Exception ignored) { - } - - Thread.sleep(CONDITION_POLL); - } - - fail("Timed out after " + CONDITION_TIMEOUT + " waiting for " + conditionDescription); - } - private record MigratedFlow(String groupId, String processorId) { } From f91878775786f24b9736e3d5567e2f5d80eba5e1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alaksiej=20=C5=A0=C4=8Darbaty?= Date: Wed, 23 Sep 2026 13:21:14 +0200 Subject: [PATCH 3/3] NIFI-16310 Align flow registry test module versions with 2.13.0-SNAPSHOT --- .../nifi-system-test-flow-registry-nar/pom.xml | 6 +++--- .../nifi-system-test-flow-registry/pom.xml | 6 +++--- .../nifi-system-test-flow-registry-bundle/pom.xml | 2 +- nifi-system-tests/nifi-system-test-suite/pom.xml | 2 +- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml index e3e6c968afe7..0dc1d425d05a 100644 --- a/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml +++ b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry-nar/pom.xml @@ -19,7 +19,7 @@ org.apache.nifi nifi-system-test-flow-registry-bundle - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT nifi-system-test-flow-registry-nar @@ -34,13 +34,13 @@ org.apache.nifi nifi-framework-api - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT provided org.apache.nifi nifi-system-test-flow-registry - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT diff --git a/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml index 8e212da4cfd1..9602517e2710 100644 --- a/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml +++ b/nifi-system-tests/nifi-system-test-flow-registry-bundle/nifi-system-test-flow-registry/pom.xml @@ -17,7 +17,7 @@ nifi-system-test-flow-registry-bundle org.apache.nifi - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT 4.0.0 @@ -27,12 +27,12 @@ org.apache.nifi nifi-framework-api - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT org.apache.nifi nifi-utils - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT com.fasterxml.jackson.core diff --git a/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml b/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml index ae125c11e997..3de1f10d06c7 100644 --- a/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml +++ b/nifi-system-tests/nifi-system-test-flow-registry-bundle/pom.xml @@ -17,7 +17,7 @@ nifi-system-tests org.apache.nifi - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT 4.0.0 diff --git a/nifi-system-tests/nifi-system-test-suite/pom.xml b/nifi-system-tests/nifi-system-test-suite/pom.xml index 699d3deacfc5..e9a935adc169 100644 --- a/nifi-system-tests/nifi-system-test-suite/pom.xml +++ b/nifi-system-tests/nifi-system-test-suite/pom.xml @@ -382,7 +382,7 @@ org.apache.nifi nifi-system-test-flow-registry-nar - 2.12.0-SNAPSHOT + 2.13.0-SNAPSHOT nar