summaryrefslogtreecommitdiff
path: root/qpid/java/broker-plugins
diff options
context:
space:
mode:
authorRobert Godfrey <rgodfrey@apache.org>2014-02-26 23:27:39 +0000
committerRobert Godfrey <rgodfrey@apache.org>2014-02-26 23:27:39 +0000
commit4ae118bb7a81155a9f3d22af3a4a3f2191799c83 (patch)
tree8bebb59700d6825f997cb501c9043513e0f785c8 /qpid/java/broker-plugins
parent6339e3b1ad22e74508510e08384c4d484bd9666c (diff)
downloadqpid-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')
-rw-r--r--qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java4
-rw-r--r--qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java4
-rw-r--r--qpid/java/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSessionDelegate.java29
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/AMQChannel.java12
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/BasicGetMethodHandler.java2
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeBoundHandler.java4
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeclareHandler.java5
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/ExchangeDeleteHandler.java4
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueBindHandler.java13
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueDeclareHandler.java2
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/handler/QueueUnbindHandler.java13
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/AMQChannelTest.java7
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/BrokerTestHelper_0_8.java4
-rw-r--r--qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/ExtractResendAndRequeueTest.java2
-rw-r--r--qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ExchangeDestination.java8
-rw-r--r--qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/NodeReceivingDestination.java1
-rw-r--r--qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java25
-rw-r--r--qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java4
-rw-r--r--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.java4
-rw-r--r--qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNode.java66
-rw-r--r--qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementNodeConsumer.java4
-rw-r--r--qpid/java/broker-plugins/management-amqp/src/main/java/org/apache/qpid/server/management/amqp/ManagementResponse.java14
-rw-r--r--qpid/java/broker-plugins/management-http/src/main/java/org/apache/qpid/server/management/plugin/servlet/rest/MessageServlet.java2
-rw-r--r--qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/ExchangeMBean.java8
-rw-r--r--qpid/java/broker-plugins/management-jmx/src/main/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBean.java2
-rw-r--r--qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/QueueMBeanTest.java6
-rw-r--r--qpid/java/broker-plugins/management-jmx/src/test/java/org/apache/qpid/server/jmx/mbeans/VirtualHostManagerMBeanTest.java4
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