X-Git-Url: http://git.argeo.org/?a=blobdiff_plain;f=runtime%2Forg.argeo.slc.support.activemq%2Fsrc%2Fmain%2Fjava%2Forg%2Fargeo%2Fslc%2Fjms%2FJmsAgent.java;h=85bdb1451dab344d6f10e9053126758bad29f667;hb=87efa1cdb79eeaf3f203cc9bf4f3d9f8d0a299f8;hp=7490f32d1db4c4a89dab9b407fa46fca3b83b3e9;hpb=7a2f320afe9a0d3d7590365b26f3f5b0e8d9fd3b;p=gpl%2Fargeo-slc.git diff --git a/runtime/org.argeo.slc.support.activemq/src/main/java/org/argeo/slc/jms/JmsAgent.java b/runtime/org.argeo.slc.support.activemq/src/main/java/org/argeo/slc/jms/JmsAgent.java index 7490f32d1..85bdb1451 100644 --- a/runtime/org.argeo.slc.support.activemq/src/main/java/org/argeo/slc/jms/JmsAgent.java +++ b/runtime/org.argeo.slc.support.activemq/src/main/java/org/argeo/slc/jms/JmsAgent.java @@ -9,9 +9,7 @@ import java.util.UUID; import javax.jms.Destination; import javax.jms.JMSException; import javax.jms.Message; -import javax.jms.MessageProducer; -import javax.jms.Session; -import javax.jms.TextMessage; +import javax.jms.MessageListener; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -19,34 +17,32 @@ import org.argeo.slc.SlcException; import org.argeo.slc.core.runtime.AbstractAgent; import org.argeo.slc.execution.ExecutionModule; import org.argeo.slc.execution.ExecutionModuleDescriptor; +import org.argeo.slc.msg.ExecutionAnswer; +import org.argeo.slc.msg.ReferenceList; import org.argeo.slc.process.SlcExecution; import org.argeo.slc.runtime.SlcAgent; import org.argeo.slc.runtime.SlcAgentDescriptor; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; +import org.springframework.jms.JmsException; import org.springframework.jms.core.JmsTemplate; -import org.springframework.jms.listener.SessionAwareMessageListener; -import org.springframework.jms.support.converter.MessageConversionException; +import org.springframework.jms.core.MessagePostProcessor; /** JMS based implementation of SLC Agent. */ public class JmsAgent extends AbstractAgent implements SlcAgent, - InitializingBean, DisposableBean, SessionAwareMessageListener { + InitializingBean, DisposableBean, MessageListener { public final static String PROPERTY_QUERY = "query"; - public final static String PROPERTY_SLC_AGENT_ID = "slc_agentId"; + public final static String QUERY_PING_ALL = "pingAll"; private final static Log log = LogFactory.getLog(JmsAgent.class); private final SlcAgentDescriptor agentDescriptor; - // private ConnectionFactory connectionFactory; private JmsTemplate jmsTemplate; private Destination agentRegister; private Destination agentUnregister; - // private Destination requestDestination; private Destination responseDestination; - // private MessageConverter messageConverter; - public JmsAgent() { try { agentDescriptor = new SlcAgentDescriptor(); @@ -58,19 +54,33 @@ public class JmsAgent extends AbstractAgent implements SlcAgent, } public void afterPropertiesSet() throws Exception { - // Initialize JMS Template - // jmsTemplate = new JmsTemplate(connectionFactory); - // jmsTemplate.setMessageConverter(messageConverter); - - jmsTemplate.convertAndSend(agentRegister, agentDescriptor); - log.info("Agent #" + agentDescriptor.getUuid() + " registered to " - + agentRegister); + try { + jmsTemplate.convertAndSend(agentRegister, agentDescriptor); + log.info("Agent #" + agentDescriptor.getUuid() + " registered to " + + agentRegister); + } catch (JmsException e) { + log + .warn("Could not register agent " + + agentDescriptor.getUuid() + + " to server: " + + e.getMessage() + + ". The agent will stay offline but will keep listening for a ping all sent by server."); + if (log.isTraceEnabled()) + log.debug("Original error.", e); + } } public void destroy() throws Exception { - jmsTemplate.convertAndSend(agentUnregister, agentDescriptor); - log.info("Agent #" + agentDescriptor.getUuid() + " unregistered from " - + agentUnregister); + try { + jmsTemplate.convertAndSend(agentUnregister, agentDescriptor); + log.info("Agent #" + agentDescriptor.getUuid() + + " unregistered from " + agentUnregister); + } catch (JmsException e) { + log.warn("Could not unregister agent " + agentDescriptor.getUuid() + + ": " + e.getMessage()); + if (log.isTraceEnabled()) + log.debug("Original error.", e); + } } public void setAgentRegister(Destination agentRegister) { @@ -81,19 +91,6 @@ public class JmsAgent extends AbstractAgent implements SlcAgent, this.agentUnregister = agentUnregister; } - /* - * public void onMessage(Message message) { // FIXME: we filter the messages - * on the client side, // because of a weird problem with selector since - * moving to OSGi try { if (message.getStringProperty("slc-agentId").equals( - * agentDescriptor.getUuid())) { runSlcExecution((SlcExecution) - * messageConverter .fromMessage(message)); } else { if - * (log.isDebugEnabled()) log.debug("Filtered out message " + message); } } - * catch (JMSException e) { throw new SlcException("Cannot convert message " - * + message, e); } - * - * } - */ - public String getMessageSelector() { String messageSelector = "slc_agentId='" + agentDescriptor.getUuid() + "'"; @@ -102,61 +99,6 @@ public class JmsAgent extends AbstractAgent implements SlcAgent, return messageSelector; } - public void onMessage(Message message, Session session) throws JMSException { - MessageProducer producer = session.createProducer(responseDestination); - String query = message.getStringProperty(PROPERTY_QUERY); - String correlationId = message.getJMSCorrelationID(); - if (log.isDebugEnabled()) - log.debug("Received query " + query + " with correlationId " - + correlationId); - - Message responseMsg = null; - if ("getExecutionModuleDescriptor".equals(query)) { - String moduleName = message.getStringProperty("moduleName"); - String version = message.getStringProperty("version"); - ExecutionModuleDescriptor emd = getExecutionModuleDescriptor( - moduleName, version); - responseMsg = jmsTemplate.getMessageConverter().toMessage(emd, - session); - } else if ("listExecutionModuleDescriptors".equals(query)) { - - List lst = listExecutionModuleDescriptors(); - SlcAgentDescriptor agentDescriptorToSend = new SlcAgentDescriptor( - agentDescriptor); - agentDescriptorToSend.setModuleDescriptors(lst); - responseMsg = jmsTemplate.getMessageConverter().toMessage( - agentDescriptorToSend, session); - } else if ("newExecution".equals(query)) { - - SlcExecution slcExecution = (SlcExecution) jmsTemplate - .getMessageConverter().fromMessage(message); - runSlcExecution(slcExecution); - } else { - // try { - // // FIXME: generalize - // SlcExecution slcExecution = (SlcExecution) jmsTemplate - // .getMessageConverter().fromMessage(message); - // runSlcExecution(slcExecution); - // } catch (MessageConversionException e) { - // if (log.isDebugEnabled()) - // log.debug("Unsupported query " + query, e); - // } - if (log.isDebugEnabled()) - log.debug("Unsupported query " + query); - return; - } - - if (responseMsg != null) { - responseMsg.setJMSCorrelationID(correlationId); - producer.send(responseMsg); - if (log.isDebugEnabled()) - log.debug("Sent response to query " + query - + " with correlationId " + correlationId + ": " - + responseMsg); - } - - } - public ExecutionModuleDescriptor getExecutionModuleDescriptor( String moduleName, String version) { return getModulesManager().getExecutionModuleDescriptor(moduleName, @@ -177,6 +119,98 @@ public class JmsAgent extends AbstractAgent implements SlcAgent, return descriptors; } + public boolean ping() { + return true; + } + + public void onMessage(final Message message) { + final String query; + final String correlationId; + try { + query = message.getStringProperty(PROPERTY_QUERY); + correlationId = message.getJMSCorrelationID(); + } catch (JMSException e1) { + throw new SlcException("Cannot analyze incoming message " + message); + } + + final Object response; + final Destination destinationSend; + if (QUERY_PING_ALL.equals(query)) { + ReferenceList refList = (ReferenceList) convertFrom(message); + if (!refList.getReferences().contains(agentDescriptor.getUuid())) { + response = agentDescriptor; + destinationSend = agentRegister; + log.info("Agent #" + agentDescriptor.getUuid() + + " registering to " + agentRegister + + " in reply to a " + QUERY_PING_ALL + " query"); + } else { + return; + } + } else { + response = process(query, message); + destinationSend = responseDestination; + } + + // Send response + jmsTemplate.convertAndSend(destinationSend, response, + new MessagePostProcessor() { + public Message postProcessMessage(Message messageToSend) + throws JMSException { + messageToSend.setStringProperty(PROPERTY_QUERY, query); + messageToSend.setStringProperty(AbstractAgent.PROPERTY_SLC_AGENT_ID, + agentDescriptor.getUuid()); + messageToSend.setJMSCorrelationID(correlationId); + return messageToSend; + } + }); + if (log.isTraceEnabled()) + log.debug("Sent response to query '" + query + + "' with correlationId " + correlationId); + } + + /** @return response */ + public Object process(String query, Message message) { + try { + if ("getExecutionModuleDescriptor".equals(query)) { + String moduleName = message.getStringProperty("moduleName"); + String version = message.getStringProperty("version"); + return getExecutionModuleDescriptor(moduleName, version); + } else if ("listExecutionModuleDescriptors".equals(query)) { + + List lst = listExecutionModuleDescriptors(); + SlcAgentDescriptor agentDescriptorToSend = new SlcAgentDescriptor( + agentDescriptor); + agentDescriptorToSend.setModuleDescriptors(lst); + return agentDescriptorToSend; + } else if ("runSlcExecution".equals(query)) { + final SlcExecution slcExecution = (SlcExecution) convertFrom(message); + new Thread() { + public void run() { + runSlcExecution(slcExecution); + } + }.start(); + return ExecutionAnswer.ok("Execution started on agent " + + agentDescriptor.getUuid()); + } else if ("ping".equals(query)) { + return ExecutionAnswer.ok("Agent " + agentDescriptor.getUuid() + + " is alive."); + } else { + throw new SlcException("Unsupported query " + query); + } + } catch (Exception e) { + log.error("Processing of query " + query + " failed", e); + return ExecutionAnswer.error(e); + } + } + + protected Object convertFrom(Message message) { + try { + return jmsTemplate.getMessageConverter().fromMessage(message); + } catch (JMSException e) { + throw new SlcException("Cannot convert message", e); + } + } + public void setResponseDestination(Destination responseDestination) { this.responseDestination = responseDestination; }