From 7b08633eec54ecd50855cf07b34267b1630379b6 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Tue, 10 Feb 2015 16:17:16 +0000 Subject: QPID-6378 : Applying patch from Xin Chen git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658750 13f79535-47bb-0310-9956-ffa450edef68 --- .../src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java | 2 +- .../src/main/java/org/apache/qpid/amqp_1_0/framing/FrameHandler.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) (limited to 'qpid') diff --git a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java index 69b4939070..cd31974e7f 100644 --- a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java +++ b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java @@ -179,7 +179,7 @@ public class Connection implements ExceptionHandler boolean ssl, int channelMax) throws ConnectionException { - this(ssl?"amqp":"amqps",address,port,username,password,maxFrameSize,container, + this(ssl?"amqps":"amqp",address,port,username,password,maxFrameSize,container, remoteHostname, getSslContext(ssl), null, diff --git a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/framing/FrameHandler.java b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/framing/FrameHandler.java index 3e9dca683e..38f28667d6 100644 --- a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/framing/FrameHandler.java +++ b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/framing/FrameHandler.java @@ -225,8 +225,8 @@ public class FrameHandler implements ProtocolHandler // PARSE HERE try { - Object val = _typeHandler.parse(in); - + Object val = in.hasRemaining() ? _typeHandler.parse(in) : null; + if(in.hasRemaining()) { if(val instanceof Transfer) -- cgit v1.2.1 From e4d9cbe7b63b862914696d605303c5401049bb23 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Tue, 10 Feb 2015 16:23:27 +0000 Subject: QPID-6352 : Always check TCP transports without looking for service loader git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658752 13f79535-47bb-0310-9956-ffa450edef68 --- .../src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java | 6 ++++++ 1 file changed, 6 insertions(+) (limited to 'qpid') diff --git a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java index cd31974e7f..a4f9ac5a3a 100644 --- a/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java +++ b/qpid/java/amqp-1-0-client/src/main/java/org/apache/qpid/amqp_1_0/client/Connection.java @@ -291,6 +291,12 @@ public class Connection implements ExceptionHandler private TransportProvider getTransportProvider(final String protocol) throws ConnectionException { + TCPTransportProviderFactory tcpTransportProviderFactory = new TCPTransportProviderFactory(); + if(tcpTransportProviderFactory.getSupportedTransports().contains(protocol)) + { + return tcpTransportProviderFactory.getProvider(protocol); + } + ClassLoader classLoader = Thread.currentThread().getContextClassLoader(); ServiceLoader providerFactories = ServiceLoader.load(TransportProviderFactory.class, classLoader); -- cgit v1.2.1 From 1556a7b49a91a921a1110c12c0ed493c4a74bd5f Mon Sep 17 00:00:00 2001 From: "Stephen D. Huston" Date: Tue, 10 Feb 2015 16:32:52 +0000 Subject: NO-JIRA - refer to OutgoingMessage as 'class' not 'struct' to match its actual definition. Fixes VC12 compile warning. git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658755 13f79535-47bb-0310-9956-ffa450edef68 --- qpid/cpp/src/qpid/client/amqp0_10/SenderImpl.h | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) (limited to 'qpid') diff --git a/qpid/cpp/src/qpid/client/amqp0_10/SenderImpl.h b/qpid/cpp/src/qpid/client/amqp0_10/SenderImpl.h index 3ed3b457ba..35ce82cf5d 100644 --- a/qpid/cpp/src/qpid/client/amqp0_10/SenderImpl.h +++ b/qpid/cpp/src/qpid/client/amqp0_10/SenderImpl.h @@ -36,7 +36,7 @@ namespace amqp0_10 { class AddressResolution; class MessageSink; -struct OutgoingMessage; +class OutgoingMessage; /** * -- cgit v1.2.1 From 4072a4925990eedbaa36844902ca584132b66806 Mon Sep 17 00:00:00 2001 From: Justin Ross Date: Tue, 10 Feb 2015 19:10:49 +0000 Subject: QPID-5703: Quiet the code generators git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658784 13f79535-47bb-0310-9956-ffa450edef68 --- qpid/cpp/managementgen/qmfgen/generate.py | 6 +++++- qpid/cpp/rubygen/amqpgen.rb | 9 +++++++-- 2 files changed, 12 insertions(+), 3 deletions(-) (limited to 'qpid') diff --git a/qpid/cpp/managementgen/qmfgen/generate.py b/qpid/cpp/managementgen/qmfgen/generate.py index a7ad43cc30..22c53aa2e0 100755 --- a/qpid/cpp/managementgen/qmfgen/generate.py +++ b/qpid/cpp/managementgen/qmfgen/generate.py @@ -257,6 +257,8 @@ class CMakeLists(Makefile): class Generator: + verbose = False + """ This class manages code generation using template files. It is instantiated once for an entire code generation session. @@ -350,7 +352,9 @@ class Generator: pass os.rename (tempFile, target) - print "Generated:", target + + if self.verbose: + print "Generated:", target def targetPackageFile (self, schema, templateFile): dot = templateFile.find(".") diff --git a/qpid/cpp/rubygen/amqpgen.rb b/qpid/cpp/rubygen/amqpgen.rb index f42e177225..7eee953c04 100755 --- a/qpid/cpp/rubygen/amqpgen.rb +++ b/qpid/cpp/rubygen/amqpgen.rb @@ -489,6 +489,7 @@ class Generator @prefix=[''] # For indentation or comments. @indentstr=' ' # One indent level. @outdent=2 + @verbose=false end # Declare next file to be public API @@ -504,10 +505,14 @@ class Generator @out=String.new # Generate in memory first yield if block if @path.exist? and @path.read == @out - puts "Skipped #{@path} - unchanged" # Dont generate if unchanged + if @verbose + puts "Skipped #{@path} - unchanged" # Dont generate if unchanged + end else @path.open('w') { |f| f << @out } - puts "Generated #{@path}" + if @verbose + puts "Generated #{@path}" + end end end end -- cgit v1.2.1 From 77208f71328e3a62ab970401eae083980e94ab44 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Tue, 10 Feb 2015 23:52:08 +0000 Subject: QPID-6380 : close()ing a durable subscription should detach rather than close a link git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658843 13f79535-47bb-0310-9956-ffa450edef68 --- .../org/apache/qpid/amqp_1_0/jms/impl/TopicSubscriberImpl.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) (limited to 'qpid') diff --git a/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/TopicSubscriberImpl.java b/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/TopicSubscriberImpl.java index b89025a27b..4b53cfa795 100644 --- a/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/TopicSubscriberImpl.java +++ b/qpid/java/amqp-1-0-client-jms/src/main/java/org/apache/qpid/amqp_1_0/jms/impl/TopicSubscriberImpl.java @@ -131,6 +131,13 @@ public class TopicSubscriberImpl extends MessageConsumerImpl implements TopicSub protected void closeUnderlyingReceiver(Receiver receiver) { - receiver.close(); + if(isDurable()) + { + receiver.detach(); + } + else + { + receiver.close(); + } } } -- cgit v1.2.1 From 93024d74d2711c3c3cdab6e98f7158ca730abbe1 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Tue, 10 Feb 2015 23:57:54 +0000 Subject: QPID-6381 : if detach with close=true is received, then actually destroy the link git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658845 13f79535-47bb-0310-9956-ffa450edef68 --- .../main/java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) (limited to 'qpid') 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 123d6ac2fb..cdaf5f0ed6 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 @@ -408,7 +408,7 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS { //TODO // if not durable or close - if(!TerminusDurability.UNSETTLED_STATE.equals(_durability)) + if(Boolean.TRUE.equals(detach.getClosed()) || !TerminusDurability.UNSETTLED_STATE.equals(_durability)) { while(!_consumer.trySendLock()) { -- cgit v1.2.1 From 35f8db0065335d4da24de4459cf228b077218138 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 11 Feb 2015 01:04:08 +0000 Subject: QPID-6384 : fix various issues with durable links git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658849 13f79535-47bb-0310-9956-ffa450edef68 --- .../qpid/amqp_1_0/transport/LinkEndpoint.java | 21 +++++++++++++++++---- .../qpid/amqp_1_0/transport/SessionEndpoint.java | 16 ++++++++++++---- .../server/protocol/v1_0/ConsumerTarget_1_0.java | 8 +++++++- .../qpid/server/protocol/v1_0/SendingLink_1_0.java | 3 ++- 4 files changed, 38 insertions(+), 10 deletions(-) (limited to 'qpid') diff --git a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/LinkEndpoint.java b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/LinkEndpoint.java index 434f939a21..246d43d3de 100644 --- a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/LinkEndpoint.java +++ b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/LinkEndpoint.java @@ -21,15 +21,28 @@ package org.apache.qpid.amqp_1_0.transport; -import org.apache.qpid.amqp_1_0.type.*; -import org.apache.qpid.amqp_1_0.type.transport.*; -import org.apache.qpid.amqp_1_0.type.transport.Error; - import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.concurrent.TimeoutException; +import org.apache.qpid.amqp_1_0.type.Binary; +import org.apache.qpid.amqp_1_0.type.DeliveryState; +import org.apache.qpid.amqp_1_0.type.Outcome; +import org.apache.qpid.amqp_1_0.type.Source; +import org.apache.qpid.amqp_1_0.type.Symbol; +import org.apache.qpid.amqp_1_0.type.Target; +import org.apache.qpid.amqp_1_0.type.UnsignedInteger; +import org.apache.qpid.amqp_1_0.type.UnsignedLong; +import org.apache.qpid.amqp_1_0.type.transport.Attach; +import org.apache.qpid.amqp_1_0.type.transport.Detach; +import org.apache.qpid.amqp_1_0.type.transport.Error; +import org.apache.qpid.amqp_1_0.type.transport.Flow; +import org.apache.qpid.amqp_1_0.type.transport.ReceiverSettleMode; +import org.apache.qpid.amqp_1_0.type.transport.Role; +import org.apache.qpid.amqp_1_0.type.transport.SenderSettleMode; +import org.apache.qpid.amqp_1_0.type.transport.Transfer; + public abstract class LinkEndpoint { diff --git a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/SessionEndpoint.java b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/SessionEndpoint.java index 0f37518773..5a28ddcb60 100644 --- a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/SessionEndpoint.java +++ b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/SessionEndpoint.java @@ -98,6 +98,12 @@ public class SessionEndpoint private int _availableOutgoingCredit; private UnsignedInteger _lastSentIncomingLimit; + private final Error _sessionEndedLinkError = + new Error(LinkError.DETACH_FORCED, + "Force detach the link because the session is remotely ended."); + + + public SessionEndpoint(final ConnectionEndpoint connectionEndpoint) { this(connectionEndpoint, UnsignedInteger.valueOf(0)); @@ -240,19 +246,21 @@ public class SessionEndpoint private void detachLinks() { Collection handles = new ArrayList(_remoteLinkEndpoints.keySet()); - Error error = new Error(); - error.setCondition(LinkError.DETACH_FORCED); - error.setDescription("Force detach the link because the session is remotely ended."); for(UnsignedInteger handle : handles) { Detach detach = new Detach(); detach.setClosed(false); detach.setHandle(handle); - detach.setError(error); + detach.setError(_sessionEndedLinkError); detach(handle, detach); } } + public boolean isSyntheticError(Error error) + { + return error == _sessionEndedLinkError; + } + public short getSendingChannel() { return _sendingChannel; diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java index 598fce03b9..f19ce6b1be 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java @@ -46,6 +46,7 @@ import org.apache.qpid.server.message.MessageInstance; import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.plugin.MessageConverter; import org.apache.qpid.server.protocol.AMQSessionModel; +import org.apache.qpid.server.protocol.LinkRegistry; import org.apache.qpid.server.protocol.MessageConverterRegistry; import org.apache.qpid.server.txn.ServerTransaction; import org.apache.qpid.server.util.ConnectionScopedRuntimeException; @@ -283,7 +284,12 @@ class ConsumerTarget_1_0 extends AbstractConsumerTarget { //TODO getEndpoint().setSource(null); - getEndpoint().detach(); + getEndpoint().close(); + + final LinkRegistry linkReg = getSession().getConnection() + .getVirtualHost() + .getLinkRegistry(getEndpoint().getSession().getConnection().getRemoteContainerId()); + linkReg.unregisterSendingLink(getEndpoint().getName()); } public boolean allocateCredit(final ServerMessage msg) 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 cdaf5f0ed6..e3994005d6 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 @@ -464,7 +464,8 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS _consumer.releaseSendLock(); } } - else if(detach == null || detach.getError() != null) + else if(detach.getError() != null + && !_linkAttachment.getEndpoint().getSession().isSyntheticError(detach.getError())) { _linkAttachment = null; _target.flowStateChanged(); -- cgit v1.2.1 From 83393913d1a7c9bced3c6681423b55fe1d6717e9 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 11 Feb 2015 11:06:23 +0000 Subject: QPID-6383 : don't auto delete queues because the children are being removed due to broker closure git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658923 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/server/queue/AbstractQueue.java | 14 +++++++++++++- .../qpid/server/protocol/v1_0/SendingLink_1_0.java | 16 +++++++--------- 2 files changed, 20 insertions(+), 10 deletions(-) (limited to 'qpid') diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java index f905558f13..8f4c3d6df0 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/AbstractQueue.java @@ -245,6 +245,7 @@ public abstract class AbstractQueue> private final ConcurrentLinkedQueue _postRecoveryQueue = new ConcurrentLinkedQueue<>(); private final QueueRunner _queueRunner = new QueueRunner(this); + private boolean _closing; protected AbstractQueue(Map attributes, VirtualHostImpl virtualHost) { @@ -753,6 +754,15 @@ public abstract class AbstractQueue> } + @Override + protected void beforeClose() + { + _closing = true; + super.beforeClose(); + } + + + synchronized void unregisterConsumer(final QueueConsumerImpl consumer) { if (consumer == null) @@ -793,7 +803,8 @@ public abstract class AbstractQueue> if(!consumer.isTransient() && ( getLifetimePolicy() == LifetimePolicy.DELETE_ON_NO_OUTBOUND_LINKS || getLifetimePolicy() == LifetimePolicy.DELETE_ON_NO_LINKS ) - && getConsumerCount() == 0) + && getConsumerCount() == 0 + && !(consumer.isDurable() && _closing)) { if (_logger.isInfoEnabled()) @@ -1797,6 +1808,7 @@ public abstract class AbstractQueue> { ReferenceCountingExecutorService.getInstance().releaseExecutorService(); } + _closing = false; } public void checkCapacity(AMQSessionModel channel) 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 e3994005d6..d1d1227818 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 @@ -292,15 +292,7 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS actualFilters.put(entry.getKey(), entry.getValue()); } - catch (ParseException e) - { - Error error = new Error(); - error.setCondition(AmqpError.INVALID_FIELD); - error.setDescription("Invalid JMS Selector: " + selectorFilter.getValue()); - error.setInfo(Collections.singletonMap(Symbol.valueOf("field"), Symbol.valueOf("filter"))); - throw new AmqpErrorException(error); - } - catch (SelectorParsingException e) + catch (ParseException | SelectorParsingException e) { Error error = new Error(); error.setCondition(AmqpError.INVALID_FIELD); @@ -364,6 +356,12 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS options.add(ConsumerImpl.Option.NO_LOCAL); } + if(_durability == TerminusDurability.CONFIGURATION || + _durability == TerminusDurability.UNSETTLED_STATE ) + { + options.add(ConsumerImpl.Option.DURABLE); + } + try { final String name; -- cgit v1.2.1 From bcfed28760d8fe968c71de9752296b9a5f34d2c2 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 11 Feb 2015 11:20:33 +0000 Subject: QPID-6383 : Check in missing file from last commit git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658927 13f79535-47bb-0310-9956-ffa450edef68 --- .../src/main/java/org/apache/qpid/server/queue/QueueConsumerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) (limited to 'qpid') diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueConsumerImpl.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueConsumerImpl.java index c85e4058a1..6b02a84e83 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueConsumerImpl.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueConsumerImpl.java @@ -153,7 +153,7 @@ class QueueConsumerImpl attributes.put(EXCLUSIVE, optionSet.contains(Option.EXCLUSIVE)); attributes.put(NO_LOCAL, optionSet.contains(Option.NO_LOCAL)); attributes.put(DISTRIBUTION_MODE, optionSet.contains(Option.ACQUIRES) ? "MOVE" : "COPY"); - attributes.put(DURABLE,false); + attributes.put(DURABLE,optionSet.contains(Option.DURABLE)); attributes.put(LIFETIME_POLICY, LifetimePolicy.DELETE_ON_SESSION_END); if(filters != null) { -- cgit v1.2.1 From cba338185d3c3f9bdc2e0b490df20d07ffade454 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 11 Feb 2015 12:04:30 +0000 Subject: QPID-6383 : Add missing file - really this time git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658941 13f79535-47bb-0310-9956-ffa450edef68 --- .../src/main/java/org/apache/qpid/server/consumer/ConsumerImpl.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) (limited to 'qpid') diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/consumer/ConsumerImpl.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/consumer/ConsumerImpl.java index b15b01ede5..c0db72d498 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/consumer/ConsumerImpl.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/consumer/ConsumerImpl.java @@ -37,7 +37,8 @@ public interface ConsumerImpl SEES_REQUEUES, TRANSIENT, EXCLUSIVE, - NO_LOCAL + NO_LOCAL, + DURABLE } long getBytesOut(); -- cgit v1.2.1 From 51f24c14d4efffede82db506aff5796e4d048118 Mon Sep 17 00:00:00 2001 From: Ken Giusti Date: Wed, 11 Feb 2015 15:12:54 +0000 Subject: QPID-5799: provide notification callback for received messages. git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1658984 13f79535-47bb-0310-9956-ffa450edef68 --- qpid/python/qpid/messaging/driver.py | 3 ++- qpid/python/qpid/messaging/endpoints.py | 18 ++++++++++++++++++ qpid/python/qpid/tests/messaging/endpoints.py | 25 +++++++++++++++++++++++++ 3 files changed, 45 insertions(+), 1 deletion(-) (limited to 'qpid') diff --git a/qpid/python/qpid/messaging/driver.py b/qpid/python/qpid/messaging/driver.py index e7d564f555..7c30e5d4ba 100644 --- a/qpid/python/qpid/messaging/driver.py +++ b/qpid/python/qpid/messaging/driver.py @@ -1368,7 +1368,8 @@ class Engine: assert rcv.received < rcv.impending, "%s, %s" % (rcv.received, rcv.impending) rcv.received += 1 log.debug("RCVD[%s]: %s", ssn.log_id, msg) - ssn.incoming.append(msg) + ssn.message_received(msg) + def _decode(self, xfr): dp = EMPTY_DP diff --git a/qpid/python/qpid/messaging/endpoints.py b/qpid/python/qpid/messaging/endpoints.py index 7d353e1cb4..6d58b4ac25 100644 --- a/qpid/python/qpid/messaging/endpoints.py +++ b/qpid/python/qpid/messaging/endpoints.py @@ -569,6 +569,7 @@ class Session(Endpoint): self.closed = False self._lock = connection._lock + self._msg_received = None def __repr__(self): return "" % self.name @@ -600,6 +601,11 @@ class Session(Endpoint): if self.closed: raise SessionClosed() + def message_received(self, msg): + self.incoming.append(msg) + if self._msg_received: + self._msg_received() + @synchronized def sender(self, target, **options): """ @@ -684,6 +690,18 @@ class Session(Endpoint): return msg return None + @synchronized + def set_message_received_handler(self, handler): + """Register a callback that will be invoked when a message arrives on the + session. Use with caution: since this callback is invoked in the context + of the driver thread, it is not safe to call any of the public messaging + APIs from within this callback. The intent of the handler is to provide + an efficient way to notify the application that a message has arrived. + This can be useful for those applications that need to schedule a task + to poll for received messages without blocking in the messaging API. + """ + self._msg_received = handler + @synchronized def next_receiver(self, timeout=None): if self._ecwait(lambda: self.incoming, timeout): diff --git a/qpid/python/qpid/tests/messaging/endpoints.py b/qpid/python/qpid/tests/messaging/endpoints.py index 247d6e9a29..56722374e5 100644 --- a/qpid/python/qpid/tests/messaging/endpoints.py +++ b/qpid/python/qpid/tests/messaging/endpoints.py @@ -660,6 +660,31 @@ class SessionTests(Base): except Detached: pass + def testRxCallback(self): + """Verify that the callback is invoked when a message is received. + """ + ADDR = 'test-rx_callback-queue; {create: always, delete: receiver}' + class CallbackHandler: + def __init__(self): + self.handler_called = False + def __call__(self): + self.handler_called = True + cb = CallbackHandler() + self.ssn.set_message_received_handler(cb) + rcv = self.ssn.receiver(ADDR) + rcv.capacity = UNLIMITED + snd = self.ssn.sender(ADDR) + assert not cb.handler_called + snd.send("Ping") + deadline = time.time() + self.timeout() + while time.time() < deadline: + if cb.handler_called: + break; + assert cb.handler_called + snd.close() + rcv.close() + + RECEIVER_Q = 'test-receiver-queue; {create: always, delete: always}' class ReceiverTests(Base): -- cgit v1.2.1 From e31d29d9127d2861b818978e324104d2cca64133 Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 11 Feb 2015 18:26:39 +0000 Subject: QPID-6381 : don't delete link endpoints which are detached and have terminus durability CONFIGURATION git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659037 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/server/protocol/v1_0/SendingLink_1_0.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) (limited to 'qpid') 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 d1d1227818..85a0b559c9 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 @@ -406,7 +406,8 @@ public class SendingLink_1_0 implements SendingLinkListener, Link_1_0, DeliveryS { //TODO // if not durable or close - if(Boolean.TRUE.equals(detach.getClosed()) || !TerminusDurability.UNSETTLED_STATE.equals(_durability)) + if(Boolean.TRUE.equals(detach.getClosed()) + || !(TerminusDurability.UNSETTLED_STATE.equals(_durability)|| TerminusDurability.CONFIGURATION.equals(_durability))) { while(!_consumer.trySendLock()) { -- cgit v1.2.1 From 647d21b7ce2bae2174a22019c6cc2e2fd0a2e13d Mon Sep 17 00:00:00 2001 From: Robert Godfrey Date: Wed, 11 Feb 2015 18:48:22 +0000 Subject: QPID-6240 : Rollback of transactions with settled acks should release the queue entries git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659039 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) (limited to 'qpid') diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java index f19ce6b1be..829b3bf336 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java @@ -426,7 +426,7 @@ class ConsumerTarget_1_0 extends AbstractConsumerTarget modified.setDeliveryFailed(true); _link.getEndpoint().updateDisposition(_deliveryTag, modified, true); _link.getEndpoint().sendFlowConditional(); - _queueEntry.unlockAcquisition(); + _queueEntry.release(); } } }); -- cgit v1.2.1 From 08f5f85f8e306c4dc20e75d976270c59753f54a4 Mon Sep 17 00:00:00 2001 From: Justin Ross Date: Wed, 11 Feb 2015 20:43:53 +0000 Subject: QPID-6347: Remove the now obsolete queue_event_generation option; this is a patch from Irina Boverman git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659063 13f79535-47bb-0310-9956-ffa450edef68 --- qpid/cpp/src/qpid/client/QueueOptions.cpp | 7 ------ qpid/cpp/src/qpid/client/QueueOptions.h | 22 ------------------ .../Cheat-Sheet-for-configuring-Queue-Options.xml | 27 ---------------------- .../client/messaging/address/QpidQueueOptions.java | 13 ----------- .../org/apache/qpid/qmf2/tools/QpidConfig.java | 18 --------------- qpid/tools/src/py/qpid-config | 6 +---- 6 files changed, 1 insertion(+), 92 deletions(-) (limited to 'qpid') diff --git a/qpid/cpp/src/qpid/client/QueueOptions.cpp b/qpid/cpp/src/qpid/client/QueueOptions.cpp index e589bd76cd..f7705afeb0 100644 --- a/qpid/cpp/src/qpid/client/QueueOptions.cpp +++ b/qpid/cpp/src/qpid/client/QueueOptions.cpp @@ -40,8 +40,6 @@ const std::string QueueOptions::strRING_STRICT("ring_strict"); const std::string QueueOptions::strLastValueQueue("qpid.last_value_queue"); const std::string QueueOptions::strLVQMatchProperty("qpid.LVQ_key"); const std::string QueueOptions::strLastValueQueueNoBrowse("qpid.last_value_queue_no_browse"); -const std::string QueueOptions::strQueueEventMode("qpid.queue_event_generation"); - QueueOptions::~QueueOptions() {} @@ -101,11 +99,6 @@ void QueueOptions::clearOrdering() erase(strLastValueQueue); } -void QueueOptions::enableQueueEvents(bool enqueueOnly) -{ - setInt(strQueueEventMode, enqueueOnly ? ENQUEUE_ONLY : ENQUEUE_AND_DEQUEUE); -} - } } diff --git a/qpid/cpp/src/qpid/client/QueueOptions.h b/qpid/cpp/src/qpid/client/QueueOptions.h index a2f30a50b5..f7ace15ff9 100644 --- a/qpid/cpp/src/qpid/client/QueueOptions.h +++ b/qpid/cpp/src/qpid/client/QueueOptions.h @@ -75,28 +75,6 @@ class QPID_CLIENT_CLASS_EXTERN QueueOptions: public framing::FieldTable */ QPID_CLIENT_EXTERN void clearOrdering(); - /** - * Turns on event generation for this queue (either enqueue only - * or for enqueue and dequeue events); the events can then be - * processed by a regsitered broker plugin. - * - * DEPRECATED - * - * This is confusing to anyone who sees only the function call - * and not the variable name / doxygen. Consider the following call: - * - * options.enableQueueEvents(false); - * - * It looks like it disables queue events, but what it really does is - * enable both enqueue and dequeue events. - * - * Use setInt() instead: - * - * options.setInt("qpid.queue_event_generation", 2); - */ - - QPID_CLIENT_EXTERN void enableQueueEvents(bool enqueueOnly); - static QPID_CLIENT_EXTERN const std::string strMaxCountKey; static QPID_CLIENT_EXTERN const std::string strMaxSizeKey; static QPID_CLIENT_EXTERN const std::string strTypeKey; diff --git a/qpid/doc/book/src/cpp-broker/Cheat-Sheet-for-configuring-Queue-Options.xml b/qpid/doc/book/src/cpp-broker/Cheat-Sheet-for-configuring-Queue-Options.xml index e693ee463b..125372e463 100644 --- a/qpid/doc/book/src/cpp-broker/Cheat-Sheet-for-configuring-Queue-Options.xml +++ b/qpid/doc/book/src/cpp-broker/Cheat-Sheet-for-configuring-Queue-Options.xml @@ -179,34 +179,7 @@
Setting additional behaviors -
- Queue - event generation - - - This option is used to determine whether enqueue/dequeue events - representing changes made to queue state are generated. These - events can then be processed by plugins such as that used for - . - - Example: - - -#include "qpid/client/QueueOptions.h" - - QueueOptions options; - options.enableQueueEvents(1); - session.queueDeclare(arg::queue="my-queue", arg::arguments=options); - - - The boolean option indicates whether only enqueue events should - be generated. The key set by this is - 'qpid.queue_event_generation' and the value is and integer value - of 1 (to replicate only enqueue events) or 2 (to replicate both - enqueue and dequeue events). -
-
Other Clients diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/messaging/address/QpidQueueOptions.java b/qpid/java/client/src/main/java/org/apache/qpid/client/messaging/address/QpidQueueOptions.java index 5b6c027f4a..24295a0832 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/messaging/address/QpidQueueOptions.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/messaging/address/QpidQueueOptions.java @@ -30,7 +30,6 @@ public class QpidQueueOptions extends HashMap<String,Object> public static final String QPID_LVQ_KEY = "qpid.LVQ_key"; public static final String QPID_LAST_VALUE_QUEUE = "qpid.last_value_queue"; public static final String QPID_LAST_VALUE_QUEUE_NO_BROWSE = "qpid.last_value_queue_no_browse"; - public static final String QPID_QUEUE_EVENT_GENERATION = "qpid.queue_event_generation"; public void validatePolicyType(String type) { @@ -83,16 +82,4 @@ public class QpidQueueOptions extends HashMap<String,Object> this.put(QPID_LVQ_KEY, key); } - public void setQueueEvents(String value) - { - if (value != null && (value.equals("1") || value.equals("2"))) - { - this.put(QPID_QUEUE_EVENT_GENERATION, value); - } - else - { - throw new IllegalArgumentException("Invalid value for " + - QPID_QUEUE_EVENT_GENERATION + " should be one of {1|2}"); - } - } } diff --git a/qpid/tools/src/java/qpid-qmf2-tools/src/main/java/org/apache/qpid/qmf2/tools/QpidConfig.java b/qpid/tools/src/java/qpid-qmf2-tools/src/main/java/org/apache/qpid/qmf2/tools/QpidConfig.java index 04a0642e2e..8cca763e4c 100644 --- a/qpid/tools/src/java/qpid-qmf2-tools/src/main/java/org/apache/qpid/qmf2/tools/QpidConfig.java +++ b/qpid/tools/src/java/qpid-qmf2-tools/src/main/java/org/apache/qpid/qmf2/tools/QpidConfig.java @@ -344,7 +344,6 @@ public final class QpidConfig private String _order = "fifo"; private boolean _msgSequence = false; private boolean _ive = false; - private long _eventGeneration = 0; private String _file = null; // New to Qpid 0.10 qpid-config @@ -364,7 +363,6 @@ public final class QpidConfig private static final String LVQNB = "qpid.last_value_queue_no_browse"; private static final String MSG_SEQUENCE = "qpid.msg_sequence"; private static final String IVE = "qpid.ive"; - private static final String QUEUE_EVENT_GENERATION = "qpid.queue_event_generation"; private static final String FLOW_STOP_COUNT = "qpid.flow_stop_count"; private static final String FLOW_RESUME_COUNT = "qpid.flow_resume_count"; private static final String FLOW_STOP_SIZE = "qpid.flow_stop_size"; @@ -387,7 +385,6 @@ public final class QpidConfig SPECIAL_ARGS.add(LVQNB); SPECIAL_ARGS.add(MSG_SEQUENCE); SPECIAL_ARGS.add(IVE); - SPECIAL_ARGS.add(QUEUE_EVENT_GENERATION); SPECIAL_ARGS.add(FLOW_STOP_COUNT); SPECIAL_ARGS.add(FLOW_RESUME_COUNT); SPECIAL_ARGS.add(FLOW_STOP_SIZE); @@ -714,11 +711,6 @@ for (Map.Entry<String, Object> entry : args.entrySet()) { System.out.printf("--order lvq-no-browse "); } - if (args.containsKey(QUEUE_EVENT_GENERATION)) - { - System.out.printf("--generate-queue-events=%d ", QmfData.getLong(args.get(QUEUE_EVENT_GENERATION))); - } - if (queue.hasValue("altExchange")) { ObjectId altExchangeRef = queue.getRefValue("altExchange"); @@ -936,11 +928,6 @@ for (Map.Entry<String, Object> entry : args.entrySet()) { properties.put(LVQNB, 1l); } - if (_eventGeneration > 0) - { - properties.put(QUEUE_EVENT_GENERATION, _eventGeneration); - } - if (_altExchange != null) { properties.put("alternate-exchange", _altExchange); @@ -1393,11 +1380,6 @@ for (Map.Entry<String, Object> entry : args.entrySet()) { _ive = true; } - if (opt[0].equals("--generate-queue-events")) - { - _eventGeneration = Long.parseLong(opt[1]); - } - if (opt[0].equals("--force")) { _ifEmpty = false; diff --git a/qpid/tools/src/py/qpid-config b/qpid/tools/src/py/qpid-config index 8d38b1a342..816e0f0a08 100755 --- a/qpid/tools/src/py/qpid-config +++ b/qpid/tools/src/py/qpid-config @@ -147,7 +147,6 @@ POLICY_TYPE = "qpid.policy_type" LVQ_KEY = "qpid.last_value_queue_key" MSG_SEQUENCE = "qpid.msg_sequence" IVE = "qpid.ive" -QUEUE_EVENT_GENERATION = "qpid.queue_event_generation" FLOW_STOP_COUNT = "qpid.flow_stop_count" FLOW_RESUME_COUNT = "qpid.flow_resume_count" FLOW_STOP_SIZE = "qpid.flow_stop_size" @@ -164,7 +163,7 @@ REPLICATE = "qpid.replicate" SPECIAL_ARGS=[ FILECOUNT,FILESIZE,EFP_PARTITION_NUM,EFP_POOL_FILE_SIZE, MAX_QUEUE_SIZE,MAX_QUEUE_COUNT,POLICY_TYPE, - LVQ_KEY,MSG_SEQUENCE,IVE,QUEUE_EVENT_GENERATION, + LVQ_KEY,MSG_SEQUENCE,IVE, FLOW_STOP_COUNT,FLOW_RESUME_COUNT,FLOW_STOP_SIZE,FLOW_RESUME_SIZE, MSG_GROUP_HDR_KEY,SHARED_MSG_GROUP,REPLICATE] @@ -541,7 +540,6 @@ class BrokerManager: if MAX_QUEUE_COUNT in args: print "--max-queue-count=%s" % args[MAX_QUEUE_COUNT], if POLICY_TYPE in args: print "--limit-policy=%s" % args[POLICY_TYPE].replace("_", "-"), if LVQ_KEY in args: print "--lvq-key=%s" % args[LVQ_KEY], - if QUEUE_EVENT_GENERATION in args: print "--generate-queue-events=%s" % args[QUEUE_EVENT_GENERATION], if q.altExchange: print "--alternate-exchange=%s" % q.altExchange, if FLOW_STOP_SIZE in args: print "--flow-stop-size=%s" % args[FLOW_STOP_SIZE], @@ -638,8 +636,6 @@ class BrokerManager: if config._lvq_key: declArgs[LVQ_KEY] = config._lvq_key - if config._eventGeneration: - declArgs[QUEUE_EVENT_GENERATION] = config._eventGeneration if config._flowStopSize is not None: declArgs[FLOW_STOP_SIZE] = config._flowStopSize -- cgit v1.2.1 From 90fcef0d551f0defd22a60b447446856cc39e750 Mon Sep 17 00:00:00 2001 From: Keith Wall <kwall@apache.org> Date: Wed, 11 Feb 2015 22:27:52 +0000 Subject: QPID-6387: [Java Client] Remove array optimisation from session/consumer maps git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659103 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/client/AMQSession.java | 92 ++------------------- .../org/apache/qpid/client/AMQSession_0_10.java | 6 +- .../org/apache/qpid/client/AMQSession_0_8.java | 7 +- .../apache/qpid/client/ChannelToSessionMap.java | 93 +++------------------- .../org/apache/qpid/client/XAConnectionImpl.java | 2 +- 5 files changed, 27 insertions(+), 173 deletions(-) (limited to 'qpid') diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java index 12e9285af8..86e1bb0a8b 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java @@ -133,7 +133,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic /** * Flag indicating to start dispatcher as a daemon thread */ - protected final boolean DEAMON_DISPATCHER_THREAD = Boolean.getBoolean(ClientProperties.DAEMON_DISPATCHER); + protected final boolean DAEMON_DISPATCHER_THREAD = Boolean.getBoolean(ClientProperties.DAEMON_DISPATCHER); /** The connection to which this session belongs. */ private AMQConnection _connection; @@ -187,7 +187,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic private MessageFactoryRegistry _messageFactoryRegistry; /** Holds all of the producers created by this session, keyed by their unique identifiers. */ - private Map<Long, MessageProducer> _producers = new ConcurrentHashMap<Long, MessageProducer>(); + private final Map<Long, MessageProducer> _producers = new ConcurrentHashMap<Long, MessageProducer>(); /** * Used as a source of unique identifiers so that the consumers can be tagged to match them to BasicConsume @@ -195,7 +195,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic */ private int _nextTag = 1; - private final IdToConsumerMap<C> _consumers = new IdToConsumerMap<C>(); + private final Map<Integer,C> _consumers = new ConcurrentHashMap<>(); /** * Contains a list of consumers which have been removed but which might still have @@ -294,12 +294,11 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic } /** - * Maps from identifying tags to message consumers, in order to pass dispatch incoming messages to the right - * consumer. + * Consumers associated with this session */ - protected IdToConsumerMap<C> getConsumers() + protected Collection<C> getConsumers() { - return _consumers; + return new ArrayList(_consumers.values()); } protected void setUsingDispatcherForCleanup(boolean usingDispatcherForCleanup) @@ -317,83 +316,6 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic abstract void handleLinkDelete(final AMQDestination dest) throws AMQException; - public static final class IdToConsumerMap<C extends BasicMessageConsumer> - { - private final BasicMessageConsumer[] _fastAccessConsumers = new BasicMessageConsumer[16]; - private final ConcurrentMap<Integer, C> _slowAccessConsumers = new ConcurrentHashMap<Integer, C>(); - - public C get(int id) - { - if ((id & 0xFFFFFFF0) == 0) - { - return (C) _fastAccessConsumers[id]; - } - else - { - return _slowAccessConsumers.get(id); - } - } - - public C put(int id, C consumer) - { - C oldVal; - if ((id & 0xFFFFFFF0) == 0) - { - oldVal = (C) _fastAccessConsumers[id]; - _fastAccessConsumers[id] = consumer; - } - else - { - oldVal = _slowAccessConsumers.put(id, consumer); - } - - return oldVal; - - } - - public C remove(int id) - { - C consumer; - if ((id & 0xFFFFFFF0) == 0) - { - consumer = (C) _fastAccessConsumers[id]; - _fastAccessConsumers[id] = null; - } - else - { - consumer = _slowAccessConsumers.remove(id); - } - - return consumer; - - } - - public Collection<C> values() - { - ArrayList<C> values = new ArrayList<C>(); - - for (int i = 0; i < 16; i++) - { - if (_fastAccessConsumers[i] != null) - { - values.add((C) _fastAccessConsumers[i]); - } - } - values.addAll(_slowAccessConsumers.values()); - - return values; - } - - public void clear() - { - _slowAccessConsumers.clear(); - for (int i = 0; i < 16; i++) - { - _fastAccessConsumers[i] = null; - } - } - } - /** * Creates a new session on a connection. * @@ -2490,7 +2412,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic String dispatcherThreadName = "Dispatcher-" + _channelId + "-Conn-" + _connection.getConnectionNumber(); _dispatcherThread.setName(dispatcherThreadName); - _dispatcherThread.setDaemon(DEAMON_DISPATCHER_THREAD); + _dispatcherThread.setDaemon(DAEMON_DISPATCHER_THREAD); _dispatcher.setConnectionStopped(initiallyStopped); _dispatcherThread.start(); if (_dispatcherLogger.isDebugEnabled()) diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java index dc1f9a719e..206ca15c82 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java @@ -833,7 +833,7 @@ public class AMQSession_0_10 extends AMQSession<BasicMessageConsumer_0_10, Basic { if (suspend) { - for (BasicMessageConsumer consumer : getConsumers().values()) + for (BasicMessageConsumer consumer : getConsumers()) { getQpidSession().messageStop(String.valueOf(consumer.getConsumerTag()), Option.UNRELIABLE); @@ -842,7 +842,7 @@ public class AMQSession_0_10 extends AMQSession<BasicMessageConsumer_0_10, Basic } else { - for (BasicMessageConsumer_0_10 consumer : getConsumers().values()) + for (BasicMessageConsumer_0_10 consumer : getConsumers()) { String consumerTag = String.valueOf(consumer.getConsumerTag()); //only set if msg list is null @@ -1320,7 +1320,7 @@ public class AMQSession_0_10 extends AMQSession<BasicMessageConsumer_0_10, Basic drainDispatchQueue(); setUsingDispatcherForCleanup(false); - for (BasicMessageConsumer consumer : getConsumers().values()) + for (BasicMessageConsumer consumer : getConsumers()) { List<Long> tags = consumer.drainReceiverQueueAndRetrieveDeliveryTags(); getPrefetchedMessageTags().addAll(tags); diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_8.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_8.java index 143de271a1..5fb9329af7 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_8.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_8.java @@ -27,6 +27,7 @@ import static org.apache.qpid.configuration.ClientProperties.QPID_FLOW_CONTROL_W import static org.apache.qpid.configuration.ClientProperties.QPID_FLOW_CONTROL_WAIT_NOTIFY_PERIOD; import java.util.ArrayList; +import java.util.Collection; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; @@ -330,10 +331,9 @@ public class AMQSession_0_8 extends AMQSession<BasicMessageConsumer_0_8, BasicMe { _logger.debug("Prefetched message: _unacknowledgedMessageTags :" + getUnacknowledgedMessageTags()); } - ArrayList<BasicMessageConsumer_0_8> consumersToCheck = new ArrayList<BasicMessageConsumer_0_8>(getConsumers().values()); boolean messageListenerFound = false; boolean serverRejectBehaviourFound = false; - for(BasicMessageConsumer_0_8 consumer : consumersToCheck) + for(BasicMessageConsumer_0_8 consumer : getConsumers()) { if (consumer.isMessageListenerSet()) { @@ -344,7 +344,6 @@ public class AMQSession_0_8 extends AMQSession<BasicMessageConsumer_0_8, BasicMe serverRejectBehaviourFound = true; } } - _logger.debug("about to pre-reject messages for " + consumersToCheck.size() + " consumer(s)"); if (serverRejectBehaviourFound) { @@ -376,7 +375,7 @@ public class AMQSession_0_8 extends AMQSession<BasicMessageConsumer_0_8, BasicMe // consumer on the queue. Whilst this is within the JMS spec it is not // user friendly and avoidable. boolean normalRejectBehaviour = true; - for (BasicMessageConsumer_0_8 consumer : getConsumers().values()) + for (BasicMessageConsumer_0_8 consumer : getConsumers()) { if(RejectBehaviour.SERVER.equals(consumer.getRejectBehaviour())) { diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java b/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java index 2fdb35de49..f46c61daa7 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java @@ -22,106 +22,45 @@ package org.apache.qpid.client; import java.util.ArrayList; import java.util.Collection; -import java.util.LinkedHashMap; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; public final class ChannelToSessionMap { - private final AMQSession[] _fastAccessSessions = new AMQSession[16]; - private final LinkedHashMap<Integer, AMQSession> _slowAccessSessions = new LinkedHashMap<Integer, AMQSession>(); - private int _size = 0; - private static final int FAST_CHANNEL_ACCESS_MASK = 0xFFFFFFF0; + private final Map<Integer, AMQSession> _sessionMap = new ConcurrentHashMap<>(); private AtomicInteger _idFactory = new AtomicInteger(0); private int _maxChannelID; private int _minChannelID; public AMQSession get(int channelId) { - if ((channelId & FAST_CHANNEL_ACCESS_MASK) == 0) - { - return _fastAccessSessions[channelId]; - } - else - { - return _slowAccessSessions.get(channelId); - } + return _sessionMap.get(channelId); } - public AMQSession put(int channelId, AMQSession session) + public void put(int channelId, AMQSession session) { - AMQSession oldVal; - if ((channelId & FAST_CHANNEL_ACCESS_MASK) == 0) - { - oldVal = _fastAccessSessions[channelId]; - _fastAccessSessions[channelId] = session; - } - else - { - oldVal = _slowAccessSessions.put(channelId, session); - } - if ((oldVal != null) && (session == null)) - { - _size--; - } - else if ((oldVal == null) && (session != null)) - { - _size++; - } - - return session; - + _sessionMap.put(channelId, session); } - public AMQSession remove(int channelId) + public void remove(int channelId) { - AMQSession session; - if ((channelId & FAST_CHANNEL_ACCESS_MASK) == 0) - { - session = _fastAccessSessions[channelId]; - _fastAccessSessions[channelId] = null; - } - else - { - session = _slowAccessSessions.remove(channelId); - } - - if (session != null) - { - _size--; - } - return session; - + _sessionMap.remove(channelId); } public Collection<AMQSession> values() { - ArrayList<AMQSession> values = new ArrayList<AMQSession>(size()); - - for (int i = 0; i < 16; i++) - { - if (_fastAccessSessions[i] != null) - { - values.add(_fastAccessSessions[i]); - } - } - values.addAll(_slowAccessSessions.values()); - - return values; + return new ArrayList<>(_sessionMap.values()); } public int size() { - return _size; + return _sessionMap.size(); } public void clear() { - _size = 0; - _slowAccessSessions.clear(); - for (int i = 0; i < 16; i++) - { - _fastAccessSessions[i] = null; - } + _sessionMap.clear(); } /* @@ -141,14 +80,8 @@ public final class ChannelToSessionMap //go back to the start _idFactory.set(_minChannelID); } - if ((id & FAST_CHANNEL_ACCESS_MASK) == 0) - { - done = (_fastAccessSessions[id] == null); - } - else - { - done = (!_slowAccessSessions.keySet().contains(id)); - } + + done = (!_sessionMap.keySet().contains(id)); } return id; diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/XAConnectionImpl.java b/qpid/java/client/src/main/java/org/apache/qpid/client/XAConnectionImpl.java index d9514338ce..d625a9ae69 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/XAConnectionImpl.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/XAConnectionImpl.java @@ -29,7 +29,7 @@ import javax.jms.XATopicConnection; import javax.jms.XATopicSession; /** - * This class implements the javax.njms.XAConnection interface + * This class implements the javax.jms.XAConnection interface */ public class XAConnectionImpl extends AMQConnection implements XAConnection, XAQueueConnection, XATopicConnection { -- cgit v1.2.1 From a3c5d96d61fdaf76f4cf9dd4e3f543fb87ee95d6 Mon Sep 17 00:00:00 2001 From: Robert Gemmell <robbie@apache.org> Date: Thu, 12 Feb 2015 11:42:03 +0000 Subject: QPID-6240: increment the delivery count when applying the new state and releasing the queue entry git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659229 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java | 1 + 1 file changed, 1 insertion(+) (limited to 'qpid') diff --git a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java index 829b3bf336..3b9521866c 100644 --- a/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java +++ b/qpid/java/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/ConsumerTarget_1_0.java @@ -426,6 +426,7 @@ class ConsumerTarget_1_0 extends AbstractConsumerTarget modified.setDeliveryFailed(true); _link.getEndpoint().updateDisposition(_deliveryTag, modified, true); _link.getEndpoint().sendFlowConditional(); + _queueEntry.incrementDeliveryCount(); _queueEntry.release(); } } -- cgit v1.2.1 From 822ef8b38a63e3b351b30c382bfc77de39904c77 Mon Sep 17 00:00:00 2001 From: Robert Godfrey <rgodfrey@apache.org> Date: Thu, 12 Feb 2015 17:57:11 +0000 Subject: QPID-6374 : avoid taking a lock when not modifying a value git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659341 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/client/AMQSession.java | 27 ++++++++++++---------- 1 file changed, 15 insertions(+), 12 deletions(-) (limited to 'qpid') diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java index 86e1bb0a8b..8f5e9524b6 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java @@ -224,7 +224,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic private volatile boolean _usingDispatcherForCleanup; /** Used to indicates that the connection to which this session belongs, has been stopped. */ - private boolean _connectionStopped; + private final AtomicBoolean _connectionStopped = new AtomicBoolean(); /** Used to indicate that this session has a message listener attached to it. */ private boolean _hasMessageListeners; @@ -3410,25 +3410,28 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic // only call while holding lock final boolean connectionStopped() { - return _connectionStopped; + return _connectionStopped.get(); } boolean setConnectionStopped(boolean connectionStopped) { - boolean currently; - synchronized (_lock) + boolean currently = _connectionStopped.get(); + if(connectionStopped != currently) { - currently = _connectionStopped; - _connectionStopped = connectionStopped; - _lock.notify(); - - if (_dispatcherLogger.isDebugEnabled()) + synchronized (_lock) { - _dispatcherLogger.debug("Set Dispatcher Connection " + (connectionStopped ? "Stopped" : "Started") - + ": Currently " + (currently ? "Stopped" : "Started")); + _connectionStopped.set(connectionStopped); + _lock.notify(); + + if (_dispatcherLogger.isDebugEnabled()) + { + _dispatcherLogger.debug("Set Dispatcher Connection " + (connectionStopped + ? "Stopped" + : "Started") + + ": Currently " + (currently ? "Stopped" : "Started")); + } } } - return currently; } -- cgit v1.2.1 From 176675ef4914dbc548e92700232ec4429614430a Mon Sep 17 00:00:00 2001 From: Robert Godfrey <rgodfrey@apache.org> Date: Thu, 12 Feb 2015 18:24:35 +0000 Subject: QPID-6386 : [AMQP 1.0 Common] close the sender if the TCP connection is terminated before connection.open has occurred git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659348 13f79535-47bb-0310-9956-ffa450edef68 --- .../main/java/org/apache/qpid/amqp_1_0/transport/ConnectionEndpoint.java | 1 + 1 file changed, 1 insertion(+) (limited to 'qpid') diff --git a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/ConnectionEndpoint.java b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/ConnectionEndpoint.java index 766c9705a1..17f334153d 100644 --- a/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/ConnectionEndpoint.java +++ b/qpid/java/amqp-1-0-common/src/main/java/org/apache/qpid/amqp_1_0/transport/ConnectionEndpoint.java @@ -457,6 +457,7 @@ public class ConnectionEndpoint implements DescribedTypeConstructorRegistry.Sour case AWAITING_OPEN: case CLOSE_SENT: _state = ConnectionState.CLOSED; + closeSender(); break; case OPEN: _state = ConnectionState.CLOSE_RECEIVED; -- cgit v1.2.1 From 8d84886a1324a42db1992a4d567487821894d691 Mon Sep 17 00:00:00 2001 From: Robert Godfrey <rgodfrey@apache.org> Date: Thu, 12 Feb 2015 18:50:09 +0000 Subject: QPID-6388 : Treat terminus with durability of "configuration" as durable git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659359 13f79535-47bb-0310-9956-ffa450edef68 --- .../main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) (limited to 'qpid') 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 1820de9d3a..b9ee0ad498 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 @@ -214,7 +214,7 @@ public class Session_1_0 implements SessionEventListener, AMQSessionModel<Sessio registerConsumer(sendingLink.getConsumer()); link = sendingLink; - if(TerminusDurability.UNSETTLED_STATE.equals(source.getDurable())) + if(TerminusDurability.UNSETTLED_STATE.equals(source.getDurable()) || TerminusDurability.CONFIGURATION.equals(source.getDurable())) { linkRegistry.registerSendingLink(endpoint.getName(), sendingLink); } @@ -376,7 +376,8 @@ public class Session_1_0 implements SessionEventListener, AMQSessionModel<Sessio receivingLinkEndpoint.setLinkEventListener(new SubjectSpecificReceivingLinkListener(receivingLink)); link = receivingLink; - if(TerminusDurability.UNSETTLED_STATE.equals(target.getDurable())) + if(TerminusDurability.UNSETTLED_STATE.equals(target.getDurable()) + || TerminusDurability.CONFIGURATION.equals(target.getDurable())) { linkRegistry.registerReceivingLink(endpoint.getName(), receivingLink); } -- cgit v1.2.1 From fbb2c460dfa60e63712f616a3e45c75c9735d5c7 Mon Sep 17 00:00:00 2001 From: Keith Wall <kwall@apache.org> Date: Fri, 13 Feb 2015 17:01:59 +0000 Subject: QPID-6374: [Java Broker] 0-10 Failover: the thread performing the failover prep now syncs the dispatch queue (avoids possibility of app level dead lock) git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1659605 13f79535-47bb-0310-9956-ffa450edef68 --- .../java/org/apache/qpid/client/AMQConnection.java | 47 +++++++++++++++++-- .../apache/qpid/client/AMQConnectionDelegate.java | 2 - .../qpid/client/AMQConnectionDelegate_0_10.java | 24 ++++++++-- .../qpid/client/AMQConnectionDelegate_8_0.java | 5 -- .../java/org/apache/qpid/client/AMQSession.java | 44 +++++++++++------- .../org/apache/qpid/client/AMQSession_0_10.java | 2 +- .../qpid/client/BasicMessageConsumer_0_10.java | 2 +- .../apache/qpid/client/ChannelToSessionMap.java | 7 ++- .../client/util/FlowControllingBlockingQueue.java | 53 ++++++++++++++++++---- 9 files changed, 140 insertions(+), 46 deletions(-) (limited to 'qpid') diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnection.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnection.java index 4c596b88a0..8e7b5b90d8 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnection.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnection.java @@ -1216,11 +1216,6 @@ public class AMQConnection extends Closeable implements Connection, QueueConnect return _failoverMutex; } - public void failoverPrep() - { - _delegate.failoverPrep(); - } - public void resubscribeSessions() throws JMSException, AMQException, FailoverException { _delegate.resubscribeSessions(); @@ -1653,4 +1648,46 @@ public class AMQConnection extends Closeable implements Connection, QueueConnect { return _messageCompressionThresholdSize; } + + void doWithAllLocks(Runnable r) + { + doWithAllLocks(r, _sessions.values()); + + } + + private void doWithAllLocks(final Runnable r, final List<AMQSession> sessions) + { + if (!sessions.isEmpty()) + { + AMQSession session = sessions.remove(0); + + final Object dispatcherLock = session.getDispatcherLock(); + if (dispatcherLock != null) + { + synchronized (dispatcherLock) + { + synchronized (session.getMessageDeliveryLock()) + { + doWithAllLocks(r, sessions); + } + } + } + else + { + synchronized (session.getMessageDeliveryLock()) + { + doWithAllLocks(r, sessions); + } + } + } + else + { + synchronized (getFailoverMutex()) + { + r.run(); + } + } + } + + } diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate.java index 74ca1ed74f..c359fbcc84 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate.java @@ -52,8 +52,6 @@ public interface AMQConnectionDelegate XASession createXASession(int ackMode) throws JMSException; - void failoverPrep(); - void resubscribeSessions() throws JMSException, AMQException, FailoverException; void closeConnection(long timeout) throws JMSException, AMQException; diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_0_10.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_0_10.java index fdeab7ae70..e22a341205 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_0_10.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_0_10.java @@ -27,6 +27,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicBoolean; import javax.jms.ExceptionListener; import javax.jms.JMSException; @@ -249,7 +250,7 @@ public class AMQConnectionDelegate_0_10 implements AMQConnectionDelegate, Connec List<AMQSession> sessions = new ArrayList<AMQSession>(_conn.getSessions().values()); for (AMQSession s : sessions) { - s.failoverPrep(); + ((AMQSession_0_10)s).failoverPrep(); } } @@ -306,16 +307,21 @@ public class AMQConnectionDelegate_0_10 implements AMQConnectionDelegate, Connec _qpidConnection.notifyFailoverRequired(); - synchronized (_conn.getFailoverMutex()) + final AtomicBoolean failoverDone = new AtomicBoolean(); + + _conn.doWithAllLocks(new Runnable() { + @Override + public void run() + { try { if (_conn.firePreFailover(false) && _conn.attemptReconnection()) { - _conn.failoverPrep(); + failoverPrep(); _conn.resubscribeSessions(); _conn.fireFailoverComplete(); - return; + failoverDone.set(true); } } catch (Exception e) @@ -327,9 +333,19 @@ public class AMQConnectionDelegate_0_10 implements AMQConnectionDelegate, Connec _conn.getProtocolHandler().getFailoverLatch().countDown(); _conn.getProtocolHandler().setFailoverLatch(null); } + + } + }); + + + if (failoverDone.get()) + { + return; } + } + _conn.setClosed(); final ExceptionListener listener = _conn.getExceptionListenerNoCheck(); diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_8_0.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_8_0.java index 35582d92b7..ae83b6ab48 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_8_0.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQConnectionDelegate_8_0.java @@ -350,11 +350,6 @@ public class AMQConnectionDelegate_8_0 implements AMQConnectionDelegate } } - public void failoverPrep() - { - // do nothing - } - /** * For all sessions, and for all consumers in those sessions, resubscribe. This is called during failover handling. * The caller must hold the failover mutex before calling this method. diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java index 8f5e9524b6..3966e75423 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession.java @@ -169,7 +169,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic private final Lock _subscriberDetails = new ReentrantLock(true); private final Lock _subscriberAccess = new ReentrantLock(true); - private final FlowControllingBlockingQueue _queue; + private final FlowControllingBlockingQueue<Dispatchable> _queue; private final AtomicLong _highestDeliveryTag = new AtomicLong(-1); private final AtomicLong _rollbackMark = new AtomicLong(-1); @@ -358,7 +358,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic if (_acknowledgeMode == NO_ACKNOWLEDGE) { _queue = - new FlowControllingBlockingQueue(_prefetchHighMark, _prefetchLowMark, + new FlowControllingBlockingQueue<Dispatchable>(_prefetchHighMark, _prefetchLowMark, new FlowControllingBlockingQueue.ThresholdListener() { private final AtomicBoolean _suspendState = new AtomicBoolean(); @@ -423,7 +423,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic } else { - _queue = new FlowControllingBlockingQueue(_prefetchHighMark, null); + _queue = new FlowControllingBlockingQueue<Dispatchable>(_prefetchHighMark, null); } // Add creation logging to tie in with the existing close logging @@ -1789,7 +1789,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic //in the pre-dispatch queue. _usingDispatcherForCleanup = true; - syncDispatchQueue(); + syncDispatchQueue(false); // Set to false before sending the recover as 0-8/9/9-1 will //send messages back before the recover completes, and we @@ -1881,7 +1881,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic setRollbackMark(); - syncDispatchQueue(); + syncDispatchQueue(false); _dispatcher.rollback(); @@ -2201,21 +2201,17 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic } - void failoverPrep() - { - syncDispatchQueue(); - } - void syncDispatchQueue() + void syncDispatchQueue(final boolean holdDispatchLock) { - if (Thread.currentThread() == _dispatcherThread) + if (Thread.currentThread() == _dispatcherThread || holdDispatchLock) { while (!super.isClosed() && !_queue.isEmpty()) { Dispatchable disp; try { - disp = (Dispatchable) _queue.take(); + disp = _queue.take(); } catch (InterruptedException e) { @@ -2267,7 +2263,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic Dispatchable disp; try { - disp = (Dispatchable) _queue.take(); + disp = _queue.take(); } catch (InterruptedException e) { @@ -3086,7 +3082,7 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic private void rejectMessagesForConsumerTag(int consumerTag, boolean requeue, boolean rejectAllConsumers) { - Iterator messages = _queue.iterator(); + Iterator<Dispatchable> messages = _queue.iterator(); if (_logger.isDebugEnabled()) { _logger.debug("Rejecting messages from _queue for Consumer tag(" + consumerTag + ") (PDispatchQ) requeue:" @@ -3237,6 +3233,12 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic public abstract void setFlowControl(final boolean active); + Object getDispatcherLock() + { + Dispatcher dispatcher = _dispatcher; + return dispatcher == null ? null : dispatcher._lock; + } + public interface Dispatchable { void dispatch(AMQSession ssn); @@ -3389,10 +3391,18 @@ public abstract class AMQSession<C extends BasicMessageConsumer, P extends Basic try { - Dispatchable disp; - while (((disp = (Dispatchable) _queue.take()) != null) && !_closed.get()) + + while (((_queue.blockingPeek()) != null) && !_closed.get()) { - disp.dispatch(AMQSession.this); + synchronized (_lock) + { + Dispatchable disp = _queue.nonBlockingTake(); + + if(disp != null) + { + disp.dispatch(AMQSession.this); + } + } } } catch (InterruptedException e) diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java index 206ca15c82..08d7ea3f67 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/AMQSession_0_10.java @@ -235,7 +235,7 @@ public class AMQSession_0_10 extends AMQSession<BasicMessageConsumer_0_10, Basic void failoverPrep() { - super.failoverPrep(); + syncDispatchQueue(true); clearUnacked(); } diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/BasicMessageConsumer_0_10.java b/qpid/java/client/src/main/java/org/apache/qpid/client/BasicMessageConsumer_0_10.java index e0d8ac3702..0cb103f0cb 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/BasicMessageConsumer_0_10.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/BasicMessageConsumer_0_10.java @@ -412,7 +412,7 @@ public class BasicMessageConsumer_0_10 extends BasicMessageConsumer<UnprocessedM _capacity, Option.UNRELIABLE); } - _0_10session.syncDispatchQueue(); + _0_10session.syncDispatchQueue(false); o = super.getMessageFromQueue(-1); } if (_capacity == 0) diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java b/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java index f46c61daa7..0ba5cfdacb 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/ChannelToSessionMap.java @@ -22,13 +22,16 @@ package org.apache.qpid.client; import java.util.ArrayList; import java.util.Collection; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; public final class ChannelToSessionMap { - private final Map<Integer, AMQSession> _sessionMap = new ConcurrentHashMap<>(); + private final Map<Integer, AMQSession> _sessionMap = Collections.synchronizedMap(new LinkedHashMap<Integer, AMQSession>()); private AtomicInteger _idFactory = new AtomicInteger(0); private int _maxChannelID; private int _minChannelID; @@ -48,7 +51,7 @@ public final class ChannelToSessionMap _sessionMap.remove(channelId); } - public Collection<AMQSession> values() + public List<AMQSession> values() { return new ArrayList<>(_sessionMap.values()); } diff --git a/qpid/java/client/src/main/java/org/apache/qpid/client/util/FlowControllingBlockingQueue.java b/qpid/java/client/src/main/java/org/apache/qpid/client/util/FlowControllingBlockingQueue.java index b194ac88de..df54b7066b 100644 --- a/qpid/java/client/src/main/java/org/apache/qpid/client/util/FlowControllingBlockingQueue.java +++ b/qpid/java/client/src/main/java/org/apache/qpid/client/util/FlowControllingBlockingQueue.java @@ -20,13 +20,13 @@ */ package org.apache.qpid.client.util; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import java.util.Iterator; import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** * A blocking queue that emits events above a user specified threshold allowing the caller to take action (e.g. flow * control) to try to prevent the queue growing (much) further. The underlying queue itself is not bounded therefore the @@ -37,12 +37,12 @@ import java.util.concurrent.ConcurrentLinkedQueue; * <p> * TODO Make this implement java.util.Queue and hide the implementation. Then different queue types can be substituted. */ -public class FlowControllingBlockingQueue +public class FlowControllingBlockingQueue<T> { private static final Logger _logger = LoggerFactory.getLogger(FlowControllingBlockingQueue.class); /** This queue is bounded and is used to store messages before being dispatched to the consumer */ - private final Queue _queue = new ConcurrentLinkedQueue(); + private final Queue<T> _queue = new ConcurrentLinkedQueue<T>(); private final int _flowControlHighThreshold; private final int _flowControlLowThreshold; @@ -82,9 +82,44 @@ public class FlowControllingBlockingQueue } } - public Object take() throws InterruptedException + public T blockingPeek() throws InterruptedException + { + T o = _queue.peek(); + if (o == null) + { + synchronized (this) + { + while ((o = _queue.peek()) == null) + { + wait(); + } + } + } + return o; + } + + public T nonBlockingTake() throws InterruptedException + { + T o = _queue.poll(); + + if (o != null && !disableFlowControl && _listener != null) + { + synchronized (_listener) + { + if (_count-- == _flowControlLowThreshold) + { + _listener.underThreshold(_count); + } + } + + } + + return o; + } + + public T take() throws InterruptedException { - Object o = _queue.poll(); + T o = _queue.poll(); if(o == null) { synchronized(this) @@ -110,7 +145,7 @@ public class FlowControllingBlockingQueue return o; } - public void add(Object o) + public void add(T o) { synchronized(this) { @@ -130,7 +165,7 @@ public class FlowControllingBlockingQueue } } - public Iterator iterator() + public Iterator<T> iterator() { return _queue.iterator(); } -- cgit v1.2.1