Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -372,29 +372,13 @@ private int pullAndDispatchPartition(String topic, int partition, int maxEvents,
// Frame architecture: the POP check key rides in frame attributes (empopck) and the
// deferred broker ACK fires on client ACK — same at-least-once goal, Frame-native.
long offset = nextOffset(topic);
// P2 fix (PR #5316 by zhang-arvin, fixes #5295): if the frame carries a POP check
// key (RocketMQ 5.x deferred ACK), build a callback that ACKs the broker on
// client ACK (restoring at-least-once).
// P2 fix: if the frame carries a POP check key (RocketMQ 5.x deferred ACK), build a
// callback that ACKs the broker on client ACK (restoring at-least-once).
String popCk = f.attributes().get("empopck");
List<Subscription> targets = subscriptionManager.targetsFor(topic, f);
if (popCk != null && !targets.isEmpty()) {
java.util.concurrent.atomic.AtomicInteger pending =
new java.util.concurrent.atomic.AtomicInteger(targets.size());
Runnable mqAck = () -> {
if (pending.decrementAndGet() == 0) {
storage.ackPulledMessage(topic, popCk);
}
};
for (Subscription target : targets) {
dispatcher.deliver(topic, partition, offset, f, target.getClientId(),
channelFor(target.getClientId()), mqAck);
}
} else {
for (Subscription target : targets) {
dispatcher.deliver(topic, partition, offset, f, target.getClientId(),
channelFor(target.getClientId()), null);
}
Runnable mqAck = (popCk != null) ? () -> storage.ackPulledMessage(topic, popCk) : null;
for (Subscription target : subscriptionManager.targetsFor(topic, f)) {
dispatcher.deliver(topic, partition, offset, f, target.getClientId(),
channelFor(target.getClientId()), mqAck);
}
}
UniTrace.end(dispatchSpan);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,16 +22,6 @@
*
* <p>Policy: Public facade -- boot wires ingress into UniRuntime, but no engine sub-package may bypass it.
*
* <p>Broker-ACK barrier (RocketMQ 5.x POP, see PR #5316 by zhang-arvin,
* fixes #5295): when the ingress frame carries a POP check key (the
* {@code empopck} attribute), deliveries to all matched subscriptions
* share an {@link java.util.concurrent.atomic.AtomicInteger} counter
* initialized to the target count; the broker is ACKed (via
* {@code storage.ackPulledMessage}) only when the last required delivery
* ACKs. This restores at-least-once semantics across LOAD_BALANCE,
* BROADCAST, and MULTICAST distribution modes. Frames without
* {@code empopck} bypass the barrier (no broker ACK to defer).
*
* <p>Marked {@link org.apache.eventmesh.common.Internal @Internal} as a
* whole package; public types must carry {@link org.apache.eventmesh.common.Public @Public}.
*/
Expand Down