diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/UniIngressService.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/UniIngressService.java index 4f34e02574..f6e617abba 100644 --- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/UniIngressService.java +++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/UniIngressService.java @@ -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 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); diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/package-info.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/package-info.java index a903e6ba8e..1d267445b9 100644 --- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/package-info.java +++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/ingress/package-info.java @@ -22,16 +22,6 @@ * *

Policy: Public facade -- boot wires ingress into UniRuntime, but no engine sub-package may bypass it. * - *

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). - * *

Marked {@link org.apache.eventmesh.common.Internal @Internal} as a * whole package; public types must carry {@link org.apache.eventmesh.common.Public @Public}. */