http://git-wip-us.apache.org/repos/asf/nifi/blob/f4c94e34/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java b/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java index da253ca..c77778d 100644 --- a/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java +++ b/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java @@ -16,7 +16,29 @@ */ package org.apache.nifi.web; -import com.google.common.collect.Sets; +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Date; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.ListIterator; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.UUID; +import java.util.function.Function; +import java.util.function.Supplier; +import java.util.stream.Collectors; + +import javax.ws.rs.WebApplicationException; +import javax.ws.rs.core.Response; + import org.apache.nifi.action.Action; import org.apache.nifi.action.Component; import org.apache.nifi.action.FlowChangeAction; @@ -161,6 +183,8 @@ import org.apache.nifi.web.api.entity.TemplateEntity; import org.apache.nifi.web.api.entity.TenantEntity; import org.apache.nifi.web.api.entity.UserEntity; import org.apache.nifi.web.api.entity.UserGroupEntity; +import org.apache.nifi.web.concurrent.DistributedLockingManager; +import org.apache.nifi.web.concurrent.LockExpiredException; import org.apache.nifi.web.controller.ControllerFacade; import org.apache.nifi.web.dao.AccessPolicyDAO; import org.apache.nifi.web.dao.ConnectionDAO; @@ -178,7 +202,6 @@ import org.apache.nifi.web.dao.UserDAO; import org.apache.nifi.web.dao.UserGroupDAO; import org.apache.nifi.web.revision.DeleteRevisionTask; import org.apache.nifi.web.revision.ExpiredRevisionClaimException; -import org.apache.nifi.web.revision.ReadOnlyRevisionCallback; import org.apache.nifi.web.revision.RevisionClaim; import org.apache.nifi.web.revision.RevisionManager; import org.apache.nifi.web.revision.RevisionUpdate; @@ -189,28 +212,7 @@ import org.apache.nifi.web.util.SnippetUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import javax.ws.rs.WebApplicationException; -import javax.ws.rs.core.Response; -import java.io.IOException; -import java.nio.charset.StandardCharsets; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.Collection; -import java.util.Date; -import java.util.HashMap; -import java.util.HashSet; -import java.util.LinkedHashMap; -import java.util.LinkedHashSet; -import java.util.List; -import java.util.ListIterator; -import java.util.Map; -import java.util.Optional; -import java.util.Set; -import java.util.UUID; -import java.util.function.Function; -import java.util.function.Supplier; -import java.util.stream.Collectors; -import java.util.stream.Stream; +import com.google.common.collect.Sets; /** * Implementation of NiFiServiceFacade that performs revision checking. @@ -256,50 +258,79 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { private Authorizer authorizer; private AuthorizableLookup authorizableLookup; + private DistributedLockingManager lockManager; // ----------------------------------------- // Synchronization methods // ----------------------------------------- + @Override + public String obtainReadLock() { + return lockManager.getReadLock().lock(); + } @Override - public void authorizeAccess(final AuthorizeAccess authorizeAccess) { - authorizeAccess.authorize(authorizableLookup); + public String obtainReadLock(final String versionIdSeed) { + return lockManager.getReadLock().lock(versionIdSeed); + } + + @Override + public <T> T withReadLock(final String versionIdentifier, final Supplier<T> action) throws LockExpiredException { + return lockManager.getReadLock().withLock(versionIdentifier, action); } @Override - public void claimRevision(final Revision revision, final NiFiUser user) { - revisionManager.requestClaim(revision, user); + public void releaseReadLock(final String versionIdentifier) throws LockExpiredException { + lockManager.getReadLock().unlock(versionIdentifier); } @Override - public void claimRevisions(final Set<Revision> revisions, final NiFiUser user) { - revisionManager.requestClaim(revisions, user); + public String obtainWriteLock() { + return lockManager.getWriteLock().lock(); } @Override - public void cancelRevision(final Revision revision) { - revisionManager.cancelClaim(revision); + public String obtainWriteLock(final String versionIdSeed) { + return lockManager.getWriteLock().lock(versionIdSeed); } @Override - public void cancelRevisions(final Set<Revision> revisions) { - revisionManager.cancelClaims(revisions); + public <T> T withWriteLock(final String versionIdentifier, final Supplier<T> action) throws LockExpiredException { + return lockManager.getWriteLock().withLock(versionIdentifier, action); } @Override - public void releaseRevisionClaim(final Revision revision, final NiFiUser user) throws InvalidRevisionException { - revisionManager.releaseClaim(new StandardRevisionClaim(revision), user); + public void releaseWriteLock(final String versionIdentifier) throws LockExpiredException { + lockManager.getWriteLock().unlock(versionIdentifier); } + + @Override - public void releaseRevisionClaims(final Set<Revision> revisions, final NiFiUser user) throws InvalidRevisionException { - revisionManager.releaseClaim(new StandardRevisionClaim(revisions), user); + public void authorizeAccess(final AuthorizeAccess authorizeAccess) { + authorizeAccess.authorize(authorizableLookup); + } + + @Override + public void verifyRevision(final Revision revision, final NiFiUser user) { + final Revision curRevision = revisionManager.getRevision(revision.getComponentId()); + if (revision.equals(curRevision)) { + return; + } + + throw new InvalidRevisionException(revision + " is not the most up-to-date revision. This component appears to have been modified"); + } + + @Override + public void verifyRevisions(final Set<Revision> revisions, final NiFiUser user) { + for (final Revision revision : revisions) { + verifyRevision(revision, user); + } } @Override public Set<Revision> getRevisionsFromGroup(final String groupId, final Function<ProcessGroup, Set<String>> getComponents) { final ProcessGroup group = processGroupDAO.getProcessGroup(groupId); - final Set<String> componentIds = revisionManager.get(group.getIdentifier(), rev -> getComponents.apply(group)); + final Set<String> componentIds = getComponents.apply(group); return componentIds.stream().map(id -> revisionManager.getRevision(id)).collect(Collectors.toSet()); } @@ -334,17 +365,12 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateConnection(final ConnectionDTO connectionDTO) { - try { - // if connection does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (connectionDAO.hasConnection(connectionDTO.getId())) { - connectionDAO.verifyUpdate(connectionDTO); - } else { - connectionDAO.verifyCreate(connectionDTO.getParentGroupId(), connectionDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(connectionDTO.getId()); - throw e; + // if connection does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (connectionDAO.hasConnection(connectionDTO.getId())) { + connectionDAO.verifyUpdate(connectionDTO); + } else { + connectionDAO.verifyCreate(connectionDTO.getParentGroupId(), connectionDTO); } } @@ -360,15 +386,10 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateInputPort(final PortDTO inputPortDTO) { - try { - // if connection does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (inputPortDAO.hasPort(inputPortDTO.getId())) { - inputPortDAO.verifyUpdate(inputPortDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(inputPortDTO.getId()); - throw e; + // if connection does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (inputPortDAO.hasPort(inputPortDTO.getId())) { + inputPortDAO.verifyUpdate(inputPortDTO); } } @@ -379,15 +400,10 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateOutputPort(final PortDTO outputPortDTO) { - try { - // if connection does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (outputPortDAO.hasPort(outputPortDTO.getId())) { - outputPortDAO.verifyUpdate(outputPortDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(outputPortDTO.getId()); - throw e; + // if connection does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (outputPortDAO.hasPort(outputPortDTO.getId())) { + outputPortDAO.verifyUpdate(outputPortDTO); } } @@ -398,15 +414,10 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateProcessor(final ProcessorDTO processorDTO) { - try { - // if group does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (processorDAO.hasProcessor(processorDTO.getId())) { - processorDAO.verifyUpdate(processorDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(processorDTO.getId()); - throw e; + // if group does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (processorDAO.hasProcessor(processorDTO.getId())) { + processorDAO.verifyUpdate(processorDTO); } } @@ -417,12 +428,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyScheduleComponents(final String groupId, final ScheduledState state, final Set<String> componentIds) { - try { - processGroupDAO.verifyScheduleComponents(groupId, state, componentIds); - } catch (final Exception e) { - componentIds.forEach(id -> revisionManager.cancelClaim(id)); - throw e; - } + processGroupDAO.verifyScheduleComponents(groupId, state, componentIds); } @Override @@ -432,36 +438,21 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateRemoteProcessGroup(final RemoteProcessGroupDTO remoteProcessGroupDTO) { - try { - // if remote group does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (remoteProcessGroupDAO.hasRemoteProcessGroup(remoteProcessGroupDTO.getId())) { - remoteProcessGroupDAO.verifyUpdate(remoteProcessGroupDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(remoteProcessGroupDTO.getId()); - throw e; + // if remote group does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (remoteProcessGroupDAO.hasRemoteProcessGroup(remoteProcessGroupDTO.getId())) { + remoteProcessGroupDAO.verifyUpdate(remoteProcessGroupDTO); } } @Override public void verifyUpdateRemoteProcessGroupInputPort(final String remoteProcessGroupId, final RemoteProcessGroupPortDTO remoteProcessGroupPortDTO) { - try { - remoteProcessGroupDAO.verifyUpdateInputPort(remoteProcessGroupId, remoteProcessGroupPortDTO); - } catch (final Exception e) { - revisionManager.cancelClaim(remoteProcessGroupId); - throw e; - } + remoteProcessGroupDAO.verifyUpdateInputPort(remoteProcessGroupId, remoteProcessGroupPortDTO); } @Override public void verifyUpdateRemoteProcessGroupOutputPort(final String remoteProcessGroupId, final RemoteProcessGroupPortDTO remoteProcessGroupPortDTO) { - try { - remoteProcessGroupDAO.verifyUpdateOutputPort(remoteProcessGroupId, remoteProcessGroupPortDTO); - } catch (final Exception e) { - revisionManager.cancelClaim(remoteProcessGroupId); - throw e; - } + remoteProcessGroupDAO.verifyUpdateOutputPort(remoteProcessGroupId, remoteProcessGroupPortDTO); } @Override @@ -471,26 +462,16 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateControllerService(final ControllerServiceDTO controllerServiceDTO) { - try { - // if service does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (controllerServiceDAO.hasControllerService(controllerServiceDTO.getId())) { - controllerServiceDAO.verifyUpdate(controllerServiceDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(controllerServiceDTO.getId()); - throw e; + // if service does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (controllerServiceDAO.hasControllerService(controllerServiceDTO.getId())) { + controllerServiceDAO.verifyUpdate(controllerServiceDTO); } } @Override public void verifyUpdateControllerServiceReferencingComponents(final String controllerServiceId, final ScheduledState scheduledState, final ControllerServiceState controllerServiceState) { - try { - controllerServiceDAO.verifyUpdateReferencingComponents(controllerServiceId, scheduledState, controllerServiceState); - } catch (final Exception e) { - revisionManager.cancelClaim(controllerServiceId); - throw e; - } + controllerServiceDAO.verifyUpdateReferencingComponents(controllerServiceId, scheduledState, controllerServiceState); } @Override @@ -500,15 +481,10 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateReportingTask(final ReportingTaskDTO reportingTaskDTO) { - try { - // if tasks does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (reportingTaskDAO.hasReportingTask(reportingTaskDTO.getId())) { - reportingTaskDAO.verifyUpdate(reportingTaskDTO); - } - } catch (final Exception e) { - revisionManager.cancelClaim(reportingTaskDTO.getId()); - throw e; + // if tasks does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (reportingTaskDAO.hasReportingTask(reportingTaskDTO.getId())) { + reportingTaskDAO.verifyUpdate(reportingTaskDTO); } } @@ -655,15 +631,10 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyUpdateSnippet(final SnippetDTO snippetDto, final Set<String> affectedComponentIds) { - try { - // if snippet does not exist, then the update request is likely creating it - // so we don't verify since it will fail - if (snippetDAO.hasSnippet(snippetDto.getId())) { - snippetDAO.verifyUpdateSnippetComponent(snippetDto); - } - } catch (final Exception e) { - affectedComponentIds.forEach(id -> revisionManager.cancelClaim(snippetDto.getId())); - throw e; + // if snippet does not exist, then the update request is likely creating it + // so we don't verify since it will fail + if (snippetDAO.hasSnippet(snippetDto.getId())) { + snippetDAO.verifyUpdateSnippetComponent(snippetDto); } } @@ -875,54 +846,32 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyCanClearProcessorState(final String processorId) { - try { - processorDAO.verifyClearState(processorId); - } catch (final Exception e) { - revisionManager.cancelClaim(processorId); - throw e; - } + processorDAO.verifyClearState(processorId); } @Override public void clearProcessorState(final String processorId) { - clearComponentState(processorId, () -> processorDAO.clearState(processorId)); - } - - private void clearComponentState(final String componentId, final Runnable clearState) { - revisionManager.get(componentId, rev -> { - clearState.run(); - return null; - }); + processorDAO.clearState(processorId); } @Override public void verifyCanClearControllerServiceState(final String controllerServiceId) { - try { - controllerServiceDAO.verifyClearState(controllerServiceId); - } catch (final Exception e) { - revisionManager.cancelClaim(controllerServiceId); - throw e; - } + controllerServiceDAO.verifyClearState(controllerServiceId); } @Override public void clearControllerServiceState(final String controllerServiceId) { - clearComponentState(controllerServiceId, () -> controllerServiceDAO.clearState(controllerServiceId)); + controllerServiceDAO.clearState(controllerServiceId); } @Override public void verifyCanClearReportingTaskState(final String reportingTaskId) { - try { - reportingTaskDAO.verifyClearState(reportingTaskId); - } catch (final Exception e) { - revisionManager.cancelClaim(reportingTaskId); - throw e; - } + reportingTaskDAO.verifyClearState(reportingTaskId); } @Override public void clearReportingTaskState(final String reportingTaskId) { - clearComponentState(reportingTaskId, () -> reportingTaskDAO.clearState(reportingTaskId)); + reportingTaskDAO.clearState(reportingTaskId); } @Override @@ -1066,12 +1015,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public void verifyDeleteSnippet(final String snippetId, final Set<String> affectedComponentIds) { - try { - snippetDAO.verifyDeleteSnippetComponents(snippetId); - } catch (final Exception e) { - affectedComponentIds.forEach(id -> revisionManager.cancelClaim(id)); - throw e; - } + snippetDAO.verifyDeleteSnippetComponents(snippetId); } @Override @@ -1229,29 +1173,22 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { */ private <D, C> RevisionUpdate<D> createComponent(final Revision revision, final ComponentDTO componentDto, final Supplier<C> daoCreation, final Function<C, D> dtoCreation) { final NiFiUser user = NiFiUserUtils.getNiFiUser(); - final String groupId = componentDto.getParentGroupId(); // read lock on the containing group - return revisionManager.get(groupId, rev -> { - // request claim for component to be created... revision already verified (version == 0) - final RevisionClaim claim = revisionManager.requestClaim(revision, user); - try { - // update revision through revision manager - return revisionManager.updateRevision(claim, user, () -> { - // add the component - final C component = daoCreation.get(); + // request claim for component to be created... revision already verified (version == 0) + final RevisionClaim claim = new StandardRevisionClaim(revision); - // save the flow - controllerFacade.save(); + // update revision through revision manager + return revisionManager.updateRevision(claim, user, () -> { + // add the component + final C component = daoCreation.get(); - final D dto = dtoCreation.apply(component); - final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); - return new StandardRevisionUpdate<D>(dto, lastMod); - }); - } finally { - // cancel in case of exception... noop if successful - revisionManager.cancelClaim(revision.getComponentId()); - } + // save the flow + controllerFacade.save(); + + final D dto = dtoCreation.apply(component); + final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); + return new StandardRevisionUpdate<D>(dto, lastMod); }); } @@ -1319,7 +1256,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { // validate any processors if (flow.getProcessors() != null) { for (final ProcessorDTO processorDTO : flow.getProcessors()) { - final ProcessorNode processorNode = revisionManager.get(processorDTO.getId(), rev -> processorDAO.getProcessor(processorDTO.getId())); + final ProcessorNode processorNode = processorDAO.getProcessor(processorDTO.getId()); final Collection<ValidationResult> validationErrors = processorNode.getValidationErrors(); if (validationErrors != null && !validationErrors.isEmpty()) { final List<String> errors = new ArrayList<>(); @@ -1333,7 +1270,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { if (flow.getInputPorts() != null) { for (final PortDTO portDTO : flow.getInputPorts()) { - final Port port = revisionManager.get(portDTO.getId(), rev -> inputPortDAO.getPort(portDTO.getId())); + final Port port = inputPortDAO.getPort(portDTO.getId()); final Collection<ValidationResult> validationErrors = port.getValidationErrors(); if (validationErrors != null && !validationErrors.isEmpty()) { final List<String> errors = new ArrayList<>(); @@ -1347,7 +1284,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { if (flow.getOutputPorts() != null) { for (final PortDTO portDTO : flow.getOutputPorts()) { - final Port port = revisionManager.get(portDTO.getId(), rev -> outputPortDAO.getPort(portDTO.getId())); + final Port port = outputPortDAO.getPort(portDTO.getId()); final Collection<ValidationResult> validationErrors = port.getValidationErrors(); if (validationErrors != null && !validationErrors.isEmpty()) { final List<String> errors = new ArrayList<>(); @@ -1362,8 +1299,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { // get any remote process group issues if (flow.getRemoteProcessGroups() != null) { for (final RemoteProcessGroupDTO remoteProcessGroupDTO : flow.getRemoteProcessGroups()) { - final RemoteProcessGroup remoteProcessGroup = revisionManager.get( - remoteProcessGroupDTO.getId(), rev -> remoteProcessGroupDAO.getRemoteProcessGroup(remoteProcessGroupDTO.getId())); + final RemoteProcessGroup remoteProcessGroup = remoteProcessGroupDAO.getRemoteProcessGroup(remoteProcessGroupDTO.getId()); if (remoteProcessGroup.getAuthorizationIssue() != null) { remoteProcessGroupDTO.setAuthorizationIssues(Arrays.asList(remoteProcessGroup.getAuthorizationIssue())); @@ -1374,20 +1310,17 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public FlowEntity copySnippet(final String groupId, final String snippetId, final Double originX, final Double originY, final String idGenerationSeed) { - final FlowDTO flowDto = revisionManager.get(groupId, - rev -> { - // create the new snippet - final FlowSnippetDTO snippet = snippetDAO.copySnippet(groupId, snippetId, originX, originY, idGenerationSeed); + // create the new snippet + final FlowSnippetDTO snippet = snippetDAO.copySnippet(groupId, snippetId, originX, originY, idGenerationSeed); - // save the flow - controllerFacade.save(); + // save the flow + controllerFacade.save(); - // drop the snippet - snippetDAO.dropSnippet(snippetId); + // drop the snippet + snippetDAO.dropSnippet(snippetId); - // post process new flow snippet - return postProcessNewFlowSnippet(groupId, snippet); - }); + // post process new flow snippet + final FlowDTO flowDto = postProcessNewFlowSnippet(groupId, snippet); final FlowEntity flowEntity = new FlowEntity(); flowEntity.setFlow(flowDto); @@ -1396,17 +1329,14 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public SnippetEntity createSnippet(final SnippetDTO snippetDTO) { - final String groupId = snippetDTO.getParentGroupId(); - final RevisionUpdate<SnippetDTO> snapshot = revisionManager.get(groupId, rev -> { - // add the component - final Snippet snippet = snippetDAO.createSnippet(snippetDTO); + // add the component + final Snippet snippet = snippetDAO.createSnippet(snippetDTO); - // save the flow - controllerFacade.save(); + // save the flow + controllerFacade.save(); - final SnippetDTO dto = dtoFactory.createSnippetDto(snippet); - return new StandardRevisionUpdate<SnippetDTO>(dto, null); - }); + final SnippetDTO dto = dtoFactory.createSnippetDto(snippet); + final RevisionUpdate<SnippetDTO> snapshot = new StandardRevisionUpdate<SnippetDTO>(dto, null); return entityFactory.createSnippetEntity(snapshot.getComponent()); } @@ -1558,27 +1488,22 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { .map(label -> label.getId()) .forEach(id -> identifiers.add(id)); - return revisionManager.get(identifiers, - () -> { - final ProcessGroup group = processGroupDAO.getProcessGroup(groupId); - final ProcessGroupStatus groupStatus = controllerFacade.getProcessGroupStatus(groupId); - return dtoFactory.createFlowDto(group, groupStatus, snippet, revisionManager); - }); + final ProcessGroup group = processGroupDAO.getProcessGroup(groupId); + final ProcessGroupStatus groupStatus = controllerFacade.getProcessGroupStatus(groupId); + return dtoFactory.createFlowDto(group, groupStatus, snippet, revisionManager); } @Override public FlowEntity createTemplateInstance(final String groupId, final Double originX, final Double originY, final String templateId, final String idGenerationSeed) { - final FlowDTO flowDto = revisionManager.get(groupId, rev -> { - // instantiate the template - there is no need to make another copy of the flow snippet since the actual template - // was copied and this dto is only used to instantiate it's components (which as already completed) - final FlowSnippetDTO snippet = templateDAO.instantiateTemplate(groupId, originX, originY, templateId, idGenerationSeed); + // instantiate the template - there is no need to make another copy of the flow snippet since the actual template + // was copied and this dto is only used to instantiate it's components (which as already completed) + final FlowSnippetDTO snippet = templateDAO.instantiateTemplate(groupId, originX, originY, templateId, idGenerationSeed); - // save the flow - controllerFacade.save(); + // save the flow + controllerFacade.save(); - // post process the new flow snippet - return postProcessNewFlowSnippet(groupId, snippet); - }); + // post process the new flow snippet + final FlowDTO flowDto = postProcessNewFlowSnippet(groupId, snippet); final FlowEntity flowEntity = new FlowEntity(); flowEntity.setFlow(flowDto); @@ -1603,11 +1528,10 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { final NiFiUser user = NiFiUserUtils.getNiFiUser(); // request claim for component to be created... revision already verified (version == 0) - final RevisionClaim claim = revisionManager.requestClaim(revision, user); + final RevisionClaim claim = new StandardRevisionClaim(revision); final RevisionUpdate<ControllerServiceDTO> snapshot; if (groupId == null) { - try { // update revision through revision manager snapshot = revisionManager.updateRevision(claim, user, () -> { // Unfortunately, we can not use the createComponent() method here because createComponent() wants to obtain the read lock @@ -1620,26 +1544,15 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); return new StandardRevisionUpdate<ControllerServiceDTO>(dto, lastMod); }); - } finally { - // cancel in case of exception... noop if successful - revisionManager.cancelClaim(revision.getComponentId()); - } } else { - snapshot = revisionManager.get(groupId, groupRevision -> { - try { - return revisionManager.updateRevision(claim, user, () -> { - final ControllerServiceNode controllerService = controllerServiceDAO.createControllerService(controllerServiceDTO); - final ControllerServiceDTO dto = dtoFactory.createControllerServiceDto(controllerService); - - controllerFacade.save(); - - final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); - return new StandardRevisionUpdate<ControllerServiceDTO>(dto, lastMod); - }); - } finally { - // cancel in case of exception... noop if successful - revisionManager.cancelClaim(revision.getComponentId()); - } + snapshot = revisionManager.updateRevision(claim, user, () -> { + final ControllerServiceNode controllerService = controllerServiceDAO.createControllerService(controllerServiceDTO); + final ControllerServiceDTO dto = dtoFactory.createControllerServiceDto(controllerService); + + controllerFacade.save(); + + final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); + return new StandardRevisionUpdate<ControllerServiceDTO>(dto, lastMod); }); } @@ -1729,13 +1642,11 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { // TODO remove once we can update a read lock referencingIds.removeAll(lockedIds); - return revisionManager.get(referencingIds, () -> { - final Map<String, Revision> referencingRevisions = new HashMap<>(); - for (final ConfiguredComponent component : reference.getReferencingComponents()) { - referencingRevisions.put(component.getIdentifier(), revisionManager.getRevision(component.getIdentifier())); - } - return createControllerServiceReferencingComponentsEntity(reference, referencingRevisions); - }); + final Map<String, Revision> referencingRevisions = new HashMap<>(); + for (final ConfiguredComponent component : reference.getReferencingComponents()) { + referencingRevisions.put(component.getIdentifier(), revisionManager.getRevision(component.getIdentifier())); + } + return createControllerServiceReferencingComponentsEntity(reference, referencingRevisions); } /** @@ -1818,29 +1729,25 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { final NiFiUser user = NiFiUserUtils.getNiFiUser(); // request claim for component to be created... revision already verified (version == 0) - final RevisionClaim claim = revisionManager.requestClaim(revision, user); - try { - // update revision through revision manager - final RevisionUpdate<ReportingTaskDTO> snapshot = revisionManager.updateRevision(claim, user, () -> { - // create the reporting task - final ReportingTaskNode reportingTask = reportingTaskDAO.createReportingTask(reportingTaskDTO); + final RevisionClaim claim = new StandardRevisionClaim(revision); - // save the update - controllerFacade.save(); + // update revision through revision manager + final RevisionUpdate<ReportingTaskDTO> snapshot = revisionManager.updateRevision(claim, user, () -> { + // create the reporting task + final ReportingTaskNode reportingTask = reportingTaskDAO.createReportingTask(reportingTaskDTO); - final ReportingTaskDTO dto = dtoFactory.createReportingTaskDto(reportingTask); - final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); - return new StandardRevisionUpdate<ReportingTaskDTO>(dto, lastMod); - }); + // save the update + controllerFacade.save(); - final ReportingTaskNode reportingTask = reportingTaskDAO.getReportingTask(reportingTaskDTO.getId()); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(reportingTask); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(reportingTask.getIdentifier())); - return entityFactory.createReportingTaskEntity(snapshot.getComponent(), dtoFactory.createRevisionDTO(snapshot.getLastModification()), accessPolicy, bulletins); - } finally { - // cancel in case of exception... noop if successful - revisionManager.cancelClaim(revision.getComponentId()); - } + final ReportingTaskDTO dto = dtoFactory.createReportingTaskDto(reportingTask); + final FlowModification lastMod = new FlowModification(revision.incrementRevision(revision.getClientId()), user.getIdentity()); + return new StandardRevisionUpdate<ReportingTaskDTO>(dto, lastMod); + }); + + final ReportingTaskNode reportingTask = reportingTaskDAO.getReportingTask(reportingTaskDTO.getId()); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(reportingTask); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(reportingTask.getIdentifier())); + return entityFactory.createReportingTaskEntity(snapshot.getComponent(), dtoFactory.createRevisionDTO(snapshot.getLastModification()), accessPolicy, bulletins); } @Override @@ -1971,47 +1878,32 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public ComponentStateDTO getProcessorState(final String processorId) { - return revisionManager.get(processorId, new ReadOnlyRevisionCallback<ComponentStateDTO>() { - @Override - public ComponentStateDTO withRevision(final Revision revision) { - final StateMap clusterState = isClustered() ? processorDAO.getState(processorId, Scope.CLUSTER) : null; - final StateMap localState = processorDAO.getState(processorId, Scope.LOCAL); + final StateMap clusterState = isClustered() ? processorDAO.getState(processorId, Scope.CLUSTER) : null; + final StateMap localState = processorDAO.getState(processorId, Scope.LOCAL); - // processor will be non null as it was already found when getting the state - final ProcessorNode processor = processorDAO.getProcessor(processorId); - return dtoFactory.createComponentStateDTO(processorId, processor.getProcessor().getClass(), localState, clusterState); - } - }); + // processor will be non null as it was already found when getting the state + final ProcessorNode processor = processorDAO.getProcessor(processorId); + return dtoFactory.createComponentStateDTO(processorId, processor.getProcessor().getClass(), localState, clusterState); } @Override public ComponentStateDTO getControllerServiceState(final String controllerServiceId) { - return revisionManager.get(controllerServiceId, new ReadOnlyRevisionCallback<ComponentStateDTO>() { - @Override - public ComponentStateDTO withRevision(final Revision revision) { - final StateMap clusterState = isClustered() ? controllerServiceDAO.getState(controllerServiceId, Scope.CLUSTER) : null; - final StateMap localState = controllerServiceDAO.getState(controllerServiceId, Scope.LOCAL); + final StateMap clusterState = isClustered() ? controllerServiceDAO.getState(controllerServiceId, Scope.CLUSTER) : null; + final StateMap localState = controllerServiceDAO.getState(controllerServiceId, Scope.LOCAL); - // controller service will be non null as it was already found when getting the state - final ControllerServiceNode controllerService = controllerServiceDAO.getControllerService(controllerServiceId); - return dtoFactory.createComponentStateDTO(controllerServiceId, controllerService.getControllerServiceImplementation().getClass(), localState, clusterState); - } - }); + // controller service will be non null as it was already found when getting the state + final ControllerServiceNode controllerService = controllerServiceDAO.getControllerService(controllerServiceId); + return dtoFactory.createComponentStateDTO(controllerServiceId, controllerService.getControllerServiceImplementation().getClass(), localState, clusterState); } @Override public ComponentStateDTO getReportingTaskState(final String reportingTaskId) { - return revisionManager.get(reportingTaskId, new ReadOnlyRevisionCallback<ComponentStateDTO>() { - @Override - public ComponentStateDTO withRevision(final Revision revision) { - final StateMap clusterState = isClustered() ? reportingTaskDAO.getState(reportingTaskId, Scope.CLUSTER) : null; - final StateMap localState = reportingTaskDAO.getState(reportingTaskId, Scope.LOCAL); + final StateMap clusterState = isClustered() ? reportingTaskDAO.getState(reportingTaskId, Scope.CLUSTER) : null; + final StateMap localState = reportingTaskDAO.getState(reportingTaskId, Scope.LOCAL); - // reporting task will be non null as it was already found when getting the state - final ReportingTaskNode reportingTask = reportingTaskDAO.getReportingTask(reportingTaskId); - return dtoFactory.createComponentStateDTO(reportingTaskId, reportingTask.getReportingTask().getClass(), localState, clusterState); - } - }); + // reporting task will be non null as it was already found when getting the state + final ReportingTaskNode reportingTask = reportingTaskDAO.getReportingTask(reportingTaskId); + return dtoFactory.createComponentStateDTO(reportingTaskId, reportingTask.getReportingTask().getClass(), localState, clusterState); } @Override @@ -2029,33 +1921,25 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { return countersDto; } + private ConnectionEntity createConnectionEntity(final Connection connection) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(connection.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(connection); + final ConnectionStatusDTO status = dtoFactory.createConnectionStatusDto(controllerFacade.getConnectionStatus(connection.getIdentifier())); + return entityFactory.createConnectionEntity(dtoFactory.createConnectionDto(connection), revision, accessPolicy, status); + } + @Override public Set<ConnectionEntity> getConnections(final String groupId) { - return revisionManager.get(groupId, rev -> { - final Set<Connection> connections = connectionDAO.getConnections(groupId); - final Set<String> connectionIds = connections.stream().map(connection -> connection.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(connectionIds, () -> { - return connections.stream() - .map(connection -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(connection.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(connection); - final ConnectionStatusDTO status = dtoFactory.createConnectionStatusDto(controllerFacade.getConnectionStatus(connection.getIdentifier())); - return entityFactory.createConnectionEntity(dtoFactory.createConnectionDto(connection), revision, accessPolicy, status); - }) - .collect(Collectors.toSet()); - }); - }); + final Set<Connection> connections = connectionDAO.getConnections(groupId); + return connections.stream() + .map(connection -> createConnectionEntity(connection)) + .collect(Collectors.toSet()); } @Override public ConnectionEntity getConnection(final String connectionId) { - return revisionManager.get(connectionId, rev -> { - final Connection connection = connectionDAO.getConnection(connectionId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(connection); - final ConnectionStatusDTO status = dtoFactory.createConnectionStatusDto(controllerFacade.getConnectionStatus(connectionId)); - return entityFactory.createConnectionEntity(dtoFactory.createConnectionDto(connectionDAO.getConnection(connectionId)), revision, accessPolicy, status); - }); + final Connection connection = connectionDAO.getConnection(connectionId); + return createConnectionEntity(connection); } @Override @@ -2086,30 +1970,28 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public ConnectionStatusDTO getConnectionStatus(final String connectionId) { - return revisionManager.get(connectionId, rev -> dtoFactory.createConnectionStatusDto(controllerFacade.getConnectionStatus(connectionId))); + return dtoFactory.createConnectionStatusDto(controllerFacade.getConnectionStatus(connectionId)); } @Override public StatusHistoryDTO getConnectionStatusHistory(final String connectionId) { - return revisionManager.get(connectionId, rev -> controllerFacade.getConnectionStatusHistory(connectionId)); + return controllerFacade.getConnectionStatusHistory(connectionId); + } + + private ProcessorEntity createProcessorEntity(final ProcessorNode processor) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(processor.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(processor); + final ProcessorStatusDTO status = dtoFactory.createProcessorStatusDto(controllerFacade.getProcessorStatus(processor.getIdentifier())); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(processor.getIdentifier())); + return entityFactory.createProcessorEntity(dtoFactory.createProcessorDto(processor), revision, accessPolicy, status, bulletins); } @Override public Set<ProcessorEntity> getProcessors(final String groupId) { final Set<ProcessorNode> processors = processorDAO.getProcessors(groupId); - final Set<String> ids = processors.stream().map(proc -> proc.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - return processors.stream() - .map(processor -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(processor.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(processor); - final ProcessorStatusDTO status = dtoFactory.createProcessorStatusDto(controllerFacade.getProcessorStatus(processor.getIdentifier())); - final List<BulletinDTO> bulletins = - dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(processor.getIdentifier())); - return entityFactory.createProcessorEntity(dtoFactory.createProcessorDto(processor), revision, accessPolicy, status, bulletins); - }) - .collect(Collectors.toSet()); - }); + return processors.stream() + .map(processor -> createProcessorEntity(processor)) + .collect(Collectors.toSet()); } @Override @@ -2164,13 +2046,8 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public ProcessorEntity getProcessor(final String id) { - return revisionManager.get(id, rev -> { - final ProcessorNode processor = processorDAO.getProcessor(id); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final ProcessorStatusDTO status = dtoFactory.createProcessorStatusDto(controllerFacade.getProcessorStatus(id)); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(id)); - return entityFactory.createProcessorEntity(dtoFactory.createProcessorDto(processor), revision, dtoFactory.createAccessPolicyDto(processor), status, bulletins); - }); + final ProcessorNode processor = processorDAO.getProcessor(id); + return createProcessorEntity(processor); } @Override @@ -2188,7 +2065,7 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public ProcessorStatusDTO getProcessorStatus(final String id) { - return revisionManager.get(id, rev -> dtoFactory.createProcessorStatusDto(controllerFacade.getProcessorStatus(id))); + return dtoFactory.createProcessorStatusDto(controllerFacade.getProcessorStatus(id)); } @Override @@ -2270,46 +2147,33 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { // serialize the input ports this NiFi has access to final Set<PortDTO> inputPortDtos = new LinkedHashSet<>(); final Set<RootGroupPort> inputPorts = controllerFacade.getInputPorts(); - final Set<String> inputPortIds = inputPorts.stream().map(port -> port.getIdentifier()).collect(Collectors.toSet()); - revisionManager.get(inputPortIds, () -> { - for (final RootGroupPort inputPort : inputPorts) { - if (isUserAuthorized(user, inputPort)) { - final PortDTO dto = new PortDTO(); - dto.setId(inputPort.getIdentifier()); - dto.setName(inputPort.getName()); - dto.setComments(inputPort.getComments()); - dto.setState(inputPort.getScheduledState().toString()); - inputPortDtos.add(dto); - } + for (final RootGroupPort inputPort : inputPorts) { + if (isUserAuthorized(user, inputPort)) { + final PortDTO dto = new PortDTO(); + dto.setId(inputPort.getIdentifier()); + dto.setName(inputPort.getName()); + dto.setComments(inputPort.getComments()); + dto.setState(inputPort.getScheduledState().toString()); + inputPortDtos.add(dto); } - return null; - }); + } // serialize the output ports this NiFi has access to final Set<PortDTO> outputPortDtos = new LinkedHashSet<>(); - final Set<RootGroupPort> outputPorts = controllerFacade.getOutputPorts(); - final Set<String> outputPortIds = outputPorts.stream().map(port -> port.getIdentifier()).collect(Collectors.toSet()); - revisionManager.get(outputPortIds, () -> { - for (final RootGroupPort outputPort : controllerFacade.getOutputPorts()) { - if (isUserAuthorized(user, outputPort)) { - final PortDTO dto = new PortDTO(); - dto.setId(outputPort.getIdentifier()); - dto.setName(outputPort.getName()); - dto.setComments(outputPort.getComments()); - dto.setState(outputPort.getScheduledState().toString()); - outputPortDtos.add(dto); - } + for (final RootGroupPort outputPort : controllerFacade.getOutputPorts()) { + if (isUserAuthorized(user, outputPort)) { + final PortDTO dto = new PortDTO(); + dto.setId(outputPort.getIdentifier()); + dto.setName(outputPort.getName()); + dto.setComments(outputPort.getComments()); + dto.setState(outputPort.getScheduledState().toString()); + outputPortDtos.add(dto); } - - return null; - }); + } // get the root group - final String rootGroupId = controllerFacade.getRootGroupId(); - final ProcessGroupCounts counts = revisionManager.get(rootGroupId, rev -> { - final ProcessGroup rootGroup = processGroupDAO.getProcessGroup(controllerFacade.getRootGroupId()); - return rootGroup.getCounts(); - }); + final ProcessGroup rootGroup = processGroupDAO.getProcessGroup(controllerFacade.getRootGroupId()); + final ProcessGroupCounts counts = rootGroup.getCounts(); // create the controller dto final ControllerDTO controllerDTO = new ControllerDTO(); @@ -2340,12 +2204,11 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public ControllerConfigurationEntity getControllerConfiguration() { - return revisionManager.get(FlowController.class.getSimpleName(), rev -> { - final ControllerConfigurationDTO dto = dtoFactory.createControllerConfigurationDto(controllerFacade); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(controllerFacade); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - return entityFactory.createControllerConfigurationEntity(dto, revision, accessPolicy); - }); + final Revision rev = revisionManager.getRevision(FlowController.class.getSimpleName()); + final ControllerConfigurationDTO dto = dtoFactory.createControllerConfigurationDto(controllerFacade); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(controllerFacade); + final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); + return entityFactory.createControllerConfigurationEntity(dto, revision, accessPolicy); } @Override @@ -2358,246 +2221,198 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public AccessPolicyEntity getAccessPolicy(final String accessPolicyId) { - AccessPolicy preRevisionRequestAccessPolicy = accessPolicyDAO.getAccessPolicy(accessPolicyId); - Set<String> ids = Stream.concat(Stream.of(accessPolicyId), - Stream.concat(preRevisionRequestAccessPolicy.getUsers().stream(), preRevisionRequestAccessPolicy.getGroups().stream())).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - final RevisionDTO requestedAccessPolicyRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(accessPolicyId)); - final AccessPolicy requestedAccessPolicy = accessPolicyDAO.getAccessPolicy(accessPolicyId); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(authorizableLookup.getAccessPolicyAuthorizable(accessPolicyId)); - return entityFactory.createAccessPolicyEntity( - dtoFactory.createAccessPolicyDto(requestedAccessPolicy, - requestedAccessPolicy.getGroups().stream().map(mapUserGroupIdToTenantEntity()).collect(Collectors.toSet()), - requestedAccessPolicy.getUsers().stream().map(mapUserIdToTenantEntity()).collect(Collectors.toSet())), - requestedAccessPolicyRevision, accessPolicy); - }); + final RevisionDTO requestedAccessPolicyRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(accessPolicyId)); + final AccessPolicy requestedAccessPolicy = accessPolicyDAO.getAccessPolicy(accessPolicyId); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(authorizableLookup.getAccessPolicyAuthorizable(accessPolicyId)); + return entityFactory.createAccessPolicyEntity( + dtoFactory.createAccessPolicyDto(requestedAccessPolicy, + requestedAccessPolicy.getGroups().stream().map(mapUserGroupIdToTenantEntity()).collect(Collectors.toSet()), + requestedAccessPolicy.getUsers().stream().map(mapUserIdToTenantEntity()).collect(Collectors.toSet())), + requestedAccessPolicyRevision, accessPolicy); } @Override public UserEntity getUser(final String userId) { final Authorizable usersAuthorizable = authorizableLookup.getTenantAuthorizable(); - Set<String> ids = Stream.concat(Stream.of(userId), userDAO.getUser(userId).getGroups().stream()).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - final RevisionDTO userRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(userId)); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(usersAuthorizable); - final User user = userDAO.getUser(userId); - final Set<TenantEntity> userGroups = user.getGroups().stream() - .map(mapUserGroupIdToTenantEntity()).collect(Collectors.toSet()); - return entityFactory.createUserEntity(dtoFactory.createUserDto(user, userGroups), userRevision, accessPolicy); - }); + final RevisionDTO userRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(userId)); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(usersAuthorizable); + final User user = userDAO.getUser(userId); + final Set<TenantEntity> userGroups = user.getGroups().stream() + .map(mapUserGroupIdToTenantEntity()).collect(Collectors.toSet()); + return entityFactory.createUserEntity(dtoFactory.createUserDto(user, userGroups), userRevision, accessPolicy); + } + + private UserEntity createUserEntity(final User user) { + final RevisionDTO userRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(user.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(authorizableLookup.getTenantAuthorizable()); + final Set<TenantEntity> userGroups = user.getGroups().stream().map(mapUserGroupIdToTenantEntity()).collect(Collectors.toSet()); + return entityFactory.createUserEntity(dtoFactory.createUserDto(user, userGroups), userRevision, accessPolicy); } @Override public Set<UserEntity> getUsers() { final Set<User> users = userDAO.getUsers(); - final Set<String> ids = users.stream().flatMap(user -> Stream.concat(Stream.of(user.getIdentifier()), user.getGroups().stream())).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - return users.stream() - .map(user -> { - final RevisionDTO userRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(user.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(authorizableLookup.getTenantAuthorizable()); - final Set<TenantEntity> userGroups = user.getGroups().stream().map(mapUserGroupIdToTenantEntity()).collect(Collectors.toSet()); - return entityFactory.createUserEntity(dtoFactory.createUserDto(user, userGroups), userRevision, accessPolicy); - }).collect(Collectors.toSet()); - }); + return users.stream() + .map(user -> createUserEntity(user)) + .collect(Collectors.toSet()); + } + + private UserGroupEntity createUserGroupEntity(final Group userGroup) { + final RevisionDTO userGroupRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(userGroup.getIdentifier())); + final Set<TenantEntity> users = userGroup.getUsers().stream().map(mapUserIdToTenantEntity()).collect(Collectors.toSet()); + return entityFactory.createUserGroupEntity(dtoFactory.createUserGroupDto(userGroup, users), userGroupRevision, + dtoFactory.createAccessPolicyDto(authorizableLookup.getTenantAuthorizable())); } @Override public UserGroupEntity getUserGroup(final String userGroupId) { - Set<String> ids = Stream.concat(Stream.of(userGroupId), userGroupDAO.getUserGroup(userGroupId).getUsers().stream()).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - final RevisionDTO userGroupRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(userGroupId)); - final Group userGroup = userGroupDAO.getUserGroup(userGroupId); - final Set<TenantEntity> users = userGroup.getUsers().stream().map(mapUserIdToTenantEntity()).collect(Collectors.toSet()); - return entityFactory.createUserGroupEntity(dtoFactory.createUserGroupDto(userGroup, users), userGroupRevision, - dtoFactory.createAccessPolicyDto(authorizableLookup.getTenantAuthorizable())); - }); + final Group userGroup = userGroupDAO.getUserGroup(userGroupId); + return createUserGroupEntity(userGroup); } @Override public Set<UserGroupEntity> getUserGroups() { - final Authorizable userGroupAuthorizable = authorizableLookup.getTenantAuthorizable(); final Set<Group> userGroups = userGroupDAO.getUserGroups(); - final Set<String> ids = userGroups.stream().flatMap(userGroup -> Stream.concat(Stream.of(userGroup.getIdentifier()), userGroup.getUsers().stream())).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - return userGroups.stream() - .map(userGroup -> { - final RevisionDTO userGroupRevision = dtoFactory.createRevisionDTO(revisionManager.getRevision(userGroup.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(userGroupAuthorizable); - final Set<TenantEntity> users = userGroup.getUsers().stream().map(mapUserIdToTenantEntity()).collect(Collectors.toSet()); - return entityFactory.createUserGroupEntity(dtoFactory.createUserGroupDto(userGroup, users), userGroupRevision, accessPolicy); - }).collect(Collectors.toSet()); - }); + return userGroups.stream() + .map(userGroup -> createUserGroupEntity(userGroup)) + .collect(Collectors.toSet()); + } + + private LabelEntity createLabelEntity(final Label label) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(label.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(label); + return entityFactory.createLabelEntity(dtoFactory.createLabelDto(label), revision, accessPolicy); } @Override public Set<LabelEntity> getLabels(final String groupId) { final Set<Label> labels = labelDAO.getLabels(groupId); - final Set<String> ids = labels.stream().map(label -> label.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - return labels.stream() - .map(label -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(label.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(label); - return entityFactory.createLabelEntity(dtoFactory.createLabelDto(label), revision, accessPolicy); - }) - .collect(Collectors.toSet()); - }); + return labels.stream() + .map(label -> createLabelEntity(label)) + .collect(Collectors.toSet()); } @Override public LabelEntity getLabel(final String labelId) { - return revisionManager.get(labelId, rev -> { - final Label label = labelDAO.getLabel(labelId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(label); - return entityFactory.createLabelEntity(dtoFactory.createLabelDto(label), revision, accessPolicy); - }); + final Label label = labelDAO.getLabel(labelId); + return createLabelEntity(label); + } + + private FunnelEntity createFunnelEntity(final Funnel funnel) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(funnel.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(funnel); + return entityFactory.createFunnelEntity(dtoFactory.createFunnelDto(funnel), revision, accessPolicy); } @Override public Set<FunnelEntity> getFunnels(final String groupId) { final Set<Funnel> funnels = funnelDAO.getFunnels(groupId); - final Set<String> funnelIds = funnels.stream().map(funnel -> funnel.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(funnelIds, () -> { - return funnels.stream() - .map(funnel -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(funnel.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(funnel); - return entityFactory.createFunnelEntity(dtoFactory.createFunnelDto(funnel), revision, accessPolicy); - }) - .collect(Collectors.toSet()); - }); + return funnels.stream() + .map(funnel -> createFunnelEntity(funnel)) + .collect(Collectors.toSet()); } @Override public FunnelEntity getFunnel(final String funnelId) { - return revisionManager.get(funnelId, rev -> { - final Funnel funnel = funnelDAO.getFunnel(funnelId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(funnel); - return entityFactory.createFunnelEntity(dtoFactory.createFunnelDto(funnel), revision, accessPolicy); - }); + final Funnel funnel = funnelDAO.getFunnel(funnelId); + return createFunnelEntity(funnel); + } + + private PortEntity createInputPortEntity(final Port port) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(port.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(port); + final PortStatusDTO status = dtoFactory.createPortStatusDto(controllerFacade.getInputPortStatus(port.getIdentifier())); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(port.getIdentifier())); + return entityFactory.createPortEntity(dtoFactory.createPortDto(port), revision, accessPolicy, status, bulletins); + } + + private PortEntity createOutputPortEntity(final Port port) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(port.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(port); + final PortStatusDTO status = dtoFactory.createPortStatusDto(controllerFacade.getOutputPortStatus(port.getIdentifier())); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(port.getIdentifier())); + return entityFactory.createPortEntity(dtoFactory.createPortDto(port), revision, accessPolicy, status, bulletins); } @Override public Set<PortEntity> getInputPorts(final String groupId) { final Set<Port> inputPorts = inputPortDAO.getPorts(groupId); - final Set<String> portIds = inputPorts.stream().map(port -> port.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(portIds, () -> { - return inputPorts.stream() - .map(port -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(port.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(port); - final PortStatusDTO status = dtoFactory.createPortStatusDto(controllerFacade.getInputPortStatus(port.getIdentifier())); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(port.getIdentifier())); - return entityFactory.createPortEntity(dtoFactory.createPortDto(port), revision, accessPolicy, status, bulletins); - }) - .collect(Collectors.toSet()); - }); + return inputPorts.stream() + .map(port -> createInputPortEntity(port)) + .collect(Collectors.toSet()); } @Override public Set<PortEntity> getOutputPorts(final String groupId) { final Set<Port> ports = outputPortDAO.getPorts(groupId); - final Set<String> ids = ports.stream().map(port -> port.getIdentifier()).collect(Collectors.toSet()); - - return revisionManager.get(ids, () -> { - return ports.stream() - .map(port -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(port.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(port); - final PortStatusDTO status = dtoFactory.createPortStatusDto(controllerFacade.getOutputPortStatus(port.getIdentifier())); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(port.getIdentifier())); - return entityFactory.createPortEntity(dtoFactory.createPortDto(port), revision, accessPolicy, status, bulletins); - }) - .collect(Collectors.toSet()); - }); + return ports.stream() + .map(port -> createOutputPortEntity(port)) + .collect(Collectors.toSet()); + } + + private ProcessGroupEntity createProcessGroupEntity(final ProcessGroup group) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(group.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(group); + final ProcessGroupStatusDTO status = dtoFactory.createConciseProcessGroupStatusDto(controllerFacade.getProcessGroupStatus(group.getIdentifier())); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(group.getIdentifier())); + return entityFactory.createProcessGroupEntity(dtoFactory.createProcessGroupDto(group), revision, accessPolicy, status, bulletins); } @Override public Set<ProcessGroupEntity> getProcessGroups(final String parentGroupId) { final Set<ProcessGroup> groups = processGroupDAO.getProcessGroups(parentGroupId); - final Set<String> ids = groups.stream().map(group -> group.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - return groups.stream() - .map(group -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(group.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(group); - final ProcessGroupStatusDTO status = dtoFactory.createConciseProcessGroupStatusDto(controllerFacade.getProcessGroupStatus(group.getIdentifier())); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(group.getIdentifier())); - return entityFactory.createProcessGroupEntity(dtoFactory.createProcessGroupDto(group), revision, accessPolicy, status, bulletins); - }) - .collect(Collectors.toSet()); - }); + return groups.stream() + .map(group -> createProcessGroupEntity(group)) + .collect(Collectors.toSet()); + } + + private RemoteProcessGroupEntity createRemoteGroupEntity(final RemoteProcessGroup rpg) { + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(rpg.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(rpg); + final RemoteProcessGroupStatusDTO status = dtoFactory.createRemoteProcessGroupStatusDto(controllerFacade.getRemoteProcessGroupStatus(rpg.getIdentifier())); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(rpg.getIdentifier())); + return entityFactory.createRemoteProcessGroupEntity(dtoFactory.createRemoteProcessGroupDto(rpg), revision, accessPolicy, status, bulletins); } @Override public Set<RemoteProcessGroupEntity> getRemoteProcessGroups(final String groupId) { final Set<RemoteProcessGroup> rpgs = remoteProcessGroupDAO.getRemoteProcessGroups(groupId); - final Set<String> ids = rpgs.stream().map(rpg -> rpg.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(ids, () -> { - return rpgs.stream() - .map(rpg -> { - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(rpg.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(rpg); - final RemoteProcessGroupStatusDTO status = dtoFactory.createRemoteProcessGroupStatusDto(controllerFacade.getRemoteProcessGroupStatus(rpg.getIdentifier())); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(rpg.getIdentifier())); - return entityFactory.createRemoteProcessGroupEntity(dtoFactory.createRemoteProcessGroupDto(rpg), revision, accessPolicy, status, bulletins); - }) - .collect(Collectors.toSet()); - }); + return rpgs.stream() + .map(rpg -> createRemoteGroupEntity(rpg)) + .collect(Collectors.toSet()); } @Override public PortEntity getInputPort(final String inputPortId) { - return revisionManager.get(inputPortId, rev -> { - final Port port = inputPortDAO.getPort(inputPortId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(port); - final PortStatusDTO status = dtoFactory.createPortStatusDto(controllerFacade.getInputPortStatus(inputPortId)); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(inputPortId)); - return entityFactory.createPortEntity(dtoFactory.createPortDto(port), revision, accessPolicy, status, bulletins); - }); + final Port port = inputPortDAO.getPort(inputPortId); + return createInputPortEntity(port); } @Override public PortStatusDTO getInputPortStatus(final String inputPortId) { - return revisionManager.get(inputPortId, rev -> dtoFactory.createPortStatusDto(controllerFacade.getInputPortStatus(inputPortId))); + return dtoFactory.createPortStatusDto(controllerFacade.getInputPortStatus(inputPortId)); } @Override public PortEntity getOutputPort(final String outputPortId) { - return revisionManager.get(outputPortId, rev -> { - final Port port = outputPortDAO.getPort(outputPortId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(port); - final PortStatusDTO status = dtoFactory.createPortStatusDto(controllerFacade.getOutputPortStatus(outputPortId)); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(outputPortId)); - return entityFactory.createPortEntity(dtoFactory.createPortDto(port), revision, accessPolicy, status, bulletins); - }); + final Port port = outputPortDAO.getPort(outputPortId); + return createOutputPortEntity(port); } @Override public PortStatusDTO getOutputPortStatus(final String outputPortId) { - return revisionManager.get(outputPortId, rev -> dtoFactory.createPortStatusDto(controllerFacade.getOutputPortStatus(outputPortId))); + return dtoFactory.createPortStatusDto(controllerFacade.getOutputPortStatus(outputPortId)); } @Override public RemoteProcessGroupEntity getRemoteProcessGroup(final String remoteProcessGroupId) { - return revisionManager.get(remoteProcessGroupId, rev -> { - final RemoteProcessGroup rpg = remoteProcessGroupDAO.getRemoteProcessGroup(remoteProcessGroupId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(rpg); - final RemoteProcessGroupStatusDTO status = dtoFactory.createRemoteProcessGroupStatusDto(controllerFacade.getRemoteProcessGroupStatus(rpg.getIdentifier())); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(rpg.getIdentifier())); - return entityFactory.createRemoteProcessGroupEntity(dtoFactory.createRemoteProcessGroupDto(rpg), revision, accessPolicy, status, bulletins); - }); + final RemoteProcessGroup rpg = remoteProcessGroupDAO.getRemoteProcessGroup(remoteProcessGroupId); + return createRemoteGroupEntity(rpg); } @Override public RemoteProcessGroupStatusDTO getRemoteProcessGroupStatus(final String id) { - return revisionManager.get(id, rev -> dtoFactory.createRemoteProcessGroupStatusDto(controllerFacade.getRemoteProcessGroupStatus(id))); + return dtoFactory.createRemoteProcessGroupStatusDto(controllerFacade.getRemoteProcessGroupStatus(id)); } @Override @@ -2620,58 +2435,59 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { @Override public ProcessGroupFlowEntity getProcessGroupFlow(final String groupId, final boolean recurse) { - return revisionManager.get(groupId, - rev -> { - // get all identifiers for every child component - final Set<String> identifiers = new HashSet<>(); - final ProcessGroup processGroup = processGroupDAO.getProcessGroup(groupId); - processGroup.getProcessors().stream() - .map(proc -> proc.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getConnections().stream() - .map(conn -> conn.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getInputPorts().stream() - .map(port -> port.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getOutputPorts().stream() - .map(port -> port.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getProcessGroups().stream() - .map(group -> group.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getRemoteProcessGroups().stream() - .map(remoteGroup -> remoteGroup.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getRemoteProcessGroups().stream() - .flatMap(remoteGroup -> remoteGroup.getInputPorts().stream()) - .map(remoteInputPort -> remoteInputPort.getIdentifier()) - .forEach(id -> identifiers.add(id)); - processGroup.getRemoteProcessGroups().stream() - .flatMap(remoteGroup -> remoteGroup.getOutputPorts().stream()) - .map(remoteOutputPort -> remoteOutputPort.getIdentifier()) - .forEach(id -> identifiers.add(id)); - - // read lock on every component being accessed in the dto conversion - return revisionManager.get(identifiers, - () -> { - final ProcessGroupStatus groupStatus = controllerFacade.getProcessGroupStatus(groupId); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(processGroup); - return entityFactory.createProcessGroupFlowEntity(dtoFactory.createProcessGroupFlowDto(processGroup, groupStatus, revisionManager), accessPolicy); - }); - }); + // get all identifiers for every child component + final Set<String> identifiers = new HashSet<>(); + final ProcessGroup processGroup = processGroupDAO.getProcessGroup(groupId); + processGroup.getProcessors().stream() + .map(proc -> proc.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getConnections().stream() + .map(conn -> conn.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getInputPorts().stream() + .map(port -> port.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getOutputPorts().stream() + .map(port -> port.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getProcessGroups().stream() + .map(group -> group.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getRemoteProcessGroups().stream() + .map(remoteGroup -> remoteGroup.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getRemoteProcessGroups().stream() + .flatMap(remoteGroup -> remoteGroup.getInputPorts().stream()) + .map(remoteInputPort -> remoteInputPort.getIdentifier()) + .forEach(id -> identifiers.add(id)); + processGroup.getRemoteProcessGroups().stream() + .flatMap(remoteGroup -> remoteGroup.getOutputPorts().stream()) + .map(remoteOutputPort -> remoteOutputPort.getIdentifier()) + .forEach(id -> identifiers.add(id)); + + // read lock on every component being accessed in the dto conversion + final ProcessGroupStatus groupStatus = controllerFacade.getProcessGroupStatus(groupId); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(processGroup); + return entityFactory.createProcessGroupFlowEntity(dtoFactory.createProcessGroupFlowDto(processGroup, groupStatus, revisionManager), accessPolicy); } @Override public ProcessGroupEntity getProcessGroup(final String groupId) { - return revisionManager.get(groupId, rev -> { - final ProcessGroup processGroup = processGroupDAO.getProcessGroup(groupId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(processGroup); - final ProcessGroupStatusDTO status = dtoFactory.createConciseProcessGroupStatusDto(controllerFacade.getProcessGroupStatus(groupId)); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(groupId)); - return entityFactory.createProcessGroupEntity(dtoFactory.createProcessGroupDto(processGroup), revision, accessPolicy, status, bulletins); - }); + final ProcessGroup processGroup = processGroupDAO.getProcessGroup(groupId); + return createProcessGroupEntity(processGroup); + } + + private ControllerServiceEntity createControllerServiceEntity(final ControllerServiceNode serviceNode, final Set<String> serviceIds) { + final ControllerServiceDTO dto = dtoFactory.createControllerServiceDto(serviceNode); + + final ControllerServiceReference ref = serviceNode.getReferences(); + final ControllerServiceReferencingComponentsEntity referencingComponentsEntity = createControllerServiceReferencingComponentsEntity(ref, serviceIds); + dto.setReferencingComponents(referencingComponentsEntity.getControllerServiceReferencingComponents()); + + final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(serviceNode.getIdentifier())); + final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(serviceNode); + final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(serviceNode.getIdentifier())); + return entityFactory.createControllerServiceEntity(dto, revision, accessPolicy, bulletins); } @Override @@ -2679,93 +2495,57 @@ public class StandardNiFiServiceFacade implements NiFiServiceFacade { final Set<ControllerServiceNode> serviceNodes = controllerServiceDAO.getControllerServices(groupId); final Set<String> serviceIds = serviceNodes.stream().map(service -> service.getIdentifier()).collect(Collectors.toSet()); - return revisionManager.get(serviceIds, () -> { - return serviceNodes.stream() - .map(serviceNode -> { - final ControllerServiceDTO dto = dtoFactory.createControllerServiceDto(serviceNode); - - final ControllerServiceReference ref = serviceNode.getReferences(); - final ControllerServiceReferencingComponentsEntity referencingComponentsEntity = createControllerServiceReferencingComponentsEntity(ref, serviceIds); - dto.setReferencingComponents(referencingComponentsEntity.getControllerServiceReferencingComponents()); - - final RevisionDTO revision = dtoFactory.createRevisionDTO(revisionManager.getRevision(serviceNode.getIdentifier())); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(serviceNode); - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(serviceNode.getIdentifier())); - return entityFactory.createControllerServiceEntity(dto, revision, accessPolicy, bulletins); - }) - .collect(Collectors.toSet()); - }); + return serviceNodes.stream() + .map(serviceNode -> createControllerServiceEntity(serviceNode, serviceIds)) + .collect(Collectors.toSet()); } @Override public ControllerServiceEntity getControllerService(final String controllerServiceId) { - return revisionManager.get(controllerServiceId, rev -> { - final ControllerServiceNode controllerService = controllerServiceDAO.getControllerService(controllerServiceId); - final RevisionDTO revision = dtoFactory.createRevisionDTO(rev); - final AccessPolicyDTO accessPolicy = dtoFactory.createAccessPolicyDto(controllerService); - final ControllerServiceDTO dto = dtoFactory.createControllerServiceDto(controllerService); - - final ControllerServiceReference ref = controllerService.getReferences(); - final ControllerServiceReferencingComponentsEntity referencingComponentsEntity = createControllerServiceReferencingComponentsEntity(ref, Sets.newHashSet(controllerServiceId)); - dto.setReferencingComponents(referencingComponentsEntity.getControllerServiceReferencingComponents()); - - final List<BulletinDTO> bulletins = dtoFactory.createBulletinDtos(bulletinRepository.findBulletinsForSource(controllerServiceId)); - - return entityFactory.createControllerServiceEntity(dto, revision, accessPolicy, bulletins); - }); + final ControllerServiceNode controllerService = controllerServiceDAO.getControllerService(controllerServiceId); + return createControllerServiceEntity(controllerService, Sets.newHashSet(controllerServiceId)); } @Override public PropertyDescriptorDTO getControllerServicePropertyDescriptor(final String id, final String property) { - return revisionManager.get(id, rev -> { - final ControllerServiceNode controllerService = controllerServiceDAO.getControllerService(id); - PropertyDescriptor descriptor = controllerServ
<TRUNCATED>
