summaryrefslogtreecommitdiff
path: root/qpid/java/systests
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2014-01-22 17:16:44 +0000
committerRobert Godfrey <rgodfrey@apache.org>2014-01-22 17:16:44 +0000
commit7cd2a924fba6b0eb78c1c7487e647f5d298a280e (patch)
tree7d20f7890ac788ce96ed4e9f72951a475e25bd7b /qpid/java/systests
parent1c7a129ba58a45726a7d14377fb8ebe447457319 (diff)
downloadqpid-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.java68
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()