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();
}
};
}