This is an automated email from the ASF dual-hosted git repository.

dakirily pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/qpid-broker-j.git


The following commit(s) were added to refs/heads/main by this push:
     new 4218c35cef QPID-8672: [Broker-J] High CPU Usage because of unnecessary 
flushCreditState() calls for all consumers of session (#436)
4218c35cef is described below

commit 4218c35ceff0ac61b491ea9bd3c30ef44c326f17
Author: Daniil Kirilyuk <[email protected]>
AuthorDate: Sat Sep 19 14:19:43 2026 +0200

    QPID-8672: [Broker-J] High CPU Usage because of unnecessary 
flushCreditState() calls for all consumers of session (#436)
---
 .../server/protocol/v0_10/ConsumerTarget_0_10.java |  38 ++--
 .../qpid/server/protocol/v0_10/ServerSession.java  |  41 ++++-
 .../protocol/v0_10/ConsumerTarget_0_10Test.java    | 193 +++++++++++++++++++++
 .../server/protocol/v0_10/ServerSessionTest.java   | 180 ++++++++++++++++++-
 .../qpid/tests/protocol/v0_10/MessageTest.java     | 126 ++++++++++++++
 5 files changed, 555 insertions(+), 23 deletions(-)

diff --git 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
index 45f9ea942a..dd72512c3c 100644
--- 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
+++ 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10.java
@@ -67,11 +67,12 @@ public class ConsumerTarget_0_10 extends 
AbstractConsumerTarget<ConsumerTarget_0
     private final MessageAcceptMode _acceptMode;
     private final MessageAcquireMode _acquireMode;
     private final ServerSession _session;
+    private final Object _deferredCreditLock = new Object();
 
     private volatile MessageFlowMode _flowMode;
     private volatile FlowCreditManager_0_10 _creditManager;
-    private volatile int _deferredMessageCredit;
-    private volatile long _deferredSizeCredit;
+    private int _deferredMessageCredit;
+    private long _deferredSizeCredit;
 
     private final StateChangeListener<MessageInstance, EntryState> 
_unacknowledgedMessageListener =
             new StateChangeListener<>()
@@ -347,21 +348,34 @@ public class ConsumerTarget_0_10 extends 
AbstractConsumerTarget<ConsumerTarget_0
 
     private void deferredAddCredit(final int deferredMessageCredit, final long 
deferredSizeCredit)
     {
-        _deferredMessageCredit += deferredMessageCredit;
-        _deferredSizeCredit += deferredSizeCredit;
-
+        synchronized (_deferredCreditLock)
+        {
+            _deferredMessageCredit += deferredMessageCredit;
+            _deferredSizeCredit += deferredSizeCredit;
+            _session.addConsumerTargetNeedingFlush(this);
+        }
     }
 
-    public void flushCreditState(boolean strict)
+    public boolean flushCreditState(final boolean strict)
     {
-        if(strict || !isSuspended() || _deferredMessageCredit >= 200
-           || !(_creditManager instanceof WindowCreditManager)
-           || ((WindowCreditManager)_creditManager).getMessageCreditLimit() < 
400 )
+        synchronized (_deferredCreditLock)
         {
-            restoreCredit(_deferredMessageCredit, _deferredSizeCredit);
+            if (_deferredMessageCredit == 0 && _deferredSizeCredit == 0L)
+            {
+                return true;
+            }
+
+            if (strict || !isSuspended() || _deferredMessageCredit >= 200
+                || !(_creditManager instanceof WindowCreditManager)
+                || ((WindowCreditManager) 
_creditManager).getMessageCreditLimit() < 400)
+            {
+                restoreCredit(_deferredMessageCredit, _deferredSizeCredit);
 
-            _deferredMessageCredit = 0;
-            _deferredSizeCredit = 0L;
+                _deferredMessageCredit = 0;
+                _deferredSizeCredit = 0L;
+                return true;
+            }
+            return false;
         }
     }
 
diff --git 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
index 1c4827f741..27d321b967 100644
--- 
a/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
+++ 
b/broker-plugins/amqp-0-10-protocol/src/main/java/org/apache/qpid/server/protocol/v0_10/ServerSession.java
@@ -64,6 +64,7 @@ import javax.security.auth.Subject;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import org.apache.qpid.server.consumer.ConsumerTarget;
 import org.apache.qpid.server.model.Session;
 import org.apache.qpid.server.txn.AsyncCommand;
 import org.apache.qpid.server.logging.LogMessage;
@@ -139,7 +140,9 @@ public class ServerSession extends SessionInvoker
     private final int commandLimit = 
Integer.getInteger("qpid.session.command_limit", 64 * 1024);
     private final Object commandsLock = new Object();
     private final Object stateLock = new Object();
+    private final Object _subscriptionCreditFlushLock = new Object();
     private final Map<String, ConsumerTarget_0_10> _subscriptions = new 
ConcurrentHashMap<>();
+    private final Set<ConsumerTarget_0_10> _subscriptionsNeedingCreditFlush = 
ConcurrentHashMap.newKeySet();
     private final AtomicReference<LogMessage> _forcedCloseLogMessage = new 
AtomicReference<>();
     private final long _blockingTimeout;
     private final ServerConnection connection;
@@ -1226,7 +1229,10 @@ public class ServerSession extends SessionInvoker
     {
         _subscriptions.remove(sub.getName());
         sub.close();
-
+        synchronized (_subscriptionCreditFlushLock)
+        {
+            _subscriptionsNeedingCreditFlush.remove(sub);
+        }
     }
 
     public boolean isTransactional()
@@ -1579,16 +1585,43 @@ public class ServerSession extends SessionInvoker
     {
         runAsSubject(() ->
         {
-            final Collection<ConsumerTarget_0_10> subscriptions = 
getSubscriptions();
-            for (ConsumerTarget_0_10 subscription_0_10 : subscriptions)
+            if (!_subscriptionsNeedingCreditFlush.isEmpty())
             {
-                subscription_0_10.flushCreditState(false);
+                final Set<ConsumerTarget_0_10> subscriptionsNeedingCreditFlush 
=
+                        new HashSet<>(_subscriptionsNeedingCreditFlush);
+                
_subscriptionsNeedingCreditFlush.removeAll(subscriptionsNeedingCreditFlush);
+                for (final ConsumerTarget_0_10 subscription : 
subscriptionsNeedingCreditFlush)
+                {
+                    boolean creditFlushed = false;
+                    try
+                    {
+                        creditFlushed = subscription.flushCreditState(false);
+                    }
+                    finally
+                    {
+                        if (!creditFlushed)
+                        {
+                            addConsumerTargetNeedingFlush(subscription);
+                        }
+                    }
+                }
             }
             awaitCommandCompletion();
             return null;
         });
     }
 
+    void addConsumerTargetNeedingFlush(final ConsumerTarget_0_10 
consumerTarget)
+    {
+        synchronized (_subscriptionCreditFlushLock)
+        {
+            if (consumerTarget.getState() != ConsumerTarget.State.CLOSED)
+            {
+                _subscriptionsNeedingCreditFlush.add(consumerTarget);
+            }
+        }
+    }
+
     public int getUnacknowledgedMessageCount()
     {
         return _messageDispositionListenerMap.size();
diff --git 
a/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10Test.java
 
b/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10Test.java
new file mode 100644
index 0000000000..ba70ca0985
--- /dev/null
+++ 
b/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ConsumerTarget_0_10Test.java
@@ -0,0 +1,193 @@
+/*
+ *
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ *
+ */
+package org.apache.qpid.server.protocol.v0_10;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import org.apache.qpid.server.message.MessageInstance;
+import org.apache.qpid.server.message.MessageInstanceConsumer;
+import org.apache.qpid.server.model.Consumer;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageAcceptMode;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageAcquireMode;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageCreditUnit;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageFlowMode;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageTransfer;
+import org.apache.qpid.test.utils.UnitTestBase;
+
+class ConsumerTarget_0_10Test extends UnitTestBase
+{
+    private static final long MESSAGE_SIZE = 10L;
+
+    private ServerSession _session;
+    private MessageInstanceConsumer _consumer;
+    private MessageInstance _entry;
+    private List<MessageTransfer> _transfers;
+
+    @BeforeEach
+    void setUp()
+    {
+        _session = mock(ServerSession.class);
+        final AMQPConnection_0_10 amqpConnection = 
mock(AMQPConnection_0_10.class);
+        when(_session.getAMQPConnection()).thenReturn(amqpConnection);
+        when(amqpConnection.getContextValue(eq(Long.class), 
eq(Consumer.SUSPEND_NOTIFICATION_PERIOD)))
+                .thenReturn(0L);
+
+        final ServerConnection connection = mock(ServerConnection.class);
+        when(_session.getConnection()).thenReturn(connection);
+        
when(connection.getConnectionDelegate()).thenReturn(mock(ServerConnectionDelegate.class));
+        when(_session.getModelObject()).thenReturn(mock(Session_0_10.class));
+
+        final MessageTransferMessage message = 
mock(MessageTransferMessage.class);
+        when(message.getSize()).thenReturn(MESSAGE_SIZE);
+        _entry = mock(MessageInstance.class);
+        when(_entry.getMessage()).thenReturn(message);
+        _consumer = mock(MessageInstanceConsumer.class);
+
+        _transfers = new ArrayList<>();
+        doAnswer(invocation ->
+        {
+            _transfers.add(invocation.getArgument(0));
+            return null;
+        }).when(_session).sendMessage(any(MessageTransfer.class), 
any(Runnable.class));
+    }
+
+    @Test
+    void completedTransferMarksTargetAndRestoresWindowOneCredit()
+    {
+        final WindowCreditManager creditManager = spy(new 
WindowCreditManager(-1L, 1L));
+        final ConsumerTarget_0_10 target = createTarget(creditManager);
+        target.updateNotifyWorkDesired();
+        exhaustCredit(target, 1);
+
+        completeTransfer(target);
+
+        verify(_session).addConsumerTargetNeedingFlush(target);
+        assertTrue(target.flushCreditState(false));
+        verify(creditManager).restoreCredit(1L, MESSAGE_SIZE);
+        assertFalse(target.isSuspended());
+    }
+
+    @Test
+    void suspendedLargeWindowRetainsCreditUntilTargetBecomesActive()
+    {
+        final WindowCreditManager creditManager = spy(new 
WindowCreditManager(-1L, 400L));
+        final ConsumerTarget_0_10 target = createTarget(creditManager);
+        target.updateNotifyWorkDesired();
+        exhaustCredit(target, 400);
+        assertTrue(target.isSuspended());
+        completeTransfer(target);
+
+        assertFalse(target.flushCreditState(false));
+        verify(creditManager, never()).restoreCredit(anyLong(), anyLong());
+
+        target.addCredit(MessageCreditUnit.MESSAGE, 1);
+
+        assertFalse(target.isSuspended());
+        assertTrue(target.flushCreditState(false));
+        verify(creditManager).restoreCredit(1L, MESSAGE_SIZE);
+    }
+
+    @Test
+    void deferredCreditThresholdRemainsUnchanged()
+    {
+        final WindowCreditManager creditManager = spy(new 
WindowCreditManager(-1L, 400L));
+        final ConsumerTarget_0_10 target = createTarget(creditManager);
+        target.updateNotifyWorkDesired();
+        exhaustCredit(target, 400);
+
+        for (int i = 0; i < 199; i++)
+        {
+            completeTransfer(target);
+        }
+        assertFalse(target.flushCreditState(false));
+        verify(creditManager, never()).restoreCredit(anyLong(), anyLong());
+
+        completeTransfer(target);
+
+        assertTrue(target.flushCreditState(false));
+        verify(creditManager).restoreCredit(200L, 200L * MESSAGE_SIZE);
+    }
+
+    @Test
+    void strictFlushRestoresSuspendedCredit()
+    {
+        final WindowCreditManager creditManager = spy(new 
WindowCreditManager(-1L, 400L));
+        final ConsumerTarget_0_10 target = createTarget(creditManager);
+        target.updateNotifyWorkDesired();
+        exhaustCredit(target, 400);
+        completeTransfer(target);
+
+        assertTrue(target.flushCreditState(true));
+        verify(creditManager).restoreCredit(1L, MESSAGE_SIZE);
+    }
+
+    @Test
+    void noPendingCreditDoesNotRestoreZeroCredit()
+    {
+        final WindowCreditManager creditManager = spy(new 
WindowCreditManager(-1L, 1L));
+        final ConsumerTarget_0_10 target = createTarget(creditManager);
+
+        assertTrue(target.flushCreditState(false));
+        verify(creditManager, never()).restoreCredit(anyLong(), anyLong());
+    }
+
+    private ConsumerTarget_0_10 createTarget(final WindowCreditManager 
creditManager)
+    {
+        return new ConsumerTarget_0_10(_session, "destination", 
MessageAcceptMode.EXPLICIT,
+                MessageAcquireMode.PRE_ACQUIRED, MessageFlowMode.WINDOW, 
creditManager, Map.of(), false);
+    }
+
+    private void exhaustCredit(final ConsumerTarget_0_10 target, final int 
credit)
+    {
+        for (int i = 0; i < credit; i++)
+        {
+            assertTrue(target.allocateCredit(_entry.getMessage()));
+        }
+        target.updateNotifyWorkDesired();
+    }
+
+    private void completeTransfer(final ConsumerTarget_0_10 target)
+    {
+        final int transferCount = _transfers.size();
+        target.doSend(_consumer, _entry, false);
+
+        final MessageTransfer transfer = _transfers.get(transferCount);
+        assertTrue(transfer.hasCompletionListener());
+        transfer.complete();
+    }
+}
diff --git 
a/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ServerSessionTest.java
 
b/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ServerSessionTest.java
index b220b41160..aefb00cd44 100644
--- 
a/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ServerSessionTest.java
+++ 
b/broker-plugins/amqp-0-10-protocol/src/test/java/org/apache/qpid/server/protocol/v0_10/ServerSessionTest.java
@@ -23,8 +23,12 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.anyBoolean;
 import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.clearInvocations;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 import java.util.ArrayList;
@@ -32,6 +36,7 @@ import java.util.Date;
 import java.util.List;
 import java.util.Set;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicReference;
 
 import javax.security.auth.Subject;
@@ -43,7 +48,7 @@ import org.junit.jupiter.api.Test;
 import org.apache.qpid.server.bytebuffer.QpidByteBuffer;
 import org.apache.qpid.server.configuration.updater.CurrentThreadTaskExecutor;
 import org.apache.qpid.server.configuration.updater.TaskExecutor;
-import org.apache.qpid.server.configuration.updater.TaskExecutorImpl;
+import org.apache.qpid.server.consumer.ConsumerTarget;
 import org.apache.qpid.server.logging.EventLogger;
 import org.apache.qpid.server.model.AuthenticationProvider;
 import org.apache.qpid.server.model.Broker;
@@ -63,6 +68,7 @@ import org.apache.qpid.server.protocol.v0_10.transport.Method;
 import org.apache.qpid.server.security.SubjectExecutionContext;
 import org.apache.qpid.server.security.auth.AuthenticatedPrincipal;
 import org.apache.qpid.server.security.auth.UsernamePrincipal;
+import org.apache.qpid.server.txn.ServerTransaction;
 import org.apache.qpid.test.utils.UnitTestBase;
 
 @SuppressWarnings({"rawtypes"})
@@ -206,21 +212,181 @@ class ServerSessionTest extends UnitTestBase
         doAnswer(invocation ->
         {
             capturedSubject.set(SubjectExecutionContext.currentSubject());
-            return null;
+            return true;
         }).when(subscription).flushCreditState(anyBoolean());
 
         session.register("dest", subscription);
+        session.addConsumerTargetNeedingFlush(subscription);
         session.receivedComplete();
 
         assertEquals(session.getSubject(), capturedSubject.get(), "Unexpected 
subject in receivedComplete");
     }
 
-    AmqpPort<?> createMockPort()
+    @Test
+    void receivedCompleteFlushesOnlyMarkedSubscriptions()
+    {
+        final ServerSession session = createServerSession();
+        final List<ConsumerTarget_0_10> subscriptions = new ArrayList<>();
+
+        for (int i = 0; i < 100; i++)
+        {
+            final ConsumerTarget_0_10 subscription = 
mock(ConsumerTarget_0_10.class);
+            when(subscription.flushCreditState(false)).thenReturn(true);
+            session.register("destination-" + i, subscription);
+            subscriptions.add(subscription);
+        }
+
+        final ConsumerTarget_0_10 markedSubscription = subscriptions.get(37);
+        session.addConsumerTargetNeedingFlush(markedSubscription);
+        session.receivedComplete();
+
+        verify(markedSubscription).flushCreditState(false);
+        for (final ConsumerTarget_0_10 subscription : subscriptions)
+        {
+            if (subscription != markedSubscription)
+            {
+                verify(subscription, never()).flushCreditState(anyBoolean());
+            }
+        }
+    }
+
+    @Test
+    void duplicateMarksAreDeduplicatedAndSuccessfulFlushRemovesMark()
+    {
+        final ServerSession session = createServerSession();
+        final ConsumerTarget_0_10 subscription = 
mock(ConsumerTarget_0_10.class);
+        when(subscription.flushCreditState(false)).thenReturn(true);
+        session.register("destination", subscription);
+
+        session.addConsumerTargetNeedingFlush(subscription);
+        session.addConsumerTargetNeedingFlush(subscription);
+        session.receivedComplete();
+        session.receivedComplete();
+
+        verify(subscription).flushCreditState(false);
+    }
+
+    @Test
+    void receivedCompleteWithoutPendingCreditFlushCompletesAsyncCommands()
+    {
+        final ServerSession session = createServerSession();
+        final ServerTransaction.Action action = 
mock(ServerTransaction.Action.class);
+        session.recordFuture(CompletableFuture.completedFuture(null), action);
+
+        session.receivedComplete();
+
+        verify(action).postCommit();
+    }
+
+    @Test
+    void retainedCreditIsRetriedUntilFlushed()
+    {
+        final ServerSession session = createServerSession();
+        final ConsumerTarget_0_10 subscription = 
mock(ConsumerTarget_0_10.class);
+        when(subscription.flushCreditState(false)).thenReturn(false, true);
+        session.register("destination", subscription);
+
+        session.addConsumerTargetNeedingFlush(subscription);
+        session.receivedComplete();
+        session.receivedComplete();
+        session.receivedComplete();
+
+        verify(subscription, times(2)).flushCreditState(false);
+    }
+
+    @Test
+    void markAddedDuringFlushIsProcessedOnNextReceivedComplete()
+    {
+        final ServerSession session = createServerSession();
+        final ConsumerTarget_0_10 subscription = 
mock(ConsumerTarget_0_10.class);
+        final AtomicBoolean firstFlush = new AtomicBoolean(true);
+        doAnswer(invocation ->
+        {
+            if (firstFlush.getAndSet(false))
+            {
+                session.addConsumerTargetNeedingFlush(subscription);
+            }
+            return true;
+        }).when(subscription).flushCreditState(false);
+        session.register("destination", subscription);
+
+        session.addConsumerTargetNeedingFlush(subscription);
+        session.receivedComplete();
+        session.receivedComplete();
+        session.receivedComplete();
+
+        verify(subscription, times(2)).flushCreditState(false);
+    }
+
+    @Test
+    void unregisterRemovesPendingCreditFlush()
+    {
+        final ServerSession session = createServerSession();
+        final ConsumerTarget_0_10 subscription = 
mock(ConsumerTarget_0_10.class);
+        final AtomicReference<ConsumerTarget.State> state = new 
AtomicReference<>(ConsumerTarget.State.OPEN);
+        when(subscription.getName()).thenReturn("destination");
+        when(subscription.getState()).thenAnswer(invocation -> state.get());
+        doAnswer(invocation ->
+        {
+            state.set(ConsumerTarget.State.CLOSED);
+            return true;
+        }).when(subscription).close();
+        session.register("destination", subscription);
+        session.addConsumerTargetNeedingFlush(subscription);
+
+        session.unregister(subscription);
+        clearInvocations(subscription);
+        session.addConsumerTargetNeedingFlush(subscription);
+        session.receivedComplete();
+
+        assertEquals(ConsumerTarget.State.CLOSED, state.get(), "Unexpected 
consumer target state");
+        verify(subscription, never()).flushCreditState(anyBoolean());
+    }
+
+    private ServerSession createServerSession()
+    {
+        final Broker<?> broker = mock(Broker.class);
+        when(broker.getContextValue(eq(Long.class), 
eq(Broker.CHANNEL_FLOW_CONTROL_ENFORCEMENT_TIMEOUT)))
+                .thenReturn(0L);
+
+        final AmqpPort<?> port = createMockPort();
+        final AMQPConnection_0_10 modelConnection = 
mock(AMQPConnection_0_10.class);
+        when(modelConnection.getCategoryClass()).thenReturn(Connection.class);
+        
when(modelConnection.getTypeClass()).thenReturn(AMQPConnection_0_10.class);
+        
when(modelConnection.closeAsync()).thenReturn(CompletableFuture.completedFuture(null));
+        when(modelConnection.getAddressSpace()).thenReturn(_virtualHost);
+        when(modelConnection.getContextProvider()).thenReturn(_virtualHost);
+        when(modelConnection.getBroker()).thenReturn(broker);
+        
when(modelConnection.getEventLogger()).thenReturn(mock(EventLogger.class));
+        when(modelConnection.getContextValue(Long.class, 
Session.PRODUCER_AUTH_CACHE_TIMEOUT))
+                .thenReturn(Session.PRODUCER_AUTH_CACHE_TIMEOUT_DEFAULT);
+        when(modelConnection.getContextValue(Integer.class, 
Session.PRODUCER_AUTH_CACHE_SIZE))
+                .thenReturn(Session.PRODUCER_AUTH_CACHE_SIZE_DEFAULT);
+        when(modelConnection.getContextValue(Long.class, 
Connection.MAX_UNCOMMITTED_IN_MEMORY_SIZE))
+                .thenReturn(Connection.DEFAULT_MAX_UNCOMMITTED_IN_MEMORY_SIZE);
+        when(modelConnection.getChildExecutor()).thenReturn(_taskExecutor);
+        when(modelConnection.getModel()).thenReturn(BrokerModel.getInstance());
+        when(modelConnection.getPort()).thenReturn(port);
+
+        final AuthenticatedPrincipal principal =
+                new AuthenticatedPrincipal(new 
UsernamePrincipal(getTestName(), mock(AuthenticationProvider.class)));
+        final Subject subject = new Subject(false, Set.of(principal), 
Set.of(), Set.of());
+        when(modelConnection.getSubject()).thenReturn(subject);
+
+        final ServerConnection connection = new ServerConnection(1, broker, 
port, Transport.TCP, modelConnection);
+        connection.setVirtualHost(_virtualHost);
+
+        final ServerSession session =
+                new ServerSession(connection, new ServerSessionDelegate(), new 
Binary(getTestName().getBytes()), 0);
+        final Session_0_10 modelSession = new Session_0_10(modelConnection, 1, 
session, getTestName());
+        session.setModelObject(modelSession);
+        return session;
+    }
+
+    private AmqpPort<?> createMockPort()
     {
-        AmqpPort port = mock(AmqpPort.class);
-        TaskExecutor childExecutor = new TaskExecutorImpl();
-        childExecutor.start();
-        when(port.getChildExecutor()).thenReturn(childExecutor);
+        final AmqpPort port = mock(AmqpPort.class);
+        when(port.getChildExecutor()).thenReturn(_taskExecutor);
         when(port.getCategoryClass()).thenReturn(Port.class);
         when(port.getModel()).thenReturn(BrokerModel.getInstance());
         return port;
diff --git 
a/systests/protocol-tests-amqp-0-10/src/test/java/org/apache/qpid/tests/protocol/v0_10/MessageTest.java
 
b/systests/protocol-tests-amqp-0-10/src/test/java/org/apache/qpid/tests/protocol/v0_10/MessageTest.java
index 019a59fb87..d1afbe0d33 100644
--- 
a/systests/protocol-tests-amqp-0-10/src/test/java/org/apache/qpid/tests/protocol/v0_10/MessageTest.java
+++ 
b/systests/protocol-tests-amqp-0-10/src/test/java/org/apache/qpid/tests/protocol/v0_10/MessageTest.java
@@ -35,6 +35,8 @@ import 
org.apache.qpid.server.protocol.v0_10.transport.ExecutionResult;
 import org.apache.qpid.server.protocol.v0_10.transport.MessageAcceptMode;
 import org.apache.qpid.server.protocol.v0_10.transport.MessageAcquireMode;
 import org.apache.qpid.server.protocol.v0_10.transport.MessageCreditUnit;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageFlowMode;
+import org.apache.qpid.server.protocol.v0_10.transport.MessageSetFlowMode;
 import org.apache.qpid.server.protocol.v0_10.transport.MessageTransfer;
 import org.apache.qpid.server.protocol.v0_10.transport.Range;
 import org.apache.qpid.server.protocol.v0_10.transport.RangeSet;
@@ -163,6 +165,130 @@ public class MessageTest extends BrokerAdminUsingTestBase
         }
     }
 
+    @Test
+    @SpecificationTest(section = "10.message.flow",
+            description = "In window mode, completing a transfer restores 
credit for another transfer.")
+    public void completedTransferRestoresWindowCredit() throws Exception
+    {
+        final String firstMessageBody = "firstMessage";
+        final String secondMessageBody = "secondMessage";
+        getBrokerAdmin().putMessageOnQueue(BrokerAdmin.TEST_QUEUE_NAME, 
firstMessageBody, secondMessageBody);
+
+        try (final FrameTransport transport = new 
FrameTransport(getBrokerAdmin()).connect())
+        {
+            final Interaction interaction = transport.newInteraction();
+            final byte[] sessionName = "testSession".getBytes(UTF_8);
+            final String subscriberName = "testSubscriber";
+            interaction.negotiateOpen()
+                       .channelId(1)
+                       .attachSession(sessionName)
+                       .message()
+                       .subscribeAcceptMode(MessageAcceptMode.EXPLICIT)
+                       .subscribeAcquireMode(MessageAcquireMode.PRE_ACQUIRED)
+                       .subscribeDestination(subscriberName)
+                       .subscribeQueue(BrokerAdmin.TEST_QUEUE_NAME)
+                       .subscribeId(0)
+                       .subscribe();
+
+            final MessageSetFlowMode setFlowMode = new 
MessageSetFlowMode(subscriberName, MessageFlowMode.WINDOW);
+            setFlowMode.setId(1);
+            interaction.sendPerformative(setFlowMode)
+                       .message()
+                       .flowId(2)
+                       .flowDestination(subscriberName)
+                       .flowUnit(MessageCreditUnit.MESSAGE)
+                       .flowValue(1)
+                       .flow()
+                       .message()
+                       .flowId(3)
+                       .flowDestination(subscriberName)
+                       .flowUnit(MessageCreditUnit.BYTE)
+                       .flowValue(-1)
+                       .flow();
+
+            final MessageTransfer firstTransfer = 
interaction.consume(MessageTransfer.class, SessionCompleted.class,
+                    SessionCommandPoint.class, SessionConfirmed.class, 
SessionFlush.class);
+            try (final QpidByteBuffer buffer = firstTransfer.getBody())
+            {
+                final byte[] body = new byte[buffer.remaining()];
+                buffer.get(body);
+                assertThat(new String(body, UTF_8), 
is(equalTo(firstMessageBody)));
+            }
+
+            interaction.sendPerformative(new 
SessionCompleted(Range.newInstance(firstTransfer.getId())));
+
+            final MessageTransfer secondTransfer = 
interaction.consume(MessageTransfer.class, SessionCompleted.class,
+                    SessionCommandPoint.class, SessionConfirmed.class, 
SessionFlush.class);
+            try (final QpidByteBuffer buffer = secondTransfer.getBody())
+            {
+                final byte[] body = new byte[buffer.remaining()];
+                buffer.get(body);
+                assertThat(new String(body, UTF_8), 
is(equalTo(secondMessageBody)));
+            }
+        }
+    }
+
+    @Test
+    @SpecificationTest(section = "10.message.flow",
+            description = "In credit mode, granting new credit re-enables 
delivery without a periodic credit refresh.")
+    public void creditModeDeliveryContinuesWithoutCompletionRefresh() throws 
Exception
+    {
+        final String[] messages = new String[20];
+        for (int i = 0; i < messages.length; i++)
+        {
+            messages[i] = "message-" + i;
+        }
+        getBrokerAdmin().putMessageOnQueue(BrokerAdmin.TEST_QUEUE_NAME, 
messages);
+
+        try (final FrameTransport transport = new 
FrameTransport(getBrokerAdmin()).connect())
+        {
+            final Interaction interaction = transport.newInteraction();
+            final byte[] sessionName = "testSession".getBytes(UTF_8);
+            final String subscriberName = "testSubscriber";
+            interaction.negotiateOpen()
+                       .channelId(1)
+                       .attachSession(sessionName)
+                       .message()
+                       .subscribeAcceptMode(MessageAcceptMode.EXPLICIT)
+                       .subscribeAcquireMode(MessageAcquireMode.PRE_ACQUIRED)
+                       .subscribeDestination(subscriberName)
+                       .subscribeQueue(BrokerAdmin.TEST_QUEUE_NAME)
+                       .subscribeId(0)
+                       .subscribe();
+
+            final MessageSetFlowMode setFlowMode = new 
MessageSetFlowMode(subscriberName, MessageFlowMode.CREDIT);
+            setFlowMode.setId(1);
+            interaction.sendPerformative(setFlowMode)
+                       .message()
+                       .flowId(2)
+                       .flowDestination(subscriberName)
+                       .flowUnit(MessageCreditUnit.BYTE)
+                       .flowValue(-1)
+                       .flow();
+
+            for (int i = 0; i < messages.length; i++)
+            {
+                interaction.message()
+                           .flowId(3 + i)
+                           .flowDestination(subscriberName)
+                           .flowUnit(MessageCreditUnit.MESSAGE)
+                           .flowValue(1)
+                           .flow();
+
+                final MessageTransfer transfer = 
interaction.consume(MessageTransfer.class, SessionCompleted.class,
+                        SessionCommandPoint.class, SessionConfirmed.class, 
SessionFlush.class);
+                try (final QpidByteBuffer buffer = transfer.getBody())
+                {
+                    final byte[] body = new byte[buffer.remaining()];
+                    buffer.get(body);
+                    assertThat(new String(body, UTF_8), 
is(equalTo(messages[i])));
+                }
+
+                interaction.sendPerformative(new 
SessionCompleted(Range.newInstance(transfer.getId())));
+            }
+        }
+    }
+
     @Test
     @SpecificationTest(section = "10.message.accept",
             description = "Accepts the message.")


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to