diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2013-08-18 09:13:02 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2013-08-18 09:13:02 +0000 |
| commit | ab6fffad2230229810c995253a6f021e42e03aaf (patch) | |
| tree | fdee7a99130750af8d7c71d25c358a282e17e405 /qpid/java/broker/src/main | |
| parent | 35b5c7fd8c761d41caa88505e8c2fee319e92a84 (diff) | |
| download | qpid-python-ab6fffad2230229810c995253a6f021e42e03aaf.tar.gz | |
QPID-5081 : [Java Broker] Refactor Queue Creation
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1515079 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker/src/main')
24 files changed, 618 insertions, 257 deletions
diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java index 53dd6df599..631490ab5f 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java @@ -167,11 +167,6 @@ public abstract class AbstractExchange implements Exchange return _virtualHost; } - public QueueRegistry getQueueRegistry() - { - return getVirtualHost().getQueueRegistry(); - } - public final boolean isBound(AMQShortString routingKey, FieldTable ft, AMQQueue queue) { return isBound(routingKey == null ? "" : routingKey.asString(), FieldTable.convertToMap(ft), queue); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java index 2873eb31e8..8e9f980e6b 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java @@ -47,6 +47,7 @@ import org.apache.qpid.server.virtualhost.VirtualHost; public class DefaultExchange implements Exchange { + private final QueueRegistry _queueRegistry; private UUID _id; private VirtualHost _virtualHost; private static final Logger _logger = Logger.getLogger(DefaultExchange.class); @@ -55,6 +56,11 @@ public class DefaultExchange implements Exchange private LogSubject _logSubject; private Map<ExchangeReferrer,Object> _referrers = new ConcurrentHashMap<ExchangeReferrer,Object>(); + public DefaultExchange(QueueRegistry queueRegistry) + { + _queueRegistry = queueRegistry; + } + @Override public void initialise(UUID id, @@ -82,7 +88,7 @@ public class DefaultExchange implements Exchange @Override public long getBindingCount() { - return _virtualHost.getQueueRegistry().getQueues().size(); + return _virtualHost.getQueues().size(); } @Override @@ -146,7 +152,7 @@ public class DefaultExchange implements Exchange @Override public Binding getBinding(String bindingKey, AMQQueue queue, Map<String, Object> arguments) { - if(_virtualHost.getQueueRegistry().getQueue(bindingKey) == queue && (arguments == null || arguments.isEmpty())) + if(_virtualHost.getQueue(bindingKey) == queue && (arguments == null || arguments.isEmpty())) { return convertToBinding(queue); } @@ -207,7 +213,7 @@ public class DefaultExchange implements Exchange @Override public List<AMQQueue> route(InboundMessage message) { - AMQQueue q = _virtualHost.getQueueRegistry().getQueue(message.getRoutingKey()); + AMQQueue q = _virtualHost.getQueue(message.getRoutingKey()); if(q == null) { List<AMQQueue> noQueues = Collections.emptyList(); @@ -235,13 +241,13 @@ public class DefaultExchange implements Exchange @Override public boolean isBound(AMQShortString routingKey) { - return _virtualHost.getQueueRegistry().getQueue(routingKey) != null; + return _virtualHost.getQueue(routingKey == null ? null : routingKey.toString()) != null; } @Override public boolean isBound(AMQQueue queue) { - return _virtualHost.getQueueRegistry().getQueue(queue.getName()) == queue; + return _virtualHost.getQueue(queue.getName()) == queue; } @Override @@ -283,7 +289,7 @@ public class DefaultExchange implements Exchange @Override public boolean isBound(String bindingKey) { - return _virtualHost.getQueueRegistry().getQueue(bindingKey) != null; + return _virtualHost.getQueue(bindingKey) != null; } @Override @@ -320,7 +326,7 @@ public class DefaultExchange implements Exchange public Collection<Binding> getBindings() { List<Binding> bindings = new ArrayList<Binding>(); - for(AMQQueue q : _virtualHost.getQueueRegistry().getQueues()) + for(AMQQueue q : _virtualHost.getQueues()) { bindings.add(convertToBinding(q)); } @@ -330,7 +336,7 @@ public class DefaultExchange implements Exchange @Override public void addBindingListener(BindingListener listener) { - _virtualHost.getQueueRegistry().addRegistryChangeListener(convertListener(listener));//To change body of implemented methods use File | Settings | File Templates. + _queueRegistry.addRegistryChangeListener(convertListener(listener)); } private QueueRegistry.RegistryChangeListener convertListener(final BindingListener listener) diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeRegistry.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeRegistry.java index 75c489c731..d8263a3c80 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeRegistry.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/exchange/DefaultExchangeRegistry.java @@ -27,6 +27,7 @@ import org.apache.qpid.exchange.ExchangeDefaults; import org.apache.qpid.protocol.AMQConstant; import org.apache.qpid.server.model.UUIDGenerator; import org.apache.qpid.server.plugin.ExchangeType; +import org.apache.qpid.server.queue.QueueRegistry; import org.apache.qpid.server.store.DurableConfigurationStore; import org.apache.qpid.server.virtualhost.VirtualHost; @@ -40,20 +41,23 @@ import java.util.concurrent.ConcurrentMap; public class DefaultExchangeRegistry implements ExchangeRegistry { private static final Logger LOGGER = Logger.getLogger(DefaultExchangeRegistry.class); - /** * Maps from exchange name to exchange instance */ private ConcurrentMap<String, Exchange> _exchangeMap = new ConcurrentHashMap<String, Exchange>(); private Exchange _defaultExchange; - private VirtualHost _host; + + private final VirtualHost _host; + private final QueueRegistry _queueRegistry; + private final Collection<RegistryChangeListener> _listeners = Collections.synchronizedCollection(new ArrayList<RegistryChangeListener>()); - public DefaultExchangeRegistry(VirtualHost host) + public DefaultExchangeRegistry(VirtualHost host, QueueRegistry queueRegistry) { _host = host; + _queueRegistry = queueRegistry; } public void initialise(ExchangeFactory exchangeFactory) throws AMQException @@ -61,7 +65,7 @@ public class DefaultExchangeRegistry implements ExchangeRegistry //create 'standard' exchanges: new ExchangeInitialiser().initialise(exchangeFactory, this, getDurableConfigurationStore()); - _defaultExchange = new DefaultExchange(); + _defaultExchange = new DefaultExchange(_queueRegistry); UUID defaultExchangeId = UUIDGenerator.generateExchangeUUID(ExchangeDefaults.DEFAULT_EXCHANGE_NAME.asString(), _host.getName()); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/Queue.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/Queue.java index 6fe0607ab2..ae2031bd71 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/Queue.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/Queue.java @@ -89,6 +89,7 @@ public interface Queue extends ConfiguredObject public static final String EXCLUSIVE = "exclusive"; public static final String MESSAGE_GROUP_KEY = "messageGroupKey"; public static final String MESSAGE_GROUP_SHARED_GROUPS = "messageGroupSharedGroups"; + public static final String MESSAGE_GROUP_DEFAULT_GROUP = "messageGroupDefaultGroup"; public static final String LVQ_KEY = "lvqKey"; public static final String MAXIMUM_DELIVERY_ATTEMPTS = "maximumDeliveryAttempts"; public static final String NO_LOCAL = "noLocal"; @@ -100,6 +101,10 @@ public interface Queue extends ConfiguredObject public static final String TYPE = "type"; public static final String PRIORITIES = "priorities"; + public static final String CREATE_DLQ_ON_CREATION = "x-qpid-dlq-enabled"; // TODO - this value should change + + public static final String FEDERATION_EXCLUDES = "federationExcludes"; + public static final String FEDERATION_ID = "federationId"; public static final Collection<String> AVAILABLE_ATTRIBUTES = @@ -134,6 +139,7 @@ public interface Queue extends ConfiguredObject PRIORITIES )); + //children Collection<Binding> getBindings(); Collection<Consumer> getConsumers(); @@ -144,6 +150,6 @@ public interface Queue extends ConfiguredObject void visit(QueueEntryVisitor visitor); void delete(); - + void setNotificationListener(QueueNotificationListener listener); } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/VirtualHost.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/VirtualHost.java index a84a041b72..26ac99d5bd 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/VirtualHost.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/VirtualHost.java @@ -120,7 +120,7 @@ public interface VirtualHost extends ConfiguredObject QUEUE_ALERT_THRESHOLD_QUEUE_DEPTH_MESSAGES, CONFIG_PATH)); - int CURRENT_CONFIG_VERSION = 2; + int CURRENT_CONFIG_VERSION = 3; //children Collection<VirtualHostAlias> getAliases(); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/QueueAdapter.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/QueueAdapter.java index 157b97cc07..96a7eacb92 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/QueueAdapter.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/QueueAdapter.java @@ -66,25 +66,6 @@ final class QueueAdapter extends AbstractAdapter implements Queue, AMQQueue.Subs put(DESCRIPTION, String.class); }}); - static final Map<String, String> ATTRIBUTE_MAPPINGS = new HashMap<String, String>(); - static - { - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.ALERT_REPEAT_GAP, AMQQueueFactory.X_QPID_MINIMUM_ALERT_REPEAT_GAP); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.ALERT_THRESHOLD_MESSAGE_AGE, AMQQueueFactory.X_QPID_MAXIMUM_MESSAGE_AGE); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.ALERT_THRESHOLD_MESSAGE_SIZE, AMQQueueFactory.X_QPID_MAXIMUM_MESSAGE_SIZE); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.ALERT_THRESHOLD_QUEUE_DEPTH_MESSAGES, AMQQueueFactory.X_QPID_MAXIMUM_MESSAGE_COUNT); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.ALERT_THRESHOLD_QUEUE_DEPTH_BYTES, AMQQueueFactory.X_QPID_MAXIMUM_QUEUE_DEPTH); - - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.MAXIMUM_DELIVERY_ATTEMPTS, AMQQueueFactory.X_QPID_MAXIMUM_DELIVERY_COUNT); - - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.QUEUE_FLOW_CONTROL_SIZE_BYTES, AMQQueueFactory.X_QPID_CAPACITY); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.QUEUE_FLOW_RESUME_SIZE_BYTES, AMQQueueFactory.X_QPID_FLOW_RESUME_CAPACITY); - - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.SORT_KEY, AMQQueueFactory.QPID_QUEUE_SORT_KEY); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.LVQ_KEY, AMQQueueFactory.QPID_LAST_VALUE_QUEUE_KEY); - QueueAdapter.ATTRIBUTE_MAPPINGS.put(Queue.PRIORITIES, AMQQueueFactory.X_QPID_PRIORITIES); - } - private final AMQQueue _queue; private final Map<Binding, BindingAdapter> _bindingAdapters = new HashMap<Binding, BindingAdapter>(); @@ -190,15 +171,7 @@ final class QueueAdapter extends AbstractAdapter implements Queue, AMQQueue.Subs { try { - QueueRegistry queueRegistry = _queue.getVirtualHost().getQueueRegistry(); - synchronized(queueRegistry) - { - _queue.delete(); - if (_queue.isDurable()) - { - DurableConfigurationStoreHelper.removeQueue(_queue.getVirtualHost().getDurableConfigurationStore(), _queue); - } - } + _queue.getVirtualHost().removeQueue(_queue); } catch(AMQException e) { @@ -414,13 +387,12 @@ final class QueueAdapter extends AbstractAdapter implements Queue, AMQQueue.Subs } else if(MESSAGE_GROUP_KEY.equals(name)) { - return _queue.getArguments().get(SimpleAMQQueue.QPID_GROUP_HEADER_KEY); + return _queue.getAttribute(MESSAGE_GROUP_KEY); } else if(MESSAGE_GROUP_SHARED_GROUPS.equals(name)) { //We only return the boolean value if message groups are actually in use - return getAttribute(MESSAGE_GROUP_KEY) == null ? null : - SimpleAMQQueue.SHARED_MSG_GROUP_ARG_VALUE.equals(_queue.getArguments().get(SimpleAMQQueue.QPID_SHARED_MSG_GROUP)); + return getAttribute(MESSAGE_GROUP_KEY) == null ? null : _queue.getAttribute(MESSAGE_GROUP_SHARED_GROUPS); } else if(LVQ_KEY.equals(name)) { diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java index c09dd9449e..977fd5ae56 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/model/adapter/VirtualHostAdapter.java @@ -41,12 +41,9 @@ import org.apache.commons.configuration.PropertiesConfiguration; import org.apache.commons.configuration.SystemConfiguration; import org.apache.log4j.Logger; import org.apache.qpid.AMQException; -import org.apache.qpid.framing.FieldTable; import org.apache.qpid.server.configuration.IllegalConfigurationException; import org.apache.qpid.server.configuration.VirtualHostConfiguration; import org.apache.qpid.server.configuration.XmlConfigurationUtilities.MyConfiguration; -import org.apache.qpid.server.connection.IConnectionRegistry; -import org.apache.qpid.server.exchange.ExchangeRegistry; import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.model.Broker; import org.apache.qpid.server.model.ConfiguredObject; @@ -68,6 +65,7 @@ import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.protocol.AMQConnectionModel; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.AMQQueueFactory; +import org.apache.qpid.server.queue.QueueArgumentsConverter; import org.apache.qpid.server.queue.QueueEntry; import org.apache.qpid.server.queue.QueueRegistry; import org.apache.qpid.server.queue.SimpleAMQQueue; @@ -86,6 +84,7 @@ import org.apache.qpid.server.virtualhost.ReservedExchangeNameException; import org.apache.qpid.server.virtualhost.UnknownExchangeException; import org.apache.qpid.server.virtualhost.VirtualHostListener; import org.apache.qpid.server.virtualhost.VirtualHostRegistry; +import org.apache.qpid.server.virtualhost.plugins.QueueExistsException; public final class VirtualHostAdapter extends AbstractAdapter implements VirtualHost, VirtualHostListener { @@ -203,7 +202,7 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual private void populateQueues() { - Collection<AMQQueue> actualQueues = _virtualHost.getQueueRegistry().getQueues(); + Collection<AMQQueue> actualQueues = _virtualHost.getQueues(); if ( actualQueues != null ) { synchronized(_queueAdapters) @@ -399,7 +398,7 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual } if (queueType == QueueType.LVQ && attributes.get(Queue.LVQ_KEY) == null) { - attributes.put(Queue.LVQ_KEY, AMQQueueFactory.QPID_LVQ_KEY); + attributes.put(Queue.LVQ_KEY, AMQQueueFactory.QPID_DEFAULT_LVQ_KEY); } else if (queueType == QueueType.PRIORITY && attributes.get(Queue.PRIORITIES) == null) { @@ -415,7 +414,7 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual { String key = MapValueConverter.getStringAttribute(Queue.MESSAGE_GROUP_KEY, attributes); attributes.remove(Queue.MESSAGE_GROUP_KEY); - attributes.put(SimpleAMQQueue.QPID_GROUP_HEADER_KEY, key); + attributes.put(QueueArgumentsConverter.QPID_GROUP_HEADER_KEY, key); } if (attributes.containsKey(Queue.MESSAGE_GROUP_SHARED_GROUPS)) @@ -423,7 +422,7 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual if(MapValueConverter.getBooleanAttribute(Queue.MESSAGE_GROUP_SHARED_GROUPS, attributes)) { attributes.remove(Queue.MESSAGE_GROUP_SHARED_GROUPS); - attributes.put(SimpleAMQQueue.QPID_SHARED_MSG_GROUP, SimpleAMQQueue.SHARED_MSG_GROUP_ARG_VALUE); + attributes.put(QueueArgumentsConverter.QPID_SHARED_MSG_GROUP, SimpleAMQQueue.SHARED_MSG_GROUP_ARG_VALUE); } } @@ -440,15 +439,6 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual attributes.remove(Queue.LIFETIME_POLICY); attributes.remove(Queue.TIME_TO_LIVE); - List<String> attrNames = new ArrayList<String>(attributes.keySet()); - for(String attr : attrNames) - { - if(QueueAdapter.ATTRIBUTE_MAPPINGS.containsKey(attr)) - { - attributes.put(QueueAdapter.ATTRIBUTE_MAPPINGS.get(attr),attributes.remove(attr)); - } - } - return createQueue(name, state, durable, exclusive, lifetime, ttl, attributes); } @@ -472,33 +462,26 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual owner = authenticatedPrincipal.getName(); } } + + final boolean autoDelete = lifetime == LifetimePolicy.AUTO_DELETE; + try { - QueueRegistry queueRegistry = _virtualHost.getQueueRegistry(); - synchronized (queueRegistry) - { - if(_virtualHost.getQueueRegistry().getQueue(name)!=null) - { - throw new IllegalArgumentException("Queue with name "+name+" already exists"); - } - AMQQueue queue = - AMQQueueFactory.createAMQQueueImpl(UUIDGenerator.generateQueueUUID(name, _virtualHost.getName()), name, - durable, owner, lifetime == LifetimePolicy.AUTO_DELETE, - exclusive, _virtualHost, attributes); - if(durable) - { - DurableConfigurationStoreHelper.createQueue(_virtualHost.getDurableConfigurationStore(), - queue, - FieldTable.convertToFieldTable(attributes)); - } - synchronized (_queueAdapters) - { - return _queueAdapters.get(queue); - } + AMQQueue queue = + _virtualHost.createQueue(UUIDGenerator.generateQueueUUID(name, _virtualHost.getName()), name, + durable, owner, autoDelete, exclusive, autoDelete && exclusive, attributes); + + synchronized (_queueAdapters) + { + return _queueAdapters.get(queue); } } + catch(QueueExistsException qe) + { + throw new IllegalArgumentException("Queue with name "+name+" already exists"); + } catch(AMQException e) { throw new IllegalArgumentException(e); @@ -1057,7 +1040,7 @@ public final class VirtualHostAdapter extends AbstractAdapter implements Virtual { if(VirtualHost.QUEUE_COUNT.equals(name)) { - return _vhost.getQueueRegistry().getQueues().size(); + return _vhost.getQueues().size(); } else if(VirtualHost.EXCHANGE_COUNT.equals(name)) { diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueue.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueue.java index 4f610cc925..cb6a9249d3 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueue.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueue.java @@ -225,7 +225,8 @@ public interface AMQQueue extends Comparable<AMQQueue>, ExchangeReferrer, Transa void setAlternateExchange(Exchange exchange); - Map<String, Object> getArguments(); + Collection<String> getAvailableAttributes(); + Object getAttribute(String attrName); void checkCapacity(AMQSessionModel channel); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueueFactory.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueueFactory.java index 1eeb6dccf3..5001c2fd2b 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueueFactory.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/AMQQueueFactory.java @@ -29,42 +29,31 @@ import org.apache.qpid.AMQException; import org.apache.qpid.AMQSecurityException; import org.apache.qpid.exchange.ExchangeDefaults; import org.apache.qpid.framing.AMQShortString; -import org.apache.qpid.framing.FieldTable; import org.apache.qpid.server.configuration.BrokerProperties; import org.apache.qpid.server.configuration.QueueConfiguration; import org.apache.qpid.server.exchange.DefaultExchangeFactory; import org.apache.qpid.server.exchange.Exchange; -import org.apache.qpid.server.exchange.ExchangeFactory; -import org.apache.qpid.server.exchange.ExchangeRegistry; +import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.model.UUIDGenerator; import org.apache.qpid.server.store.DurableConfigurationStoreHelper; import org.apache.qpid.server.virtualhost.ExchangeExistsException; import org.apache.qpid.server.virtualhost.VirtualHost; -public class AMQQueueFactory +public class AMQQueueFactory implements QueueFactory { - public static final String X_QPID_FLOW_RESUME_CAPACITY = "x-qpid-flow-resume-capacity"; - public static final String X_QPID_CAPACITY = "x-qpid-capacity"; - public static final String X_QPID_MINIMUM_ALERT_REPEAT_GAP = "x-qpid-minimum-alert-repeat-gap"; - public static final String X_QPID_MAXIMUM_MESSAGE_COUNT = "x-qpid-maximum-message-count"; - public static final String X_QPID_MAXIMUM_MESSAGE_SIZE = "x-qpid-maximum-message-size"; - public static final String X_QPID_MAXIMUM_MESSAGE_AGE = "x-qpid-maximum-message-age"; - public static final String X_QPID_MAXIMUM_QUEUE_DEPTH = "x-qpid-maximum-queue-depth"; - - public static final String X_QPID_PRIORITIES = "x-qpid-priorities"; - public static final String X_QPID_DESCRIPTION = "x-qpid-description"; - public static final String QPID_LVQ_KEY = "qpid.LVQ_key"; - public static final String QPID_LAST_VALUE_QUEUE = "qpid.last_value_queue"; - public static final String QPID_LAST_VALUE_QUEUE_KEY = "qpid.last_value_queue_key"; - public static final String QPID_QUEUE_SORT_KEY = "qpid.queue_sort_key"; + public static final String QPID_DEFAULT_LVQ_KEY = "qpid.LVQ_key"; + - public static final String DLQ_ROUTING_KEY = "dlq"; - public static final String X_QPID_DLQ_ENABLED = "x-qpid-dlq-enabled"; - public static final String X_QPID_MAXIMUM_DELIVERY_COUNT = "x-qpid-maximum-delivery-count"; public static final String DEFAULT_DLQ_NAME_SUFFIX = "_DLQ"; + public static final String DLQ_ROUTING_KEY = "dlq"; + + private final VirtualHost _virtualHost; + private final QueueRegistry _queueRegistry; - private AMQQueueFactory() + public AMQQueueFactory(VirtualHost virtualHost, QueueRegistry queueRegistry) { + _virtualHost = virtualHost; + _queueRegistry = queueRegistry; } private abstract static class QueueProperty @@ -129,56 +118,56 @@ public class AMQQueueFactory } private static final QueueProperty[] DECLAREABLE_PROPERTIES = { - new QueueLongProperty(X_QPID_MAXIMUM_MESSAGE_AGE) + new QueueLongProperty(Queue.ALERT_THRESHOLD_MESSAGE_AGE) { public void setPropertyValue(AMQQueue queue, long value) { queue.setMaximumMessageAge(value); } }, - new QueueLongProperty(X_QPID_MAXIMUM_MESSAGE_SIZE) + new QueueLongProperty(Queue.ALERT_THRESHOLD_MESSAGE_SIZE) { public void setPropertyValue(AMQQueue queue, long value) { queue.setMaximumMessageSize(value); } }, - new QueueLongProperty(X_QPID_MAXIMUM_MESSAGE_COUNT) + new QueueLongProperty(Queue.ALERT_THRESHOLD_QUEUE_DEPTH_MESSAGES) { public void setPropertyValue(AMQQueue queue, long value) { queue.setMaximumMessageCount(value); } }, - new QueueLongProperty(X_QPID_MAXIMUM_QUEUE_DEPTH) + new QueueLongProperty(Queue.ALERT_THRESHOLD_QUEUE_DEPTH_BYTES) { public void setPropertyValue(AMQQueue queue, long value) { queue.setMaximumQueueDepth(value); } }, - new QueueLongProperty(X_QPID_MINIMUM_ALERT_REPEAT_GAP) + new QueueLongProperty(Queue.ALERT_REPEAT_GAP) { public void setPropertyValue(AMQQueue queue, long value) { queue.setMinimumAlertRepeatGap(value); } }, - new QueueLongProperty(X_QPID_CAPACITY) + new QueueLongProperty(Queue.QUEUE_FLOW_CONTROL_SIZE_BYTES) { public void setPropertyValue(AMQQueue queue, long value) { queue.setCapacity(value); } }, - new QueueLongProperty(X_QPID_FLOW_RESUME_CAPACITY) + new QueueLongProperty(Queue.QUEUE_FLOW_RESUME_SIZE_BYTES) { public void setPropertyValue(AMQQueue queue, long value) { queue.setFlowResumeCapacity(value); } }, - new QueueIntegerProperty(X_QPID_MAXIMUM_DELIVERY_COUNT) + new QueueIntegerProperty(Queue.MAXIMUM_DELIVERY_ATTEMPTS) { public void setPropertyValue(AMQQueue queue, int value) { @@ -189,13 +178,17 @@ public class AMQQueueFactory /** * @param id the id to use. + * @param deleteOnNoConsumer */ - public static AMQQueue createAMQQueueImpl(UUID id, - String queueName, - boolean durable, - String owner, - boolean autoDelete, - boolean exclusive, VirtualHost virtualHost, Map<String, Object> arguments) throws AMQSecurityException, AMQException + @Override + public AMQQueue createAMQQueueImpl(UUID id, + String queueName, + boolean durable, + String owner, + boolean autoDelete, + boolean exclusive, + boolean deleteOnNoConsumer, + Map<String, Object> arguments) throws AMQSecurityException, AMQException { if (id == null) { @@ -206,16 +199,11 @@ public class AMQQueueFactory throw new IllegalArgumentException("Queue name must not be null"); } - // Access check - if (!virtualHost.getSecurityManager().authoriseCreateQueue(autoDelete, durable, exclusive, null, null, new AMQShortString(queueName), owner)) - { - String description = "Permission denied: queue-name '" + queueName + "'"; - throw new AMQSecurityException(description); - } - QueueConfiguration queueConfiguration = virtualHost.getConfiguration().getQueueConfiguration(queueName); - boolean isDLQEnabled = isDLQEnabled(autoDelete, arguments, queueConfiguration); - if (isDLQEnabled) + QueueConfiguration queueConfiguration = _virtualHost.getConfiguration().getQueueConfiguration(queueName); + + boolean createDLQ = createDLQ(autoDelete, arguments, queueConfiguration); + if (createDLQ) { validateDLNames(queueName); } @@ -226,17 +214,17 @@ public class AMQQueueFactory if(arguments != null) { - if(arguments.containsKey(QPID_LAST_VALUE_QUEUE) || arguments.containsKey(QPID_LAST_VALUE_QUEUE_KEY)) + if(arguments.containsKey(Queue.LVQ_KEY)) { - conflationKey = (String) arguments.get(QPID_LAST_VALUE_QUEUE_KEY); + conflationKey = (String) arguments.get(Queue.LVQ_KEY); if(conflationKey == null) { - conflationKey = QPID_LVQ_KEY; + conflationKey = QPID_DEFAULT_LVQ_KEY; } } - else if(arguments.containsKey(X_QPID_PRIORITIES)) + else if(arguments.containsKey(Queue.PRIORITIES)) { - Object prioritiesObj = arguments.get(X_QPID_PRIORITIES); + Object prioritiesObj = arguments.get(Queue.PRIORITIES); if(prioritiesObj instanceof Number) { priorities = ((Number)prioritiesObj).intValue(); @@ -257,33 +245,36 @@ public class AMQQueueFactory // TODO - should warn here of invalid format } } - else if(arguments.containsKey(QPID_QUEUE_SORT_KEY)) + else if(arguments.containsKey(Queue.SORT_KEY)) { - sortingKey = (String)arguments.get(QPID_QUEUE_SORT_KEY); + sortingKey = (String)arguments.get(Queue.SORT_KEY); } } AMQQueue q; if(sortingKey != null) { - q = new SortedQueue(id, queueName, durable, owner, autoDelete, exclusive, virtualHost, arguments, sortingKey); + q = new SortedQueue(id, queueName, durable, owner, autoDelete, exclusive, _virtualHost, arguments, sortingKey); } else if(conflationKey != null) { - q = new ConflationQueue(id, queueName, durable, owner, autoDelete, exclusive, virtualHost, arguments, conflationKey); + q = new ConflationQueue(id, queueName, durable, owner, autoDelete, exclusive, _virtualHost, arguments, conflationKey); } else if(priorities > 1) { - q = new AMQPriorityQueue(id, queueName, durable, owner, autoDelete, exclusive, virtualHost, arguments, priorities); + q = new AMQPriorityQueue(id, queueName, durable, owner, autoDelete, exclusive, _virtualHost, arguments, priorities); } else { - q = new SimpleAMQQueue(id, queueName, durable, owner, autoDelete, exclusive, virtualHost, arguments); + q = new SimpleAMQQueue(id, queueName, durable, owner, autoDelete, exclusive, _virtualHost, arguments); } + q.setDeleteOnNoConsumers(deleteOnNoConsumer); + //Register the new queue - virtualHost.getQueueRegistry().registerQueue(q); - q.configure(virtualHost.getConfiguration().getQueueConfiguration(queueName)); + _queueRegistry.registerQueue(q); + + q.configure(_virtualHost.getConfiguration().getQueueConfiguration(queueName)); if(arguments != null) { @@ -294,21 +285,25 @@ public class AMQQueueFactory p.setPropertyValue(q, arguments.get(p.getArgumentName().toString())); } } + + if(arguments.get(Queue.NO_LOCAL) instanceof Boolean) + { + q.setNoLocal((Boolean)arguments.get(Queue.NO_LOCAL)); + } + } - if(isDLQEnabled) + if(createDLQ) { final String dlExchangeName = getDeadLetterExchangeName(queueName); final String dlQueueName = getDeadLetterQueueName(queueName); - final QueueRegistry queueRegistry = virtualHost.getQueueRegistry(); - Exchange dlExchange = null; - final UUID dlExchangeId = UUIDGenerator.generateExchangeUUID(dlExchangeName, virtualHost.getName()); + final UUID dlExchangeId = UUIDGenerator.generateExchangeUUID(dlExchangeName, _virtualHost.getName()); try { - dlExchange = virtualHost.createExchange(dlExchangeId, + dlExchange = _virtualHost.createExchange(dlExchangeId, dlExchangeName, ExchangeDefaults.FANOUT_EXCHANGE_CLASS.toString(), true, false, null); @@ -321,23 +316,19 @@ public class AMQQueueFactory AMQQueue dlQueue = null; - synchronized(queueRegistry) + synchronized(_queueRegistry) { - dlQueue = queueRegistry.getQueue(dlQueueName); + dlQueue = _queueRegistry.getQueue(dlQueueName); if(dlQueue == null) { //set args to disable DLQ'ing/MDC from the DLQ itself, preventing loops etc final Map<String, Object> args = new HashMap<String, Object>(); - args.put(X_QPID_DLQ_ENABLED, false); - args.put(X_QPID_MAXIMUM_DELIVERY_COUNT, 0); - - dlQueue = createAMQQueueImpl(UUIDGenerator.generateQueueUUID(dlQueueName, virtualHost.getName()), dlQueueName, true, owner, false, exclusive, virtualHost, args); + args.put(Queue.CREATE_DLQ_ON_CREATION, false); + args.put(Queue.MAXIMUM_DELIVERY_ATTEMPTS, 0); - //enter the dlq in the persistent store - DurableConfigurationStoreHelper.createQueue(virtualHost.getDurableConfigurationStore(), - dlQueue, - FieldTable.convertToFieldTable(args)); + dlQueue = _virtualHost.createQueue(UUIDGenerator.generateQueueUUID(dlQueueName, _virtualHost.getName()), dlQueueName, true, owner, false, exclusive, + false, args); } } @@ -350,11 +341,31 @@ public class AMQQueueFactory } q.setAlternateExchange(dlExchange); } + else if(arguments != null && arguments.get(Queue.ALTERNATE_EXCHANGE) instanceof String) + { + + final String altExchangeAttr = (String) arguments.get(Queue.ALTERNATE_EXCHANGE); + Exchange altExchange; + try + { + altExchange = _virtualHost.getExchange(UUID.fromString(altExchangeAttr)); + } + catch(IllegalArgumentException e) + { + altExchange = _virtualHost.getExchange(altExchangeAttr); + } + q.setAlternateExchange(altExchange); + } + + if (q.isDurable() && !q.isAutoDelete()) + { + DurableConfigurationStoreHelper.createQueue(_virtualHost.getDurableConfigurationStore(), q); + } return q; } - public static AMQQueue createAMQQueueImpl(QueueConfiguration config, VirtualHost host) throws AMQException + public AMQQueue createAMQQueueImpl(QueueConfiguration config) throws AMQException { String queueName = config.getName(); @@ -365,9 +376,9 @@ public class AMQQueueFactory Map<String, Object> arguments = createQueueArgumentsFromConfig(config); // we need queues that are defined in config to have deterministic ids. - UUID id = UUIDGenerator.generateQueueUUID(queueName, host.getName()); + UUID id = UUIDGenerator.generateQueueUUID(queueName, _virtualHost.getName()); - AMQQueue q = createAMQQueueImpl(id, queueName, durable, owner, autodelete, exclusive, host, arguments); + AMQQueue q = createAMQQueueImpl(id, queueName, durable, owner, autodelete, exclusive, false, arguments); q.configure(config); return q; } @@ -414,21 +425,23 @@ public class AMQQueueFactory * queue configuration * @return true if DLQ enabled */ - protected static boolean isDLQEnabled(boolean autoDelete, Map<String, Object> arguments, QueueConfiguration qConfig) + protected static boolean createDLQ(boolean autoDelete, Map<String, Object> arguments, QueueConfiguration qConfig) { //feature is not to be enabled for temporary queues or when explicitly disabled by argument - if (!autoDelete) + if (!(autoDelete || (arguments != null && arguments.containsKey(Queue.ALTERNATE_EXCHANGE)))) { - boolean dlqArgumentPresent = arguments != null && arguments.containsKey(X_QPID_DLQ_ENABLED); + boolean dlqArgumentPresent = arguments != null + && arguments.containsKey(Queue.CREATE_DLQ_ON_CREATION); if (dlqArgumentPresent || qConfig.isDeadLetterQueueEnabled()) { boolean dlqEnabled = true; if (dlqArgumentPresent) { - Object argument = arguments.get(X_QPID_DLQ_ENABLED); - dlqEnabled = argument instanceof Boolean && ((Boolean)argument).booleanValue(); + Object argument = arguments.get(Queue.CREATE_DLQ_ON_CREATION); + dlqEnabled = (argument instanceof Boolean && ((Boolean)argument).booleanValue()) + || (argument instanceof String && Boolean.parseBoolean(argument.toString())); } - return dlqEnabled; + return dlqEnabled ; } } return false; @@ -464,31 +477,30 @@ public class AMQQueueFactory if(config.getArguments() != null && !config.getArguments().isEmpty()) { - arguments.putAll(config.getArguments()); + arguments.putAll(QueueArgumentsConverter.convertWireArgsToModel(new HashMap<String, Object>(config.getArguments()))); } if(config.isLVQ() || config.getLVQKey() != null) { - arguments.put(QPID_LAST_VALUE_QUEUE, 1); - arguments.put(QPID_LAST_VALUE_QUEUE_KEY, config.getLVQKey() == null ? QPID_LVQ_KEY : config.getLVQKey()); + arguments.put(Queue.LVQ_KEY, config.getLVQKey() == null ? QPID_DEFAULT_LVQ_KEY : config.getLVQKey()); } else if (config.getPriority() || config.getPriorities() > 0) { - arguments.put(X_QPID_PRIORITIES, config.getPriorities() < 0 ? 10 : config.getPriorities()); + arguments.put(Queue.PRIORITIES, config.getPriorities() < 0 ? 10 : config.getPriorities()); } else if (config.getQueueSortKey() != null && !"".equals(config.getQueueSortKey())) { - arguments.put(QPID_QUEUE_SORT_KEY, config.getQueueSortKey()); + arguments.put(Queue.SORT_KEY, config.getQueueSortKey()); } if (!config.getAutoDelete() && config.isDeadLetterQueueEnabled()) { - arguments.put(X_QPID_DLQ_ENABLED, true); + arguments.put(Queue.CREATE_DLQ_ON_CREATION, true); } if (config.getDescription() != null && !"".equals(config.getDescription())) { - arguments.put(X_QPID_DESCRIPTION, config.getDescription()); + arguments.put(Queue.DESCRIPTION, config.getDescription()); } if (arguments.isEmpty()) diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/DefaultQueueRegistry.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/DefaultQueueRegistry.java index 27a9e13617..7308433759 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/DefaultQueueRegistry.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/DefaultQueueRegistry.java @@ -59,8 +59,9 @@ public class DefaultQueueRegistry implements QueueRegistry } } - public void unregisterQueue(AMQShortString name) + public void unregisterQueue(String nameString) { + AMQShortString name = new AMQShortString(nameString); AMQQueue q = _queueMap.remove(name); if(q != null) { @@ -74,16 +75,11 @@ public class DefaultQueueRegistry implements QueueRegistry } } - public AMQQueue getQueue(AMQShortString name) + private AMQQueue getQueue(AMQShortString name) { return _queueMap.get(name); } - public Collection<AMQShortString> getQueueNames() - { - return _queueMap.keySet(); - } - public Collection<AMQQueue> getQueues() { return _queueMap.values(); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueArgumentsConverter.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueArgumentsConverter.java new file mode 100644 index 0000000000..f5bee850c2 --- /dev/null +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueArgumentsConverter.java @@ -0,0 +1,154 @@ +/* + * + * 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 java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.Map; +import org.apache.qpid.server.model.Queue; + +public class QueueArgumentsConverter +{ + public static final String X_QPID_FLOW_RESUME_CAPACITY = "x-qpid-flow-resume-capacity"; + public static final String X_QPID_CAPACITY = "x-qpid-capacity"; + public static final String X_QPID_MINIMUM_ALERT_REPEAT_GAP = "x-qpid-minimum-alert-repeat-gap"; + public static final String X_QPID_MAXIMUM_MESSAGE_COUNT = "x-qpid-maximum-message-count"; + public static final String X_QPID_MAXIMUM_MESSAGE_SIZE = "x-qpid-maximum-message-size"; + public static final String X_QPID_MAXIMUM_MESSAGE_AGE = "x-qpid-maximum-message-age"; + public static final String X_QPID_MAXIMUM_QUEUE_DEPTH = "x-qpid-maximum-queue-depth"; + + public static final String QPID_ALERT_COUNT = "qpid.alert_count"; + public static final String QPID_ALERT_SIZE = "qpid.alert_size"; + public static final String QPID_ALERT_REPEAT_GAP = "qpid.alert_repeat_gap"; + + public static final String X_QPID_PRIORITIES = "x-qpid-priorities"; + + public static final String X_QPID_DESCRIPTION = "x-qpid-description"; + /* public static final String QPID_LVQ_KEY = "qpid.LVQ_key"; + public static final String QPID_LAST_VALUE_QUEUE = "qpid.last_value_queue"; + */ + public static final String QPID_LAST_VALUE_QUEUE_KEY = "qpid.last_value_queue_key"; + + public static final String QPID_QUEUE_SORT_KEY = "qpid.queue_sort_key"; + public static final String X_QPID_DLQ_ENABLED = "x-qpid-dlq-enabled"; + public static final String X_QPID_MAXIMUM_DELIVERY_COUNT = "x-qpid-maximum-delivery-count"; + public static final String QPID_GROUP_HEADER_KEY = "qpid.group_header_key"; + public static final String QPID_SHARED_MSG_GROUP = "qpid.shared_msg_group"; + public static final String QPID_DEFAULT_MESSAGE_GROUP_ARG = "qpid.default-message-group"; + public static final String QPID_TRACE_EXCLUDE = "qpid.trace.exclude"; + public static final String QPID_TRACE_ID = "qpid.trace.id"; + + public static final String QPID_LAST_VALUE_QUEUE = "qpid.last_value_queue"; + + /** + * No-local queue argument is used to support the no-local feature of Durable Subscribers. + */ + public static final String QPID_NO_LOCAL = "no-local"; + static final Map<String, String> ATTRIBUTE_MAPPINGS = new LinkedHashMap<String, String>(); + static + { + ATTRIBUTE_MAPPINGS.put(X_QPID_MINIMUM_ALERT_REPEAT_GAP, Queue.ALERT_REPEAT_GAP); + ATTRIBUTE_MAPPINGS.put(X_QPID_MAXIMUM_MESSAGE_AGE, Queue.ALERT_THRESHOLD_MESSAGE_AGE); + ATTRIBUTE_MAPPINGS.put(X_QPID_MAXIMUM_MESSAGE_SIZE, Queue.ALERT_THRESHOLD_MESSAGE_SIZE); + + ATTRIBUTE_MAPPINGS.put(X_QPID_MAXIMUM_MESSAGE_COUNT, Queue.ALERT_THRESHOLD_QUEUE_DEPTH_MESSAGES); + ATTRIBUTE_MAPPINGS.put(X_QPID_MAXIMUM_QUEUE_DEPTH, Queue.ALERT_THRESHOLD_QUEUE_DEPTH_BYTES); + ATTRIBUTE_MAPPINGS.put(QPID_ALERT_COUNT, Queue.ALERT_THRESHOLD_QUEUE_DEPTH_MESSAGES); + ATTRIBUTE_MAPPINGS.put(QPID_ALERT_SIZE, Queue.ALERT_THRESHOLD_QUEUE_DEPTH_BYTES); + ATTRIBUTE_MAPPINGS.put(QPID_ALERT_REPEAT_GAP, Queue.ALERT_REPEAT_GAP); + + ATTRIBUTE_MAPPINGS.put(X_QPID_MAXIMUM_DELIVERY_COUNT, Queue.MAXIMUM_DELIVERY_ATTEMPTS); + + ATTRIBUTE_MAPPINGS.put(X_QPID_CAPACITY, Queue.QUEUE_FLOW_CONTROL_SIZE_BYTES); + ATTRIBUTE_MAPPINGS.put(X_QPID_FLOW_RESUME_CAPACITY, Queue.QUEUE_FLOW_RESUME_SIZE_BYTES); + + ATTRIBUTE_MAPPINGS.put(QPID_QUEUE_SORT_KEY, Queue.SORT_KEY); + ATTRIBUTE_MAPPINGS.put(QPID_LAST_VALUE_QUEUE_KEY, Queue.LVQ_KEY); + ATTRIBUTE_MAPPINGS.put(X_QPID_PRIORITIES, Queue.PRIORITIES); + + ATTRIBUTE_MAPPINGS.put(X_QPID_DESCRIPTION, Queue.DESCRIPTION); + + ATTRIBUTE_MAPPINGS.put(X_QPID_DLQ_ENABLED, Queue.CREATE_DLQ_ON_CREATION); + ATTRIBUTE_MAPPINGS.put(QPID_GROUP_HEADER_KEY, Queue.MESSAGE_GROUP_KEY); + //ATTRIBUTE_MAPPINGS.put(QPID_SHARED_MSG_GROUP, Queue.MESSAGE_GROUP_SHARED_GROUPS); + ATTRIBUTE_MAPPINGS.put(QPID_DEFAULT_MESSAGE_GROUP_ARG, Queue.MESSAGE_GROUP_DEFAULT_GROUP); + ATTRIBUTE_MAPPINGS.put(QPID_TRACE_EXCLUDE, Queue.FEDERATION_EXCLUDES); + ATTRIBUTE_MAPPINGS.put(QPID_TRACE_ID, Queue.FEDERATION_ID); + ATTRIBUTE_MAPPINGS.put(QPID_NO_LOCAL, Queue.NO_LOCAL); + + } + + + public static Map<String,Object> convertWireArgsToModel(Map<String,Object> wireArguments) + { + Map<String,Object> modelArguments = new HashMap<String, Object>(); + if(wireArguments != null) + { + for(Map.Entry<String,String> entry : ATTRIBUTE_MAPPINGS.entrySet()) + { + if(wireArguments.containsKey(entry.getKey())) + { + modelArguments.put(entry.getValue(), wireArguments.get(entry.getKey())); + } + } + if(wireArguments.containsKey(QPID_LAST_VALUE_QUEUE) && !wireArguments.containsKey(QPID_LAST_VALUE_QUEUE_KEY)) + { + modelArguments.put(Queue.LVQ_KEY, AMQQueueFactory.QPID_DEFAULT_LVQ_KEY); + } + if(wireArguments.containsKey(QPID_SHARED_MSG_GROUP)) + { + modelArguments.put(Queue.MESSAGE_GROUP_SHARED_GROUPS, + SimpleAMQQueue.SHARED_MSG_GROUP_ARG_VALUE.equals(wireArguments.get(QPID_SHARED_MSG_GROUP))); + } + if(wireArguments.get(X_QPID_DLQ_ENABLED) != null) + { + modelArguments.put(Queue.CREATE_DLQ_ON_CREATION, Boolean.parseBoolean(wireArguments.get(X_QPID_DLQ_ENABLED).toString())); + } + + if(wireArguments.get(QPID_NO_LOCAL) != null) + { + modelArguments.put(Queue.NO_LOCAL, Boolean.parseBoolean(wireArguments.get(QPID_NO_LOCAL).toString())); + } + + } + return modelArguments; + } + + + public static Map<String,Object> convertModelArgsToWire(Map<String,Object> modelArguments) + { + Map<String,Object> wireArguments = new HashMap<String, Object>(); + for(Map.Entry<String,String> entry : ATTRIBUTE_MAPPINGS.entrySet()) + { + if(modelArguments.containsKey(entry.getValue())) + { + wireArguments.put(entry.getKey(), modelArguments.get(entry.getValue())); + } + } + + if(Boolean.TRUE.equals(modelArguments.get(Queue.MESSAGE_GROUP_SHARED_GROUPS))) + { + wireArguments.put(QPID_SHARED_MSG_GROUP, SimpleAMQQueue.SHARED_MSG_GROUP_ARG_VALUE); + } + + return wireArguments; + } +} diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueFactory.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueFactory.java new file mode 100644 index 0000000000..5411a2bc9c --- /dev/null +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueFactory.java @@ -0,0 +1,38 @@ +/* + * + * 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 java.util.Map; +import java.util.UUID; +import org.apache.qpid.AMQException; +import org.apache.qpid.AMQSecurityException; + +public interface QueueFactory +{ + AMQQueue createAMQQueueImpl(UUID id, + String queueName, + boolean durable, + String owner, + boolean autoDelete, + boolean exclusive, + boolean deleteOnNoConsumer, + Map<String, Object> arguments) throws AMQSecurityException, AMQException; +} diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueRegistry.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueRegistry.java index e8c34128e9..bc1d5942bd 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueRegistry.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/QueueRegistry.java @@ -20,7 +20,6 @@ */ package org.apache.qpid.server.queue; -import org.apache.qpid.framing.AMQShortString; import org.apache.qpid.server.virtualhost.VirtualHost; import java.util.Collection; @@ -32,11 +31,7 @@ public interface QueueRegistry void registerQueue(AMQQueue queue); - void unregisterQueue(AMQShortString name); - - AMQQueue getQueue(AMQShortString name); - - Collection<AMQShortString> getQueueNames(); + void unregisterQueue(String name); Collection<AMQQueue> getQueues(); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java index b0ab93162a..e3dbd62b6c 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java @@ -22,7 +22,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.EnumSet; -import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -51,6 +51,7 @@ import org.apache.qpid.server.logging.actors.QueueActor; import org.apache.qpid.server.logging.messages.QueueMessages; import org.apache.qpid.server.logging.subjects.QueueLogSubject; import org.apache.qpid.server.message.ServerMessage; +import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.protocol.AMQSessionModel; import org.apache.qpid.server.security.AuthorizationHolder; import org.apache.qpid.server.subscription.AssignedSubscriptionMessageGroupManager; @@ -68,12 +69,10 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes private static final Logger _logger = Logger.getLogger(SimpleAMQQueue.class); - public static final String QPID_GROUP_HEADER_KEY = "qpid.group_header_key"; - public static final String QPID_SHARED_MSG_GROUP = "qpid.shared_msg_group"; public static final String SHARED_MSG_GROUP_ARG_VALUE = "1"; - private static final String QPID_DEFAULT_MESSAGE_GROUP_ARG = "qpid.default-message-group"; private static final String QPID_NO_GROUP = "qpid.no-group"; private static final String DEFAULT_SHARED_MESSAGE_GROUP = System.getProperty(BrokerProperties.PROPERTY_DEFAULT_SHARED_MESSAGE_GROUP, QPID_NO_GROUP); + // TODO - should make this configurable at the vhost / broker level private static final int DEFAULT_MAX_GROUPS = 255; @@ -237,7 +236,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes _exclusive = exclusive; _virtualHost = virtualHost; _entries = entryListFactory.createQueueEntryList(this); - _arguments = arguments == null ? new HashMap<String, Object>() : new HashMap<String, Object>(arguments); + _arguments = Collections.synchronizedMap(arguments == null ? new LinkedHashMap<String, Object>() : new LinkedHashMap<String, Object>(arguments)); _id = id; _asyncDelivery = ReferenceCountingExecutorService.getInstance().acquireExecutorService(); @@ -255,19 +254,21 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes durable, !durable, _entries.getPriorities() > 0)); - if(arguments != null && arguments.containsKey(QPID_GROUP_HEADER_KEY)) + if(arguments != null && arguments.containsKey(Queue.MESSAGE_GROUP_KEY)) { - if(arguments.containsKey(QPID_SHARED_MSG_GROUP) && String.valueOf(arguments.get(QPID_SHARED_MSG_GROUP)).equals(SHARED_MSG_GROUP_ARG_VALUE)) + if(arguments.get(Queue.MESSAGE_GROUP_SHARED_GROUPS) != null + && (Boolean)(arguments.get(Queue.MESSAGE_GROUP_SHARED_GROUPS))) { - Object defaultGroup = arguments.get(QPID_DEFAULT_MESSAGE_GROUP_ARG); + Object defaultGroup = arguments.get(Queue.MESSAGE_GROUP_DEFAULT_GROUP); _messageGroupManager = - new DefinedGroupMessageGroupManager(String.valueOf(arguments.get(QPID_GROUP_HEADER_KEY)), + new DefinedGroupMessageGroupManager(String.valueOf(arguments.get(Queue.MESSAGE_GROUP_KEY)), defaultGroup == null ? DEFAULT_SHARED_MESSAGE_GROUP : defaultGroup.toString(), this); } else { - _messageGroupManager = new AssignedSubscriptionMessageGroupManager(String.valueOf(arguments.get(QPID_GROUP_HEADER_KEY)), DEFAULT_MAX_GROUPS); + _messageGroupManager = new AssignedSubscriptionMessageGroupManager(String.valueOf(arguments.get( + Queue.MESSAGE_GROUP_KEY)), DEFAULT_MAX_GROUPS); } } else @@ -358,13 +359,17 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes _alternateExchange = exchange; } - /** - * Arguments used to create this queue. The caller is assured - * that null will never be returned. - */ - public Map<String, Object> getArguments() + + @Override + public Collection<String> getAvailableAttributes() + { + return new ArrayList<String>(_arguments.keySet()); + } + + @Override + public Object getAttribute(String attrName) { - return _arguments; + return _arguments.get(attrName); } public boolean isAutoDelete() @@ -511,7 +516,7 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes _logger.info("Auto-deleteing queue:" + this); } - delete(); + getVirtualHost().removeQueue(this); // we need to manually fire the event to the removed subscription (which was the last one left for this // queue. This is because the delete method uses the subscription set which has just been cleared @@ -1340,7 +1345,6 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes } } - _virtualHost.getQueueRegistry().unregisterQueue(_name); List<QueueEntry> entries = getMessagesOnTheQueue(new QueueEntryFilter() { @@ -2282,18 +2286,18 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes { if (description == null) { - _arguments.remove(AMQQueueFactory.X_QPID_DESCRIPTION); + _arguments.remove(Queue.DESCRIPTION); } else { - _arguments.put(AMQQueueFactory.X_QPID_DESCRIPTION, description); + _arguments.put(Queue.DESCRIPTION, description); } } @Override public String getDescription() { - return (String) _arguments.get(AMQQueueFactory.X_QPID_DESCRIPTION); + return (String) _arguments.get(Queue.DESCRIPTION); } } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/security/SecurityManager.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/security/SecurityManager.java index 931368cb97..960986ec45 100755 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/security/SecurityManager.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/security/SecurityManager.java @@ -167,12 +167,12 @@ public class SecurityManager implements ConfigurationChangeListener { String pluginTypeName = getPluginTypeName(accessControl); _hostPlugins.put(pluginTypeName, accessControl); - + if(_logger.isDebugEnabled()) { _logger.debug("Added access control to host plugins with name: " + vhostName); } - + break; } } @@ -366,7 +366,7 @@ public class SecurityManager implements ConfigurationChangeListener } public boolean authoriseCreateQueue(final Boolean autoDelete, final Boolean durable, final Boolean exclusive, - final Boolean nowait, final Boolean passive, final AMQShortString queueName, final String owner) + final Boolean nowait, final Boolean passive, final String queueName, final String owner) { return checkAllPlugins(new AccessCheck() { diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/security/access/ObjectProperties.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/security/access/ObjectProperties.java index 6c631fc360..893b371d11 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/security/access/ObjectProperties.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/security/access/ObjectProperties.java @@ -212,7 +212,7 @@ public class ObjectProperties } public ObjectProperties(Boolean autoDelete, Boolean durable, Boolean exclusive, Boolean nowait, Boolean passive, - AMQShortString queueName, String owner) + String queueName, String owner) { super(); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/DurableConfigurationStoreHelper.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/DurableConfigurationStoreHelper.java index efb1e95e99..e9181c0e12 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/store/DurableConfigurationStoreHelper.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/store/DurableConfigurationStoreHelper.java @@ -20,6 +20,7 @@ */ package org.apache.qpid.server.store; +import java.util.Collection; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.Map; @@ -32,6 +33,7 @@ import org.apache.qpid.server.model.Exchange; import org.apache.qpid.server.model.LifetimePolicy; import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.queue.AMQQueue; +import org.apache.qpid.server.queue.QueueArgumentsConverter; public class DurableConfigurationStoreHelper { @@ -46,28 +48,23 @@ public class DurableConfigurationStoreHelper attributesMap.put(Queue.NAME, queue.getName()); attributesMap.put(Queue.OWNER, AMQShortString.toString(queue.getOwner())); attributesMap.put(Queue.EXCLUSIVE, queue.isExclusive()); + if (queue.getAlternateExchange() != null) { attributesMap.put(Queue.ALTERNATE_EXCHANGE, queue.getAlternateExchange().getId()); } - else - { - attributesMap.remove(Queue.ALTERNATE_EXCHANGE); - } - if (attributesMap.containsKey(Queue.ARGUMENTS)) - { - // We wouldn't need this if createQueueConfiguredObject took only AMQQueue - Map<String, Object> currentArgs = (Map<String, Object>) attributesMap.get(Queue.ARGUMENTS); - currentArgs.putAll(queue.getArguments()); - } - else + + Collection<String> availableAttrs = queue.getAvailableAttributes(); + + for(String attrName : availableAttrs) { - attributesMap.put(Queue.ARGUMENTS, queue.getArguments()); + attributesMap.put(attrName, queue.getAttribute(attrName)); } + store.update(queue.getId(), QUEUE, attributesMap); } - public static void createQueue(DurableConfigurationStore store, AMQQueue queue, FieldTable arguments) + public static void createQueue(DurableConfigurationStore store, AMQQueue queue) throws AMQStoreException { Map<String, Object> attributesMap = new HashMap<String, Object>(); @@ -78,11 +75,9 @@ public class DurableConfigurationStoreHelper { attributesMap.put(Queue.ALTERNATE_EXCHANGE, queue.getAlternateExchange().getId()); } - // TODO KW i think the arguments could come from the queue itself removing the need for the parameter arguments. - // It would also do away with the need for the if/then/else within updateQueueConfiguredObject - if (arguments != null) + for(String attrName : queue.getAvailableAttributes()) { - attributesMap.put(Queue.ARGUMENTS, FieldTable.convertToMap(arguments)); + attributesMap.put(attrName, queue.getAttribute(attrName)); } store.create(queue.getId(), QUEUE,attributesMap); } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java index 4e27a008dd..d87431a415 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/AbstractVirtualHost.java @@ -34,6 +34,7 @@ import java.util.concurrent.TimeUnit; import org.apache.commons.configuration.ConfigurationException; import org.apache.log4j.Logger; import org.apache.qpid.AMQException; +import org.apache.qpid.AMQSecurityException; import org.apache.qpid.server.configuration.ExchangeConfiguration; import org.apache.qpid.server.configuration.QueueConfiguration; import org.apache.qpid.server.configuration.VirtualHostConfiguration; @@ -58,11 +59,13 @@ import org.apache.qpid.server.queue.QueueRegistry; import org.apache.qpid.server.security.SecurityManager; import org.apache.qpid.server.stats.StatisticsCounter; import org.apache.qpid.server.stats.StatisticsGatherer; +import org.apache.qpid.server.store.DurableConfigurationStore; import org.apache.qpid.server.store.DurableConfigurationStoreHelper; import org.apache.qpid.server.store.DurableConfiguredObjectRecoverer; import org.apache.qpid.server.store.Event; import org.apache.qpid.server.store.EventListener; import org.apache.qpid.server.txn.DtxRegistry; +import org.apache.qpid.server.virtualhost.plugins.QueueExistsException; public abstract class AbstractVirtualHost implements VirtualHost, IConnectionRegistry.RegistryChangeListener, EventListener { @@ -95,6 +98,7 @@ public abstract class AbstractVirtualHost implements VirtualHost, IConnectionReg private final ConnectionRegistry _connectionRegistry; private final DtxRegistry _dtxRegistry; + private final AMQQueueFactory _queueFactory; private volatile State _state = State.INITIALISING; @@ -136,11 +140,14 @@ public abstract class AbstractVirtualHost implements VirtualHost, IConnectionReg _houseKeepingTasks = new ScheduledThreadPoolExecutor(_vhostConfig.getHouseKeepingThreadCount()); + _queueRegistry = new DefaultQueueRegistry(this); + _queueFactory = new AMQQueueFactory(this, _queueRegistry); + _exchangeFactory = new DefaultExchangeFactory(this); - _exchangeRegistry = new DefaultExchangeRegistry(this); + _exchangeRegistry = new DefaultExchangeRegistry(this, _queueRegistry); initialiseStatistics(); @@ -298,12 +305,12 @@ public abstract class AbstractVirtualHost implements VirtualHost, IConnectionReg private void configureQueue(QueueConfiguration queueConfiguration) throws AMQException, ConfigurationException { - AMQQueue queue = AMQQueueFactory.createAMQQueueImpl(queueConfiguration, this); + AMQQueue queue = _queueFactory.createAMQQueueImpl(queueConfiguration); String queueName = queue.getName(); if (queue.isDurable()) { - DurableConfigurationStoreHelper.createQueue(getDurableConfigurationStore(), queue, null); + DurableConfigurationStoreHelper.createQueue(getDurableConfigurationStore(), queue); } //get the exchange name (returns default exchange name if none was specified) @@ -428,12 +435,102 @@ public abstract class AbstractVirtualHost implements VirtualHost, IConnectionReg } @Override + public AMQQueue getQueue(String name) + { + return _queueRegistry.getQueue(name); + } + + @Override + public AMQQueue getQueue(UUID id) + { + return _queueRegistry.getQueue(id); + } + + @Override + public Collection<AMQQueue> getQueues() + { + return _queueRegistry.getQueues(); + } + + @Override + public int removeQueue(AMQQueue queue) throws AMQException + { + synchronized (getQueueRegistry()) + { + int purged = queue.delete(); + + getQueueRegistry().unregisterQueue(queue.getName()); + if (queue.isDurable() && !queue.isAutoDelete()) + { + DurableConfigurationStore store = getDurableConfigurationStore(); + DurableConfigurationStoreHelper.removeQueue(store, queue); + } + return purged; + } + } + + @Override + public AMQQueue createQueue(UUID id, + String queueName, + boolean durable, + String owner, + boolean autoDelete, + boolean exclusive, + boolean deleteOnNoConsumer, + Map<String, Object> arguments) throws AMQException + { + // Access check + if (!getSecurityManager().authoriseCreateQueue(autoDelete, + durable, + exclusive, + null, + null, + queueName, + owner)) + { + String description = "Permission denied: queue-name '" + queueName + "'"; + throw new AMQSecurityException(description); + } + + synchronized (_queueRegistry) + { + if(_queueRegistry.getQueue(queueName) != null) + { + throw new QueueExistsException("Queue with name " + queueName + " already exists", _queueRegistry.getQueue(queueName)); + } + if(id == null) + { + + id = UUIDGenerator.generateExchangeUUID(queueName, getName()); + while(_queueRegistry.getQueue(id) != null) + { + id = UUID.randomUUID(); + } + + } + else if(_queueRegistry.getQueue(id) != null) + { + throw new QueueExistsException("Queue with id " + id + " already exists", _queueRegistry.getQueue(queueName)); + } + return _queueFactory.createAMQQueueImpl(id, queueName, durable, owner, autoDelete, exclusive, deleteOnNoConsumer, + arguments); + } + + } + + @Override public Exchange getExchange(String name) { return _exchangeRegistry.getExchange(name); } @Override + public Exchange getExchange(UUID id) + { + return _exchangeRegistry.getExchange(id); + } + + @Override public Exchange getDefaultExchange() { return _exchangeRegistry.getDefaultExchange(); @@ -747,7 +844,7 @@ public abstract class AbstractVirtualHost implements VirtualHost, IConnectionReg protected Map<String, DurableConfiguredObjectRecoverer> getDurableConfigurationRecoverers() { DurableConfiguredObjectRecoverer[] recoverers = { - new QueueRecoverer(this, getExchangeRegistry()), + new QueueRecoverer(this, getExchangeRegistry(), _queueFactory), new ExchangeRecoverer(getExchangeRegistry(), getExchangeFactory()), new BindingRecoverer(this, getExchangeRegistry()) }; diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/BindingRecoverer.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/BindingRecoverer.java index 7cfadbcadf..2d3a620e91 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/BindingRecoverer.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/BindingRecoverer.java @@ -91,7 +91,7 @@ public class BindingRecoverer extends AbstractDurableConfiguredObjectRecoverer<B { _unresolvedDependencies.add(new ExchangeDependency()); } - _queue = _virtualHost.getQueueRegistry().getQueue(_queueId); + _queue = _virtualHost.getQueue(_queueId); if(_queue == null) { _unresolvedDependencies.add(new QueueDependency()); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java index 3526551073..8d05e719ee 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/DefaultUpgraderProvider.java @@ -30,6 +30,7 @@ import org.apache.qpid.server.exchange.TopicExchange; import org.apache.qpid.server.model.Binding; import org.apache.qpid.server.model.Exchange; import org.apache.qpid.server.model.Queue; +import org.apache.qpid.server.queue.QueueArgumentsConverter; import org.apache.qpid.server.store.ConfiguredObjectRecord; import org.apache.qpid.server.store.DurableConfigurationRecoverer; import org.apache.qpid.server.store.DurableConfigurationStoreUpgrader; @@ -60,6 +61,9 @@ public class DefaultUpgraderProvider implements UpgraderProvider currentUpgrader = addUpgrader(currentUpgrader, new Version0Upgrader()); case 1: currentUpgrader = addUpgrader(currentUpgrader, new Version1Upgrader()); + case 2: + currentUpgrader = addUpgrader(currentUpgrader, new Version2Upgrader()); + case CURRENT_CONFIG_VERSION: currentUpgrader = addUpgrader(currentUpgrader, new NullUpgrader(recoverer)); break; @@ -213,7 +217,7 @@ public class DefaultUpgraderProvider implements UpgraderProvider UUID queueId = UUID.fromString(queueIdString); ConfiguredObjectRecord localRecord = getUpdateMap().get(queueId); return !((localRecord != null && localRecord.getType().equals(Queue.class.getSimpleName())) - || _virtualHost.getQueueRegistry().getQueue(queueId) != null); + || _virtualHost.getQueue(queueId) != null); } private boolean isBinding(final String type) @@ -224,4 +228,39 @@ public class DefaultUpgraderProvider implements UpgraderProvider } + /* + * Convert the storage of queue attributes to remove the separate "ARGUMENT" attribute, and flatten the + * attributes into the map using the model attribute names rather than the wire attribute names + */ + private class Version2Upgrader extends NonNullUpgrader + { + + private static final String ARGUMENTS = "arguments"; + + @Override + public void configuredObject(UUID id, String type, Map<String, Object> attributes) + { + if(Queue.class.getSimpleName().equals(type)) + { + Map<String, Object> newAttributes = new LinkedHashMap<String, Object>(); + if(attributes.get(ARGUMENTS) instanceof Map) + { + newAttributes.putAll(QueueArgumentsConverter.convertWireArgsToModel((Map<String, Object>) attributes + .get(ARGUMENTS))); + } + newAttributes.putAll(attributes); + attributes = newAttributes; + getUpdateMap().put(id, new ConfiguredObjectRecord(id,type,attributes)); + } + + getNextUpgrader().configuredObject(id,type,attributes); + } + + @Override + public void complete() + { + getNextUpgrader().complete(); + } + } + } diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/QueueRecoverer.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/QueueRecoverer.java index 7929cd3e39..b4fbdf7544 100644 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/QueueRecoverer.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/QueueRecoverer.java @@ -20,18 +20,19 @@ */ package org.apache.qpid.server.virtualhost; +import java.util.LinkedHashMap; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.UUID; import org.apache.log4j.Logger; import org.apache.qpid.AMQException; -import org.apache.qpid.server.configuration.IllegalConfigurationException; import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.exchange.ExchangeRegistry; import org.apache.qpid.server.model.Queue; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.AMQQueueFactory; +import org.apache.qpid.server.queue.QueueFactory; import org.apache.qpid.server.store.AbstractDurableConfiguredObjectRecoverer; import org.apache.qpid.server.store.UnresolvedDependency; import org.apache.qpid.server.store.UnresolvedObject; @@ -41,11 +42,15 @@ public class QueueRecoverer extends AbstractDurableConfiguredObjectRecoverer<AMQ private static final Logger _logger = Logger.getLogger(QueueRecoverer.class); private final VirtualHost _virtualHost; private final ExchangeRegistry _exchangeRegistry; + private final QueueFactory _queueFactory; - public QueueRecoverer(final VirtualHost virtualHost, final ExchangeRegistry exchangeRegistry) + public QueueRecoverer(final VirtualHost virtualHost, + final ExchangeRegistry exchangeRegistry, + final QueueFactory queueFactory) { _virtualHost = virtualHost; _exchangeRegistry = exchangeRegistry; + _queueFactory = queueFactory; } @Override @@ -101,26 +106,24 @@ public class QueueRecoverer extends AbstractDurableConfiguredObjectRecoverer<AMQ String queueName = (String) _attributes.get(Queue.NAME); String owner = (String) _attributes.get(Queue.OWNER); boolean exclusive = (Boolean) _attributes.get(Queue.EXCLUSIVE); - @SuppressWarnings("unchecked") - Map<String, Object> queueArgumentsMap = (Map<String, Object>) _attributes.get(Queue.ARGUMENTS); + + Map<String, Object> queueArgumentsMap = new LinkedHashMap<String, Object>(_attributes); + queueArgumentsMap.remove(Queue.NAME); + queueArgumentsMap.remove(Queue.OWNER); + queueArgumentsMap.remove(Queue.EXCLUSIVE); + try { - _queue = _virtualHost.getQueueRegistry().getQueue(_id); + _queue = _virtualHost.getQueue(_id); if(_queue == null) { - _queue = _virtualHost.getQueueRegistry().getQueue(queueName); + _queue = _virtualHost.getQueue(queueName); } if (_queue == null) { - _queue = AMQQueueFactory.createAMQQueueImpl(_id, queueName, true, owner, false, exclusive, _virtualHost, - queueArgumentsMap); - _virtualHost.getQueueRegistry().registerQueue(_queue); - - if (_alternateExchange != null) - { - _queue.setAlternateExchange(_alternateExchange); - } + _queue = _queueFactory.createAMQQueueImpl(_id, queueName, true, owner, false, exclusive, + false, queueArgumentsMap); } } catch (AMQException e) diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java index e06e785338..2ebbedccd4 100755 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHost.java @@ -21,15 +21,18 @@ package org.apache.qpid.server.virtualhost; import java.util.Collection; +import java.util.Map; import java.util.UUID; import java.util.concurrent.ScheduledFuture; import org.apache.qpid.AMQException; +import org.apache.qpid.AMQSecurityException; import org.apache.qpid.common.Closeable; import org.apache.qpid.server.configuration.VirtualHostConfiguration; import org.apache.qpid.server.connection.IConnectionRegistry; import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.protocol.LinkRegistry; +import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.QueueRegistry; import org.apache.qpid.server.security.SecurityManager; import org.apache.qpid.server.stats.StatisticsGatherer; @@ -45,7 +48,23 @@ public interface VirtualHost extends DurableConfigurationStore.Source, Closeable String getName(); - QueueRegistry getQueueRegistry(); + AMQQueue getQueue(String name); + + AMQQueue getQueue(UUID id); + + Collection<AMQQueue> getQueues(); + + int removeQueue(AMQQueue queue) throws AMQException; + + AMQQueue createQueue(UUID id, + String queueName, + boolean durable, + String owner, + boolean autoDelete, + boolean exclusive, + boolean deleteOnNoConsumer, + Map<String, Object> arguments) throws AMQException; + Exchange createExchange(UUID id, String exchange, @@ -58,6 +77,8 @@ public interface VirtualHost extends DurableConfigurationStore.Source, Closeable void removeExchange(Exchange exchange, boolean force) throws AMQException; Exchange getExchange(String name); + Exchange getExchange(UUID id); + Exchange getDefaultExchange(); diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostConfigRecoveryHandler.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostConfigRecoveryHandler.java index 3738306f6a..39ca3197b4 100755 --- a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostConfigRecoveryHandler.java +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/VirtualHostConfigRecoveryHandler.java @@ -119,7 +119,7 @@ public class VirtualHostConfigRecoveryHandler implements } for(Transaction.Record record : enqueues) { - final AMQQueue queue = _virtualHost.getQueueRegistry().getQueue(record.getQueue().getId()); + final AMQQueue queue = _virtualHost.getQueue(record.getQueue().getId()); if(queue != null) { final long messageId = record.getMessage().getMessageNumber(); @@ -179,7 +179,7 @@ public class VirtualHostConfigRecoveryHandler implements } for(Transaction.Record record : dequeues) { - final AMQQueue queue = _virtualHost.getQueueRegistry().getQueue(record.getQueue().getId()); + final AMQQueue queue = _virtualHost.getQueue(record.getQueue().getId()); if(queue != null) { final long messageId = record.getMessage().getMessageNumber(); @@ -268,7 +268,7 @@ public class VirtualHostConfigRecoveryHandler implements public void queueEntry(final UUID queueId, long messageId) { - AMQQueue queue = _virtualHost.getQueueRegistry().getQueue(queueId); + AMQQueue queue = _virtualHost.getQueue(queueId); try { if(queue != null) diff --git a/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/plugins/QueueExistsException.java b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/plugins/QueueExistsException.java new file mode 100644 index 0000000000..54f7d0d172 --- /dev/null +++ b/qpid/java/broker/src/main/java/org/apache/qpid/server/virtualhost/plugins/QueueExistsException.java @@ -0,0 +1,40 @@ +/* + * + * 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.virtualhost.plugins; + +import org.apache.qpid.AMQException; +import org.apache.qpid.server.queue.AMQQueue; + +public class QueueExistsException extends AMQException +{ + private final AMQQueue _existing; + + public QueueExistsException(String name, AMQQueue existing) + { + super(name); + _existing = existing; + } + + public AMQQueue getExistingQueue() + { + return _existing; + } +} |
