From 6713bfc5ddc1ff6202dad0d950a252273f73f795 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 5 Mar 2014 16:04:16 +0000 Subject: QPID-4000 , QPID-5601 : Improve conversion of reply-to between different protocols. Add functionality to the default exchange to understand AMQP 1.0 addresses. git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1574551 13f79535-47bb-0310-9956-ffa450edef68 --- .../server/protocol/v0_10/ConsumerTarget_0_10.java | 2 +- .../v0_10/MessageConverter_Internal_to_v0_10.java | 15 +- .../protocol/v0_10/MessageConverter_v0_10.java | 3 +- .../protocol/v0_10/MessageTransferMessage.java | 2 +- .../qpid/server/protocol/v0_10/ServerSession.java | 7 +- .../qpid/server/protocol/v0_8/AMQChannel.java | 9 +- .../qpid/server/protocol/v0_8/AMQMessage.java | 2 +- .../v0_8/MessageConverter_Internal_to_v0_8.java | 8 +- .../v0_8/MessageConverter_v0_8_to_Internal.java | 210 ++++++++++++++++++++- .../server/protocol/v1_0/ExchangeDestination.java | 2 +- .../server/protocol/v1_0/MessageMetaData_1_0.java | 5 + .../qpid/server/protocol/v1_0/Message_1_0.java | 10 +- .../protocol/v1_0/NodeReceivingDestination.java | 2 +- .../v0_10_v1_0/MessageConverter_0_10_to_1_0.java | 2 +- .../v0_10_v1_0/MessageConverter_1_0_to_v0_10.java | 51 ++++- .../v0_8_v0_10/MessageConverter_0_8_to_0_10.java | 2 +- .../v0_8_v1_0/MessageConverter_0_8_to_1_0.java | 42 ++++- .../server/management/amqp/ManagementNode.java | 16 +- 18 files changed, 337 insertions(+), 53 deletions(-) (limited to 'qpid/java/broker-plugins') diff --git a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java index 69c625d41d..120ba2d951 100644 --- a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java +++ b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java @@ -417,7 +417,7 @@ public class ConsumerTarget_0_10 extends AbstractConsumerTarget implements FlowC { logActor.message(ChannelMessages.DISCARDMSG_NOALTEXCH(msg.getMessageNumber(), queue.getName(), - msg.getRoutingKey())); + msg.getInitialRoutingAddress())); } } } diff --git a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/MessageConverter_Internal_to_v0_10.java b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/MessageConverter_Internal_to_v0_10.java index 37bbd810b4..7ff3873856 100644 --- a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/MessageConverter_Internal_to_v0_10.java +++ b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/MessageConverter_Internal_to_v0_10.java @@ -30,17 +30,8 @@ import org.apache.qpid.transport.DeliveryProperties; import org.apache.qpid.transport.Header; import org.apache.qpid.transport.MessageDeliveryPriority; import org.apache.qpid.transport.MessageProperties; -import org.apache.qpid.transport.codec.BBDecoder; -import org.apache.qpid.typedmessage.TypedBytesContentReader; -import org.apache.qpid.typedmessage.TypedBytesFormatException; -import java.io.EOFException; import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.ListIterator; -import java.util.Map; public class MessageConverter_Internal_to_v0_10 implements MessageConverter { @@ -123,7 +114,7 @@ public class MessageConverter_Internal_to_v0_10 implements MessageConverter> } }; - int enqueues = _currentMessage.getDestination().send(amqMessage, instanceProperties, _transaction, - immediate ? _immediateAction : _capacityCheckAction); + int enqueues = _currentMessage.getDestination().send(amqMessage, + amqMessage.getInitialRoutingAddress(), + instanceProperties, _transaction, + immediate ? _immediateAction : _capacityCheckAction + ); if(enqueues == 0) { handleUnroutableMessage(amqMessage); @@ -1574,7 +1577,7 @@ public class AMQChannel> if (altExchange == null) { _logger.debug("No alternate exchange configured for queue, must discard the message as unable to DLQ: delivery tag: " + deliveryTag); - _actor.message(_logSubject, ChannelMessages.DISCARDMSG_NOALTEXCH(msg.getMessageNumber(), queue.getName(), msg.getRoutingKey())); + _actor.message(_logSubject, ChannelMessages.DISCARDMSG_NOALTEXCH(msg.getMessageNumber(), queue.getName(), msg.getInitialRoutingAddress())); } else diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQMessage.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQMessage.java index 833f5fb06f..0ed63daf7c 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQMessage.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQMessage.java @@ -71,7 +71,7 @@ public class AMQMessage extends AbstractServerMessageImpl headerProps = new LinkedHashMap(); for(String headerName : serverMsg.getMessageHeader().getHeaderNames()) @@ -184,6 +185,7 @@ public class MessageConverter_Internal_to_v0_8 implements MessageConverter { @@ -58,9 +65,210 @@ public class MessageConverter_v0_8_to_Internal implements MessageConverter names) + { + return _delegate.containsHeaders(names); + } + + @Override + public boolean containsHeader(final String name) + { + return _delegate.containsHeader(name); + } + + @Override + public Collection getHeaderNames() + { + return _delegate.getHeaderNames(); + } + } + + private static Object convertMessageBody(String mimeType, byte[] data) { if("text/plain".equals(mimeType) || "text/xml".equals(mimeType)) diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java index d83665ad39..fc2c0d93d0 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java @@ -76,7 +76,7 @@ public class ExchangeDestination implements ReceivingDestination, SendingDestina return null; }}; - int enqueues = _exchange.send(message, instanceProperties, txn, null); + int enqueues = _exchange.send(message, message.getInitialRoutingAddress(), instanceProperties, txn, null); return enqueues == 0 ? REJECTED : ACCEPTED; diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageMetaData_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageMetaData_1_0.java index d5a349304c..4540308f61 100755 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageMetaData_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/MessageMetaData_1_0.java @@ -563,6 +563,11 @@ public class MessageMetaData_1_0 implements StorableMessageMetaData { return _properties == null ? null : _properties.getTo(); } + + public Map getHeadersAsMap() + { + return new HashMap(_appProperties); + } } } diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Message_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Message_1_0.java index 66094f52f0..36796851e0 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Message_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Message_1_0.java @@ -69,7 +69,7 @@ public class Message_1_0 extends AbstractServerMessageImpl convertToStoredMessage(final Message_1_0 serverMsg) + private StoredMessage convertToStoredMessage(final Message_1_0 serverMsg, + final VirtualHost vhost) { Object bodyObject = MessageConverter_from_1_0.convertBodyToObject(serverMsg); final byte[] messageContent = MessageConverter_from_1_0.convertToBody(bodyObject); final MessageMetaData_0_10 messageMetaData_0_10 = convertMetaData(serverMsg, + vhost, MessageConverter_from_1_0.getBodyMimeType(bodyObject), messageContent.length); @@ -119,25 +123,54 @@ public class MessageConverter_1_0_to_v0_10 implements MessageConverter { @@ -102,9 +104,45 @@ public class MessageConverter_0_8_to_1_0 extends MessageConverter_to_1_0> int send(final M message, - final InstanceProperties instanceProperties, - final ServerTransaction txn, - final Action postEnqueueAction) + final String routingAddress, + final InstanceProperties instanceProperties, + final ServerTransaction txn, + final Action postEnqueueAction) { @SuppressWarnings("unchecked") @@ -361,11 +363,19 @@ class ManagementNode implements MessageSource, MessageDestination ManagementNodeConsumer consumer = _consumers.get(message.getMessageHeader().getReplyTo()); + response.setInitialRoutingAddress(message.getMessageHeader().getReplyTo()); if(consumer != null) { // TODO - check same owner consumer.send(response); } + else + { + _virtualHost.getDefaultDestination().send(response, + message.getMessageHeader().getReplyTo(), InstanceProperties.EMPTY, + new AutoCommitTransaction(_virtualHost.getMessageStore()), + null); + } // TODO - route to a queue } -- cgit v1.2.1