diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-02-26 23:27:39 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-02-26 23:27:39 +0000 |
| commit | 4ae118bb7a81155a9f3d22af3a4a3f2191799c83 (patch) | |
| tree | 8bebb59700d6825f997cb501c9043513e0f785c8 /qpid/java/broker-plugins | |
| parent | 6339e3b1ad22e74508510e08384c4d484bd9666c (diff) | |
| download | qpid-python-4ae118bb7a81155a9f3d22af3a4a3f2191799c83.tar.gz | |
QPID-5577 : [Java Broker] Change Exchange,Queue,Binding,Consumer to implement ConfiguredObject and remove adapter classes
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1572343 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker-plugins')
27 files changed, 127 insertions, 126 deletions
diff --git a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java index ff4bd1dc2e..69c625d41d 100644 --- a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java +++ b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java @@ -20,7 +20,7 @@ */ package org.apache.qpid.server.protocol.v0_10; -import org.apache.qpid.server.exchange.Exchange; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.flow.FlowCreditManager; import org.apache.qpid.server.logging.LogActor; import org.apache.qpid.server.logging.actors.CurrentActor; @@ -406,7 +406,7 @@ public class ConsumerTarget_0_10 extends AbstractConsumerTarget implements FlowC if(owningResource instanceof AMQQueue) { final AMQQueue queue = (AMQQueue)owningResource; - final Exchange alternateExchange = queue.getAlternateExchange(); + final ExchangeImpl alternateExchange = queue.getAlternateExchange(); if(alternateExchange != null) { diff --git a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java index 12c10240bb..236a955ea9 100644 --- a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java +++ b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java @@ -936,10 +936,10 @@ public class ServerSession extends Session return getId().compareTo(o.getId()); } - private class CheckCapacityAction<C extends Consumer> implements Action<MessageInstance<?,C>> + private class CheckCapacityAction implements Action<MessageInstance> { @Override - public void performAction(final MessageInstance<?,C> entry) + public void performAction(final MessageInstance entry) { TransactionLogResource queue = entry.getOwningResource(); if(queue instanceof CapacityChecker) diff --git a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSessionDelegate.java b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSessionDelegate.java index e70b34a426..1fb82efd2d 100644 --- a/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSessionDelegate.java +++ b/qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSessionDelegate.java @@ -27,13 +27,11 @@ import java.util.LinkedHashMap; import java.util.UUID; import org.apache.log4j.Logger; -import org.apache.qpid.server.binding.*; -import org.apache.qpid.server.binding.Binding; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.model.ExclusivityPolicy; import org.apache.qpid.server.model.LifetimePolicy; import org.apache.qpid.server.store.StoreException; import org.apache.qpid.server.exchange.AMQUnknownExchangeType; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.exchange.HeadersExchange; import org.apache.qpid.server.filter.AMQInvalidArgumentException; import org.apache.qpid.server.filter.FilterManager; @@ -688,7 +686,7 @@ public class ServerSessionDelegate extends SessionDelegate if(method.getPassive()) { - Exchange exchange = getExchange(session, exchangeName); + ExchangeImpl exchange = getExchange(session, exchangeName); if(exchange == null) { @@ -736,7 +734,7 @@ public class ServerSessionDelegate extends SessionDelegate } catch(ExchangeExistsException e) { - Exchange exchange = e.getExistingExchange(); + ExchangeImpl exchange = e.getExistingExchange(); if(!exchange.getTypeName().equals(method.getType())) { exception(session, method, ExecutionErrorCode.NOT_ALLOWED, @@ -776,7 +774,7 @@ public class ServerSessionDelegate extends SessionDelegate ((ServerSession)session).close(errorCode.getValue(), description); } - private Exchange getExchange(Session session, String exchangeName) + private ExchangeImpl getExchange(Session session, String exchangeName) { return getVirtualHost(session).getExchange(exchangeName); } @@ -827,7 +825,7 @@ public class ServerSessionDelegate extends SessionDelegate return; } - Exchange exchange = getExchange(session, method.getExchange()); + ExchangeImpl exchange = getExchange(session, method.getExchange()); if(exchange == null) { @@ -862,7 +860,7 @@ public class ServerSessionDelegate extends SessionDelegate return false; } - private boolean isStandardExchange(Exchange exchange, Collection<ExchangeType<? extends Exchange>> registeredTypes) + private boolean isStandardExchange(ExchangeImpl exchange, Collection<ExchangeType<? extends ExchangeImpl>> registeredTypes) { for(ExchangeType type : registeredTypes) { @@ -880,7 +878,7 @@ public class ServerSessionDelegate extends SessionDelegate ExchangeQueryResult result = new ExchangeQueryResult(); - Exchange exchange = getExchange(session, method.getName()); + ExchangeImpl exchange = getExchange(session, method.getName()); if(exchange != null) { @@ -919,7 +917,7 @@ public class ServerSessionDelegate extends SessionDelegate method.setBindingKey(method.getQueue()); } AMQQueue queue = virtualHost.getQueue(method.getQueue()); - Exchange exchange = virtualHost.getExchange(method.getExchange()); + ExchangeImpl exchange = virtualHost.getExchange(method.getExchange()); if(queue == null) { exception(session, method, ExecutionErrorCode.NOT_FOUND, "Queue: '" + method.getQueue() + "' not found"); @@ -978,7 +976,7 @@ public class ServerSessionDelegate extends SessionDelegate else { AMQQueue queue = virtualHost.getQueue(method.getQueue()); - Exchange exchange = virtualHost.getExchange(method.getExchange()); + ExchangeImpl exchange = virtualHost.getExchange(method.getExchange()); if(queue == null) { exception(session, method, ExecutionErrorCode.NOT_FOUND, "Queue: '" + method.getQueue() + "' not found"); @@ -991,10 +989,9 @@ public class ServerSessionDelegate extends SessionDelegate { try { - Binding binding = exchange.getBinding(method.getBindingKey(), queue); - if(binding != null) + if(exchange.hasBinding(method.getBindingKey(), queue)) { - binding.delete(); + exchange.deleteBinding(method.getBindingKey(), queue); } } catch (AccessControlException e) @@ -1011,7 +1008,7 @@ public class ServerSessionDelegate extends SessionDelegate ExchangeBoundResult result = new ExchangeBoundResult(); VirtualHost virtualHost = getVirtualHost(session); - Exchange exchange; + ExchangeImpl exchange; AMQQueue queue; if(method.hasExchange()) { @@ -1378,7 +1375,7 @@ public class ServerSessionDelegate extends SessionDelegate arguments.put(attrName, queue.getAttribute(attrName)); } result.setArguments(QueueArgumentsConverter.convertModelArgsToWire(arguments)); - result.setMessageCount(queue.getMessageCount()); + result.setMessageCount(queue.getQueueDepthMessages()); result.setSubscriberCount(queue.getConsumerCount()); } diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java index 453f3035e1..9e0c5b6be6 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java @@ -32,6 +32,7 @@ import org.apache.log4j.Logger; import org.apache.qpid.AMQConnectionException; import org.apache.qpid.AMQException; import org.apache.qpid.server.connection.SessionPrincipal; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.filter.AMQInvalidArgumentException; import org.apache.qpid.server.filter.Filterable; import org.apache.qpid.server.filter.MessageFilter; @@ -48,7 +49,6 @@ import org.apache.qpid.protocol.AMQConstant; import org.apache.qpid.server.TransactionTimeoutHelper; import org.apache.qpid.server.TransactionTimeoutHelper.CloseAction; import org.apache.qpid.server.configuration.BrokerProperties; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.filter.FilterManager; import org.apache.qpid.server.filter.FilterManagerFactory; import org.apache.qpid.server.filter.SimpleFilterManager; @@ -1267,14 +1267,14 @@ public class AMQChannel<T extends AMQProtocolSession<T>> } - private class ImmediateAction<C extends Consumer> implements Action<MessageInstance<?,C>> + private class ImmediateAction implements Action<MessageInstance> { public ImmediateAction() { } - public void performAction(MessageInstance<?,C> entry) + public void performAction(MessageInstance entry) { TransactionLogResource queue = entry.getOwningResource(); @@ -1332,10 +1332,10 @@ public class AMQChannel<T extends AMQProtocolSession<T>> } } - private final class CapacityCheckAction<C extends Consumer> implements Action<MessageInstance<?,C>> + private final class CapacityCheckAction implements Action<MessageInstance> { @Override - public void performAction(final MessageInstance<?,C> entry) + public void performAction(final MessageInstance entry) { TransactionLogResource queue = entry.getOwningResource(); if(queue instanceof CapacityChecker) @@ -1569,7 +1569,7 @@ public class AMQChannel<T extends AMQProtocolSession<T>> { final AMQQueue queue = (AMQQueue) owningResource; - final Exchange altExchange = queue.getAlternateExchange(); + final ExchangeImpl altExchange = queue.getAlternateExchange(); if (altExchange == null) { diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/BasicGetMethodHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/BasicGetMethodHandler.java index 611999d8c6..f620abf30f 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/BasicGetMethodHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/BasicGetMethodHandler.java @@ -210,7 +210,7 @@ public class BasicGetMethodHandler implements StateAwareMethodListener<BasicGetB props, _channel.getChannelId(), deliveryTag, - _queue.getMessageCount()); + _queue.getQueueDepthMessages()); _deliveredMessage = true; } diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeBoundHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeBoundHandler.java index 4ebddb0f68..27837844ff 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeBoundHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeBoundHandler.java @@ -25,8 +25,8 @@ import org.apache.qpid.framing.AMQShortString; import org.apache.qpid.framing.ExchangeBoundBody; import org.apache.qpid.framing.ExchangeBoundOkBody; import org.apache.qpid.framing.MethodRegistry; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.protocol.v0_8.AMQChannel; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.protocol.v0_8.AMQProtocolSession; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.protocol.v0_8.state.AMQStateManager; @@ -82,7 +82,7 @@ public class ExchangeBoundHandler implements StateAwareMethodListener<ExchangeBo AMQShortString exchangeName = body.getExchange() == null ? AMQShortString.EMPTY_STRING : body.getExchange(); AMQShortString queueName = body.getQueue(); AMQShortString routingKey = body.getRoutingKey(); - Exchange exchange = virtualHost.getExchange(exchangeName.toString()); + ExchangeImpl exchange = virtualHost.getExchange(exchangeName.toString()); ExchangeBoundOkBody response; if (exchange == null) { diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeclareHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeclareHandler.java index 9446f53188..3b630c684c 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeclareHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeclareHandler.java @@ -30,13 +30,12 @@ import org.apache.qpid.framing.AMQShortString; import org.apache.qpid.framing.ExchangeDeclareBody; import org.apache.qpid.framing.MethodRegistry; import org.apache.qpid.protocol.AMQConstant; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.model.LifetimePolicy; import org.apache.qpid.server.protocol.v0_8.AMQChannel; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.protocol.v0_8.AMQProtocolSession; import org.apache.qpid.server.protocol.v0_8.state.AMQStateManager; import org.apache.qpid.server.protocol.v0_8.state.StateAwareMethodListener; -import org.apache.qpid.server.virtualhost.AbstractVirtualHost; import org.apache.qpid.server.virtualhost.ExchangeExistsException; import org.apache.qpid.server.virtualhost.ReservedExchangeNameException; import org.apache.qpid.server.virtualhost.UnknownExchangeException; @@ -77,7 +76,7 @@ public class ExchangeDeclareHandler implements StateAwareMethodListener<Exchange _logger.debug("Request to declare exchange of type " + body.getType() + " with name " + exchangeName); } - Exchange exchange; + ExchangeImpl exchange; if (body.getPassive()) { diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeleteHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeleteHandler.java index bbe6028a63..720677064b 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeleteHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeleteHandler.java @@ -24,8 +24,8 @@ import org.apache.qpid.AMQException; import org.apache.qpid.framing.ExchangeDeleteBody; import org.apache.qpid.framing.ExchangeDeleteOkBody; import org.apache.qpid.protocol.AMQConstant; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.protocol.v0_8.AMQChannel; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.protocol.v0_8.AMQProtocolSession; import org.apache.qpid.server.protocol.v0_8.state.AMQStateManager; import org.apache.qpid.server.protocol.v0_8.state.StateAwareMethodListener; @@ -62,7 +62,7 @@ public class ExchangeDeleteHandler implements StateAwareMethodListener<ExchangeD { final String exchangeName = body.getExchange() == null ? null : body.getExchange().toString(); - final Exchange exchange = virtualHost.getExchange(exchangeName); + final ExchangeImpl exchange = virtualHost.getExchange(exchangeName); if(exchange == null) { throw body.getChannelException(AMQConstant.NOT_FOUND, "No such exchange: " + body.getExchange()); diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueBindHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueBindHandler.java index 7ac71babf3..1e0382f456 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueBindHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueBindHandler.java @@ -29,10 +29,9 @@ import org.apache.qpid.framing.FieldTable; import org.apache.qpid.framing.MethodRegistry; import org.apache.qpid.framing.QueueBindBody; import org.apache.qpid.protocol.AMQConstant; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.exchange.TopicExchange; import org.apache.qpid.server.protocol.v0_8.AMQChannel; -import org.apache.qpid.server.binding.Binding; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.protocol.v0_8.AMQProtocolSession; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.protocol.v0_8.state.AMQStateManager; @@ -103,7 +102,7 @@ public class QueueBindHandler implements StateAwareMethodListener<QueueBindBody> throw body.getChannelException(AMQConstant.NOT_FOUND, "Queue " + queueName + " does not exist."); } final String exchangeName = body.getExchange() == null ? null : body.getExchange().toString(); - final Exchange exch = virtualHost.getExchange(exchangeName); + final ExchangeImpl exch = virtualHost.getExchange(exchangeName); if (exch == null) { throw body.getChannelException(AMQConstant.NOT_FOUND, "Exchange " + exchangeName + " does not exist."); @@ -121,13 +120,7 @@ public class QueueBindHandler implements StateAwareMethodListener<QueueBindBody> if(!exch.addBinding(bindingKey, queue, arguments) && TopicExchange.TYPE.equals(exch.getExchangeType())) { - Binding oldBinding = exch.getBinding(bindingKey, queue); - - Map<String, Object> oldArgs = oldBinding.getArguments(); - if((oldArgs == null && !arguments.isEmpty()) || (oldArgs != null && !oldArgs.equals(arguments))) - { - exch.replaceBinding(oldBinding.getId(), bindingKey, queue, arguments); - } + exch.replaceBinding(bindingKey, queue, arguments); } } } diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueDeclareHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueDeclareHandler.java index 4a3bd1e921..97aac8424a 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueDeclareHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueDeclareHandler.java @@ -163,7 +163,7 @@ public class QueueDeclareHandler implements StateAwareMethodListener<QueueDeclar MethodRegistry methodRegistry = protocolConnection.getMethodRegistry(); QueueDeclareOkBody responseBody = methodRegistry.createQueueDeclareOkBody(queueName, - queue.getMessageCount(), + queue.getQueueDepthMessages(), queue.getConsumerCount()); protocolConnection.writeFrame(responseBody.generateFrame(channelId)); diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueUnbindHandler.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueUnbindHandler.java index f6dbd0cee0..a828ca323d 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueUnbindHandler.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueUnbindHandler.java @@ -30,9 +30,8 @@ import org.apache.qpid.framing.QueueUnbindBody; import org.apache.qpid.framing.amqp_0_9.MethodRegistry_0_9; import org.apache.qpid.framing.amqp_0_91.MethodRegistry_0_91; import org.apache.qpid.protocol.AMQConstant; -import org.apache.qpid.server.binding.Binding; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.protocol.v0_8.AMQChannel; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.protocol.v0_8.AMQProtocolSession; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.protocol.v0_8.state.AMQStateManager; @@ -94,13 +93,13 @@ public class QueueUnbindHandler implements StateAwareMethodListener<QueueUnbindB { throw body.getChannelException(AMQConstant.NOT_FOUND, "Queue " + body.getQueue() + " does not exist."); } - final Exchange exch = virtualHost.getExchange(body.getExchange() == null ? null : body.getExchange().toString()); + final ExchangeImpl exch = virtualHost.getExchange(body.getExchange() == null ? null : body.getExchange().toString()); if (exch == null) { throw body.getChannelException(AMQConstant.NOT_FOUND, "Exchange " + body.getExchange() + " does not exist."); } - if(exch.getBinding(String.valueOf(routingKey), queue) == null) + if(!exch.hasBinding(String.valueOf(routingKey), queue)) { throw body.getChannelException(AMQConstant.NOT_FOUND,"No such binding"); } @@ -108,11 +107,7 @@ public class QueueUnbindHandler implements StateAwareMethodListener<QueueUnbindB { try { - Binding binding = exch.getBinding(String.valueOf(routingKey), queue); - if(binding != null) - { - binding.delete(); - } + exch.deleteBinding(String.valueOf(routingKey), queue); } catch (AccessControlException e) { diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/AMQChannelTest.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/AMQChannelTest.java index 317a544b6c..face730195 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/AMQChannelTest.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/AMQChannelTest.java @@ -26,13 +26,12 @@ import static org.mockito.Mockito.when; import java.util.HashMap; import java.util.Map; -import org.apache.qpid.AMQException; import org.apache.qpid.framing.AMQShortString; import org.apache.qpid.framing.BasicContentHeaderProperties; import org.apache.qpid.framing.ContentHeaderBody; import org.apache.qpid.framing.abstraction.MessagePublishInfo; import org.apache.qpid.server.configuration.BrokerProperties; -import org.apache.qpid.server.exchange.Exchange; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.message.MessageContentSource; import org.apache.qpid.server.model.Broker; import org.apache.qpid.server.util.BrokerTestHelper; @@ -100,7 +99,7 @@ public class AMQChannelTest extends QpidTestCase channel.setLocalTransactional(); MessagePublishInfo info = mock(MessagePublishInfo.class); - Exchange e = mock(Exchange.class); + ExchangeImpl e = mock(ExchangeImpl.class); ContentHeaderBody contentHeaderBody= mock(ContentHeaderBody.class); BasicContentHeaderProperties properties = mock(BasicContentHeaderProperties.class); @@ -123,7 +122,7 @@ public class AMQChannelTest extends QpidTestCase channel.setLocalTransactional(); MessagePublishInfo info = mock(MessagePublishInfo.class); - Exchange e = mock(Exchange.class); + ExchangeImpl e = mock(ExchangeImpl.class); ContentHeaderBody contentHeaderBody= mock(ContentHeaderBody.class); BasicContentHeaderProperties properties = mock(BasicContentHeaderProperties.class); diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/BrokerTestHelper_0_8.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/BrokerTestHelper_0_8.java index e5a3475feb..86adc585c3 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/BrokerTestHelper_0_8.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/BrokerTestHelper_0_8.java @@ -25,7 +25,7 @@ import org.apache.qpid.framing.AMQShortString; import org.apache.qpid.framing.BasicContentHeaderProperties; import org.apache.qpid.framing.ContentHeaderBody; import org.apache.qpid.framing.abstraction.MessagePublishInfo; -import org.apache.qpid.server.exchange.Exchange; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.util.BrokerTestHelper; import org.apache.qpid.server.virtualhost.VirtualHost; @@ -73,7 +73,7 @@ public class BrokerTestHelper_0_8 extends BrokerTestHelper when(info.getExchange()).thenReturn(exchangeNameAsShortString); when(info.getRoutingKey()).thenReturn(routingKey); - Exchange exchange = channel.getVirtualHost().getExchange(exchangeName); + ExchangeImpl exchange = channel.getVirtualHost().getExchange(exchangeName); for (int count = 0; count < numberOfMessages; count++) { channel.setPublishFrame(info, exchange); diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/ExtractResendAndRequeueTest.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/ExtractResendAndRequeueTest.java index aa5a75396a..f18da87d09 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/ExtractResendAndRequeueTest.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/ExtractResendAndRequeueTest.java @@ -75,7 +75,7 @@ public class ExtractResendAndRequeueTest extends TestCase when(_queue.getName()).thenReturn(getName()); when(_queue.isDeleted()).thenReturn(_queueDeleted); _consumer = mock(Consumer.class); - when(_consumer.getId()).thenReturn(Consumer.SUB_ID_GENERATOR.getAndIncrement()); + when(_consumer.getConsumerNumber()).thenReturn(Consumer.CONSUMER_NUMBER_GENERATOR.getAndIncrement()); long id = 0; diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java index 5356a6e6a3..d83665ad39 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java @@ -25,7 +25,7 @@ import org.apache.qpid.amqp_1_0.type.messaging.Accepted; import org.apache.qpid.amqp_1_0.type.messaging.Rejected; import org.apache.qpid.amqp_1_0.type.messaging.TerminusDurability; import org.apache.qpid.amqp_1_0.type.messaging.TerminusExpiryPolicy; -import org.apache.qpid.server.exchange.Exchange; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.message.InstanceProperties; import org.apache.qpid.server.txn.ServerTransaction; @@ -35,11 +35,11 @@ public class ExchangeDestination implements ReceivingDestination, SendingDestina public static final Rejected REJECTED = new Rejected(); private static final Outcome[] OUTCOMES = { ACCEPTED, REJECTED}; - private Exchange _exchange; + private ExchangeImpl _exchange; private TerminusDurability _durability; private TerminusExpiryPolicy _expiryPolicy; - public ExchangeDestination(Exchange exchange, TerminusDurability durable, TerminusExpiryPolicy expiryPolicy) + public ExchangeDestination(ExchangeImpl exchange, TerminusDurability durable, TerminusExpiryPolicy expiryPolicy) { _exchange = exchange; _durability = durable; @@ -98,7 +98,7 @@ public class ExchangeDestination implements ReceivingDestination, SendingDestina return 20000; } - public Exchange getExchange() + public ExchangeImpl getExchange() { return _exchange; } diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/NodeReceivingDestination.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/NodeReceivingDestination.java index 70f659b546..f7f049831e 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/NodeReceivingDestination.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/NodeReceivingDestination.java @@ -25,7 +25,6 @@ import org.apache.qpid.amqp_1_0.type.messaging.Accepted; import org.apache.qpid.amqp_1_0.type.messaging.Rejected; import org.apache.qpid.amqp_1_0.type.messaging.TerminusDurability; import org.apache.qpid.amqp_1_0.type.messaging.TerminusExpiryPolicy; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.message.InstanceProperties; import org.apache.qpid.server.message.MessageDestination; import org.apache.qpid.server.txn.ServerTransaction; diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java index 94120371fb..394ab69990 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java @@ -21,16 +21,12 @@ package org.apache.qpid.server.protocol.v1_0; import java.security.AccessControlException; -import java.util.ArrayList; -import java.util.Collections; -import java.util.EnumSet; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.UUID; +import java.util.*; import java.util.concurrent.ConcurrentHashMap; import org.apache.log4j.Logger; +import org.apache.qpid.server.binding.BindingImpl; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.model.ExclusivityPolicy; import org.apache.qpid.server.model.LifetimePolicy; import org.apache.qpid.server.model.Queue; @@ -51,9 +47,7 @@ import org.apache.qpid.amqp_1_0.type.transport.Error; import org.apache.qpid.amqp_1_0.type.transport.Transfer; import org.apache.qpid.filter.SelectorParsingException; import org.apache.qpid.filter.selector.ParseException; -import org.apache.qpid.server.binding.Binding; import org.apache.qpid.server.exchange.DirectExchange; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.exchange.TopicExchange; import org.apache.qpid.server.filter.JMSSelectorFilter; import org.apache.qpid.server.filter.SimpleFilterManager; @@ -205,7 +199,7 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS } AMQQueue queue = _vhost.getQueue(name); - Exchange exchange = exchangeDestination.getExchange(); + ExchangeImpl exchange = exchangeDestination.getExchange(); if(queue == null) { @@ -220,17 +214,16 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS } else { - List<Binding> bindings = queue.getBindings(); - List<Binding> bindingsToRemove = new ArrayList<Binding>(); - for(Binding existingBinding : bindings) + Collection<BindingImpl> bindings = queue.getBindings(); + List<BindingImpl> bindingsToRemove = new ArrayList<BindingImpl>(); + for(BindingImpl existingBinding : bindings) { - if(existingBinding.getExchangeImpl() != _vhost.getDefaultExchange() - && existingBinding.getExchangeImpl() != exchange) + if(existingBinding.getExchange() != exchange) { bindingsToRemove.add(existingBinding); } } - for(Binding existingBinding : bindingsToRemove) + for(BindingImpl existingBinding : bindingsToRemove) { existingBinding.delete(); } diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java index ac1517aaf5..6132b48722 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java @@ -43,9 +43,9 @@ import org.apache.qpid.amqp_1_0.type.transport.*; import org.apache.qpid.amqp_1_0.type.transport.Error; import org.apache.qpid.server.connection.SessionPrincipal; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.model.*; import org.apache.qpid.protocol.AMQConstant; -import org.apache.qpid.server.exchange.Exchange; import org.apache.qpid.server.logging.LogSubject; import org.apache.qpid.server.message.MessageDestination; import org.apache.qpid.server.message.MessageSource; @@ -138,7 +138,7 @@ public class Session_1_0 implements SessionEventListener, AMQSessionModel<Sessio } else { - Exchange exchg = getVirtualHost().getExchange(addr); + ExchangeImpl exchg = getVirtualHost().getExchange(addr); if(exchg != null) { destination = new ExchangeDestination(exchg, source.getDurable(), source.getExpiryPolicy()); diff --git a/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-0-10/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v0_10/MessageConverter_0_10_to_0_8.java b/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-0-10/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v0_10/MessageConverter_0_10_to_0_8.java index 0c83c31ad4..d385f684ea 100644 --- a/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-0-10/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v0_10/MessageConverter_0_10_to_0_8.java +++ b/qpid/java/broker-plugins/amqp-msg-conv-0-8-to-0-10/src/main/java/org/apache/qpid/server/protocol/converter/v0_8_v0_10/MessageConverter_0_10_to_0_8.java @@ -30,7 +30,7 @@ import org.apache.qpid.framing.BasicContentHeaderProperties; import org.apache.qpid.framing.ContentHeaderBody; import org.apache.qpid.framing.FieldTable; import org.apache.qpid.framing.abstraction.MessagePublishInfo; -import org.apache.qpid.server.exchange.Exchange; +import org.apache.qpid.server.exchange.ExchangeImpl; import org.apache.qpid.server.plugin.MessageConverter; import org.apache.qpid.server.protocol.v0_10.MessageTransferMessage; import org.apache.qpid.server.protocol.v0_8.AMQMessage; @@ -110,7 +110,7 @@ public class MessageConverter_0_10_to_0_8 implements MessageConverter<MessageTra exchangeName = ""; } - Exchange exchange = vhost.getExchange(exchangeName); + ExchangeImpl exchange = vhost.getExchange(exchangeName); String exchangeClass = exchange == null ? ExchangeDefaults.DIRECT_EXCHANGE_CLASS : exchange.getTypeName(); diff --git a/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNode.java b/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNode.java index c163bfa238..6029b09466 100644 --- a/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNode.java +++ b/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNode.java @@ -56,7 +56,7 @@ import java.util.*; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; -class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementNode>, MessageDestination +class ManagementNode implements MessageSource, MessageDestination { public static final String NAME_ATTRIBUTE = "name"; @@ -93,8 +93,8 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN private final UUID _id; - private final CopyOnWriteArrayList<ConsumerRegistrationListener<ManagementNode>> _consumerRegistrationListeners = - new CopyOnWriteArrayList<ConsumerRegistrationListener<ManagementNode>>(); + private final CopyOnWriteArrayList<ConsumerRegistrationListener<? super MessageSource>> _consumerRegistrationListeners = + new CopyOnWriteArrayList<ConsumerRegistrationListener<? super MessageSource>>(); private final SystemNodeCreator.SystemNodeRegistry _registry; private final ConfiguredObject<?> _managedObject; @@ -139,18 +139,41 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN private Class getManagementClass(Class objectClass) { - List<Class> allClasses = new ArrayList<Class>(); - allClasses.add(objectClass); - allClasses.addAll(Arrays.asList(objectClass.getInterfaces())); - allClasses.add(objectClass.getSuperclass()); - for(Class clazz : allClasses) + + if(objectClass.getAnnotation(ManagedObject.class)!=null) { - ManagedObject annotation = (ManagedObject) clazz.getAnnotation(ManagedObject.class); - if(annotation != null) + return objectClass; + } + List<Class> allClasses = Collections.singletonList(objectClass); + List<Class> testedClasses = new ArrayList<Class>(); + do + { + testedClasses.addAll( allClasses ); + allClasses = new ArrayList<Class>(); + for(Class c : testedClasses) { - return clazz; + for(Class i : c.getInterfaces()) + { + if(!allClasses.contains(i)) + { + allClasses.add(i); + } + } + if(c.getSuperclass() != null && !allClasses.contains(c.getSuperclass())) + { + allClasses.add(c.getSuperclass()); + } + } + allClasses.removeAll(testedClasses); + for(Class c : allClasses) + { + if(c.getAnnotation(ManagedObject.class) != null) + { + return c; + } } } + while(!allClasses.isEmpty()); return null; } @@ -240,7 +263,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN public <M extends ServerMessage<? extends StorableMessageMetaData>> int send(final M message, final InstanceProperties instanceProperties, final ServerTransaction txn, - final Action<? super MessageInstance<?, ? extends Consumer>> postEnqueueAction) + final Action<? super MessageInstance> postEnqueueAction) { @SuppressWarnings("unchecked") @@ -286,7 +309,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN return header.containsHeader(name) && header.getHeader(name) instanceof String; } - synchronized void enqueue(InternalMessage message, InstanceProperties properties, Action<? super MessageInstance<?, ? extends Consumer>> postEnqueueAction) + synchronized void enqueue(InternalMessage message, InstanceProperties properties, Action<? super MessageInstance> postEnqueueAction) { if(postEnqueueAction != null) { @@ -925,7 +948,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN } @Override - public synchronized <T extends ConsumerTarget> ManagementNodeConsumer addConsumer(final T target, + public synchronized ManagementNodeConsumer addConsumer(final ConsumerTarget target, final FilterManager filters, final Class<? extends ServerMessage> messageClass, final String consumerName, @@ -935,7 +958,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN final ManagementNodeConsumer managementNodeConsumer = new ManagementNodeConsumer(consumerName,this, target); target.consumerAdded(managementNodeConsumer); _consumers.put(consumerName, managementNodeConsumer); - for(ConsumerRegistrationListener<ManagementNode> listener : _consumerRegistrationListeners) + for(ConsumerRegistrationListener<? super MessageSource> listener : _consumerRegistrationListeners) { listener.consumerAdded(this, managementNodeConsumer); } @@ -949,7 +972,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN } @Override - public void addConsumerRegistrationListener(final ConsumerRegistrationListener<ManagementNode> listener) + public void addConsumerRegistrationListener(final ConsumerRegistrationListener<? super MessageSource> listener) { _consumerRegistrationListeners.add(listener); } @@ -984,7 +1007,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN return false; } - private class ConsumedMessageInstance implements MessageInstance<ConsumedMessageInstance,Consumer> + private class ConsumedMessageInstance implements MessageInstance { private final ServerMessage _message; private final InstanceProperties _properties; @@ -1015,13 +1038,13 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN } @Override - public void addStateChangeListener(final StateChangeListener<? super ConsumedMessageInstance, State> listener) + public void addStateChangeListener(final StateChangeListener<? super MessageInstance, State> listener) { } @Override - public boolean removeStateChangeListener(final StateChangeListener<? super ConsumedMessageInstance, State> listener) + public boolean removeStateChangeListener(final StateChangeListener<? super MessageInstance, State> listener) { return false; } @@ -1094,7 +1117,7 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN } @Override - public int routeToAlternate(final Action<? super MessageInstance<?, ? extends Consumer>> action, + public int routeToAlternate(final Action<? super MessageInstance> action, final ServerTransaction txn) { return 0; @@ -1182,7 +1205,8 @@ class ManagementNode implements MessageSource<ManagementNodeConsumer,ManagementN @Override public void childAdded(final ConfiguredObject object, final ConfiguredObject child) { - final ManagedEntityType entityType = _entityTypes.get(getManagementClass(child.getClass()).getName()); + final Class managementClass = getManagementClass(child.getClass()); + final ManagedEntityType entityType = _entityTypes.get(managementClass.getName()); if(entityType != null) { _entities.get(entityType).put(child.getName(), child); diff --git a/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNodeConsumer.java b/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNodeConsumer.java index 1e2c7b0652..8a1f39fdfe 100644 --- a/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNodeConsumer.java +++ b/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNodeConsumer.java @@ -35,7 +35,7 @@ import java.util.concurrent.locks.ReentrantLock; class ManagementNodeConsumer implements Consumer { - private final long _id = Consumer.SUB_ID_GENERATOR.getAndIncrement(); + private final long _id = Consumer.CONSUMER_NUMBER_GENERATOR.getAndIncrement(); private final ManagementNode _managementNode; private final List<ManagementResponse> _queue = Collections.synchronizedList(new ArrayList<ManagementResponse>()); private final ConsumerTarget _target; @@ -95,7 +95,7 @@ class ManagementNodeConsumer implements Consumer } @Override - public long getId() + public long getConsumerNumber() { return _id; } diff --git a/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementResponse.java b/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementResponse.java index 59ab849848..18c68bd198 100644 --- a/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementResponse.java +++ b/qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementResponse.java @@ -31,7 +31,7 @@ import org.apache.qpid.server.txn.ServerTransaction; import org.apache.qpid.server.util.Action; import org.apache.qpid.server.util.StateChangeListener; -class ManagementResponse implements MessageInstance<ManagementResponse,ManagementNodeConsumer> +class ManagementResponse implements MessageInstance { private final ManagementNodeConsumer _consumer; private int _deliveryCount; @@ -65,13 +65,13 @@ class ManagementResponse implements MessageInstance<ManagementResponse,Managemen } @Override - public void addStateChangeListener(final StateChangeListener<? super ManagementResponse, State> listener) + public void addStateChangeListener(final StateChangeListener<? super MessageInstance, State> listener) { } @Override - public boolean removeStateChangeListener(final StateChangeListener<? super ManagementResponse, State> listener) + public boolean removeStateChangeListener(final StateChangeListener<? super MessageInstance, State> listener) { return false; } @@ -84,7 +84,7 @@ class ManagementResponse implements MessageInstance<ManagementResponse,Managemen } @Override - public boolean isAcquiredBy(final ManagementNodeConsumer consumer) + public boolean isAcquiredBy(final Consumer consumer) { return consumer == _consumer && !isDeleted(); } @@ -114,7 +114,7 @@ class ManagementResponse implements MessageInstance<ManagementResponse,Managemen } @Override - public boolean isRejectedBy(final ManagementNodeConsumer consumer) + public boolean isRejectedBy(final Consumer consumer) { return false; } @@ -132,7 +132,7 @@ class ManagementResponse implements MessageInstance<ManagementResponse,Managemen } @Override - public boolean acquire(final ManagementNodeConsumer sub) + public boolean acquire(final Consumer sub) { return false; } @@ -144,7 +144,7 @@ class ManagementResponse implements MessageInstance<ManagementResponse,Managemen } @Override - public int routeToAlternate(final Action<? super MessageInstance<?, ? extends Consumer>> action, + public int routeToAlternate(final Action<? super MessageInstance> action, final ServerTransaction txn) { return 0; diff --git a/qpid/java/broker-plugins/management-http/src/main/java/org/apache/qpid/server/management/plugin/servlet/rest/MessageServlet.java b/qpid/java/broker-plugins/management-http/src/main/java/org/apache/qpid/server/management/plugin/servlet/rest/MessageServlet.java index d28338b354..baf92e8522 100644 --- a/qpid/java/broker-plugins/management-http/src/main/java/org/apache/qpid/server/management/plugin/servlet/rest/MessageServlet.java +++ b/qpid/java/broker-plugins/management-http/src/main/java/org/apache/qpid/server/management/plugin/servlet/rest/MessageServlet.java @@ -328,7 +328,7 @@ public class MessageServlet extends AbstractServlet ? "Acquired" : ""); final Consumer deliveredConsumer = entry.getDeliveredConsumer(); - object.put("deliveredTo", deliveredConsumer == null ? null : deliveredConsumer.getId()); + object.put("deliveredTo", deliveredConsumer == null ? null : deliveredConsumer.getConsumerNumber()); ServerMessage message = entry.getMessage(); if(message != null) diff --git a/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/ExchangeMBean.java b/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/ExchangeMBean.java index 27901bfbf7..8b247ef5fa 100644 --- a/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/ExchangeMBean.java +++ b/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/ExchangeMBean.java @@ -22,6 +22,8 @@ package org.apache.qpid.server.jmx.mbeans; import org.apache.qpid.management.common.mbeans.ManagedExchange; +import org.apache.qpid.server.binding.BindingImpl; +import org.apache.qpid.server.exchange.HeadersExchange; import org.apache.qpid.server.jmx.AMQManagedObject; import org.apache.qpid.server.jmx.ManagedObject; import org.apache.qpid.server.model.Binding; @@ -173,7 +175,7 @@ public class ExchangeMBean extends AMQManagedObject implements ManagedExchange { if(HEADERS_EXCHANGE_TYPE.equals(_exchange.getType())) { - return getHeadersBindings(_exchange.getBindings()); + return getHeadersBindings(_exchange.getBindings()); } else { @@ -181,7 +183,7 @@ public class ExchangeMBean extends AMQManagedObject implements ManagedExchange } } - private TabularData getHeadersBindings(Collection<Binding> bindings) throws OpenDataException + private TabularData getHeadersBindings(Collection<? extends Binding> bindings) throws OpenDataException { TabularType bindinglistDataType = new TabularType("Exchange Bindings", "List of exchange bindings for " + getName(), @@ -221,7 +223,7 @@ public class ExchangeMBean extends AMQManagedObject implements ManagedExchange } - private TabularData getNonHeadersBindings(Collection<Binding> bindings) throws OpenDataException + private TabularData getNonHeadersBindings(Collection<? extends Binding> bindings) throws OpenDataException { TabularType bindinglistDataType = diff --git a/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBean.java b/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBean.java index 88d68cff9a..ca8cc7eb7d 100644 --- a/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBean.java +++ b/qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBean.java @@ -166,7 +166,7 @@ public class VirtualHostManagerMBean extends AbstractStatisticsGatheringMBean<Vi try { getConfiguredObject().createExchange(name, State.ACTIVE, durable, - LifetimePolicy.PERMANENT, 0l, type, Collections.EMPTY_MAP); + LifetimePolicy.PERMANENT, type, Collections.EMPTY_MAP); } catch (IllegalArgumentException iae) { diff --git a/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/QueueMBeanTest.java b/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/QueueMBeanTest.java index 4a88884bf8..f2ca04f709 100644 --- a/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/QueueMBeanTest.java +++ b/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/QueueMBeanTest.java @@ -87,7 +87,7 @@ public class QueueMBeanTest extends QpidTestCase public void testGetMessageCount() throws Exception { - when(_mockQueue.getQueueDepthMessages()).thenReturn(1000l); + when(_mockQueue.getQueueDepthMessages()).thenReturn(1000); assertStatistic("messageCount", 1000); } @@ -105,13 +105,13 @@ public class QueueMBeanTest extends QpidTestCase public void testActiveConsumerCount() throws Exception { - when(_mockQueue.getConsumerCountWithCredit()).thenReturn(3l); + when(_mockQueue.getConsumerCountWithCredit()).thenReturn(3); assertStatistic("activeConsumerCount", 3); } public void testConsumerCount() throws Exception { - when(_mockQueue.getConsumerCount()).thenReturn(3l); + when(_mockQueue.getConsumerCount()).thenReturn(3); assertStatistic("consumerCount", 3); } diff --git a/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBeanTest.java b/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBeanTest.java index 2dc2cb8d3b..8d56e766fc 100644 --- a/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBeanTest.java +++ b/qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBeanTest.java @@ -146,7 +146,7 @@ public class VirtualHostManagerMBeanTest extends TestCase public void testCreateNewDurableExchange() throws Exception { _virtualHostManagerMBean.createNewExchange(TEST_EXCHANGE_NAME, TEST_EXCHANGE_TYPE, true); - verify(_mockVirtualHost).createExchange(TEST_EXCHANGE_NAME, State.ACTIVE, true, LifetimePolicy.PERMANENT, 0, TEST_EXCHANGE_TYPE, EMPTY_ARGUMENT_MAP); + verify(_mockVirtualHost).createExchange(TEST_EXCHANGE_NAME, State.ACTIVE, true, LifetimePolicy.PERMANENT, TEST_EXCHANGE_TYPE, EMPTY_ARGUMENT_MAP); } public void testCreateNewExchangeWithUnknownExchangeType() throws Exception @@ -161,7 +161,7 @@ public class VirtualHostManagerMBeanTest extends TestCase { // PASS } - verify(_mockVirtualHost, never()).createExchange(TEST_EXCHANGE_NAME, State.ACTIVE, true, LifetimePolicy.PERMANENT, 0, exchangeType, EMPTY_ARGUMENT_MAP); + verify(_mockVirtualHost, never()).createExchange(TEST_EXCHANGE_NAME, State.ACTIVE, true, LifetimePolicy.PERMANENT, exchangeType, EMPTY_ARGUMENT_MAP); } public void testUnregisterExchange() throws Exception |
