Commit 8d91166a authored by Brian Long's avatar Brian Long
Browse files

Merge branch 'develop' into stable

parents 370b385a 4fe113b8
Loading
Loading
Loading
Loading
+38 −2
Original line number Diff line number Diff line
@@ -1186,6 +1186,40 @@
                }
            ]
        },
        {
            "name": "mq_prioritypackage",
            "properties": [
                {
                    "id": "mq_priority",
                    "type": "String",
                    "title": "MQ Priority",
                    "value": "",
                    "description": "MQ message priority; depends on protocol specification",
                    "popular": false,
                    "custom": {
                        "includeInXML": true,
                        "xmlPropertyName": "mq_priority"
                    }
                }
            ]
        },
        {
            "name": "mq_concurrencypackage",
            "properties": [
                {
                    "id": "mq_concurrency",
                    "type": "String",
                    "title": "MQ Concurrency (positive number)",
                    "value": "",
                    "description": "Number of MQ subscription threads",
                    "popular": false,
                    "custom": {
                        "includeInXML": true,
                        "xmlPropertyName": "mq_concurrency"
                    }
                }
            ]
        },
        {
            "name": "mq_payloadpackage",
            "properties": [
@@ -3050,10 +3084,11 @@
                "mq_connectorIdpackage",
                "mq_queueNamepackage",
                "mq_messageNamepackage",
                "mq_prioritypackage",
                "mq_payloadpackage",
                "mq_replyQueueNamepackage",
                "mq_statusQueueNamepackage",
                "mq_metadataProcessScopepackage",
                "mq_payloadpackage"
                "mq_metadataProcessScopepackage"
            ],
            "hiddenPropertyPackages": [
                "multiinstance_typepackage",
@@ -3153,6 +3188,7 @@
                "mq_connectorIdpackage",
                "mq_queueNamepackage",
                "mq_messageNamepackage",
                "mq_concurrencypackage",
                "mq_metadataProcessScopepackage"
            ],
            "hiddenPropertyPackages": [
+2 −0
Original line number Diff line number Diff line
@@ -9,6 +9,8 @@ public class Constants {
	public static final String FIELD_CONNECTOR_ID = "mq_connectorId";
	public static final String FIELD_QUEUE_NAME = "mq_queueName";
	public static final String FIELD_MESSAGE_NAME = "mq_messageName";
	public static final String FIELD_CONCURRENCY = "mq_concurrency";
	public static final String FIELD_PRIORITY = "mq_priority";
	public static final String FIELD_REPLY_QUEUE_NAME = "mq_replyQueueName";
	public static final String FIELD_STATUS_QUEUE_NAME = "mq_statusQueueName";
	public static final String FIELD_PAYLOAD = "mq_payload";
+76 −17
Original line number Diff line number Diff line
@@ -12,6 +12,7 @@ import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;

import org.activiti.bpmn.model.BpmnModel;
import org.activiti.bpmn.model.FieldExtension;
import org.activiti.bpmn.model.FlowElement;
import org.activiti.bpmn.model.SequenceFlow;
import org.activiti.bpmn.model.ServiceTask;
@@ -30,6 +31,9 @@ import org.activiti.engine.impl.persistence.entity.ProcessDefinitionEntity;
import org.activiti.engine.impl.util.ProcessDefinitionUtil;
import org.activiti.engine.repository.ProcessDefinition;
import org.activiti.engine.repository.ProcessDefinitionQuery;
import org.activiti.engine.runtime.Execution;
import org.activiti.engine.runtime.ExecutionQuery;
import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
@@ -74,27 +78,83 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic
		}
		
		if (event instanceof ContextRefreshedEvent || event instanceof ContextStartedEvent) {
			String tenantId = this.findTenantId();
			List<ProcessDefinition> procDefs = this.findLatestActiveProcessDefinnitions(tenantId);
			this.logger.debug("Found {} active process definitions", procDefs.size());
			for (ProcessDefinition procDef : procDefs) {
				this.logger.trace("Inspecting process definition for qualifying MQ subscriptions: {}", procDef.getId());
				
				ServiceTask task = this.findMqStartSubscribeTask(procDef.getId());
				if (task == null)
					continue;
				
				int concurrency = this.determineConcurrency(task);
				this.logger.debug("Process definition MQ subscription is configured for concurrency: {}: {}", procDef.getId(), concurrency);

				List<Execution> execs = this.findExecutionsByServiceTask(tenantId, procDef.getId(), task);
				this.logger.debug("Process appears to have {} executions waiting on the MQ subscription: {}", execs.size(), procDef.getId());
				
				if (execs.size() < concurrency) {
					this.logger.info("Process has {} too few executions waiting on the MQ subscription; starting them: {}", (concurrency - execs.size()), procDef.getId());
					for (int thread = execs.size(); thread < concurrency; thread++) {
						this.loopMqSubscribeTask(procDef.getId(), task);
					}
				}
			}
		}
	}
	
	protected String findTenantId() {
		Tenant tenant = this.tenantFinderService.findTenant();
		return this.tenantFinderService.transform(tenant);
	}
	
			ProcessDefinitionQuery procDefQuery = this.services.getRepositoryService().createProcessDefinitionQuery().active();
			if (tenant == null) {
	protected List<ProcessDefinition> findLatestActiveProcessDefinnitions(String tenantId) {
		ProcessDefinitionQuery procDefQuery = this.services.getRepositoryService().createProcessDefinitionQuery()
				.latestVersion()
				.active();
		if (tenantId == null) {
			procDefQuery.processDefinitionWithoutTenantId();
		} else {
				String tenantId = this.tenantFinderService.transform(tenant.getId());
			procDefQuery.processDefinitionTenantId(tenantId);
		}
		
			List<ProcessDefinition> procDefs = procDefQuery.list();
			this.logger.debug("Found {} active process definitions", procDefs.size());
			for (ProcessDefinition procDef : procDefs) {
				this.logger.debug("Inspecting process definition for qualifying MQ subscriptions: {}", procDef.getId());
		return procDefQuery.list();
	}
	
				ServiceTask task = this.findMqStartSubscribeTask(procDef.getId());
				if (task == null)
					continue;
	protected List<Execution> findExecutionsByServiceTask(String tenantId, String processDefinitionId, ServiceTask task) {
		ExecutionQuery execQuery = this.services.getRuntimeService().createExecutionQuery()
				.processDefinitionId(processDefinitionId)
				.activityId(task.getId());
		if (tenantId == null) {
			execQuery.executionWithoutTenantId();
		} else {
			execQuery.executionTenantId(tenantId);
		}
		
				this.loopMqSubscribeTask(procDef.getId(), task);
		return execQuery.list();
	}
	
	protected Integer findConcurrency(ServiceTask task) {
		for (FieldExtension fieldext : task.getFieldExtensions()) {
			if (fieldext.getFieldName().equals(Constants.FIELD_CONCURRENCY)) {
				String concurrencyStr = StringUtils.trimToNull(fieldext.getStringValue());
				return concurrencyStr == null ? null : Integer.valueOf(concurrencyStr);
			}
		}
		
		return null;
	}
	
	protected int determineConcurrency(ServiceTask task) {
		Integer concurrency = this.findConcurrency(task);
		if (concurrency == null) {
			return 1;
		} else if (concurrency.intValue() < 1) {
			this.logger.warn("The task defines an illegal concurrency of {}; using 1: {}", concurrency, task.getId());
			return 1;
		} else {
			return concurrency.intValue();
		}
	}
	
@@ -212,10 +272,9 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic
		this.logger.debug("Starting process instance on process '{}' to subscribe to an MQ queue", processDefId);
		this.services.getRuntimeService().startProcessInstanceById(processDefId);
		
		if (this.activeListeners.containsKey(processDefId)) {
			this.logger.debug("The process definition already has a looping listener: {}", processDefId);
		if (this.activeListeners.containsKey(processDefId))
			// one listener, no matter how many instances are subscribed
			return;
		}
		
		AbstractActivityListener listener = new AbstractActivityListener(processDefId, task) {
			@Override
@@ -334,7 +393,7 @@ public class MQProcessDefinitionMonitor implements ActivitiEventListener, Applic
			return null;
		}
		
		this.logger.debug("Process starts with an MQ subscription: {}: {}", processDefId, task.getId(), task.getName());
		this.logger.info("Process starts with an MQ subscription: {}: {}", processDefId, task.getId(), task.getName());
		return task;
	}
	
+14 −71
Original line number Diff line number Diff line
@@ -2,92 +2,35 @@ package com.inteligr8.activiti.mq;

import java.time.OffsetDateTime;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;

import javax.annotation.OverridingMethodsMustInvokeSuper;

import org.activiti.bpmn.model.FieldExtension;
import org.activiti.bpmn.model.FlowElement;
import org.activiti.bpmn.model.ServiceTask;
import org.activiti.engine.ProcessEngine;
import org.activiti.engine.delegate.DelegateExecution;
import org.activiti.engine.delegate.Expression;
import org.activiti.engine.impl.context.Context;
import org.activiti.engine.impl.el.ExpressionManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class MqDelegateExecution {
	
	private final Logger logger = LoggerFactory.getLogger(this.getClass());
	
	protected final ProcessEngine services;
	protected final DelegateExecution execution;
	protected final ServiceTask task;
	protected final Map<String, FieldExtension> fieldMap = new HashMap<>();
	private boolean metadataToProcessScope = false;
	protected final MqServiceTask task;
	
	public MqDelegateExecution(ProcessEngine services, DelegateExecution execution) {
		this.services = services;
	public MqDelegateExecution(ProcessEngine services, MqServiceTaskService msts, DelegateExecution execution) {
		this.execution = execution;
		
		FlowElement flowElement = execution.getCurrentFlowElement();
		if (!(flowElement instanceof ServiceTask))
			throw new IllegalStateException("This should never happen");
		this.task = (ServiceTask) flowElement;
		this.logger.trace("Discovered task: {}: {}", this.task.getId(), this.task.getName());

		this.logger.trace("Indexing {} fields", this.task.getFieldExtensions().size());
		for (FieldExtension field : this.task.getFieldExtensions()) {
			this.logger.trace("Discovering field: {}: {}: {}", field.getId(), field.getFieldName(), field.getStringValue());
			
			switch (field.getFieldName()) {
				case Constants.FIELD_METADATA_PROCESS_SCOPE:
					this.metadataToProcessScope = Boolean.valueOf(field.getStringValue());
					break;
				default:
					this.fieldMap.put(field.getFieldName(), field);
			}
		}
		this.task = msts.get(execution.getProcessDefinitionId(), execution.getCurrentFlowElement());
	}
	
	@OverridingMethodsMustInvokeSuper
	public void validate() {
    	if (this.fieldMap.get(Constants.FIELD_CONNECTOR_ID) == null)
    		throw new IllegalStateException("The '" + this.execution.getCurrentActivityId() + "' activity must define an 'MQ Connector ID'");
    	if (this.fieldMap.get(Constants.FIELD_QUEUE_NAME) == null)
    		throw new IllegalStateException("The '" + this.execution.getCurrentActivityId() + "' activity must define an 'MQ Queue Name'");
	}
	
	public String getTaskId() {
		return this.task.getId();
	}
	
	public boolean doWriteToProcessScope() {
		return this.metadataToProcessScope;
	}

	/**
	 * Unlike MqServiceTask, this allows for expression expansion based on the variables
	 */
	protected <T> T getFieldValueFromModel(String fieldName, Class<T> type) {
		return this.getFieldValueFromModel(fieldName, type, false);
		return this.task.getFieldValueFromModel(fieldName, type, this.execution, false);
	}
	
	@SuppressWarnings("unchecked")
	/**
	 * Unlike MqServiceTask, this allows for expression expansion based on the variables; even partial values
	 */
	protected <T> T getFieldValueFromModel(String fieldName, Class<T> type, boolean forceExpressionProcessing) {
		FieldExtension field = this.fieldMap.get(fieldName);
		if (field == null) {
			return null;
		} else if (field.getExpression() != null && field.getExpression().length() > 0) {
			ExpressionManager exprman = Context.getProcessEngineConfiguration().getExpressionManager();
			Expression expr = exprman.createExpression(field.getExpression());
			return (T) expr.getValue(this.execution);
		} else if (forceExpressionProcessing) {
			ExpressionManager exprman = Context.getProcessEngineConfiguration().getExpressionManager();
			Expression expr = exprman.createExpression(field.getStringValue());
			return (T) expr.getValue(this.execution);
		} else {
			return (T) field.getStringValue();
		}
		return this.task.getFieldValueFromModel(fieldName, type, this.execution, forceExpressionProcessing);
	}
	
	public String getConnectorIdFromModel() {
@@ -112,7 +55,7 @@ public class MqDelegateExecution {
    
    protected void setMqVariable(String varName, Object value) {
    	varName = this.formulateVariableName(varName);
    	if (this.doWriteToProcessScope()) {
    	if (this.task.doWriteToProcessScope()) {
    		this.execution.setVariable(varName, value);
    	} else {
    		this.execution.setVariableLocal(varName, value);
+6 −1
Original line number Diff line number Diff line
@@ -30,6 +30,9 @@ public class MqPublishDelegate extends AbstractMqDelegate {
    @Autowired
    private ProcessEngine services;
    
    @Autowired
    private MqServiceTaskService msts;
    
    /**
     * This method sends a message to an MQ queue.
	 * 
@@ -46,6 +49,7 @@ public class MqPublishDelegate extends AbstractMqDelegate {
	 * @field mq_connectorId An Activiti App Tenant Endpoint ID or Java system property in the format `inteligr8.mq.connectors.{connectorId}.url`.  Using system properties support `url`, `username`, and `password`.
	 * @field mq_queueName The name of an MQ destination (queue or topic).  This is the target of the message.  If it doesn't exist, it will be created.
	 * @field mq_messageName [optional] A unique identifier to append to Activiti variables when more than one MQ message publication/subscription is supported by a process.
	 * @field mq_priority [optional] A priority of the MQ message.  May be an expression.  Value depends on MQ protocol.
	 * @field mq_payload [optional] The body of the MQ message.  May include expressions.
	 * @field mq_payloadMimeType [optional] The MIME type of the body of the MQ message.
	 * @field mq_replyQueueName [optional] The name of an MQ destination (queue or topic).  This tells the processor of the message where to send a reply.
@@ -59,7 +63,7 @@ public class MqPublishDelegate extends AbstractMqDelegate {
     */
    @Override
    public void execute(DelegateExecution execution) {
    	MqPublishDelegateExecution mqExecution = new MqPublishDelegateExecution(this.services, execution);
    	MqPublishDelegateExecution mqExecution = new MqPublishDelegateExecution(this.services, this.msts, execution);
    	mqExecution.validate();
    	
    	GenericDestination destination = new GenericDestination();
@@ -72,6 +76,7 @@ public class MqPublishDelegate extends AbstractMqDelegate {
        	if (mqExecution.getStatusQueueNameFromModel() != null)
        		message.setProperty("inteligr8.statusQueueName", mqExecution.getStatusQueueNameFromModel());
        	message.setReplyToQueueName(mqExecution.getReplyQueueNameFromModel());
        	message.setPriority(mqExecution.getPriorityFromModel());
        	
        	String payload = mqExecution.getPayloadFromModel();
        	if (payload != null) {
Loading