diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 16e19a8..e351f8a 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -994,7 +994,11 @@ public final class FleetMcp { if (isBlank(target) || isBlank(msgId)) { return error("target and msgId are required"); } - messages.ackReply(target, msgId); + if (!messages.ackReply(target, msgId)) { + return error(msgId + " is not in " + target + "'s reply inbox (wrong id, wrong target, " + + "or already acked). Held lead-to-lead (peer) mail cannot be acked this way — " + + "read it with fleet_poll{coordId}."); + } return text("acknowledged " + msgId); } @@ -1766,7 +1770,10 @@ public final class FleetMcp { + "has processed a reply and wants to confirm it, leaving other pending replies " + "in the inbox for later drain.", objectSchema(Map.of( - "target", stringProp("Worker session id whose inbox to ack from"), + "target", stringProp("Worker session id whose inbox to ack from. Must name a " + + "reply actually queued for it — an id in no inbox, or a coord-id " + + "(peer held mail, read with fleet_poll{coordId} instead), errors " + + "rather than reporting a false success"), "msgId", stringProp("The message id to acknowledge")), List.of("target", "msgId"))); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java index 95b984a..3c631f8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java @@ -394,17 +394,17 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { } @Override - public void ack(String target, String msgId) { + public boolean ack(String target, String msgId) { var perTarget = held.get(target); if (perTarget == null || perTarget == RELEASED) { - return; + return false; } Held h; synchronized (perTarget) { h = perTarget.remove(msgId); } if (h == null) { - return; // never held (or already acked) — no-op + return false; // never held (or already acked) — no-op } try { synchronized (channelLock) { @@ -418,6 +418,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { } throw new IllegalStateException("cannot ack reply " + msgId + " on " + queueName(target), e); } + return true; } private DeliverCallback deliverCallback(String target) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/InMemoryReplyInbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/InMemoryReplyInbox.java index 69d72c9..efab1ad 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/InMemoryReplyInbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/InMemoryReplyInbox.java @@ -59,16 +59,17 @@ public final class InMemoryReplyInbox implements ReplyInbox { } @Override - public void ack(String target, String msgId) { + public boolean ack(String target, String msgId) { if (!owned.contains(target)) { - return; + return false; } var perTarget = store.get(target); - if (perTarget != null) { - //noinspection SynchronizationOnLocalVariableOrMethodParameter - synchronized (perTarget) { - perTarget.remove(msgId); - } + if (perTarget == null) { + return false; + } + //noinspection SynchronizationOnLocalVariableOrMethodParameter + synchronized (perTarget) { + return perTarget.remove(msgId) != null; } } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 04ebc21..cf5a76f 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -826,9 +826,13 @@ public final class MessageService { /** * Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox * so that a subsequent drain or peek no longer returns it. + * + * @return {@code true} if an entry was actually removed, {@code false} if {@code msgId} was not + * in {@code target}'s inbox (wrong id, wrong target, or already acked). The caller — + * {@link dev.ltms.fleet.mcp.FleetMcp#ack} — must not report success on {@code false}. */ - public void ackReply(String target, String msgId) { - inbox.ack(target, msgId); + public boolean ackReply(String target, String msgId) { + return inbox.ack(target, msgId); } /** diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyInbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyInbox.java index 38ce4b9..dca7e8b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyInbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyInbox.java @@ -46,6 +46,13 @@ public interface ReplyInbox { /** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */ List peek(String target); - /** Remove the reply {@code msgId} for {@code target} once the primary has taken it. No-op if absent. */ - void ack(String target, String msgId); + /** + * Remove the reply {@code msgId} for {@code target} once the primary has taken it. + * + * @return {@code true} if an entry was actually removed, {@code false} if there was nothing to + * remove (unknown {@code target}, unowned {@code target}, or a {@code msgId} not held for + * it). A {@code false} is not an error — acking a {@code target} this daemon does not own is + * part of the normal contract, not a failure. + */ + boolean ack(String target, String msgId); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 2186371..d1aac0d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -352,7 +352,7 @@ class FleetMcpTest { fail("a lead fleet_reply must not publish to the worker inbox"); } @Override public List peek(String target) { return List.of(); } - @Override public void ack(String target, String msgId) { } + @Override public boolean ack(String target, String msgId) { return false; } }; MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(), inboxThatRejectsPublishes); @@ -1386,6 +1386,10 @@ class FleetMcpTest { @Test void bridgeAckReturnsConfirmationForValidArgs() { + // fleet_ack only reports success for a msgId actually queued in the target's inbox + // (fleetd #437) — publish one via the inbox directly rather than asserting on a + // fabricated id nothing ever queued. + inbox.publish("term_a", "msg-1", "queued reply"); McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "msg-1"); assertNotEquals(Boolean.TRUE, res.isError()); assertTrue(textOf(res).contains("msg-1"), "response should mention the msgId"); @@ -1398,6 +1402,26 @@ class FleetMcpTest { assertTrue(FleetMcp.ack(messages, " ", "msg-1").isError()); } + @Test + void bridgeAckOfAnIdInNoInboxIsAnError() { + // fleetd #437: fleet_ack used to say "acknowledged " for a message it never + // touched, because nothing in the chain reported hit vs. miss. "never-queued" is in no + // inbox at all, so this must error rather than claim success. + McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "never-queued"); + assertTrue(res.isError()); + assertTrue(textOf(res).contains("never-queued"), textOf(res)); + } + + @Test + void bridgeAckOfACoordIdTargetIsAnErrorNamingFleetPoll() { + // A coord-id names a peer lead's held mailbox (LeadChannel/LeadMailbox), never a + // worker's ReplyInbox — fleet_ack has no route to it and must say so, pointing at + // fleet_poll{coordId} instead of reporting a false "acknowledged". + McpSchema.CallToolResult res = FleetMcp.ack(messages, "coord-some-peer", "msg-1"); + assertTrue(res.isError()); + assertTrue(textOf(res).contains("fleet_poll{coordId}"), textOf(res)); + } + @Test void bridgeAckRemovesSpecificReply() { // Queue a reply and capture its msgId. @@ -1411,8 +1435,18 @@ class FleetMcpTest { var peeked = messages.drainReplies("term_a"); assertEquals(1, peeked.size(), "one fresh reply in the inbox"); - // ackReply works (no-op since published with a different UUID, but callable). - assertDoesNotThrow(() -> messages.ackReply("term_a", msgId)); + // fleetd #437: msgId was already drained above (a fresh UUID each publish), so it is no + // longer in the inbox — ackReply must now report that miss instead of pretending to ack. + assertFalse(messages.ackReply("term_a", msgId)); + } + + @Test + void bridgeAckRemovingARealQueuedReplyReportsSuccessAndRemovesIt() { + // The worker path must not change behaviour: acking a reply that IS still in the inbox + // still succeeds and still removes it (fleetd #437). + inbox.publish("term_a", "real-1", "still queued"); + assertTrue(messages.ackReply("term_a", "real-1"), "ack of a real queued reply must report true"); + assertTrue(inbox.peek("term_a").isEmpty(), "the acked reply must be gone from the inbox"); } @Test diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java index b310062..8f2e35a 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java @@ -16,6 +16,7 @@ import java.util.List; import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -221,6 +222,31 @@ class AmqpReplyInboxContractTest { } } + @Test + void ackReportsHitVsMissAgainstARealBroker() throws Exception { + // fleetd #437: fleet_ack said "acknowledged " for a message it never touched, + // because ReplyInbox.ack() (void) could not tell a hit from a miss. Pin the fixed + // boolean contract against a real broker — the adapter fleetd actually runs live. + String target = "worker-ack-contract-" + System.nanoTime(); + try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) { + inbox.own(target); + + // Never held for this target at all: must report false, not throw. + assertFalse(inbox.ack(target, "never-held"), + "acking a msgId never held for an owned target must report false"); + + // A real message: first ack removes it and reports true... + inbox.publish(target, "m1", "ack me"); + assertEquals(1, awaitPeek(inbox, target).size(), "the published reply should be held"); + assertTrue(inbox.ack(target, "m1"), "acking a held reply must report true"); + assertTrue(inbox.peek(target).isEmpty(), "an acked reply is dropped"); + + // ...and the second ack of the SAME msgId has nothing left to remove: false. + assertFalse(inbox.ack(target, "m1"), + "acking the same msgId twice must report false the second time"); + } + } + @Test void confirmedPublishDeliversNormally() throws Exception { String target = "worker-confirm-" + System.nanoTime(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/InMemoryReplyInboxTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/InMemoryReplyInboxTest.java index f48e4c8..ec4a021 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/InMemoryReplyInboxTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/InMemoryReplyInboxTest.java @@ -41,20 +41,20 @@ class InMemoryReplyInboxTest { @Test void ackRemovesTheMessage() { inbox.publish("term_a", "m1", "hello"); - inbox.ack("term_a", "m1"); + assertTrue(inbox.ack("term_a", "m1"), "fleetd #437: ack of a real entry must report true"); assertTrue(inbox.peek("term_a").isEmpty(), "after ack, the message is gone"); } @Test void ackForUnknownMsgIdIsNoOp() { inbox.publish("term_a", "m1", "hello"); - inbox.ack("term_a", "no-such-id"); // no-op + assertFalse(inbox.ack("term_a", "no-such-id"), "fleetd #437: a miss must report false"); // no-op assertEquals(1, inbox.peek("term_a").size(), "the published message is still there"); } @Test void ackForUnknownTargetIsNoOp() { - inbox.ack("no-such-target", "m1"); // no-op, should not throw + assertFalse(inbox.ack("no-such-target", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw } @Test @@ -167,7 +167,7 @@ class InMemoryReplyInboxTest { @Test void peekAndAckAreNoOpsForUnownedTarget() { assertTrue(inbox.peek("term_not_owned").isEmpty()); - inbox.ack("term_not_owned", "m1"); // no-op, should not throw + assertFalse(inbox.ack("term_not_owned", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw } @Test