NIFI-11946 Support preserving created timestamp for Registry Flows

This closes #7604

Signed-off-by: David Handermann <exceptionfactory@apache.org>
This commit is contained in:
ravisingh 2023-08-13 19:11:16 -07:00 committed by exceptionfactory
parent af365414e9
commit 6bb2d62ffd
No known key found for this signature in database
GPG Key ID: 29B6A52D2AAE8DBA
2 changed files with 76 additions and 13 deletions

View File

@ -218,7 +218,7 @@ public class RegistryService {
final List<BucketEntity> bucketsWithSameName = metadataService.getBucketsByName(bucket.getName());
if (bucketsWithSameName != null) {
for (final BucketEntity bucketWithSameName : bucketsWithSameName) {
if (!bucketWithSameName.getId().equals(existingBucketById.getId())){
if (!bucketWithSameName.getId().equals(existingBucketById.getId())) {
throw new IllegalStateException("A bucket with the same name already exists - " + bucket.getName());
}
}
@ -339,11 +339,11 @@ public class RegistryService {
if (versionedFlow.getBucketIdentifier() == null) {
versionedFlow.setBucketIdentifier(bucketIdentifier);
}
final long timestamp = System.currentTimeMillis();
if (versionedFlow.getCreatedTimestamp() <= 0) {
versionedFlow.setCreatedTimestamp(timestamp);
}
versionedFlow.setModifiedTimestamp(timestamp);
validate(versionedFlow, "Cannot create versioned flow");
// ensure the bucket exists
@ -480,7 +480,7 @@ public class RegistryService {
final List<FlowEntity> flowsWithSameName = metadataService.getFlowsByName(existingBucket.getId(), versionedFlow.getName());
if (flowsWithSameName != null) {
for (final FlowEntity flowWithSameName : flowsWithSameName) {
if(!flowWithSameName.getId().equals(existingFlow.getId())) {
if (!flowWithSameName.getId().equals(existingFlow.getId())) {
throw new IllegalStateException("A versioned flow with the same name already exists in the selected bucket");
}
}
@ -949,6 +949,7 @@ public class RegistryService {
/**
* Group the differences in the comparison by component
*
* @param flowDifferences The differences to group together by component
* @return A set of componentDifferenceGroups where each entry contains a set of differences specific to that group
*/
@ -958,9 +959,9 @@ public class RegistryService {
ComponentDifferenceGroup group;
// A component may only exist on only one version for new/removed components
VersionedComponent component = ObjectUtils.firstNonNull(diff.getComponentA(), diff.getComponentB());
if(differenceGroups.containsKey(component.getIdentifier())){
if (differenceGroups.containsKey(component.getIdentifier())) {
group = differenceGroups.get(component.getIdentifier());
}else{
} else {
group = FlowMappings.map(component);
differenceGroups.put(component.getIdentifier(), group);
}

View File

@ -57,6 +57,7 @@ import java.util.Set;
import java.util.SortedSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@ -376,6 +377,67 @@ public class TestRegistryService {
assertEquals(versionedFlow.getDescription(), createdFlow.getDescription());
}
@Test
public void testCreateFlowWithCreatedTimestamp() {
final BucketEntity existingBucket = new BucketEntity();
existingBucket.setId("b1");
existingBucket.setName("My Bucket");
existingBucket.setDescription("This is my bucket");
existingBucket.setCreated(new Date());
when(metadataService.getBucketById(existingBucket.getId())).thenReturn(existingBucket);
final long timestamp = System.currentTimeMillis()-1000; // 1 millisecond previous
final VersionedFlow versionedFlow = new VersionedFlow();
versionedFlow.setIdentifier("f1");
versionedFlow.setName("My Flow");
versionedFlow.setBucketIdentifier("b1");
versionedFlow.setCreatedTimestamp(timestamp);
versionedFlow.setModifiedTimestamp(timestamp);
doAnswer(createFlowAnswer()).when(metadataService).createFlow(any(FlowEntity.class));
final VersionedFlow createdFlow = registryService.createFlow(versionedFlow.getBucketIdentifier(), versionedFlow);
assertNotNull(createdFlow);
assertNotNull(createdFlow.getIdentifier());
assertTrue(createdFlow.getCreatedTimestamp() > 0);
assertTrue(createdFlow.getModifiedTimestamp() > 0);
assertEquals(timestamp,createdFlow.getCreatedTimestamp());
assertNotEquals(timestamp,createdFlow.getModifiedTimestamp());
assertEquals(versionedFlow.getIdentifier(), createdFlow.getIdentifier());
assertEquals(versionedFlow.getName(), createdFlow.getName());
assertEquals(versionedFlow.getBucketIdentifier(), createdFlow.getBucketIdentifier());
assertEquals(versionedFlow.getDescription(), createdFlow.getDescription());
}
@Test
public void testCreateFlowWithPreserveSourcePropertiesTrueAndInvalidTimestamps() {
final BucketEntity existingBucket = new BucketEntity();
existingBucket.setId("b1");
existingBucket.setName("My Bucket");
existingBucket.setDescription("This is my bucket");
existingBucket.setCreated(new Date());
when(metadataService.getBucketById(existingBucket.getId())).thenReturn(existingBucket);
final long timestamp = -1;
final VersionedFlow versionedFlow = new VersionedFlow();
versionedFlow.setIdentifier("f1");
versionedFlow.setName("My Flow");
versionedFlow.setBucketIdentifier("b1");
versionedFlow.setCreatedTimestamp(timestamp);
versionedFlow.setModifiedTimestamp(timestamp);
doAnswer(createFlowAnswer()).when(metadataService).createFlow(any(FlowEntity.class));
final VersionedFlow createdFlow = registryService.createFlow(versionedFlow.getBucketIdentifier(), versionedFlow);
assertNotNull(createdFlow);
assertNotNull(createdFlow.getIdentifier());
assertTrue(createdFlow.getCreatedTimestamp() > 0);
assertTrue(createdFlow.getModifiedTimestamp() > 0);
assertEquals(versionedFlow.getIdentifier(), createdFlow.getIdentifier());
assertEquals(versionedFlow.getName(), createdFlow.getName());
assertEquals(versionedFlow.getBucketIdentifier(), createdFlow.getBucketIdentifier());
assertEquals(versionedFlow.getDescription(), createdFlow.getDescription());
}
@Test
public void testGetFlowDoesNotExist() {
when(metadataService.getFlowById(any(String.class))).thenReturn(null);
@ -385,7 +447,7 @@ public class TestRegistryService {
@Test
public void testGetFlowDirectDoesNotExist() {
when(metadataService.getFlowById(any(String.class))).thenReturn(null);
assertThrows(ResourceNotFoundException.class , () -> registryService.getFlow("flow1"));
assertThrows(ResourceNotFoundException.class, () -> registryService.getFlow("flow1"));
}
@Test