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

ffang pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/cxf.git


The following commit(s) were added to refs/heads/main by this push:
     new 6c072f66868 [CXF-9255]WS-RM: resend due messages of a sequence in 
message number order (#3546)
6c072f66868 is described below

commit 6c072f66868f97ac560ca1adcfabe05c679c3d6c
Author: Freeman(Yue) Fang <[email protected]>
AuthorDate: Wed Oct 7 12:09:59 2026 -0400

    [CXF-9255]WS-RM: resend due messages of a sequence in message number order 
(#3546)
    
    java.util.Timer does not run tasks scheduled for the same time in
    scheduling order. When two messages of a sequence are cached within the
    same millisecond, the resend of the later one could run first. With
    in-order delivery the destination holds that message back until the
    earlier one arrives, so the resend blocks the resend thread (the timer
    thread with the SynchronousExecutor) until the receive timeout, and the
    earlier message is not resent in the meantime.
    
    When a resend task fires, first resend any earlier message of the same
    sequence that is due no later than it and not pending or suspended, in
    message number order, cancelling its own scheduled task. Resend times
    are unchanged.
    
    This made the in-order one-way tests in systests/ws-rm
    (DeliveryAssuranceOnewayTest, MessageCallbackOnewayTest) time out on
    fast machines.
---
 .../cxf/ws/rm/soap/RetransmissionQueueImpl.java    | 53 +++++++++++++++++++++-
 .../ws/rm/soap/RetransmissionQueueImplTest.java    | 52 +++++++++++++++++++++
 2 files changed, 104 insertions(+), 1 deletion(-)

diff --git 
a/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java 
b/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java
index 59ab60d3490..2eecdcc819c 100644
--- 
a/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java
+++ 
b/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java
@@ -26,6 +26,7 @@ import java.io.OutputStream;
 import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.Comparator;
 import java.util.Date;
 import java.util.HashMap;
 import java.util.List;
@@ -401,6 +402,39 @@ public class RetransmissionQueueImpl implements 
RetransmissionQueue {
         return suspendedCandidates.containsKey(key);
     }
 
+    /**
+     * Initiate the resend of the given candidate, after initiating the resend 
of any earlier message of the
+     * same sequence which is due no later than it and not pending yet. 
java.util.Timer does not run tasks
+     * scheduled for the same time in scheduling order, so for messages cached 
within the same millisecond
+     * the resend of a later message could otherwise run (and, with in-order 
delivery, block the resend
+     * thread until the receive timeout) before the resend of an earlier one.
+     *
+     * @param candidate the candidate whose resend is due
+     */
+    protected void initiateInOrder(ResendCandidate candidate) {
+        final Date due = candidate.getNext();
+        final List<ResendCandidate> earlier = new ArrayList<>();
+        if (null != due) {
+            String key = 
RMContextUtils.retrieveRMProperties(candidate.getMessage(), true)
+                .getSequence().getIdentifier().getValue();
+            synchronized (this) {
+                List<ResendCandidate> sequenceCandidates = 
getSequenceCandidates(key);
+                if (null != sequenceCandidates) {
+                    for (ResendCandidate c : sequenceCandidates) {
+                        if (c.getNumber() < candidate.getNumber() && 
c.takeOverDueResend(due)) {
+                            earlier.add(c);
+                        }
+                    }
+                }
+            }
+            earlier.sort(Comparator.comparingLong(ResendCandidate::getNumber));
+        }
+        for (ResendCandidate c : earlier) {
+            c.initiate(c.includeAckRequested);
+        }
+        candidate.initiate(candidate.includeAckRequested);
+    }
+
     /**
      * Represents a candidate for resend, i.e. an unacked outgoing message.
      */
@@ -560,6 +594,23 @@ public class RetransmissionQueueImpl implements 
RetransmissionQueue {
             }
         }
 
+        /**
+         * Take over the resend of this candidate if it is due no later than 
the given time and is neither
+         * pending nor suspended: its own scheduled resend is then cancelled, 
as the caller initiates it.
+         *
+         * @param due the time by which the resend must be due
+         * @return true if the caller must initiate the resend of this 
candidate
+         */
+        protected synchronized boolean takeOverDueResend(Date due) {
+            if (pending || suspended || null == next || next.after(due)) {
+                return false;
+            }
+            if (null != nextTask) {
+                nextTask.cancel();
+            }
+            return true;
+        }
+
         /**
          * Cancel further resend (although no ACK has been received).
          */
@@ -641,7 +692,7 @@ public class RetransmissionQueueImpl implements 
RetransmissionQueue {
                 @Override
                 public void run() {
                     if (!candidate.isPending()) {
-                        candidate.initiate(includeAckRequested);
+                        initiateInOrder(candidate);
                     }
                 }
             }
diff --git 
a/rt/ws/rm/src/test/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImplTest.java
 
b/rt/ws/rm/src/test/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImplTest.java
index 7a8f7c2e964..08801124881 100644
--- 
a/rt/ws/rm/src/test/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImplTest.java
+++ 
b/rt/ws/rm/src/test/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImplTest.java
@@ -21,12 +21,15 @@
 package org.apache.cxf.ws.rm.soap;
 
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Date;
 import java.util.List;
 import java.util.concurrent.Executor;
 
 import org.apache.cxf.binding.soap.SoapMessage;
+import org.apache.cxf.message.Exchange;
 import org.apache.cxf.message.Message;
+import org.apache.cxf.service.Service;
 import org.apache.cxf.ws.rm.RMConfiguration;
 import org.apache.cxf.ws.rm.RMEndpoint;
 import org.apache.cxf.ws.rm.RMException;
@@ -220,6 +223,43 @@ public class RetransmissionQueueImplTest {
                    sequence1List.get(1).getMessage());
     }
 
+    @Test
+    public void testInitiateInOrderResendsEarlierDueMessageFirst() {
+        SoapMessage message1 = setUpMessage("sequence1", ONE);
+        SoapMessage message2 = setUpMessage("sequence1", TWO);
+        setupMessagePolicies(message1);
+        setupMessagePolicies(message2);
+        setUpExecutor(message1);
+        setUpExecutor(message2);
+        ready(false);
+
+        queue.cacheUnacknowledged(message1);
+        RetransmissionQueueImpl.ResendCandidate candidate2 = 
queue.cacheUnacknowledged(message2);
+
+        // the resend task of message 2 runs first, as java.util.Timer may do 
for tasks due at the same time
+        queue.initiateInOrder(candidate2);
+        assertEquals(Arrays.asList(message1, message2), resender.resent);
+    }
+
+    @Test
+    public void testInitiateInOrderSkipsEarlierMessageNotDue() {
+        SoapMessage message1 = setUpMessage("sequence1", ONE);
+        SoapMessage message2 = setUpMessage("sequence1", TWO);
+        setupMessagePolicies(message1);
+        setupMessagePolicies(message2);
+        setUpExecutor(message1);
+        setUpExecutor(message2);
+        ready(false);
+
+        RetransmissionQueueImpl.ResendCandidate candidate1 = 
queue.cacheUnacknowledged(message1);
+        RetransmissionQueueImpl.ResendCandidate candidate2 = 
queue.cacheUnacknowledged(message2);
+
+        // message 1 was already resent, so its next resend is due after the 
one of message 2
+        candidate1.attempted();
+        queue.initiateInOrder(candidate2);
+        assertEquals(Arrays.asList(message2), resender.resent);
+    }
+
     @Test
     public void testPurgeAcknowledgedSome() {
         Long[] messageNumbers = {TEN, ONE};
@@ -368,6 +408,15 @@ public class RetransmissionQueueImplTest {
         return message;
     }
 
+    private void setUpExecutor(Message message) {
+        Exchange exchange = createMock(Exchange.class);
+        org.apache.cxf.endpoint.Endpoint ep = 
createMock(org.apache.cxf.endpoint.Endpoint.class);
+        Service service = createMock(Service.class);
+        when(message.getExchange()).thenReturn(exchange);
+        when(exchange.getEndpoint()).thenReturn(ep);
+        when(ep.getService()).thenReturn(service);
+    }
+
     private void setupMessagePolicies(Message message) {
         RMConfiguration cfg = new RMConfiguration();
         when(manager.getEffectiveConfiguration(message)).thenReturn(cfg);
@@ -451,15 +500,18 @@ public class RetransmissionQueueImplTest {
     static class TestResender implements RetransmissionQueueImpl.Resender {
         Message message;
         boolean includeAckRequested;
+        List<Message> resent = new ArrayList<>();
 
         public void resend(Message ctx, boolean requestAcknowledge) {
             message = ctx;
             includeAckRequested = requestAcknowledge;
+            resent.add(ctx);
         }
 
         void clear() {
             message = null;
             includeAckRequested = false;
+            resent.clear();
         }
     };
 }

Reply via email to