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");