Loading src/main/java/com/inteligr8/activiti/mq/AbstractMqDelegate.java +1 −1 Original line number Diff line number Diff line Loading @@ -33,7 +33,7 @@ public abstract class AbstractMqDelegate implements JavaDelegate { private Map<String, MqCommunicator> communicators = new HashMap<>(); protected synchronized MqCommunicator getConnection(String connectorId) throws JMSException, TimeoutException, IOException { protected synchronized MqCommunicator getCommunicator(String connectorId) throws JMSException, TimeoutException, IOException { MqCommunicator communicator = this.communicators.get(connectorId); if (communicator == null) { EndpointConfiguration endpointConfig = this.endpointService.getConfigurationByName(connectorId); Loading src/main/java/com/inteligr8/activiti/mq/MQProcessDefinitionMonitor.java +29 −1 Original line number Diff line number Diff line Loading @@ -5,6 +5,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; import java.util.regex.Matcher; import java.util.regex.Pattern; Loading Loading @@ -63,6 +64,9 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic @Autowired private TenantFinderService tenantFinderService; @Autowired private MqSubscriptionService subscriptionService; private Map<String, AbstractActivityListener> activeListeners = new HashMap<>(); /** Loading Loading @@ -236,6 +240,8 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic protected void onProcessDefinitionAddEvent(ProcessDefinitionEntity entity) { this.logger.debug("Triggered by process definition addition: {}", entity); this.unsubscribeMqSubscribeTasks(entity.getId()); ServiceTask task = this.findMqStartSubscribeTask(entity.getId()); if (task == null) return; Loading @@ -254,6 +260,8 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic for (ProcessDefinitionEntity procDefEntity : procDefEntities) { this.logger.debug("Inspecting process definition: {}: {}: {}", procDefEntity.getId(), procDefEntity.getKey(), procDefEntity.getName()); this.unsubscribeMqSubscribeTasks(procDefEntity.getId()); ServiceTask task = this.findMqStartSubscribeTask(procDefEntity.getId()); if (task == null) return; Loading @@ -264,9 +272,29 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic protected void onProcessDefinitionRemoveEvent(ProcessDefinitionEntity entity) { this.logger.debug("Triggered by process definition removal: {}", entity); this.unsubscribeMqSubscribeTasks(entity.getId()); ServiceTask task = this.findMqStartSubscribeTask(entity.getId()); if (task == null) return; this.deloopMqSubscribeTask(entity.getId()); } protected void unsubscribeMqSubscribeTasks(String procDefId) { try { Set<String> executionIds = this.subscriptionService.clear(procDefId); if (this.logger.isDebugEnabled()) { this.logger.debug("Subscription executions ended early: {}: {}", procDefId, executionIds); } else { this.logger.info("Subscriptions ended early: {}: {}", procDefId, executionIds.size()); } } catch (Exception e) { this.logger.error("The subscriptions could not be cleared: " + procDefId, e); } } protected synchronized void loopMqSubscribeTask(String processDefId, ServiceTask task) { // start a process instance on the process this.logger.debug("Starting process instance on process '{}' to subscribe to an MQ queue", processDefId); Loading src/main/java/com/inteligr8/activiti/mq/MqCommunicator.java +27 −11 Original line number Diff line number Diff line Loading @@ -14,33 +14,49 @@ public interface MqCommunicator { <BodyType> String send(GenericDestination destination, PreparedMessage<BodyType> message) throws JMSException, IOException, TimeoutException; default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, null); return this.receive(destination, -1L, null, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, null); return this.receive(destination, timeoutInMillis, null, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, null); return this.receive(destination, -1L, correlationId, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, correlationId, null); return this.receive(destination, timeoutInMillis, correlationId, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, correlationId, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, handler); default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, listener, handler); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, handler); default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, listener, handler); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, handler); default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, listener, handler); } <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException; <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException; } src/main/java/com/inteligr8/activiti/mq/MqExecutionService.java 0 → 100644 +101 −0 Original line number Diff line number Diff line package com.inteligr8.activiti.mq; import java.util.Collection; import java.util.Collections; import java.util.HashSet; import java.util.Set; import org.activiti.engine.delegate.DelegateExecution; import org.apache.commons.collections4.MultiValuedMap; import org.apache.commons.collections4.multimap.HashSetValuedHashMap; import org.apache.commons.lang3.tuple.Pair; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; @Component public class MqExecutionService { private final Logger logger = LoggerFactory.getLogger(this.getClass()); /** * The size of the keys is limited to the number of process definitions * defined. It would actually only contain ones with an MQ subscribe task. * Even if we kept versioned or inactive process definitions in the map, it * would never be a significant memory hog. * * The size of the values is limited to the number of MQ subscribe tasks * defined in each process definition. So it would never be a significant * memory hog. * * The size of the keys/values have nothing to do with the number of * process instances or executions. * * This means it does not need to be trimmed. However, it is a good idea * to remove process definition keys that have no active executions. You * could do the same with active activities, but cleaning up process * definitions will clean those up too. */ private MultiValuedMap<String, String> processDefinitionActivityMap = new HashSetValuedHashMap<>(); /** * The size of the keys is limited to the number of MQ subscribe tasks * defined in all process definitions. So it would never be a significant * memory hog. * * The size of the values has no limit. It will grow with the number of * executions (related to process instances). * * This means the map values need to be trimmed. When an MQ subscribe task * is completed, it is paramount to remove the execution from the values of * this map. It is also a good idea to remove the activity key when it is * removed from the `processDefinitionActivityMap` map; and to propagate * the removal of executions from the `executionSubscriptionMap` map. */ private MultiValuedMap<Pair<String, String>, String> activityExecutionMap = new HashSetValuedHashMap<>(); public synchronized void executing(DelegateExecution execution) { this.processDefinitionActivityMap.put(execution.getProcessDefinitionId(), execution.getCurrentActivityId()); Pair<String, String> key = this.toKey(execution); this.activityExecutionMap.put(key, execution.getId()); } public synchronized void executed(DelegateExecution execution) { Pair<String, String> key = this.toKey(execution); this.activityExecutionMap.removeMapping(key, execution.getId()); } /** * @param processDefinitionId A process definition identifier. * @return A set of execution identifiers that were in the now cleared map. */ public synchronized Set<String> clear(String processDefinitionId) throws Exception { Collection<String> activityIds = this.processDefinitionActivityMap.get(processDefinitionId); if (activityIds == null) { this.logger.debug("No activities/executions to clear for process definition: {}", processDefinitionId); return Collections.emptySet(); } Set<String> executionIds = new HashSet<>(); for (String activityId : activityIds) { this.logger.trace("Clearing process definition activity: {}: {}", processDefinitionId, activityId); Pair<String, String> key = this.toKey(processDefinitionId, activityId); Collection<String> activityExecutionIds = this.activityExecutionMap.get(key); if (activityExecutionIds != null) executionIds.addAll(activityExecutionIds); } return executionIds; } protected Pair<String, String> toKey(DelegateExecution execution) { return this.toKey(execution.getProcessDefinitionId(), execution.getCurrentActivityId()); } protected Pair<String, String> toKey(String processDefinitionId, String activityId) { return Pair.of(processDefinitionId, activityId); } } src/main/java/com/inteligr8/activiti/mq/MqPublishDelegate.java +1 −1 Original line number Diff line number Diff line Loading @@ -70,7 +70,7 @@ public class MqPublishDelegate extends AbstractMqDelegate { destination.setQueueName(mqExecution.getQueueNameFromModel()); try { MqCommunicator communicator = this.getConnection(mqExecution.getConnectorIdFromModel()); MqCommunicator communicator = this.getCommunicator(mqExecution.getConnectorIdFromModel()); PreparedMessage<String> message = communicator.createPreparedMessage(); if (mqExecution.getStatusQueueNameFromModel() != null) Loading Loading
src/main/java/com/inteligr8/activiti/mq/AbstractMqDelegate.java +1 −1 Original line number Diff line number Diff line Loading @@ -33,7 +33,7 @@ public abstract class AbstractMqDelegate implements JavaDelegate { private Map<String, MqCommunicator> communicators = new HashMap<>(); protected synchronized MqCommunicator getConnection(String connectorId) throws JMSException, TimeoutException, IOException { protected synchronized MqCommunicator getCommunicator(String connectorId) throws JMSException, TimeoutException, IOException { MqCommunicator communicator = this.communicators.get(connectorId); if (communicator == null) { EndpointConfiguration endpointConfig = this.endpointService.getConfigurationByName(connectorId); Loading
src/main/java/com/inteligr8/activiti/mq/MQProcessDefinitionMonitor.java +29 −1 Original line number Diff line number Diff line Loading @@ -5,6 +5,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; import java.util.regex.Matcher; import java.util.regex.Pattern; Loading Loading @@ -63,6 +64,9 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic @Autowired private TenantFinderService tenantFinderService; @Autowired private MqSubscriptionService subscriptionService; private Map<String, AbstractActivityListener> activeListeners = new HashMap<>(); /** Loading Loading @@ -236,6 +240,8 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic protected void onProcessDefinitionAddEvent(ProcessDefinitionEntity entity) { this.logger.debug("Triggered by process definition addition: {}", entity); this.unsubscribeMqSubscribeTasks(entity.getId()); ServiceTask task = this.findMqStartSubscribeTask(entity.getId()); if (task == null) return; Loading @@ -254,6 +260,8 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic for (ProcessDefinitionEntity procDefEntity : procDefEntities) { this.logger.debug("Inspecting process definition: {}: {}: {}", procDefEntity.getId(), procDefEntity.getKey(), procDefEntity.getName()); this.unsubscribeMqSubscribeTasks(procDefEntity.getId()); ServiceTask task = this.findMqStartSubscribeTask(procDefEntity.getId()); if (task == null) return; Loading @@ -264,9 +272,29 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic protected void onProcessDefinitionRemoveEvent(ProcessDefinitionEntity entity) { this.logger.debug("Triggered by process definition removal: {}", entity); this.unsubscribeMqSubscribeTasks(entity.getId()); ServiceTask task = this.findMqStartSubscribeTask(entity.getId()); if (task == null) return; this.deloopMqSubscribeTask(entity.getId()); } protected void unsubscribeMqSubscribeTasks(String procDefId) { try { Set<String> executionIds = this.subscriptionService.clear(procDefId); if (this.logger.isDebugEnabled()) { this.logger.debug("Subscription executions ended early: {}: {}", procDefId, executionIds); } else { this.logger.info("Subscriptions ended early: {}: {}", procDefId, executionIds.size()); } } catch (Exception e) { this.logger.error("The subscriptions could not be cleared: " + procDefId, e); } } protected synchronized void loopMqSubscribeTask(String processDefId, ServiceTask task) { // start a process instance on the process this.logger.debug("Starting process instance on process '{}' to subscribe to an MQ queue", processDefId); Loading
src/main/java/com/inteligr8/activiti/mq/MqCommunicator.java +27 −11 Original line number Diff line number Diff line Loading @@ -14,33 +14,49 @@ public interface MqCommunicator { <BodyType> String send(GenericDestination destination, PreparedMessage<BodyType> message) throws JMSException, IOException, TimeoutException; default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, null); return this.receive(destination, -1L, null, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, null); return this.receive(destination, timeoutInMillis, null, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, null); return this.receive(destination, -1L, correlationId, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, correlationId, null); return this.receive(destination, timeoutInMillis, correlationId, null, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId, MqSubscriptionListener listener) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, correlationId, listener, null); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, handler); default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, null, listener, handler); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, handler); default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, timeoutInMillis, null, listener, handler); } default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, handler); default <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, String correlationId, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException { return this.receive(destination, -1L, correlationId, listener, handler); } <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException; <BodyType> DeliveredMessage<BodyType> receive(GenericDestination destination, long timeoutInMillis, String correlationId, MqSubscriptionListener listener, TransactionalMessageHandler<BodyType> handler) throws JMSException, IOException, TimeoutException; }
src/main/java/com/inteligr8/activiti/mq/MqExecutionService.java 0 → 100644 +101 −0 Original line number Diff line number Diff line package com.inteligr8.activiti.mq; import java.util.Collection; import java.util.Collections; import java.util.HashSet; import java.util.Set; import org.activiti.engine.delegate.DelegateExecution; import org.apache.commons.collections4.MultiValuedMap; import org.apache.commons.collections4.multimap.HashSetValuedHashMap; import org.apache.commons.lang3.tuple.Pair; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; @Component public class MqExecutionService { private final Logger logger = LoggerFactory.getLogger(this.getClass()); /** * The size of the keys is limited to the number of process definitions * defined. It would actually only contain ones with an MQ subscribe task. * Even if we kept versioned or inactive process definitions in the map, it * would never be a significant memory hog. * * The size of the values is limited to the number of MQ subscribe tasks * defined in each process definition. So it would never be a significant * memory hog. * * The size of the keys/values have nothing to do with the number of * process instances or executions. * * This means it does not need to be trimmed. However, it is a good idea * to remove process definition keys that have no active executions. You * could do the same with active activities, but cleaning up process * definitions will clean those up too. */ private MultiValuedMap<String, String> processDefinitionActivityMap = new HashSetValuedHashMap<>(); /** * The size of the keys is limited to the number of MQ subscribe tasks * defined in all process definitions. So it would never be a significant * memory hog. * * The size of the values has no limit. It will grow with the number of * executions (related to process instances). * * This means the map values need to be trimmed. When an MQ subscribe task * is completed, it is paramount to remove the execution from the values of * this map. It is also a good idea to remove the activity key when it is * removed from the `processDefinitionActivityMap` map; and to propagate * the removal of executions from the `executionSubscriptionMap` map. */ private MultiValuedMap<Pair<String, String>, String> activityExecutionMap = new HashSetValuedHashMap<>(); public synchronized void executing(DelegateExecution execution) { this.processDefinitionActivityMap.put(execution.getProcessDefinitionId(), execution.getCurrentActivityId()); Pair<String, String> key = this.toKey(execution); this.activityExecutionMap.put(key, execution.getId()); } public synchronized void executed(DelegateExecution execution) { Pair<String, String> key = this.toKey(execution); this.activityExecutionMap.removeMapping(key, execution.getId()); } /** * @param processDefinitionId A process definition identifier. * @return A set of execution identifiers that were in the now cleared map. */ public synchronized Set<String> clear(String processDefinitionId) throws Exception { Collection<String> activityIds = this.processDefinitionActivityMap.get(processDefinitionId); if (activityIds == null) { this.logger.debug("No activities/executions to clear for process definition: {}", processDefinitionId); return Collections.emptySet(); } Set<String> executionIds = new HashSet<>(); for (String activityId : activityIds) { this.logger.trace("Clearing process definition activity: {}: {}", processDefinitionId, activityId); Pair<String, String> key = this.toKey(processDefinitionId, activityId); Collection<String> activityExecutionIds = this.activityExecutionMap.get(key); if (activityExecutionIds != null) executionIds.addAll(activityExecutionIds); } return executionIds; } protected Pair<String, String> toKey(DelegateExecution execution) { return this.toKey(execution.getProcessDefinitionId(), execution.getCurrentActivityId()); } protected Pair<String, String> toKey(String processDefinitionId, String activityId) { return Pair.of(processDefinitionId, activityId); } }
src/main/java/com/inteligr8/activiti/mq/MqPublishDelegate.java +1 −1 Original line number Diff line number Diff line Loading @@ -70,7 +70,7 @@ public class MqPublishDelegate extends AbstractMqDelegate { destination.setQueueName(mqExecution.getQueueNameFromModel()); try { MqCommunicator communicator = this.getConnection(mqExecution.getConnectorIdFromModel()); MqCommunicator communicator = this.getCommunicator(mqExecution.getConnectorIdFromModel()); PreparedMessage<String> message = communicator.createPreparedMessage(); if (mqExecution.getStatusQueueNameFromModel() != null) Loading