diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-01-22 17:16:44 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-01-22 17:16:44 +0000 |
| commit | 7cd2a924fba6b0eb78c1c7487e647f5d298a280e (patch) | |
| tree | 7d20f7890ac788ce96ed4e9f72951a475e25bd7b /qpid/java/systests | |
| parent | 1c7a129ba58a45726a7d14377fb8ebe447457319 (diff) | |
| download | qpid-python-7cd2a924fba6b0eb78c1c7487e647f5d298a280e.tar.gz | |
QPID-5504 : initial refactoring to move common code into shared classes, make transports work similarly with respect to message routing
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1560424 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/systests')
| -rw-r--r-- | qpid/java/systests/src/main/java/org/apache/qpid/server/store/MessageStoreTest.java | 68 |
1 files changed, 23 insertions, 45 deletions
diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/server/store/MessageStoreTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/server/store/MessageStoreTest.java index 07f0d0c369..19dc1a5a02 100644 --- a/qpid/java/systests/src/main/java/org/apache/qpid/server/store/MessageStoreTest.java +++ b/qpid/java/systests/src/main/java/org/apache/qpid/server/store/MessageStoreTest.java @@ -50,8 +50,6 @@ import org.apache.qpid.server.queue.AMQPriorityQueue; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.BaseQueue; import org.apache.qpid.server.queue.ConflationQueue; -import org.apache.qpid.server.protocol.v0_8.IncomingMessage; -import org.apache.qpid.server.queue.QueueArgumentsConverter; import org.apache.qpid.server.queue.SimpleAMQQueue; import org.apache.qpid.server.txn.AutoCommitTransaction; import org.apache.qpid.server.txn.ServerTransaction; @@ -617,61 +615,41 @@ public class MessageStoreTest extends QpidTestCase MessagePublishInfo messageInfo = new TestMessagePublishInfo(exchange, false, false, routingKey); - final IncomingMessage currentMessage; - - - currentMessage = new IncomingMessage(messageInfo); - - currentMessage.setExchange(exchange); - ContentHeaderBody headerBody = new ContentHeaderBody(BasicConsumeBodyImpl.CLASS_ID,0,properties,0l); - try - { - currentMessage.setContentHeaderBody(headerBody); - } - catch (AMQException e) - { - fail(e.getMessage()); - } + MessageMetaData mmd = new MessageMetaData(messageInfo, headerBody, System.currentTimeMillis()); - currentMessage.setExpiration(); + final StoredMessage<MessageMetaData> storedMessage = getVirtualHost().getMessageStore().addMessage(mmd); + storedMessage.flushToStore(); + final AMQMessage currentMessage = new AMQMessage(storedMessage); - MessageMetaData mmd = currentMessage.headersReceived(System.currentTimeMillis()); - currentMessage.setStoredMessage(getVirtualHost().getMessageStore().addMessage(mmd)); - currentMessage.getStoredMessage().flushToStore(); - currentMessage.route(); + final List<? extends BaseQueue> destinationQueues = exchange.route(currentMessage); - // check and deliver if header says body length is zero - if (currentMessage.allContentReceived()) - { - ServerTransaction trans = new AutoCommitTransaction(getVirtualHost().getMessageStore()); - final List<? extends BaseQueue> destinationQueues = currentMessage.getDestinationQueues(); - trans.enqueue(currentMessage.getDestinationQueues(), currentMessage, new ServerTransaction.Action() { - public void postCommit() - { - try - { - AMQMessage message = new AMQMessage(currentMessage.getStoredMessage()); + ServerTransaction trans = new AutoCommitTransaction(getVirtualHost().getMessageStore()); - for(BaseQueue queue : destinationQueues) - { - queue.enqueue(message); - } - } - catch (AMQException e) + trans.enqueue(destinationQueues, currentMessage, new ServerTransaction.Action() { + public void postCommit() + { + try + { + for(BaseQueue queue : destinationQueues) { - _logger.error("Problem enqueing message", e); + queue.enqueue(currentMessage); } } - - public void onRollback() + catch (AMQException e) { - //To change body of implemented methods use File | Settings | File Templates. + _logger.error("Problem enqueing message", e); } - }); - } + } + + public void onRollback() + { + //To change body of implemented methods use File | Settings | File Templates. + } + }); + } private void createAllQueues() |
