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

rgodfrey 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 d919b5619b QPID-8751 : [Broker-J] AMQP 1.0 consumer is never resumed 
when the session incoming window is reopened by a flow naming another link
d919b5619b is described below

commit d919b5619b7ccd3c9a7a4b2cbdf2961f210521ed
Author: Rob Godfrey <[email protected]>
AuthorDate: Sun Aug 16 14:08:27 2026 +0200

    QPID-8751 : [Broker-J] AMQP 1.0 consumer is never resumed when the session 
incoming window is reopened by a flow naming another link
    
    Patch thanks to Marco Geri
---
 .../qpid/server/protocol/v1_0/Session_1_0.java     | 33 ++++++++--
 .../qpid/server/protocol/v1_0/Session_1_0Test.java | 70 ++++++++++++++++++++++
 2 files changed, 97 insertions(+), 6 deletions(-)

diff --git 
a/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
 
b/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
index 35019c7f7f..fd432b5cca 100644
--- 
a/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
+++ 
b/broker-plugins/amqp-1-0-protocol/src/main/java/org/apache/qpid/server/protocol/v1_0/Session_1_0.java
@@ -529,6 +529,8 @@ public class Session_1_0 extends 
AbstractAMQPSession<Session_1_0, ConsumerTarget
             return;
         }
 
+        final boolean hadCreditToSend = hasCreditToSend();
+
         _remoteIncomingWindow = flowNextIncomingId.longValue() + 
flow.getIncomingWindow().longValue()
                                 - _nextOutgoingId.longValue();
 
@@ -548,16 +550,22 @@ public class Session_1_0 extends 
AbstractAMQPSession<Session_1_0, ConsumerTarget
             else
             {
                 endpoint.receiveFlow(flow);
+
+                // Regaining session credit is a session scoped event: it 
makes every sending link on the session
+                // sendable again, not only the one named by the flow's 
optional handle. A link suspended whilst the
+                // session had no credit has no other route back, since it is 
offered no message, its own link credit
+                // therefore never changes, and a well behaved peer has no 
reason to send a flow naming it. Losing
+                // session credit needs no such sweep: a still interested 
consumer is simply refused by
+                // ConsumerTarget_1_0#allocateCredit on the next delivery 
attempt, which suspends it accurately.
+                if (!hadCreditToSend && hasCreditToSend())
+                {
+                    notifyFlowStateChanged(endpoint);
+                }
             }
         }
         else
         {
-            final Collection<LinkEndpoint<? extends BaseSource, ? extends 
BaseTarget>> allLinkEndpoints =
-                    _inputHandleToEndpoint.values();
-            for (final LinkEndpoint<? extends BaseSource, ? extends 
BaseTarget> le : allLinkEndpoints)
-            {
-                le.flowStateChanged();
-            }
+            notifyFlowStateChanged(null);
 
             if (Boolean.TRUE.equals(flow.getEcho()))
             {
@@ -566,6 +574,19 @@ public class Session_1_0 extends 
AbstractAMQPSession<Session_1_0, ConsumerTarget
         }
     }
 
+    private void notifyFlowStateChanged(final LinkEndpoint<? extends 
BaseSource, ? extends BaseTarget> alreadyNotified)
+    {
+        final Collection<LinkEndpoint<? extends BaseSource, ? extends 
BaseTarget>> allLinkEndpoints =
+                new ArrayList<>(_inputHandleToEndpoint.values());
+        for (final LinkEndpoint<? extends BaseSource, ? extends BaseTarget> le 
: allLinkEndpoints)
+        {
+            if (le != alreadyNotified && le instanceof SendingLinkEndpoint)
+            {
+                le.flowStateChanged();
+            }
+        }
+    }
+
     public void receiveDisposition(final Disposition disposition)
     {
         if (_sessionState != SessionState.ACTIVE)
diff --git 
a/broker-plugins/amqp-1-0-protocol/src/test/java/org/apache/qpid/server/protocol/v1_0/Session_1_0Test.java
 
b/broker-plugins/amqp-1-0-protocol/src/test/java/org/apache/qpid/server/protocol/v1_0/Session_1_0Test.java
index 820e437451..e054295d10 100644
--- 
a/broker-plugins/amqp-1-0-protocol/src/test/java/org/apache/qpid/server/protocol/v1_0/Session_1_0Test.java
+++ 
b/broker-plugins/amqp-1-0-protocol/src/test/java/org/apache/qpid/server/protocol/v1_0/Session_1_0Test.java
@@ -22,6 +22,7 @@ package org.apache.qpid.server.protocol.v1_0;
 
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertInstanceOf;
 import static org.junit.jupiter.api.Assertions.assertNotEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
@@ -1244,6 +1245,75 @@ class Session_1_0Test extends UnitTestBase
         verify(_connection, times(1)).sendFrame(eq(session.getChannelId()), 
any(Flow.class));
     }
 
+    @Test
+    void receiveFlowReopeningSessionWindowResumesConsumersOnOtherLinks()
+    {
+        final Queue<?> queue = _virtualHost.createChild(Queue.class, 
Map.of(Queue.NAME, QUEUE_NAME));
+
+        final Attach attachA = createQueueAttach(false, "linkA", QUEUE_NAME);
+        final Attach attachB = createQueueAttach(false, "linkB", QUEUE_NAME);
+        _session.receiveAttach(attachA);
+        _session.receiveAttach(attachB);
+        final QueueConsumer<?, ?> consumerA = consumerOnLink(queue, "linkA");
+        final QueueConsumer<?, ?> consumerB = consumerOnLink(queue, "linkB");
+
+        // Both links are granted credit whilst the session incoming window is 
open: both consumers are sendable.
+        _session.receiveFlow(createLinkFlow(attachA.getHandle(), 10, 2048));
+        _session.receiveFlow(createLinkFlow(attachB.getHandle(), 10, 2048));
+
+        assertTrue(consumerA.isNotifyWorkDesired(), "Consumer A should be 
sendable whilst the session window is open");
+        assertTrue(consumerB.isNotifyWorkDesired(), "Consumer B should be 
sendable whilst the session window is open");
+        assertEquals(2, queue.getConsumerCountWithCredit(), "Both consumers 
should hold credit");
+
+        // The session incoming window closes. The flow carrying the closure 
names link A, so consumer A is
+        // re-evaluated and suspended for a session scoped reason whilst its 
own link credit is untouched.
+        _session.receiveFlow(createLinkFlow(attachA.getHandle(), 10, 0));
+
+        assertFalse(_session.hasCreditToSend(), "Session should have no credit 
once the incoming window is closed");
+        assertFalse(consumerA.isNotifyWorkDesired(), "Consumer A should be 
suspended whilst the session window is closed");
+
+        // The window reopens on a flow naming link B - the protonj2 client 
decorates every flow with the handle of
+        // the link that triggered it, so a pure session window replenish 
routinely names some other link.
+        _session.receiveFlow(createLinkFlow(attachB.getHandle(), 10, 2048));
+
+        assertTrue(_session.hasCreditToSend(), "Session should have credit 
once the incoming window is reopened");
+        assertTrue(consumerB.isNotifyWorkDesired(), "Consumer B should be 
sendable once the session window is reopened");
+
+        // Consumer A still has link credit on an attached link, and the 
session it shares with B has credit again,
+        // so the queue must serve it again. Nothing else will ever wake it: a 
suspended link is offered no message,
+        // so its own credit never changes and the peer has no reason to send 
a flow naming its handle.
+        assertTrue(consumerA.isNotifyWorkDesired(),
+                "Consumer A should be resumed once the session window is 
reopened by a flow naming another link");
+        assertEquals(2, queue.getConsumerCountWithCredit(),
+                "Both consumers should hold credit once the session window is 
reopened");
+    }
+
+    /**
+     * Returns the queue consumer serving the named link. A sending link names 
its consumer after the link itself
+     * unless the attach carries a target address, which {@link 
#createReceiverAttach(String)} does not, and the
+     * queue then qualifies that name with the connection and channel the 
consumer belongs to.
+     */
+    private QueueConsumer<?, ?> consumerOnLink(final Queue<?> queue, final 
String linkName)
+    {
+        return queue.getConsumers().stream()
+                .filter(consumer -> consumer.getName().endsWith("|" + 
linkName))
+                .findFirst()
+                .orElseThrow(() -> new AssertionError(String.format("No 
consumer for link '%s' among %s", linkName,
+                        
queue.getConsumers().stream().map(QueueConsumer::getName).toList())));
+    }
+
+    private Flow createLinkFlow(final UnsignedInteger handle, final int 
linkCredit, final int incomingWindow)
+    {
+        final Flow flow = new Flow();
+        flow.setHandle(handle);
+        flow.setLinkCredit(UnsignedInteger.valueOf(linkCredit));
+        flow.setNextIncomingId(UnsignedInteger.ZERO);
+        flow.setIncomingWindow(UnsignedInteger.valueOf(incomingWindow));
+        flow.setNextOutgoingId(UnsignedInteger.ZERO);
+        flow.setOutgoingWindow(UnsignedInteger.valueOf(2048));
+        return flow;
+    }
+
     private static byte[] createSequentialBytes(final int length)
     {
         final byte[] bytes = new byte[length];


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

Reply via email to