ffang commented on code in PR #3546:
URL: https://github.com/apache/cxf/pull/3546#discussion_r4208607122
##########
rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java:
##########
@@ -401,6 +402,39 @@ protected boolean isSequenceSuspended(String key) {
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)
Review Comment:
Thanks @reta for looking into it! I went through the scenario and I don't
think it can happen with the current code, but please correct me if I'm missing
something:
1. **No two resend tasks run at the same time.** `initiate()` is only called
from `initiateInOrder()`, which is only called from `ResendTask.run()`
([L694-L695](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java#L694-L695)),
and that runs on the single `RMManager-Timer` thread ([`RMManager`
L212](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/RMManager.java#L212)).
`java.util.Timer` executes one task at a time, so a second `initiateInOrder()`
can't start while the first is running.
2. **We never cancel a running task.** The only task running is the caller's
own, and it is excluded by `c.getNumber() < candidate.getNumber()`
([L424](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java#L424)).
The tasks of the earlier candidates are still waiting in the timer queue, so
`cancel()` reliably prevents them from running.
3. **Resends running on an executor are left alone.** If the
endpoint/service has an executor, the resend runs there, but `initiate()` sets
`pending = true`
([L497](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java#L497))
before handing it over, and `takeOverDueResend()` skips pending candidates
([L605](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java#L605)).
A later task only picks the candidate up again after `attempted()` has cleared
`pending` and moved `next` to the next retry.
4. **The live list is read under the queue lock** (`synchronized (this)`,
[L420](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java#L420)),
the same lock `purgeCandidates()` takes when acknowledged candidates are
removed
([L180](https://github.com/apache/cxf/blob/1760ba1363/rt/ws/rm/src/main/java/org/apache/cxf/ws/rm/soap/RetransmissionQueueImpl.java#L180)),
and the selected candidates are copied out before anything is resent outside
the lock.
Best Regards
Freeman
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]