diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-04-17 01:07:34 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-04-17 01:07:34 +0000 |
| commit | 7177135ca38651943b3701b171ef29e4fa52ad86 (patch) | |
| tree | 4d221e6969f959b7c0c2a3c0736f71bdf935b459 /qpid/java/broker-plugins/amqp-0-8-protocol | |
| parent | 359bd6e75abf11027b668d33d2d733b4cd399e38 (diff) | |
| download | qpid-python-7177135ca38651943b3701b171ef29e4fa52ad86.tar.gz | |
QPID-5709 : [Java Broker] Replace exchange registry / factory with use of common configured object mechanism for registration of children
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1588126 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker-plugins/amqp-0-8-protocol')
3 files changed, 22 insertions, 3 deletions
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 b852b22abb..70094ea7c7 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 @@ -86,6 +86,7 @@ import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.model.ConfigurationChangeListener; import org.apache.qpid.server.model.ConfiguredObject; import org.apache.qpid.server.model.Consumer; +import org.apache.qpid.server.model.Session; import org.apache.qpid.server.model.State; import org.apache.qpid.server.protocol.AMQSessionModel; import org.apache.qpid.server.protocol.CapacityChecker; @@ -195,6 +196,7 @@ public class AMQChannel<T extends AMQProtocolSession<T>> private final CopyOnWriteArrayList<Consumer<?>> _consumers = new CopyOnWriteArrayList<Consumer<?>>(); private final ConfigurationChangeListener _consumerClosedListener = new ConsumerClosedListener(); private final CopyOnWriteArrayList<ConsumerListener> _consumerListeners = new CopyOnWriteArrayList<ConsumerListener>(); + private Session<?> _modelObject; public AMQChannel(T session, int channelId, final MessageStore messageStore) @@ -737,6 +739,10 @@ public class AMQChannel<T extends AMQProtocolSession<T>> _transaction.rollback(); + if(_modelObject != null) + { + _modelObject.delete(); + } try { @@ -1759,4 +1765,16 @@ public class AMQChannel<T extends AMQProtocolSession<T>> { _consumerListeners.remove(listener); } + + @Override + public void setModelObject(final Session<?> session) + { + _modelObject = session; + } + + @Override + public Session<?> getModelObject() + { + return _modelObject; + } } 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 42cb66ce7e..ef8d01d89f 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 @@ -36,7 +36,6 @@ import org.apache.qpid.protocol.AMQConstant; import org.apache.qpid.server.model.ExclusivityPolicy; import org.apache.qpid.server.model.LifetimePolicy; import org.apache.qpid.server.model.Queue; -import org.apache.qpid.server.model.UUIDGenerator; import org.apache.qpid.server.protocol.AMQSessionModel; import org.apache.qpid.server.protocol.v0_8.AMQChannel; import org.apache.qpid.server.protocol.v0_8.AMQProtocolSession; @@ -192,7 +191,7 @@ public class QueueDeclareHandler implements StateAwareMethodListener<QueueDeclar QueueArgumentsConverter.convertWireArgsToModel(FieldTable.convertToMap(body.getArguments())); final String queueNameString = AMQShortString.toString(queueName); attributes.put(Queue.NAME, queueNameString); - attributes.put(Queue.ID, UUIDGenerator.generateQueueUUID(queueNameString, virtualHost.getName())); + attributes.put(Queue.ID, UUID.randomUUID()); attributes.put(Queue.DURABLE, durable); LifetimePolicy lifetimePolicy; diff --git a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java index 2322865c80..7f4a3701cd 100644 --- a/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java +++ b/qpid/java/broker-plugins/amqp-0-8-protocol/src/test/java/org/apache/qpid/server/protocol/v0_8/InternalTestProtocolSession.java @@ -272,11 +272,13 @@ public class InternalTestProtocolSession extends AMQProtocolEngine implements Pr } } + private static final AtomicInteger portNumber = new AtomicInteger(0); + private static class TestNetworkConnection implements NetworkConnection { private String _remoteHost = "127.0.0.1"; private String _localHost = "127.0.0.1"; - private int _port = 1; + private int _port = portNumber.incrementAndGet(); private final Sender<ByteBuffer> _sender; public TestNetworkConnection() |
