From 672f3f7ab44bb4666deb95e80008a9d3e7b35806 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Thu, 17 Apr 2014 12:38:47 +0000 Subject: QPID-5580 : [Java Broker] Introduce explicit type hierarchy for queues in the ConfiguredObject model git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1588234 13f79535-47bb-0310-9956-ffa450edef68 --- .../qpid/server/queue/ConflationQueueTest.java | 574 --------------------- .../qpid/server/queue/LastValueQueueTest.java | 574 +++++++++++++++++++++ .../server/store/VirtualHostMessageStoreTest.java | 20 +- .../management/jmx/QueueManagementTest.java | 4 +- .../java/org/apache/qpid/systest/rest/Asserts.java | 13 +- .../qpid/systest/rest/VirtualHostRestTest.java | 35 +- 6 files changed, 614 insertions(+), 606 deletions(-) delete mode 100644 qpid/java/systests/src/main/java/org/apache/qpid/server/queue/ConflationQueueTest.java create mode 100644 qpid/java/systests/src/main/java/org/apache/qpid/server/queue/LastValueQueueTest.java (limited to 'qpid/java/systests/src') diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/server/queue/ConflationQueueTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/server/queue/ConflationQueueTest.java deleted file mode 100644 index 0e59e9cceb..0000000000 --- a/qpid/java/systests/src/main/java/org/apache/qpid/server/queue/ConflationQueueTest.java +++ /dev/null @@ -1,574 +0,0 @@ -/* - * - * 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.queue; - -import org.apache.log4j.Logger; -import org.apache.qpid.AMQException; -import org.apache.qpid.client.AMQConnection; -import org.apache.qpid.client.AMQDestination; -import org.apache.qpid.client.AMQQueue; -import org.apache.qpid.client.AMQSession; -import org.apache.qpid.framing.AMQShortString; -import org.apache.qpid.test.utils.QpidBrokerTestCase; -import org.apache.qpid.url.AMQBindingURL; - -import javax.jms.Connection; -import javax.jms.JMSException; -import javax.jms.Message; -import javax.jms.MessageConsumer; -import javax.jms.MessageProducer; -import javax.jms.Queue; -import javax.jms.Session; - -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; - -public class ConflationQueueTest extends QpidBrokerTestCase -{ - private static final Logger LOGGER = Logger.getLogger(ConflationQueueTest.class); - - private static final String MESSAGE_SEQUENCE_NUMBER_PROPERTY = "msg"; - private static final String KEY_PROPERTY = "key"; - - private static final int MSG_COUNT = 400; - - private String _queueName; - private Queue _queue; - private Connection _producerConnection; - private MessageProducer _producer; - private Session _producerSession; - private Connection _consumerConnection; - private Session _consumerSession; - private MessageConsumer _consumer; - - protected void setUp() throws Exception - { - super.setUp(); - - _queueName = getTestQueueName(); - _producerConnection = getConnection(); - _producerSession = _producerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - } - - public void testConflation() throws Exception - { - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - createConflationQueue(_producerSession); - _producer = _producerSession.createProducer(_queue); - - for (int msg = 0; msg < MSG_COUNT; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - } - - _producer.close(); - _producerSession.close(); - _producerConnection.close(); - - _consumer = _consumerSession.createConsumer(_queue); - _consumerConnection.start(); - Message received; - - List messages = new ArrayList(); - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - - assertEquals("Unexpected number of messages received",10,messages.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - } - - public void testConflationWithRelease() throws Exception - { - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.CLIENT_ACKNOWLEDGE); - - - createConflationQueue(_producerSession); - _producer = _producerSession.createProducer(_queue); - - for (int msg = 0; msg < MSG_COUNT/2; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - - } - - // HACK to do something synchronous - ((AMQSession)_producerSession).sync(); - - _consumer = _consumerSession.createConsumer(_queue); - _consumerConnection.start(); - Message received; - List messages = new ArrayList(); - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - - assertEquals("Unexpected number of messages received",10,messages.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received", MSG_COUNT/2 - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - - _consumerSession.close(); - _consumerConnection.close(); - - - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - - for (int msg = MSG_COUNT/2; msg < MSG_COUNT; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - } - - - // HACK to do something synchronous - ((AMQSession)_producerSession).sync(); - - _consumer = _consumerSession.createConsumer(_queue); - _consumerConnection.start(); - - messages = new ArrayList(); - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - - assertEquals("Unexpected number of messages received",10,messages.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - - } - - - public void testConflationWithReleaseAfterNewPublish() throws Exception - { - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.CLIENT_ACKNOWLEDGE); - - - createConflationQueue(_producerSession); - _producer = _producerSession.createProducer(_queue); - - for (int msg = 0; msg < MSG_COUNT/2; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - } - - // HACK to do something synchronous - ((AMQSession)_producerSession).sync(); - - _consumer = _consumerSession.createConsumer(_queue); - _consumerConnection.start(); - Message received; - List messages = new ArrayList(); - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - - assertEquals("Unexpected number of messages received",10,messages.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received", MSG_COUNT/2 - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - - _consumer.close(); - - for (int msg = MSG_COUNT/2; msg < MSG_COUNT; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - } - - // HACK to do something synchronous - ((AMQSession)_producerSession).sync(); - - - // this causes the "old" messages to be released - _consumerSession.close(); - _consumerConnection.close(); - - - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - - - _consumer = _consumerSession.createConsumer(_queue); - _consumerConnection.start(); - - messages = new ArrayList(); - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - - assertEquals("Unexpected number of messages received",10,messages.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - - } - - public void testConflatedQueueDepth() throws Exception - { - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - createConflationQueue(_producerSession); - _producer = _producerSession.createProducer(_queue); - - for (int msg = 0; msg < MSG_COUNT; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - } - - final long queueDepth = ((AMQSession)_producerSession).getQueueDepth((AMQDestination)_queue, true); - - assertEquals(10, queueDepth); - } - - public void testConflationBrowser() throws Exception - { - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - - createConflationQueue(_producerSession); - _producer = _producerSession.createProducer(_queue); - - for (int msg = 0; msg < MSG_COUNT; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - - } - - ((AMQSession)_producerSession).sync(); - - AMQBindingURL url = new AMQBindingURL("direct://amq.direct//"+_queueName+"?browse='true'&durable='true'"); - AMQQueue browseQueue = new AMQQueue(url); - - _consumer = _consumerSession.createConsumer(browseQueue); - _consumerConnection.start(); - Message received; - List messages = new ArrayList(); - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - - assertEquals("Unexpected number of messages received",10,messages.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - - messages.clear(); - - _producer.send(nextMessage(MSG_COUNT, _producerSession)); - - ((AMQSession)_producerSession).sync(); - - while((received = _consumer.receive(1000))!=null) - { - messages.add(received); - } - assertEquals("Unexpected number of messages received",1,messages.size()); - assertEquals("Unexpected message number received", MSG_COUNT, messages.get(0).getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - - - _producer.close(); - _producerSession.close(); - _producerConnection.close(); - } - - public void testConflation2Browsers() throws Exception - { - _consumerConnection = getConnection(); - _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - - createConflationQueue(_producerSession); - _producer = _producerSession.createProducer(_queue); - - for (int msg = 0; msg < MSG_COUNT; msg++) - { - _producer.send(nextMessage(msg, _producerSession)); - } - - ((AMQSession)_producerSession).sync(); - - AMQBindingURL url = new AMQBindingURL("direct://amq.direct//"+_queueName+"?browse='true'&durable='true'"); - AMQQueue browseQueue = new AMQQueue(url); - - _consumer = _consumerSession.createConsumer(browseQueue); - MessageConsumer consumer2 = _consumerSession.createConsumer(browseQueue); - _consumerConnection.start(); - List messages = new ArrayList(); - List messages2 = new ArrayList(); - Message received = _consumer.receive(1000); - Message received2 = consumer2.receive(1000); - - while(received!=null || received2!=null) - { - if(received != null) - { - messages.add(received); - } - if(received2 != null) - { - messages2.add(received2); - } - - - received = _consumer.receive(1000); - received2 = consumer2.receive(1000); - - } - - assertEquals("Unexpected number of messages received on first browser",10,messages.size()); - assertEquals("Unexpected number of messages received on second browser",10,messages2.size()); - - for(int i = 0 ; i < 10; i++) - { - Message msg = messages.get(i); - assertEquals("Unexpected message number received on first browser", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - msg = messages2.get(i); - assertEquals("Unexpected message number received on second browser", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); - } - - - _producer.close(); - _producerSession.close(); - _producerConnection.close(); - } - - public void testParallelProductionAndConsumption() throws Exception - { - createConflationQueue(_producerSession); - - // Start producing threads that send messages - BackgroundMessageProducer messageProducer1 = new BackgroundMessageProducer("Message sender1"); - messageProducer1.startSendingMessages(); - BackgroundMessageProducer messageProducer2 = new BackgroundMessageProducer("Message sender2"); - messageProducer2.startSendingMessages(); - - Map lastReceivedMessages = receiveMessages(messageProducer1); - - messageProducer1.join(); - messageProducer2.join(); - - final Map lastSentMessages1 = messageProducer1.getMessageSequenceNumbersByKey(); - assertEquals("Unexpected number of last sent messages sent by producer1", 2, lastSentMessages1.size()); - final Map lastSentMessages2 = messageProducer2.getMessageSequenceNumbersByKey(); - assertEquals(lastSentMessages1, lastSentMessages2); - - assertEquals("The last message sent for each key should match the last message received for that key", - lastSentMessages1, lastReceivedMessages); - - assertNull("Unexpected exception from background producer thread", messageProducer1.getException()); - } - - private Map receiveMessages(BackgroundMessageProducer producer) throws Exception - { - producer.waitUntilQuarterOfMessagesSentToEncourageConflation(); - - _consumerConnection = getConnection(); - int smallPrefetchToEncourageConflation = 1; - _consumerSession = ((AMQConnection)_consumerConnection).createSession(false, Session.AUTO_ACKNOWLEDGE, smallPrefetchToEncourageConflation); - - LOGGER.info("Starting to receive"); - - _consumer = _consumerSession.createConsumer(_queue); - _consumerConnection.start(); - - Map messageSequenceNumbersByKey = new HashMap(); - - Message message; - int numberOfShutdownsReceived = 0; - int numberOfMessagesReceived = 0; - while(numberOfShutdownsReceived < 2) - { - message = _consumer.receive(10000); - assertNotNull(message); - - if (message.propertyExists(BackgroundMessageProducer.SHUTDOWN)) - { - numberOfShutdownsReceived++; - } - else - { - numberOfMessagesReceived++; - putMessageInMap(message, messageSequenceNumbersByKey); - } - } - - LOGGER.info("Finished receiving. Received " + numberOfMessagesReceived + " message(s) in total"); - - return messageSequenceNumbersByKey; - } - - private void putMessageInMap(Message message, Map messageSequenceNumbersByKey) throws JMSException - { - String keyValue = message.getStringProperty(KEY_PROPERTY); - Integer messageSequenceNumber = message.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY); - messageSequenceNumbersByKey.put(keyValue, messageSequenceNumber); - } - - private class BackgroundMessageProducer - { - static final String SHUTDOWN = "SHUTDOWN"; - - private final String _threadName; - - private volatile Exception _exception; - - private Thread _thread; - private Map _messageSequenceNumbersByKey = new HashMap(); - private CountDownLatch _quarterOfMessagesSentLatch = new CountDownLatch(MSG_COUNT/4); - - public BackgroundMessageProducer(String threadName) - { - _threadName = threadName; - } - - public void waitUntilQuarterOfMessagesSentToEncourageConflation() throws InterruptedException - { - final long latchTimeout = 60000; - boolean success = _quarterOfMessagesSentLatch.await(latchTimeout, TimeUnit.MILLISECONDS); - assertTrue("Failed to be notified that 1/4 of the messages have been sent within " + latchTimeout + " ms.", success); - LOGGER.info("Quarter of messages sent"); - } - - public Exception getException() - { - return _exception; - } - - public Map getMessageSequenceNumbersByKey() - { - return Collections.unmodifiableMap(_messageSequenceNumbersByKey); - } - - public void startSendingMessages() - { - Runnable messageSender = new Runnable() - { - @Override - public void run() - { - try - { - LOGGER.info("Starting to send in background thread"); - Connection producerConnection = getConnection(); - Session producerSession = producerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); - - MessageProducer backgroundProducer = producerSession.createProducer(_queue); - for (int messageNumber = 0; messageNumber < MSG_COUNT; messageNumber++) - { - Message message = nextMessage(messageNumber, producerSession, 2); - backgroundProducer.send(message); - - putMessageInMap(message, _messageSequenceNumbersByKey); - _quarterOfMessagesSentLatch.countDown(); - } - - Message shutdownMessage = producerSession.createMessage(); - shutdownMessage.setBooleanProperty(SHUTDOWN, true); - backgroundProducer.send(shutdownMessage); - - LOGGER.info("Finished sending in background thread"); - } - catch (Exception e) - { - _exception = e; - throw new RuntimeException(e); - } - } - }; - - _thread = new Thread(messageSender); - _thread.setName(_threadName); - _thread.start(); - } - - public void join() throws InterruptedException - { - final int timeoutInMillis = 120000; - _thread.join(timeoutInMillis); - assertFalse("Expected producer thread to finish within " + timeoutInMillis + "ms", _thread.isAlive()); - } - } - - private void createConflationQueue(Session session) throws AMQException - { - final Map arguments = new HashMap(); - arguments.put("qpid.last_value_queue_key",KEY_PROPERTY); - ((AMQSession) session).createQueue(new AMQShortString(_queueName), false, true, false, arguments); - _queue = new AMQQueue("amq.direct", _queueName); - ((AMQSession) session).declareAndBind((AMQDestination)_queue); - } - - private Message nextMessage(int msg, Session producerSession) throws JMSException - { - return nextMessage(msg, producerSession, 10); - } - - private Message nextMessage(int msg, Session producerSession, int numberOfUniqueKeyValues) throws JMSException - { - Message send = producerSession.createTextMessage("Message: " + msg); - - send.setStringProperty(KEY_PROPERTY, String.valueOf(msg % numberOfUniqueKeyValues)); - send.setIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY, msg); - - return send; - } -} diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/server/queue/LastValueQueueTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/server/queue/LastValueQueueTest.java new file mode 100644 index 0000000000..dc30c02951 --- /dev/null +++ b/qpid/java/systests/src/main/java/org/apache/qpid/server/queue/LastValueQueueTest.java @@ -0,0 +1,574 @@ +/* + * + * 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.queue; + +import org.apache.log4j.Logger; +import org.apache.qpid.AMQException; +import org.apache.qpid.client.AMQConnection; +import org.apache.qpid.client.AMQDestination; +import org.apache.qpid.client.AMQQueue; +import org.apache.qpid.client.AMQSession; +import org.apache.qpid.framing.AMQShortString; +import org.apache.qpid.test.utils.QpidBrokerTestCase; +import org.apache.qpid.url.AMQBindingURL; + +import javax.jms.Connection; +import javax.jms.JMSException; +import javax.jms.Message; +import javax.jms.MessageConsumer; +import javax.jms.MessageProducer; +import javax.jms.Queue; +import javax.jms.Session; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +public class LastValueQueueTest extends QpidBrokerTestCase +{ + private static final Logger LOGGER = Logger.getLogger(LastValueQueueTest.class); + + private static final String MESSAGE_SEQUENCE_NUMBER_PROPERTY = "msg"; + private static final String KEY_PROPERTY = "key"; + + private static final int MSG_COUNT = 400; + + private String _queueName; + private Queue _queue; + private Connection _producerConnection; + private MessageProducer _producer; + private Session _producerSession; + private Connection _consumerConnection; + private Session _consumerSession; + private MessageConsumer _consumer; + + protected void setUp() throws Exception + { + super.setUp(); + + _queueName = getTestQueueName(); + _producerConnection = getConnection(); + _producerSession = _producerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + } + + public void testConflation() throws Exception + { + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + createConflationQueue(_producerSession); + _producer = _producerSession.createProducer(_queue); + + for (int msg = 0; msg < MSG_COUNT; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + } + + _producer.close(); + _producerSession.close(); + _producerConnection.close(); + + _consumer = _consumerSession.createConsumer(_queue); + _consumerConnection.start(); + Message received; + + List messages = new ArrayList(); + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + + assertEquals("Unexpected number of messages received",10,messages.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + } + + public void testConflationWithRelease() throws Exception + { + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.CLIENT_ACKNOWLEDGE); + + + createConflationQueue(_producerSession); + _producer = _producerSession.createProducer(_queue); + + for (int msg = 0; msg < MSG_COUNT/2; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + + } + + // HACK to do something synchronous + ((AMQSession)_producerSession).sync(); + + _consumer = _consumerSession.createConsumer(_queue); + _consumerConnection.start(); + Message received; + List messages = new ArrayList(); + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + + assertEquals("Unexpected number of messages received",10,messages.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received", MSG_COUNT/2 - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + + _consumerSession.close(); + _consumerConnection.close(); + + + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + + for (int msg = MSG_COUNT/2; msg < MSG_COUNT; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + } + + + // HACK to do something synchronous + ((AMQSession)_producerSession).sync(); + + _consumer = _consumerSession.createConsumer(_queue); + _consumerConnection.start(); + + messages = new ArrayList(); + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + + assertEquals("Unexpected number of messages received",10,messages.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + + } + + + public void testConflationWithReleaseAfterNewPublish() throws Exception + { + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.CLIENT_ACKNOWLEDGE); + + + createConflationQueue(_producerSession); + _producer = _producerSession.createProducer(_queue); + + for (int msg = 0; msg < MSG_COUNT/2; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + } + + // HACK to do something synchronous + ((AMQSession)_producerSession).sync(); + + _consumer = _consumerSession.createConsumer(_queue); + _consumerConnection.start(); + Message received; + List messages = new ArrayList(); + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + + assertEquals("Unexpected number of messages received",10,messages.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received", MSG_COUNT/2 - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + + _consumer.close(); + + for (int msg = MSG_COUNT/2; msg < MSG_COUNT; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + } + + // HACK to do something synchronous + ((AMQSession)_producerSession).sync(); + + + // this causes the "old" messages to be released + _consumerSession.close(); + _consumerConnection.close(); + + + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + + + _consumer = _consumerSession.createConsumer(_queue); + _consumerConnection.start(); + + messages = new ArrayList(); + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + + assertEquals("Unexpected number of messages received",10,messages.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + + } + + public void testConflatedQueueDepth() throws Exception + { + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + createConflationQueue(_producerSession); + _producer = _producerSession.createProducer(_queue); + + for (int msg = 0; msg < MSG_COUNT; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + } + + final long queueDepth = ((AMQSession)_producerSession).getQueueDepth((AMQDestination)_queue, true); + + assertEquals(10, queueDepth); + } + + public void testConflationBrowser() throws Exception + { + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + + createConflationQueue(_producerSession); + _producer = _producerSession.createProducer(_queue); + + for (int msg = 0; msg < MSG_COUNT; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + + } + + ((AMQSession)_producerSession).sync(); + + AMQBindingURL url = new AMQBindingURL("direct://amq.direct//"+_queueName+"?browse='true'&durable='true'"); + AMQQueue browseQueue = new AMQQueue(url); + + _consumer = _consumerSession.createConsumer(browseQueue); + _consumerConnection.start(); + Message received; + List messages = new ArrayList(); + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + + assertEquals("Unexpected number of messages received",10,messages.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + + messages.clear(); + + _producer.send(nextMessage(MSG_COUNT, _producerSession)); + + ((AMQSession)_producerSession).sync(); + + while((received = _consumer.receive(1000))!=null) + { + messages.add(received); + } + assertEquals("Unexpected number of messages received",1,messages.size()); + assertEquals("Unexpected message number received", MSG_COUNT, messages.get(0).getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + + + _producer.close(); + _producerSession.close(); + _producerConnection.close(); + } + + public void testConflation2Browsers() throws Exception + { + _consumerConnection = getConnection(); + _consumerSession = _consumerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + + createConflationQueue(_producerSession); + _producer = _producerSession.createProducer(_queue); + + for (int msg = 0; msg < MSG_COUNT; msg++) + { + _producer.send(nextMessage(msg, _producerSession)); + } + + ((AMQSession)_producerSession).sync(); + + AMQBindingURL url = new AMQBindingURL("direct://amq.direct//"+_queueName+"?browse='true'&durable='true'"); + AMQQueue browseQueue = new AMQQueue(url); + + _consumer = _consumerSession.createConsumer(browseQueue); + MessageConsumer consumer2 = _consumerSession.createConsumer(browseQueue); + _consumerConnection.start(); + List messages = new ArrayList(); + List messages2 = new ArrayList(); + Message received = _consumer.receive(1000); + Message received2 = consumer2.receive(1000); + + while(received!=null || received2!=null) + { + if(received != null) + { + messages.add(received); + } + if(received2 != null) + { + messages2.add(received2); + } + + + received = _consumer.receive(1000); + received2 = consumer2.receive(1000); + + } + + assertEquals("Unexpected number of messages received on first browser",10,messages.size()); + assertEquals("Unexpected number of messages received on second browser",10,messages2.size()); + + for(int i = 0 ; i < 10; i++) + { + Message msg = messages.get(i); + assertEquals("Unexpected message number received on first browser", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + msg = messages2.get(i); + assertEquals("Unexpected message number received on second browser", MSG_COUNT - 10 + i, msg.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY)); + } + + + _producer.close(); + _producerSession.close(); + _producerConnection.close(); + } + + public void testParallelProductionAndConsumption() throws Exception + { + createConflationQueue(_producerSession); + + // Start producing threads that send messages + BackgroundMessageProducer messageProducer1 = new BackgroundMessageProducer("Message sender1"); + messageProducer1.startSendingMessages(); + BackgroundMessageProducer messageProducer2 = new BackgroundMessageProducer("Message sender2"); + messageProducer2.startSendingMessages(); + + Map lastReceivedMessages = receiveMessages(messageProducer1); + + messageProducer1.join(); + messageProducer2.join(); + + final Map lastSentMessages1 = messageProducer1.getMessageSequenceNumbersByKey(); + assertEquals("Unexpected number of last sent messages sent by producer1", 2, lastSentMessages1.size()); + final Map lastSentMessages2 = messageProducer2.getMessageSequenceNumbersByKey(); + assertEquals(lastSentMessages1, lastSentMessages2); + + assertEquals("The last message sent for each key should match the last message received for that key", + lastSentMessages1, lastReceivedMessages); + + assertNull("Unexpected exception from background producer thread", messageProducer1.getException()); + } + + private Map receiveMessages(BackgroundMessageProducer producer) throws Exception + { + producer.waitUntilQuarterOfMessagesSentToEncourageConflation(); + + _consumerConnection = getConnection(); + int smallPrefetchToEncourageConflation = 1; + _consumerSession = ((AMQConnection)_consumerConnection).createSession(false, Session.AUTO_ACKNOWLEDGE, smallPrefetchToEncourageConflation); + + LOGGER.info("Starting to receive"); + + _consumer = _consumerSession.createConsumer(_queue); + _consumerConnection.start(); + + Map messageSequenceNumbersByKey = new HashMap(); + + Message message; + int numberOfShutdownsReceived = 0; + int numberOfMessagesReceived = 0; + while(numberOfShutdownsReceived < 2) + { + message = _consumer.receive(10000); + assertNotNull(message); + + if (message.propertyExists(BackgroundMessageProducer.SHUTDOWN)) + { + numberOfShutdownsReceived++; + } + else + { + numberOfMessagesReceived++; + putMessageInMap(message, messageSequenceNumbersByKey); + } + } + + LOGGER.info("Finished receiving. Received " + numberOfMessagesReceived + " message(s) in total"); + + return messageSequenceNumbersByKey; + } + + private void putMessageInMap(Message message, Map messageSequenceNumbersByKey) throws JMSException + { + String keyValue = message.getStringProperty(KEY_PROPERTY); + Integer messageSequenceNumber = message.getIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY); + messageSequenceNumbersByKey.put(keyValue, messageSequenceNumber); + } + + private class BackgroundMessageProducer + { + static final String SHUTDOWN = "SHUTDOWN"; + + private final String _threadName; + + private volatile Exception _exception; + + private Thread _thread; + private Map _messageSequenceNumbersByKey = new HashMap(); + private CountDownLatch _quarterOfMessagesSentLatch = new CountDownLatch(MSG_COUNT/4); + + public BackgroundMessageProducer(String threadName) + { + _threadName = threadName; + } + + public void waitUntilQuarterOfMessagesSentToEncourageConflation() throws InterruptedException + { + final long latchTimeout = 60000; + boolean success = _quarterOfMessagesSentLatch.await(latchTimeout, TimeUnit.MILLISECONDS); + assertTrue("Failed to be notified that 1/4 of the messages have been sent within " + latchTimeout + " ms.", success); + LOGGER.info("Quarter of messages sent"); + } + + public Exception getException() + { + return _exception; + } + + public Map getMessageSequenceNumbersByKey() + { + return Collections.unmodifiableMap(_messageSequenceNumbersByKey); + } + + public void startSendingMessages() + { + Runnable messageSender = new Runnable() + { + @Override + public void run() + { + try + { + LOGGER.info("Starting to send in background thread"); + Connection producerConnection = getConnection(); + Session producerSession = producerConnection.createSession(false, Session.AUTO_ACKNOWLEDGE); + + MessageProducer backgroundProducer = producerSession.createProducer(_queue); + for (int messageNumber = 0; messageNumber < MSG_COUNT; messageNumber++) + { + Message message = nextMessage(messageNumber, producerSession, 2); + backgroundProducer.send(message); + + putMessageInMap(message, _messageSequenceNumbersByKey); + _quarterOfMessagesSentLatch.countDown(); + } + + Message shutdownMessage = producerSession.createMessage(); + shutdownMessage.setBooleanProperty(SHUTDOWN, true); + backgroundProducer.send(shutdownMessage); + + LOGGER.info("Finished sending in background thread"); + } + catch (Exception e) + { + _exception = e; + throw new RuntimeException(e); + } + } + }; + + _thread = new Thread(messageSender); + _thread.setName(_threadName); + _thread.start(); + } + + public void join() throws InterruptedException + { + final int timeoutInMillis = 120000; + _thread.join(timeoutInMillis); + assertFalse("Expected producer thread to finish within " + timeoutInMillis + "ms", _thread.isAlive()); + } + } + + private void createConflationQueue(Session session) throws AMQException + { + final Map arguments = new HashMap(); + arguments.put("qpid.last_value_queue_key",KEY_PROPERTY); + ((AMQSession) session).createQueue(new AMQShortString(_queueName), false, true, false, arguments); + _queue = new AMQQueue("amq.direct", _queueName); + ((AMQSession) session).declareAndBind((AMQDestination)_queue); + } + + private Message nextMessage(int msg, Session producerSession) throws JMSException + { + return nextMessage(msg, producerSession, 10); + } + + private Message nextMessage(int msg, Session producerSession, int numberOfUniqueKeyValues) throws JMSException + { + Message send = producerSession.createTextMessage("Message: " + msg); + + send.setStringProperty(KEY_PROPERTY, String.valueOf(msg % numberOfUniqueKeyValues)); + send.setIntProperty(MESSAGE_SEQUENCE_NUMBER_PROPERTY, msg); + + return send; + } +} diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/server/store/VirtualHostMessageStoreTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/server/store/VirtualHostMessageStoreTest.java index d857bb4ce0..9949102af8 100644 --- a/qpid/java/systests/src/main/java/org/apache/qpid/server/store/VirtualHostMessageStoreTest.java +++ b/qpid/java/systests/src/main/java/org/apache/qpid/server/store/VirtualHostMessageStoreTest.java @@ -50,9 +50,11 @@ import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.protocol.v0_8.AMQMessage; import org.apache.qpid.server.protocol.v0_8.MessageMetaData; import org.apache.qpid.server.queue.AMQQueue; -import org.apache.qpid.server.queue.ConflationQueue; +import org.apache.qpid.server.queue.LastValueQueue; +import org.apache.qpid.server.queue.LastValueQueueImpl; import org.apache.qpid.server.queue.PriorityQueue; -import org.apache.qpid.server.queue.StandardQueue; +import org.apache.qpid.server.queue.PriorityQueueImpl; +import org.apache.qpid.server.queue.StandardQueueImpl; import org.apache.qpid.server.txn.AutoCommitTransaction; import org.apache.qpid.server.txn.ServerTransaction; import org.apache.qpid.server.util.BrokerTestHelper; @@ -553,18 +555,18 @@ public class VirtualHostMessageStoreTest extends QpidTestCase if (usePriority) { - assertEquals("Queue is no longer a Priority Queue", PriorityQueue.class, queue.getClass()); + assertEquals("Queue is no longer a Priority Queue", PriorityQueueImpl.class, queue.getClass()); assertEquals("Priority Queue does not have set priorities", - DEFAULT_PRIORTY_LEVEL, ((PriorityQueue) queue).getPriorities()); + DEFAULT_PRIORTY_LEVEL, ((PriorityQueueImpl) queue).getPriorities()); } else if (lastValueQueue) { - assertEquals("Queue is no longer a LastValue Queue", ConflationQueue.class, queue.getClass()); - assertEquals("LastValue Queue Key has changed", LVQ_KEY, ((ConflationQueue) queue).getConflationKey()); + assertEquals("Queue is no longer a LastValue Queue", LastValueQueueImpl.class, queue.getClass()); + assertEquals("LastValue Queue Key has changed", LVQ_KEY, ((LastValueQueueImpl) queue).getConflationKey()); } else { - assertEquals("Queue is not 'simple'", StandardQueue.class, queue.getClass()); + assertEquals("Queue is not 'simple'", StandardQueueImpl.class, queue.getClass()); } assertEquals("Queue owner is not as expected for queue " + queue.getName(), exclusive ? queueOwner : null, queue.getOwner()); @@ -660,12 +662,12 @@ public class VirtualHostMessageStoreTest extends QpidTestCase if (usePriority) { - queueArguments.put(Queue.PRIORITIES, DEFAULT_PRIORTY_LEVEL); + queueArguments.put(PriorityQueue.PRIORITIES, DEFAULT_PRIORTY_LEVEL); } if (lastValueQueue) { - queueArguments.put(Queue.LVQ_KEY, LVQ_KEY); + queueArguments.put(LastValueQueue.LVQ_KEY, LVQ_KEY); } queueArguments.put(Queue.ID, UUIDGenerator.generateRandomUUID()); diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/systest/management/jmx/QueueManagementTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/systest/management/jmx/QueueManagementTest.java index 931974942f..3bd91faa3e 100644 --- a/qpid/java/systests/src/main/java/org/apache/qpid/systest/management/jmx/QueueManagementTest.java +++ b/qpid/java/systests/src/main/java/org/apache/qpid/systest/management/jmx/QueueManagementTest.java @@ -28,7 +28,7 @@ import org.apache.qpid.management.common.mbeans.ManagedBroker; import org.apache.qpid.management.common.mbeans.ManagedQueue; import org.apache.qpid.server.queue.NotificationCheckTest; import org.apache.qpid.server.queue.QueueArgumentsConverter; -import org.apache.qpid.server.queue.StandardQueue; +import org.apache.qpid.server.queue.StandardQueueImpl; import org.apache.qpid.test.client.destination.AddressBasedDestinationTest; import org.apache.qpid.test.utils.JMXTestUtils; import org.apache.qpid.test.utils.QpidBrokerTestCase; @@ -660,7 +660,7 @@ public class QueueManagementTest extends QpidBrokerTestCase final Object messageGroupKey = "test"; final Map arguments = new HashMap(2); arguments.put(QueueArgumentsConverter.QPID_GROUP_HEADER_KEY, messageGroupKey); - arguments.put(QueueArgumentsConverter.QPID_SHARED_MSG_GROUP, StandardQueue.SHARED_MSG_GROUP_ARG_VALUE); + arguments.put(QueueArgumentsConverter.QPID_SHARED_MSG_GROUP, StandardQueueImpl.SHARED_MSG_GROUP_ARG_VALUE); managedBroker.createNewQueue(queueName, null, true, arguments); final ManagedQueue managedQueue = _jmxUtils.getManagedQueue(queueName); diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/Asserts.java b/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/Asserts.java index a2bb05dd00..47a25dbbec 100644 --- a/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/Asserts.java +++ b/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/Asserts.java @@ -43,6 +43,9 @@ import org.apache.qpid.server.model.Port; import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.model.State; import org.apache.qpid.server.model.VirtualHost; +import org.apache.qpid.server.queue.LastValueQueue; +import org.apache.qpid.server.queue.PriorityQueue; +import org.apache.qpid.server.queue.SortedQueue; import org.apache.qpid.test.utils.TestBrokerConfiguration; public class Asserts @@ -113,11 +116,11 @@ public class Asserts Queue.ALTERNATE_EXCHANGE, Queue.OWNER, Queue.NO_LOCAL, - Queue.LVQ_KEY, - Queue.SORT_KEY, + LastValueQueue.LVQ_KEY, + SortedQueue.SORT_KEY, Queue.MESSAGE_GROUP_KEY, Queue.MESSAGE_GROUP_SHARED_GROUPS, - Queue.PRIORITIES, + PriorityQueue.PRIORITIES, ConfiguredObject.CONTEXT); assertEquals("Unexpected value of queue attribute " + Queue.NAME, queueName, queueData.get(Queue.NAME)); @@ -127,9 +130,9 @@ public class Asserts queueData.get(Queue.STATE)); assertEquals("Unexpected value of queue attribute " + Queue.LIFETIME_POLICY, LifetimePolicy.PERMANENT.name(), queueData.get(Queue.LIFETIME_POLICY)); - assertEquals("Unexpected value of queue attribute " + Queue.QUEUE_TYPE, + assertEquals("Unexpected value of queue attribute " + Queue.TYPE, queueType, - queueData.get(Queue.QUEUE_TYPE)); + queueData.get(Queue.TYPE)); if (expectedAttributes == null) { assertEquals("Unexpected value of queue attribute " + Queue.EXCLUSIVE, diff --git a/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/VirtualHostRestTest.java b/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/VirtualHostRestTest.java index 421c609e46..4535425ea4 100644 --- a/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/VirtualHostRestTest.java +++ b/qpid/java/systests/src/main/java/org/apache/qpid/systest/rest/VirtualHostRestTest.java @@ -37,7 +37,10 @@ import org.apache.qpid.client.AMQConnection; import org.apache.qpid.server.model.Exchange; import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.model.VirtualHost; -import org.apache.qpid.server.queue.ConflationQueue; +import org.apache.qpid.server.queue.LastValueQueue; +import org.apache.qpid.server.queue.LastValueQueueImpl; +import org.apache.qpid.server.queue.PriorityQueue; +import org.apache.qpid.server.queue.SortedQueue; import org.apache.qpid.server.store.MessageStore; import org.apache.qpid.server.virtualhost.StandardVirtualHost; import org.apache.qpid.util.FileUtils; @@ -193,15 +196,15 @@ public class VirtualHostRestTest extends QpidRestTestCase createQueue(queueName + "-standard", "standard", null); Map sortedQueueAttributes = new HashMap(); - sortedQueueAttributes.put(Queue.SORT_KEY, "sortme"); + sortedQueueAttributes.put(SortedQueue.SORT_KEY, "sortme"); createQueue(queueName + "-sorted", "sorted", sortedQueueAttributes); Map priorityQueueAttributes = new HashMap(); - priorityQueueAttributes.put(Queue.PRIORITIES, 10); + priorityQueueAttributes.put(PriorityQueue.PRIORITIES, 10); createQueue(queueName + "-priority", "priority", priorityQueueAttributes); Map lvqQueueAttributes = new HashMap(); - lvqQueueAttributes.put(Queue.LVQ_KEY, "LVQ"); + lvqQueueAttributes.put(LastValueQueue.LVQ_KEY, "LVQ"); createQueue(queueName + "-lvq", "lvq", lvqQueueAttributes); Map hostDetails = getRestTestHelper().getJsonAsSingletonList("/rest/virtualhost/test"); @@ -223,9 +226,9 @@ public class VirtualHostRestTest extends QpidRestTestCase assertEquals("Unexpected value of queue attribute " + Queue.DURABLE, Boolean.TRUE, priorityQueue.get(Queue.DURABLE)); assertEquals("Unexpected value of queue attribute " + Queue.DURABLE, Boolean.TRUE, lvqQueue.get(Queue.DURABLE)); - assertEquals("Unexpected sorted key attribute", "sortme", sortedQueue.get(Queue.SORT_KEY)); - assertEquals("Unexpected lvq key attribute", "LVQ", lvqQueue.get(Queue.LVQ_KEY)); - assertEquals("Unexpected priorities key attribute", 10, priorityQueue.get(Queue.PRIORITIES)); + assertEquals("Unexpected sorted key attribute", "sortme", sortedQueue.get(SortedQueue.SORT_KEY)); + assertEquals("Unexpected lvq key attribute", "LVQ", lvqQueue.get(LastValueQueue.LVQ_KEY)); + assertEquals("Unexpected priorities key attribute", 10, priorityQueue.get(PriorityQueue.PRIORITIES)); } public void testPutCreateExchange() throws Exception @@ -271,7 +274,7 @@ public class VirtualHostRestTest extends QpidRestTestCase Asserts.assertQueue(queueName , "lvq", lvqQueue); assertEquals("Unexpected value of queue attribute " + Queue.DURABLE, Boolean.TRUE, lvqQueue.get(Queue.DURABLE)); - assertEquals("Unexpected lvq key attribute", ConflationQueue.DEFAULT_LVQ_KEY, lvqQueue.get(Queue.LVQ_KEY)); + assertEquals("Unexpected lvq key attribute", LastValueQueueImpl.DEFAULT_LVQ_KEY, lvqQueue.get(LastValueQueue.LVQ_KEY)); } public void testPutCreateSortedQueueWithoutKey() throws Exception @@ -302,7 +305,7 @@ public class VirtualHostRestTest extends QpidRestTestCase Asserts.assertQueue(queueName , "priority", priorityQueue); assertEquals("Unexpected value of queue attribute " + Queue.DURABLE, Boolean.TRUE, priorityQueue.get(Queue.DURABLE)); - assertEquals("Unexpected number of priorities", 10, priorityQueue.get(Queue.PRIORITIES)); + assertEquals("Unexpected number of priorities", 10, priorityQueue.get(PriorityQueue.PRIORITIES)); } public void testPutCreateStandardQueueWithoutType() throws Exception @@ -401,17 +404,17 @@ public class VirtualHostRestTest extends QpidRestTestCase Map sortedQueueAttributes = new HashMap(); sortedQueueAttributes.putAll(attributes); - sortedQueueAttributes.put(Queue.SORT_KEY, "sortme"); + sortedQueueAttributes.put(SortedQueue.SORT_KEY, "sortme"); createQueue(queueName + "-sorted", "sorted", sortedQueueAttributes); Map priorityQueueAttributes = new HashMap(); priorityQueueAttributes.putAll(attributes); - priorityQueueAttributes.put(Queue.PRIORITIES, 10); + priorityQueueAttributes.put(PriorityQueue.PRIORITIES, 10); createQueue(queueName + "-priority", "priority", priorityQueueAttributes); Map lvqQueueAttributes = new HashMap(); lvqQueueAttributes.putAll(attributes); - lvqQueueAttributes.put(Queue.LVQ_KEY, "LVQ"); + lvqQueueAttributes.put(LastValueQueue.LVQ_KEY, "LVQ"); createQueue(queueName + "-lvq", "lvq", lvqQueueAttributes); Map hostDetails = getRestTestHelper().getJsonAsSingletonList("/rest/virtualhost/test"); @@ -429,9 +432,9 @@ public class VirtualHostRestTest extends QpidRestTestCase Asserts.assertQueue(queueName + "-priority", "priority", priorityQueue, attributes); Asserts.assertQueue(queueName + "-lvq", "lvq", lvqQueue, attributes); - assertEquals("Unexpected sorted key attribute", "sortme", sortedQueue.get(Queue.SORT_KEY)); - assertEquals("Unexpected lvq key attribute", "LVQ", lvqQueue.get(Queue.LVQ_KEY)); - assertEquals("Unexpected priorities key attribute", 10, priorityQueue.get(Queue.PRIORITIES)); + assertEquals("Unexpected sorted key attribute", "sortme", sortedQueue.get(SortedQueue.SORT_KEY)); + assertEquals("Unexpected lvq key attribute", "LVQ", lvqQueue.get(LastValueQueue.LVQ_KEY)); + assertEquals("Unexpected priorities key attribute", 10, priorityQueue.get(PriorityQueue.PRIORITIES)); } @SuppressWarnings("unchecked") @@ -506,7 +509,7 @@ public class VirtualHostRestTest extends QpidRestTestCase queueData.put(Queue.DURABLE, Boolean.TRUE); if (queueType != null) { - queueData.put(Queue.QUEUE_TYPE, queueType); + queueData.put(Queue.TYPE, queueType); } if (attributes != null) { -- cgit v1.2.1