You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@qpid.apache.org by kw...@apache.org on 2016/05/04 22:24:02 UTC
svn commit: r1742339 - in /qpid/java/trunk/broker-plugins:
amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/
amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/
Author: kwall
Date: Wed May 4 22:24:02 2016
New Revision: 1742339
URL: http://svn.apache.org/viewvc?rev=1742339&view=rev
Log:
QPID-7223: [Java Broker] Ensure 0-10's consumer target decrements "unacknowledgedBytes" and "unacknowledgedMessages"
Ensured that when queue entries are stolen from consumers by management, the consumer accounting takes place immediately.
Modified:
qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java
qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
qpid/java/trunk/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java
Modified: qpid/java/trunk/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/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java?rev=1742339&r1=1742338&r2=1742339&view=diff
==============================================================================
--- qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java (original)
+++ qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java Wed May 4 22:24:02 2016
@@ -44,12 +44,12 @@ import org.apache.qpid.server.model.Exch
import org.apache.qpid.server.model.Queue;
import org.apache.qpid.server.plugin.MessageConverter;
import org.apache.qpid.server.protocol.MessageConverterRegistry;
-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/trunk/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/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java?rev=1742339&r1=1742338&r2=1742339&view=diff
==============================================================================
--- qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java (original)
+++ qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ImplicitAcceptDispositionChangeListener.java Wed May 4 22:24:02 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/trunk/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/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java?rev=1742339&r1=1742338&r2=1742339&view=diff
==============================================================================
--- qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java (original)
+++ qpid/java/trunk/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java Wed May 4 22:24:02 2016
@@ -535,7 +535,6 @@ public class ServerSession extends Sessi
public void postCommit()
{
- target.acquisitionRemoved(entry);
entry.delete();
}
Modified: qpid/java/trunk/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/trunk/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java?rev=1742339&r1=1742338&r2=1742339&view=diff
==============================================================================
--- qpid/java/trunk/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java (original)
+++ qpid/java/trunk/broker-plugins/amqp-0-8-protocol/src/main/java/org/apache/qpid/server/protocol/v0_8/ConsumerTarget_0_8.java Wed May 4 22:24:02 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