From fde2c156277e9e6d61a4065f141a1ec82a565341 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 06:24:35 +0700 Subject: [PATCH] t385: a redelivered lead message must not write the pane twice An AMQP recovery clears the held delivery tags, the broker redelivers with fresh ones, and the coordination loop wrote the same peer message into the lead pane again. Measured on the live daemon: one msgId reached the pane 12 times in 9 hours, across 19 recovery events. LeadCoordLoop now remembers the msgIds it has written to a pane (bounded at 1024) and acks a redelivery without a second write. LeadMailbox.ack no longer returns quietly for an unknown msgId: it throws, so an ack that never reached the broker is reported instead of hidden. A repeat ack that this connection already completed stays quiet, tracked in a bounded set. --- .../java/dev/ltms/fleet/msg/LeadChannel.java | 6 ++- .../dev/ltms/fleet/msg/LeadCoordLoop.java | 40 ++++++++++++++++-- .../java/dev/ltms/fleet/msg/LeadMailbox.java | 42 +++++++++++++++---- .../dev/ltms/fleet/msg/LeadCoordLoopTest.java | 14 +++++++ .../dev/ltms/fleet/msg/LeadMailboxTest.java | 27 ++++++++++++ 5 files changed, 118 insertions(+), 11 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java index db16aa3..d7c8931 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java @@ -32,7 +32,11 @@ public interface LeadChannel { /** Non-destructive FIFO snapshot of the messages held for this daemon's own coord-id. */ List peek(); - /** Drop {@code msgId} from the held set and ack it on the broker. A no-op if it is not held. */ + /** + * Drop {@code msgId} from the held set and ack it on the broker. A repeated ack that this + * connection already completed may be a no-op. Any other unknown msgId must throw rather than + * report an ack that did not reach the broker. + */ void ack(String msgId); /** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */ diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java index ad86a63..e7041aa 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java @@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.AgentStatus; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.concurrent.ScheduledExecutorService; @@ -43,6 +44,9 @@ public final class LeadCoordLoop { private static final Logger log = LoggerFactory.getLogger(LeadCoordLoop.class); + /** A bounded window is enough: redelivery can only follow a recent failed ack or recovery. */ + private static final int RECENT_DELIVERY_LIMIT = 1_024; + /** How an arriving peer message is rendered into the lead's pane — the sender's coord-id, then its text. */ static final String DELIVERY_FORMAT = "[lead %s] %s"; @@ -51,6 +55,8 @@ public final class LeadCoordLoop { private final Supplier> leads; private final ScheduledExecutorService scheduler; private final long intervalMs; + /** msgIds already written to the pane, so recovery redelivery is acked without another pane write. */ + private final LinkedHashMap delivered = new LinkedHashMap<>(); private volatile boolean running; @@ -117,6 +123,13 @@ public final class LeadCoordLoop { if (held.isEmpty()) { return; } + LeadMessage msg = held.getFirst(); + if (wasDelivered(msg.msgId())) { + // This lives here, rather than in LeadMailbox, because only this loop knows a pane write + // happened. The mailbox only knows broker delivery tags and must still redeliver after a crash. + ackDelivered(msg); + return; + } String lead = resolveLocalLead(); if (lead == null) { // Left unacked on purpose: the broker keeps holding it until a lead pane exists. @@ -136,7 +149,6 @@ public final class LeadCoordLoop { lead, status, held.size()); return; } - LeadMessage msg = held.getFirst(); try { agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content())); } catch (RuntimeException e) { @@ -145,6 +157,13 @@ public final class LeadCoordLoop { msg.msgId(), msg.from(), lead, e.toString()); return; } + rememberDelivered(msg.msgId()); + if (ackDelivered(msg)) { + log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead); + } + } + + private boolean ackDelivered(LeadMessage msg) { try { channel.ack(msg.msgId()); } catch (RuntimeException e) { @@ -152,9 +171,24 @@ public final class LeadCoordLoop { // deliberate direction of this trade. log.warn("lead coordination: delivered message {} but could not ack it: {}", msg.msgId(), e.toString()); - return; + return false; + } + return true; + } + + private boolean wasDelivered(String msgId) { + synchronized (delivered) { + return delivered.containsKey(msgId); + } + } + + private void rememberDelivered(String msgId) { + synchronized (delivered) { + delivered.put(msgId, Boolean.TRUE); + if (delivered.size() > RECENT_DELIVERY_LIMIT) { + delivered.remove(delivered.keySet().iterator().next()); + } } - log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead); } /** diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java index a2f3b64..a9283f1 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -59,7 +59,8 @@ import java.util.concurrent.TimeoutException; *

Recovery. The connection is opened with automatic + topology recovery * enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags, * so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on - * redelivery) and any publish still awaiting its confirm is failed rather than left to idle out + * redelivery). {@link LeadCoordLoop} separately deduplicates pane writes, since it alone knows + * which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out * the confirm timeout against a sequence number that means nothing on the new channel. */ public final class LeadMailbox implements LeadChannel, AutoCloseable { @@ -85,6 +86,10 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { private final Object channelLock = new Object(); /** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */ private final LinkedHashMap held = new LinkedHashMap<>(); + /** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */ + private final LinkedHashMap recentlyAcked = new LinkedHashMap<>(); + /** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */ + private static final int RECENT_ACK_LIMIT = 1_024; /** * A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack) @@ -162,16 +167,15 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { } // On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the // tags we were holding are now stale. Drop the held snapshot so the re-attached consumer - // repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish + // repopulates it with valid tags (dedup by msgId still prevents any double-queue). LeadCoordLoop + // remembers successful pane writes separately, so that redelivery cannot write a pane twice. Any publish // confirm still in flight when the connection dropped is equally stale — fail it now rather // than let it silently ride out CONFIRM_TIMEOUT_MS. if (connection instanceof Recoverable recoverable) { recoverable.addRecoveryListener(new RecoveryListener() { @Override public void handleRecovery(Recoverable recoverable) { - synchronized (held) { - held.clear(); - } + clearHeldForRecovery(); failPendingPublishesOnRecovery(); log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery"); } @@ -350,15 +354,26 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { return snapshot; } - /** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */ + /** + * Remove the held message {@code msgId} and ack it on the broker. + * + *

A repeated ack that this connection already completed is a no-op, tracked in the bounded + * {@link #recentlyAcked} set. Any other missing entry throws: recovery clears {@link #held} while + * the broker still owns the unacked delivery, and quiet success there would hide a required retry. + * The set is bounded because it only distinguishes a recent duplicate caller ack from an unknown + * delivery; it is not a substitute for broker state across a reconnect. + */ @Override public void ack(String msgId) { Held h; synchronized (held) { h = held.remove(msgId); + if (h == null && recentlyAcked.containsKey(msgId)) { + return; + } } if (h == null) { - return; // never held (or already acked) — no-op + throw new IllegalStateException("cannot ack lead message " + msgId + ": it is not held"); } try { synchronized (channelLock) { @@ -372,6 +387,19 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { } throw new IllegalStateException("cannot ack lead message " + msgId, e); } + synchronized (held) { + recentlyAcked.put(msgId, Boolean.TRUE); + if (recentlyAcked.size() > RECENT_ACK_LIMIT) { + recentlyAcked.remove(recentlyAcked.keySet().iterator().next()); + } + } + } + + /** Clear stale delivery tags after recovery; package-private so the recovery contract test drives this exact path. */ + void clearHeldForRecovery() { + synchronized (held) { + held.clear(); + } } private DeliverCallback deliverCallback() { diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java index e1150b9..dc35d24 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java @@ -54,6 +54,20 @@ class LeadCoordLoopTest { assertTrue(channel.peek().isEmpty(), "and is no longer held"); } + @Test + void redeliveryOfAMessageAlreadyWrittenToThePaneIsAckedWithoutAnotherPaneWrite() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me")); + var herdr = new FakeHerdr().agentStatus("idle"); + var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF)); + + loop.tick(); + channel.hold(new LeadMessage("m1", PEER, SELF, "recover me")); + loop.tick(); + + assertEquals(1, prompts(herdr).size(), "a redelivery must not consume the lead pane twice"); + assertEquals(List.of("m1", "m1"), channel.acked(), "the redelivery still needs a fresh broker ack"); + } + @Test void leavesTheMessageUnackedWhenTheLeadIsMidTurn() { var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java index ccc5144..6998b47 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java @@ -92,6 +92,33 @@ class LeadMailboxTest { } } + @Test + void ackThrowsWhenRecoveryClearedTheHeldMessage() throws Exception { + String to = coordId("lead-recovery-ack"); + try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) { + inbox.publish(to, new LeadMessage("recovery-ack", "lead-from", to, "in flight")); + assertEquals(1, awaitPeek(inbox).size(), "the broker delivery must be held before recovery clears it"); + + inbox.clearHeldForRecovery(); + + assertThrows(IllegalStateException.class, () -> inbox.ack("recovery-ack"), + "a cleared delivery has no valid tag, so ack must report that it did not reach the broker"); + } + } + + @Test + void ackOfAMessageAlreadyAckedOnThisConnectionStaysQuiet() throws Exception { + String to = coordId("lead-repeat-ack"); + try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) { + inbox.publish(to, new LeadMessage("repeat-ack", "lead-from", to, "once")); + awaitPeek(inbox); + + inbox.ack("repeat-ack"); + + inbox.ack("repeat-ack"); + } + } + @Test void duplicateMsgIdIsNotDoubleQueued() throws Exception { String to = coordId("lead-dedup");