You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@qpid.apache.org by lq...@apache.org on 2016/09/12 13:40:43 UTC

svn commit: r1760368 - in /qpid/java/branches/6.0.x: ./ broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/

Author: lquack
Date: Mon Sep 12 13:40:43 2016
New Revision: 1760368

URL: http://svn.apache.org/viewvc?rev=1760368&view=rev
Log:
QPID-7223: [Java Broker] Ensure 0-10's consumer target decrements "unacknowledgedBytes" and "unacknowledgedMessages"

Merged from trunk with command:

svn merge -c 1742339 ^/qpid/java/trunk

Modified:
    qpid/java/branches/6.0.x/   (props changed)
    qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
    qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java
    qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
    qpid/java/branches/6.0.x/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java

Propchange: qpid/java/branches/6.0.x/
------------------------------------------------------------------------------
--- svn:mergeinfo (original)
+++ svn:mergeinfo Mon Sep 12 13:40:43 2016
@@ -9,5 +9,5 @@
 /qpid/branches/java-broker-vhost-refactor/java:1493674-1494547
 /qpid/branches/java-network-refactor/qpid/java:805429-821809
 /qpid/branches/qpid-2935/qpid/java:1061302-1072333
-/qpid/java/trunk:1715445-1715447,1715586,1715940,1716086-1716087,1716127-1716128,1716141,1716153,1716155,1716194,1716204,1716209,1716227,1716277,1716357,1716368,1716370,1716374,1716432,1716444-1716445,1716455,1716461,1716474,1716489,1716497,1716515,1716555,1716602,1716606-1716610,1716619,1716636,1717269,1717299,1717401,1717446,1717449,1717626,1717691,1717735,1717780,1718744,1718889,1718893,1718918,1718922,1719026,1719028,1719033,1719037,1719047,1719051,1720340,1720664,1721151,1721198,1722019-1722020,1722246,1722339,1722416,1722674,1722678,1722683,1722711,1723064,1723194,1723563,1724216,1724251,1724257,1724292,1724375,1724397,1724432,1724582,1724603,1724780,1724843-1724844,1725295,1725569,1725760,1726176,1726244-1726246,1726249,1726358,1726436,1726449,1726456,1726646,1726653,1726755,1726778,1727532,1727555,1727608,1727951,1727954,1728089,1728167,1728302,1728497,1728501,1728524,1728639,1728651,1728772,1729215,1729297,1729347,1729356,1729406,1729408,1729412,1729515,1729638,1729656-1729
 657,1729783,1729828,1729832,1729841,1729851,1729886,1729904,1729973,1730019,1730025,1730052,1730072,1730088,1730494,1730499,1730547,1730559,1730567,1730578,1730585,1730651,1730697,1730712-1730713,1730805,1731029,1731110,1731210,1731225,1731444,1731551,1731612,1732184,1732452,1732461,1732465,1732525,1732812,1733467,1734452,1736478,1736751,1736838,1737804,1737835,1737853,1737984,1737992,1738119,1738135,1738231,1738271,1738607,1738610,1738731,1738914,1741702,1742257,1742284,1742544,1742900,1742926,1743161,1743228,1743383,1743982,1744012-1744013,1744046,1744123,1744157,1744276,1744403,1745424,1745450,1746140,1746273,1747526,1748254,1748723,1748818,1749349,1749399,1749482,1749524,1750359-1750360,1750943,1751433,1754251,1754354,1754392,1754429,1754510,1754550,1755561,1755957,1758628,1758640,1758766,1758964,1759774,1759783
+/qpid/java/trunk:1715445-1715447,1715586,1715940,1716086-1716087,1716127-1716128,1716141,1716153,1716155,1716194,1716204,1716209,1716227,1716277,1716357,1716368,1716370,1716374,1716432,1716444-1716445,1716455,1716461,1716474,1716489,1716497,1716515,1716555,1716602,1716606-1716610,1716619,1716636,1717269,1717299,1717401,1717446,1717449,1717626,1717691,1717735,1717780,1718744,1718889,1718893,1718918,1718922,1719026,1719028,1719033,1719037,1719047,1719051,1720340,1720664,1721151,1721198,1722019-1722020,1722246,1722339,1722416,1722674,1722678,1722683,1722711,1723064,1723194,1723563,1724216,1724251,1724257,1724292,1724375,1724397,1724432,1724582,1724603,1724780,1724843-1724844,1725295,1725569,1725760,1726176,1726244-1726246,1726249,1726358,1726436,1726449,1726456,1726646,1726653,1726755,1726778,1727532,1727555,1727608,1727951,1727954,1728089,1728167,1728302,1728497,1728501,1728524,1728639,1728651,1728772,1729215,1729297,1729347,1729356,1729406,1729408,1729412,1729515,1729638,1729656-1729
 657,1729783,1729828,1729832,1729841,1729851,1729886,1729904,1729973,1730019,1730025,1730052,1730072,1730088,1730494,1730499,1730547,1730559,1730567,1730578,1730585,1730651,1730697,1730712-1730713,1730805,1731029,1731110,1731210,1731225,1731444,1731551,1731612,1732184,1732452,1732461,1732465,1732525,1732812,1733467,1734452,1736478,1736751,1736838,1737804,1737835,1737853,1737984,1737992,1738119,1738135,1738231,1738271,1738607,1738610,1738731,1738914,1741702,1742257,1742284,1742339,1742544,1742900,1742926,1743161,1743228,1743383,1743982,1744012-1744013,1744046,1744123,1744157,1744276,1744403,1745424,1745450,1746140,1746273,1747526,1748254,1748723,1748818,1749349,1749399,1749482,1749524,1750359-1750360,1750943,1751433,1754251,1754354,1754392,1754429,1754510,1754550,1755561,1755957,1758628,1758640,1758766,1758964,1759774,1759783
 /qpid/trunk/qpid:796646-796653

Modified: qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
URL: http://svn.apache.org/viewvc/qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java?rev=1760368&r1=1760367&r2=1760368&view=diff
==============================================================================
--- qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java (original)
+++ qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java Mon Sep 12 13:40:43 2016
@@ -44,12 +44,12 @@ import org.apache.qpid.server.model.Exch
 import org.apache.qpid.server.plugin.MessageConverter;
 import org.apache.qpid.server.protocol.MessageConverterRegistry;
 import org.apache.qpid.server.queue.AMQQueue;
-import org.apache.qpid.server.queue.QueueConsumer;
 import org.apache.qpid.server.store.TransactionLogResource;
 import org.apache.qpid.server.txn.AutoCommitTransaction;
 import org.apache.qpid.server.txn.ServerTransaction;
 import org.apache.qpid.server.util.Action;
 import org.apache.qpid.server.util.ConnectionScopedRuntimeException;
+import org.apache.qpid.server.util.StateChangeListener;
 import org.apache.qpid.transport.DeliveryProperties;
 import org.apache.qpid.transport.Header;
 import org.apache.qpid.transport.MessageAcceptMode;
@@ -89,6 +89,18 @@ public class ConsumerTarget_0_10 extends
     private long _deferredSizeCredit;
     private final List<ConsumerImpl> _consumers = new CopyOnWriteArrayList<>();
 
+    private final StateChangeListener<MessageInstance, MessageInstance.State> _unacknowledgedMessageListener = new StateChangeListener<MessageInstance, MessageInstance.State>()
+    {
+
+        public void stateChanged(MessageInstance entry, MessageInstance.State oldState, MessageInstance.State newState)
+        {
+            if(oldState == MessageInstance.State.ACQUIRED && newState != MessageInstance.State.ACQUIRED)
+            {
+                removeUnacknowledgedMessage(entry);
+                entry.removeStateChangeListener(this);
+            }
+        }
+    };
 
     public ConsumerTarget_0_10(ServerSession session,
                                String name,
@@ -351,14 +363,30 @@ public class ConsumerTarget_0_10 extends
         }
         else if(_acquireMode == MessageAcquireMode.PRE_ACQUIRED)
         {
-            recordUnacknowledged(entry);
+            addUnacknowledgedMessage(entry);
         }
     }
 
-    void recordUnacknowledged(MessageInstance entry)
+    void addUnacknowledgedMessage(MessageInstance entry)
     {
         _unacknowledgedCount.incrementAndGet();
         _unacknowledgedBytes.addAndGet(entry.getMessage().getSize());
+        entry.addStateChangeListener(_unacknowledgedMessageListener);
+    }
+
+    private void removeUnacknowledgedMessage(MessageInstance entry)
+    {
+        _unacknowledgedBytes.addAndGet(-entry.getMessage().getSize());
+        _unacknowledgedCount.decrementAndGet();
+    }
+
+    @Override
+    public void acquisitionRemoved(final MessageInstance entry)
+    {
+        if (entry.removeStateChangeListener(_unacknowledgedMessageListener))
+        {
+            removeUnacknowledgedMessage(entry);
+        }
     }
 
     private void deferredAddCredit(final int deferredMessageCredit, final long deferredSizeCredit)
@@ -411,17 +439,6 @@ public class ConsumerTarget_0_10 extends
         }
     }
 
-    private boolean isAcquiredByConsumer(final MessageInstance entry)
-    {
-        ConsumerImpl acquiringConsumer = entry.getAcquiringConsumer();
-        if(acquiringConsumer instanceof QueueConsumer)
-        {
-            return ((QueueConsumer)acquiringConsumer).getTarget() == this;
-        }
-
-        return false;
-    }
-
     void release(final ConsumerImpl consumer,
                  final MessageInstance entry,
                  final boolean setRedelivered)
@@ -518,7 +535,6 @@ public class ConsumerTarget_0_10 extends
         return _creditManager;
     }
 
-
     public void stop()
     {
         try
@@ -590,13 +606,6 @@ public class ConsumerTarget_0_10 extends
         return _stopped.get();
     }
 
-    @Override
-    public void acquisitionRemoved(final MessageInstance entry)
-    {
-        _unacknowledgedBytes.addAndGet(-entry.getMessage().getSize());
-        _unacknowledgedCount.decrementAndGet();
-    }
-
     public void flush()
     {
         flushCreditState(true);

Modified: qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java
URL: http://svn.apache.org/viewvc/qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java?rev=1760368&r1=1760367&r2=1760368&view=diff
==============================================================================
--- qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java (original)
+++ qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java Mon Sep 12 13:40:43 2016
@@ -65,7 +65,7 @@ class ImplicitAcceptDispositionChangeLis
         boolean acquired = _entry.acquire(_consumer);
         if(acquired)
         {
-            _target.recordUnacknowledged(_entry);
+            _target.addUnacknowledgedMessage(_entry);
         }
         return acquired;
 

Modified: qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
URL: http://svn.apache.org/viewvc/qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java?rev=1760368&r1=1760367&r2=1760368&view=diff
==============================================================================
--- qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java (original)
+++ qpid/java/branches/6.0.x/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java Mon Sep 12 13:40:43 2016
@@ -535,7 +535,6 @@ public class ServerSession extends Sessi
 
                                      public void postCommit()
                                      {
-                                         target.acquisitionRemoved(entry);
                                          entry.delete();
                                      }
 

Modified: qpid/java/branches/6.0.x/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java
URL: http://svn.apache.org/viewvc/qpid/java/branches/6.0.x/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java?rev=1760368&r1=1760367&r2=1760368&view=diff
==============================================================================
--- qpid/java/branches/6.0.x/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java (original)
+++ qpid/java/branches/6.0.x/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java Mon Sep 12 13:40:43 2016
@@ -515,9 +515,23 @@ public abstract class ConsumerTarget_0_8
         entry.addStateChangeListener(_unacknowledgedMessageListener);
     }
 
+    private void removeUnacknowledgedMessage(MessageInstance entry)
+    {
+
+        final long _size = entry.getMessage().getSize();
+        _unacknowledgedBytes.addAndGet(-_size);
+        _unacknowledgedCount.decrementAndGet();
+
+        _creditManager.restoreCredit(1, _size);
+    }
+
     @Override
     public void acquisitionRemoved(final MessageInstance node)
     {
+        if (node.removeStateChangeListener(_unacknowledgedMessageListener))
+        {
+            removeUnacknowledgedMessage(node);
+        }
     }
 
     public long getUnacknowledgedBytes()
@@ -537,12 +551,7 @@ public abstract class ConsumerTarget_0_8
         {
             if(oldState == MessageInstance.State.ACQUIRED && newState != MessageInstance.State.ACQUIRED)
             {
-                final long _size = entry.getMessage().getSize();
-                _unacknowledgedBytes.addAndGet(-_size);
-                _unacknowledgedCount.decrementAndGet();
-
-                _creditManager.restoreCredit(1, _size);
-
+                removeUnacknowledgedMessage(entry);
                 entry.removeStateChangeListener(this);
             }
 



---------------------------------------------------------------------
To unsubscribe, e-mail: commits-unsubscribe@qpid.apache.org
For additional commands, e-mail: commits-help@qpid.apache.org