summaryrefslogtreecommitdiff
path: root/java/broker/src
diff options
context:
space:
mode:
authorRajith Muditha Attapattu <rajith@apache.org>2007-07-24 00:35:26 +0000
committerRajith Muditha Attapattu <rajith@apache.org>2007-07-24 00:35:26 +0000
commit586e63b99de7711689b0728d7a0c20354256c8dc (patch)
tree0559ed6348a880f8ce979a951eaff56ecab62d19 /java/broker/src
parent42238d6f0a49bd9311229752c07278329b90e05c (diff)
downloadqpid-python-586e63b99de7711689b0728d7a0c20354256c8dc.tar.gz
adding synapse exchange
git-svn-id: https://svn.apache.org/repos/asf/incubator/qpid/branches/client_restructure@558903 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'java/broker/src')
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/AMQBrokerManagerMBean.java2
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java7
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeFactory.java19
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/Exchange.java2
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/ExchangeFactory.java2
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/MessageContextCreatorForQpid.java258
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/QpidSynapseEnvironment.java74
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/SynapseExchange.java104
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/TestClassMediator.java62
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/handler/ExchangeDeclareHandler.java2
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/protocol/ExchangeInitialiser.java5
-rw-r--r--java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java9
12 files changed, 529 insertions, 17 deletions
diff --git a/java/broker/src/main/java/org/apache/qpid/server/AMQBrokerManagerMBean.java b/java/broker/src/main/java/org/apache/qpid/server/AMQBrokerManagerMBean.java
index 204b5674ce..f0cf8d37ab 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/AMQBrokerManagerMBean.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/AMQBrokerManagerMBean.java
@@ -94,7 +94,7 @@ public class AMQBrokerManagerMBean extends AMQManagedObject implements ManagedBr
Exchange exchange = _exchangeRegistry.getExchange(new AMQShortString(exchangeName));
if (exchange == null)
{
- exchange = _exchangeFactory.createExchange(new AMQShortString(exchangeName), new AMQShortString(type), durable, autoDelete, 0);
+ exchange = _exchangeFactory.createExchange(_exchangeRegistry,new AMQShortString(exchangeName), new AMQShortString(type), durable, autoDelete, 0);
_exchangeRegistry.registerExchange(exchange);
}
else
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java
index 8b4f41a7a0..579ddf64d7 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java
@@ -126,7 +126,7 @@ public abstract class AbstractExchange implements Exchange, Managable
*/
protected abstract ExchangeMBean createMBean() throws AMQException;
- public void initialise(VirtualHost host, AMQShortString name, boolean durable, int ticket, boolean autoDelete) throws AMQException
+ public void initialise(VirtualHost host, AMQShortString name, boolean durable, int ticket, boolean autoDelete, ExchangeRegistry exchangeRegistry) throws AMQException
{
_virtualHost = host;
_name = name;
@@ -134,7 +134,10 @@ public abstract class AbstractExchange implements Exchange, Managable
_autoDelete = autoDelete;
_ticket = ticket;
_exchangeMbean = createMBean();
- _exchangeMbean.register();
+ if(_exchangeMbean != null)
+ {
+ _exchangeMbean.register();
+ }
}
public boolean isDurable()
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeFactory.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeFactory.java
index 86feb46bb6..db78555b0e 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeFactory.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeFactory.java
@@ -20,19 +20,15 @@
*/
package org.apache.qpid.server.exchange;
+import java.util.HashMap;
+import java.util.Map;
+
import org.apache.log4j.Logger;
import org.apache.qpid.AMQException;
-import org.apache.qpid.AMQChannelException;
import org.apache.qpid.AMQUnknownExchangeType;
-import org.apache.qpid.server.virtualhost.VirtualHost;
-import org.apache.qpid.protocol.AMQConstant;
import org.apache.qpid.exchange.ExchangeDefaults;
import org.apache.qpid.framing.AMQShortString;
-
-import java.util.HashMap;
-import java.util.Map;
-import org.apache.qpid.AMQUnknownExchangeType;
-import org.apache.qpid.exchange.ExchangeDefaults;
+import org.apache.qpid.server.virtualhost.VirtualHost;
public class DefaultExchangeFactory implements ExchangeFactory
{
@@ -48,9 +44,12 @@ public class DefaultExchangeFactory implements ExchangeFactory
_exchangeClassMap.put(ExchangeDefaults.TOPIC_EXCHANGE_CLASS, org.apache.qpid.server.exchange.DestWildExchange.class);
_exchangeClassMap.put(ExchangeDefaults.HEADERS_EXCHANGE_CLASS, org.apache.qpid.server.exchange.HeadersExchange.class);
_exchangeClassMap.put(ExchangeDefaults.FANOUT_EXCHANGE_CLASS, org.apache.qpid.server.exchange.FanoutExchange.class);
+
+ // I'd rather allow an extention mechanism to register custom exchanges. for standard default exchanges this is fine.
+ _exchangeClassMap.put(new AMQShortString("synapse"), org.apache.qpid.server.exchange.synapse.SynapseExchange.class);
}
- public Exchange createExchange(AMQShortString exchange, AMQShortString type, boolean durable, boolean autoDelete,
+ public Exchange createExchange(ExchangeRegistry exchangeRegistry,AMQShortString exchange, AMQShortString type, boolean durable, boolean autoDelete,
int ticket)
throws AMQException
{
@@ -62,7 +61,7 @@ public class DefaultExchangeFactory implements ExchangeFactory
try
{
Exchange e = exchClass.newInstance();
- e.initialise(_host, exchange, durable, ticket, autoDelete);
+ e.initialise(_host, exchange, durable, ticket, autoDelete, exchangeRegistry);
return e;
}
catch (InstantiationException e)
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/Exchange.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/Exchange.java
index c012a1c1c9..084811df09 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/exchange/Exchange.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/Exchange.java
@@ -32,7 +32,7 @@ public interface Exchange
AMQShortString getName();
AMQShortString getType();
- void initialise(VirtualHost host, AMQShortString name, boolean durable, int ticket, boolean autoDelete) throws AMQException;
+ void initialise(VirtualHost host, AMQShortString name, boolean durable, int ticket, boolean autoDelete, ExchangeRegistry exchangeRegistry) throws AMQException;
boolean isDurable();
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/ExchangeFactory.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/ExchangeFactory.java
index e07fd0b8fc..7b57a860e4 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/exchange/ExchangeFactory.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/ExchangeFactory.java
@@ -26,7 +26,7 @@ import org.apache.qpid.framing.AMQShortString;
public interface ExchangeFactory
{
- Exchange createExchange(AMQShortString exchange, AMQShortString type, boolean durable, boolean autoDelete,
+ Exchange createExchange(ExchangeRegistry exchangeRegistry, AMQShortString exchange, AMQShortString type, boolean durable, boolean autoDelete,
int ticket)
throws AMQException;
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/MessageContextCreatorForQpid.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/MessageContextCreatorForQpid.java
new file mode 100644
index 0000000000..c7887ed99b
--- /dev/null
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/MessageContextCreatorForQpid.java
@@ -0,0 +1,258 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.qpid.server.exchange.synapse;
+
+import java.io.ByteArrayInputStream;
+import java.io.FileInputStream;
+
+import javax.activation.DataHandler;
+import javax.xml.namespace.QName;
+import javax.xml.stream.XMLStreamException;
+import javax.xml.stream.XMLStreamReader;
+
+import org.apache.axiom.attachments.ByteArrayDataSource;
+import org.apache.axiom.om.OMDocument;
+import org.apache.axiom.om.OMElement;
+import org.apache.axiom.om.OMText;
+import org.apache.axiom.om.impl.builder.StAXOMBuilder;
+import org.apache.axiom.om.util.StAXUtils;
+import org.apache.axiom.soap.SOAPEnvelope;
+import org.apache.axiom.soap.SOAPFactory;
+import org.apache.axiom.soap.impl.llom.soap11.SOAP11Factory;
+import org.apache.axis2.AxisFault;
+import org.apache.axis2.addressing.EndpointReference;
+import org.apache.axis2.addressing.RelatesTo;
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
+import org.apache.qpid.framing.AMQShortString;
+import org.apache.qpid.framing.Content;
+import org.apache.qpid.framing.MessageTransferBody;
+import org.apache.qpid.server.queue.AMQMessage;
+import org.apache.synapse.MessageContext;
+import org.apache.synapse.SynapseException;
+import org.apache.synapse.config.SynapseConfiguration;
+import org.apache.synapse.core.SynapseEnvironment;
+import org.apache.synapse.core.axis2.Axis2MessageContext;
+
+/**
+ * The MessageContext needs to be set up and then is used by the SynapseMessageReceiver to inject messages.
+ * This class is used by the SynapseMessageReceiver to find the environment. The env is stored in a Parameter to the Axis2 config
+ */
+public class MessageContextCreatorForQpid{
+
+ private static Log log = LogFactory.getLog(MessageContextCreatorForQpid.class);
+
+ private static SynapseConfiguration synCfg = null;
+ private static SynapseEnvironment synEnv = null;
+
+ final static String ORIGINAL_MESSAGE = "ORIGINAL_MESSAGE";
+ final static String AMQP_CONTENT_TYPE = "AMQP_CONTENT_TYPE";
+ final static String DEFAULT_CHAR_SET_ENCODING = "UTF-8";
+
+ enum ContentType
+ {
+ TEXT_PLAIN ("text/plain"),
+ TEXT_XML ("text/xml"),
+ APPLICATION_OCTECT ("application/octet-stream");
+
+ private final String _value;
+
+ private ContentType (String value)
+ {
+ _value = value;
+ }
+
+ public String value()
+ {
+ return _value;
+ }
+ }
+
+ private static String createURL(String exchangeName,String routingKey)
+ {
+ StringBuffer buf = new StringBuffer();
+ buf.append("amqp://");
+ buf.append(exchangeName);
+ buf.append("?");
+ buf.append("routingKey=");
+ buf.append(routingKey);
+
+ return buf.toString();
+ }
+
+ public static MessageContext getSynapseMessageContext(AMQMessage amqMsg) throws SynapseException {
+
+ if (synCfg == null || synEnv == null) {
+ String msg = "Synapse environment has not initialized properly..";
+ log.fatal(msg);
+ throw new SynapseException(msg);
+ }
+
+ org.apache.axis2.context.MessageContext axis2MC = new org.apache.axis2.context.MessageContext();
+ Axis2MessageContext synCtx = new Axis2MessageContext(axis2MC, synCfg, synEnv);
+ synCtx.setMessageID(amqMsg.getTransferBody().getMessageId().asString());
+ if(amqMsg.getTransferBody().getCorrelationId() != null)
+ {
+ synCtx.setRelatesTo(new RelatesTo[]{new RelatesTo(amqMsg.getTransferBody().getCorrelationId().asString())});
+ }
+ synCtx.setTo(new EndpointReference(createURL(amqMsg.getTransferBody().getExchange().asString(),amqMsg.getTransferBody().getRoutingKey().asString())));
+
+ if(amqMsg.getTransferBody().getReplyTo() != null)
+ {
+ synCtx.setReplyTo(new EndpointReference(createURL(amqMsg.getTransferBody().getExchange().asString(),amqMsg.getTransferBody().getReplyTo().asString())));
+ }
+ synCtx.setDoingPOX(true);
+ synCtx.setProperty(ORIGINAL_MESSAGE, amqMsg);
+
+ //Creating a fictitious SOAP envelope to support the synapse model
+
+ SOAPFactory soapFactory = new SOAP11Factory();
+ SOAPEnvelope envelope = soapFactory.getDefaultEnvelope();
+
+ String contentType = amqMsg.getTransferBody().getContentType().asString();
+ if(ContentType.TEXT_PLAIN.value().equals(contentType))
+ {
+ OMElement wrapper = soapFactory.createOMElement(new QName("payload"), null);
+ OMText textData = soapFactory.createOMText(amqMsg.getTransferBody().getBody().getContentAsString());
+ wrapper.addChild(textData);
+ envelope.getBody().addChild(wrapper);
+ }
+ else if (ContentType.TEXT_XML.value().equals(contentType))
+ {
+ XMLStreamReader parser;
+ try
+ {
+ parser = StAXUtils.createXMLStreamReader(
+ new ByteArrayInputStream(amqMsg.getTransferBody().getBody().getContentAsByteArray()),
+ DEFAULT_CHAR_SET_ENCODING);
+ }
+ catch (XMLStreamException e)
+ {
+ throw new SynapseException("Error reading the XML message",e);
+ }
+
+ StAXOMBuilder builder = new StAXOMBuilder(parser);
+ //builder.setOMBuilderFactory(soapFactory);
+
+ Object obj = builder.getDocumentElement();
+ envelope.getBody().addChild(builder.getDocumentElement());
+ }
+ else if (ContentType.APPLICATION_OCTECT.value().equals(contentType))
+ {
+ // treat binary data as an attachment
+ DataHandler dataHandler = new DataHandler(
+ new ByteArrayDataSource(amqMsg.getTransferBody().getBody().getContentAsByteArray()));
+ OMText textData = soapFactory.createOMText(dataHandler, true);
+ OMElement wrapper = soapFactory.createOMElement(new QName("payload"), null);
+ wrapper.addChild(textData);
+ synCtx.setDoingMTOM(true);
+
+ envelope.getBody().addChild(wrapper);
+ }
+ else
+ {
+ throw new SynapseException("Unsupported Content Type : " + contentType);
+ }
+
+ synCtx.setProperty(AMQP_CONTENT_TYPE, contentType);
+
+ try
+ {
+ synCtx.setEnvelope(envelope);
+ }
+ catch(AxisFault e)
+ {
+ throw new SynapseException(e);
+ }
+
+ return synCtx;
+ }
+
+ public static AMQMessage getAMQMessage(MessageContext mc)
+ {
+ AMQMessage origMsg = (AMQMessage)mc.getProperty(ORIGINAL_MESSAGE);
+ OMElement payload = mc.getEnvelope().getBody().getFirstElement();
+
+ String amqContentType = (String)mc.getProperty(AMQP_CONTENT_TYPE);
+ byte[] content = new byte[0];
+
+ if(ContentType.TEXT_PLAIN.value().equals(amqContentType))
+ {
+ // For plain text there was a wrapper element
+ content = payload.getText().getBytes();
+ }
+ else if (ContentType.TEXT_XML.value().equals(amqContentType))
+ {
+ content = payload.getText().getBytes();
+ }
+ else if (ContentType.APPLICATION_OCTECT.value().equals(amqContentType) && mc.isDoingMTOM())
+ {
+
+ }
+
+ String url = mc.getTo().getAddress();;
+ // very crude
+ // should have utility class to do this, but do it when amqp
+ // officialy converge on an addressing scheme
+ String exchangeName = url.substring(7,url.indexOf('?'));
+ String routingKey = url.substring(url.indexOf('=')+1,url.length());
+
+
+ MessageTransferBody origTransferBody = origMsg.getTransferBody();
+ MessageTransferBody transferBody = MessageTransferBody.createMethodBody(
+ origTransferBody.getMajor(),
+ origTransferBody.getMinor(),
+ origTransferBody.getAppId(), //appId
+ origTransferBody.getApplicationHeaders(), //applicationHeaders
+ new Content(Content.TypeEnum.INLINE_T, content), //body
+ origTransferBody.getContentType(), //contentEncoding,
+ origTransferBody.getContentType(), //contentType
+ origTransferBody.getCorrelationId(), //correlationId
+ origTransferBody.getDeliveryMode(), //deliveryMode non persistant
+ new AMQShortString(exchangeName),// destination
+ new AMQShortString(exchangeName),// exchange
+ origTransferBody.getExpiration(), //expiration
+ origTransferBody.getImmediate(), //immediate
+ origTransferBody.getMandatory(), //mandatory
+ origTransferBody.getMessageId(), //messageId
+ origTransferBody.getPriority(), //priority
+ origTransferBody.getRedelivered(), //redelivered
+ origTransferBody.getReplyTo(), //replyTo
+ new AMQShortString(routingKey), //routingKey,
+ "abc".getBytes(), //securityToken
+ origTransferBody.ticket, //ticket
+ System.currentTimeMillis(), //timestamp
+ origTransferBody.getTransactionId(), //transactionId
+ origTransferBody.getTtl(), //ttl,
+ origTransferBody.getUserId() //userId
+ );
+ AMQMessage newMsg = new AMQMessage(origMsg.getMessageStore(),transferBody,origMsg.getTransactionContext());
+
+ return newMsg;
+ }
+
+ public static void setSynConfig(SynapseConfiguration synCfg) {
+ MessageContextCreatorForQpid.synCfg = synCfg;
+ }
+
+ public static void setSynEnv(SynapseEnvironment synEnv) {
+ MessageContextCreatorForQpid.synEnv = synEnv;
+ }
+}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/QpidSynapseEnvironment.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/QpidSynapseEnvironment.java
new file mode 100644
index 0000000000..60fdc40788
--- /dev/null
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/QpidSynapseEnvironment.java
@@ -0,0 +1,74 @@
+package org.apache.qpid.server.exchange.synapse;
+
+import org.apache.axis2.addressing.EndpointReference;
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
+import org.apache.qpid.server.exchange.Exchange;
+import org.apache.qpid.server.queue.AMQMessage;
+import org.apache.synapse.MessageContext;
+import org.apache.synapse.SynapseException;
+import org.apache.synapse.config.SynapseConfiguration;
+import org.apache.synapse.core.SynapseEnvironment;
+import org.apache.synapse.core.axis2.Axis2MessageContext;
+import org.apache.synapse.endpoints.utils.EndpointDefinition;
+import org.apache.synapse.statistics.StatisticsCollector;
+
+public class QpidSynapseEnvironment implements SynapseEnvironment
+{
+
+ private static final Log log = LogFactory.getLog(QpidSynapseEnvironment.class);
+
+ private SynapseConfiguration synapseConfig;
+
+ private StatisticsCollector statisticsCollector;
+
+ private SynapseExchange qpidExchange;
+
+ public QpidSynapseEnvironment(SynapseConfiguration synapseConfig, SynapseExchange qpidExchange)
+ {
+ this.synapseConfig = synapseConfig;
+ this.qpidExchange = qpidExchange;
+ }
+
+ public MessageContext createMessageContext()
+ {
+ org.apache.axis2.context.MessageContext axis2MC = new org.apache.axis2.context.MessageContext();
+ MessageContext mc = new Axis2MessageContext(axis2MC, synapseConfig, this);
+ return mc;
+ }
+
+ public StatisticsCollector getStatisticsCollector()
+ {
+ return statisticsCollector;
+ }
+
+ public void injectMessage(MessageContext synCtx)
+ {
+
+ synCtx.getMainSequence().mediate(synCtx);
+ }
+
+ public void send(EndpointDefinition endpoint, MessageContext smc)
+ {
+ if(endpoint != null)
+ {
+ smc.setTo(new EndpointReference(endpoint.getAddress()));
+ AMQMessage newMessage = MessageContextCreatorForQpid.getAMQMessage(smc);
+ try
+ {
+ qpidExchange.getExchangeRegistry().routeContent(newMessage);
+ }
+ catch(Exception e)
+ {
+ throw new SynapseException("Faulty endpoint",e);
+ }
+ }
+
+ }
+
+ public void setStatisticsCollector(StatisticsCollector statisticsCollector)
+ {
+ this.statisticsCollector = statisticsCollector;
+ }
+
+}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/SynapseExchange.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/SynapseExchange.java
new file mode 100644
index 0000000000..c408529fbd
--- /dev/null
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/SynapseExchange.java
@@ -0,0 +1,104 @@
+package org.apache.qpid.server.exchange.synapse;
+
+import org.apache.qpid.AMQException;
+import org.apache.qpid.framing.AMQShortString;
+import org.apache.qpid.framing.FieldTable;
+import org.apache.qpid.server.exchange.AbstractExchange;
+import org.apache.qpid.server.exchange.ExchangeRegistry;
+import org.apache.qpid.server.queue.AMQMessage;
+import org.apache.qpid.server.queue.AMQQueue;
+import org.apache.qpid.server.virtualhost.VirtualHost;
+import org.apache.synapse.Constants;
+import org.apache.synapse.MessageContext;
+import org.apache.synapse.config.SynapseConfiguration;
+import org.apache.synapse.config.SynapseConfigurationBuilder;
+import org.apache.synapse.core.SynapseEnvironment;
+
+public class SynapseExchange extends AbstractExchange
+{
+
+ public final static AMQShortString TYPE = new AMQShortString("synapse");
+
+ private SynapseEnvironment synEnv;
+
+ private ExchangeRegistry exchangeRegistry;
+
+ public SynapseExchange()
+ {
+ super();
+ }
+
+ @Override
+ public void initialise(VirtualHost host, AMQShortString name, boolean durable, int ticket, boolean autoDelete, ExchangeRegistry exchangeRegistry) throws AMQException
+ {
+ super.initialise(host, name, durable, ticket, autoDelete, exchangeRegistry);
+
+ String config = System.getProperty(Constants.SYNAPSE_XML);
+ SynapseConfiguration synapseConfiguration = SynapseConfigurationBuilder.getConfiguration(config);
+ synEnv = new QpidSynapseEnvironment(synapseConfiguration,this);
+ MessageContextCreatorForQpid.setSynConfig(synapseConfiguration);
+ MessageContextCreatorForQpid.setSynEnv(synEnv);
+ this.exchangeRegistry = exchangeRegistry;
+ }
+
+ @Override
+ protected ExchangeMBean createMBean() throws AMQException
+ {
+ // TODO Auto-generated method stub
+ return null;
+ }
+
+ public void deregisterQueue(AMQShortString routingKey, AMQQueue queue) throws AMQException
+ {
+ throw new UnsupportedOperationException("This exchange does not take bindings");
+ }
+
+ public AMQShortString getType()
+ {
+ return TYPE;
+ }
+
+ public boolean hasBindings() throws AMQException
+ {
+ return false;
+ }
+
+ public boolean isBound(AMQShortString routingKey, AMQQueue queue) throws AMQException
+ {
+ throw new UnsupportedOperationException("This exchange does not take bindings");
+ }
+
+ public boolean isBound(AMQShortString routingKey) throws AMQException
+ {
+ throw new UnsupportedOperationException("This exchange does not take bindings");
+ }
+
+ public boolean isBound(AMQQueue queue) throws AMQException
+ {
+ throw new UnsupportedOperationException("This exchange does not take bindings");
+ }
+
+ public void registerQueue(AMQShortString routingKey, AMQQueue queue, FieldTable args) throws AMQException
+ {
+ throw new UnsupportedOperationException("This exchange does not take bindings");
+ }
+
+ public void route(AMQMessage message) throws AMQException
+ {
+ try
+ {
+ MessageContext mc = MessageContextCreatorForQpid.getSynapseMessageContext(message);
+ synEnv.injectMessage(mc);
+ }
+ catch(Exception e)
+ {
+ throw new AMQException("Error occurred while trying to mediate message through Synapse",e);
+ }
+ }
+
+ public ExchangeRegistry getExchangeRegistry()
+ {
+ return exchangeRegistry;
+ }
+
+}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/TestClassMediator.java b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/TestClassMediator.java
new file mode 100644
index 0000000000..09ef377ef7
--- /dev/null
+++ b/java/broker/src/main/java/org/apache/qpid/server/exchange/synapse/TestClassMediator.java
@@ -0,0 +1,62 @@
+package org.apache.qpid.server.exchange.synapse;
+
+import javax.activation.DataHandler;
+import javax.xml.namespace.QName;
+
+import org.apache.axiom.attachments.ByteArrayDataSource;
+import org.apache.axiom.om.OMElement;
+import org.apache.axiom.om.OMText;
+import org.apache.axiom.soap.SOAPFactory;
+import org.apache.axiom.soap.impl.llom.soap11.SOAP11Factory;
+import org.apache.synapse.Mediator;
+import org.apache.synapse.MessageContext;
+
+public class TestClassMediator implements Mediator
+{
+
+ public int getTraceState()
+ {
+ // TODO Auto-generated method stub
+ return 0;
+ }
+
+ public String getType()
+ {
+ // TODO Auto-generated method stub
+ return null;
+ }
+
+ public boolean mediate(MessageContext mc)
+ {
+ SOAPFactory soapFactory = new SOAP11Factory();
+ OMElement binaryNode = mc.getEnvelope().getBody().getFirstChildWithName(new QName("payload"));
+ byte[] source = binaryNode.getText().getBytes();
+
+ byte[] b = new byte[source.length];
+ int j = 0;
+ for(int i=source.length-1; i>0; i--)
+ {
+ b[j] = source[i];
+ j++;
+ }
+
+ mc.getEnvelope().getBody().getFirstChildWithName(new QName("payload")).detach();
+
+ DataHandler dataHandler = new DataHandler(
+ new ByteArrayDataSource(b));
+ OMText textData = soapFactory.createOMText(dataHandler, true);
+ OMElement wrapper = soapFactory.createOMElement(new QName("payload"), null);
+ wrapper.addChild(textData);
+ mc.setDoingMTOM(true);
+
+ mc.getEnvelope().getBody().addChild(wrapper);
+ return true;
+ }
+
+ public void setTraceState(int arg0)
+ {
+ // TODO Auto-generated method stub
+
+ }
+
+}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/handler/ExchangeDeclareHandler.java b/java/broker/src/main/java/org/apache/qpid/server/handler/ExchangeDeclareHandler.java
index 7b129f0187..067eef5003 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/handler/ExchangeDeclareHandler.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/handler/ExchangeDeclareHandler.java
@@ -79,7 +79,7 @@ public class ExchangeDeclareHandler implements StateAwareMethodListener<Exchange
{
try
{
- exchange = exchangeFactory.createExchange(body.exchange, body.type, body.durable,
+ exchange = exchangeFactory.createExchange(exchangeRegistry,body.exchange, body.type, body.durable,
body.passive, body.ticket);
exchangeRegistry.registerExchange(exchange);
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/protocol/ExchangeInitialiser.java b/java/broker/src/main/java/org/apache/qpid/server/protocol/ExchangeInitialiser.java
index 8b5f05e8ea..df2d840bc9 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/protocol/ExchangeInitialiser.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/protocol/ExchangeInitialiser.java
@@ -34,12 +34,15 @@ public class ExchangeInitialiser
define(registry, factory, ExchangeDefaults.HEADERS_EXCHANGE_NAME, ExchangeDefaults.HEADERS_EXCHANGE_CLASS);
define(registry, factory, ExchangeDefaults.FANOUT_EXCHANGE_NAME, ExchangeDefaults.FANOUT_EXCHANGE_CLASS);
+ //There should be an extention mechanism to register
+ define(registry,factory,new AMQShortString("amq.synapse"),new AMQShortString("synapse"));
+
registry.setDefaultExchange(registry.getExchange(ExchangeDefaults.DIRECT_EXCHANGE_NAME));
}
private void define(ExchangeRegistry r, ExchangeFactory f,
AMQShortString name, AMQShortString type) throws AMQException
{
- r.registerExchange(f.createExchange(name, type, true, false, 0));
+ r.registerExchange(f.createExchange(r,name, type, true, false, 0));
}
}
diff --git a/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java b/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java
index 711e045516..59f88e2f43 100644
--- a/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java
+++ b/java/broker/src/main/java/org/apache/qpid/server/queue/AMQMessage.java
@@ -643,4 +643,13 @@ public class AMQMessage
return _requestId;
}
+ public MessageStore getMessageStore()
+ {
+ return _store;
+ }
+
+ public TransactionalContext getTransactionContext()
+ {
+ return _txnContext;
+ }
}