Merge #448: fleet_ack errors instead of claiming success on a miss (fleetd #437)
CI / contract (push) Successful in 52s
CI / build (push) Successful in 1m33s

Verified by the lead on head 5289eb5, and then again on the MERGE.

Battery on the branch tip:

  FULL BUILD  Tests run: 1577, Failures: 0, Errors: 0, Skipped: 0  BUILD SUCCESS
              compile errors: 0
  CONTROL     contract test, real broker: Tests run: 9, Failures: 0, Skipped: 0

  M3a  AmqpReplyInbox.ack's `h == null` branch reports true
       -> KILLED  AmqpReplyInboxContractTest.ackReportsHitVsMissAgainstARealBroker
  M3b  AmqpReplyInbox.ack's `perTarget == null` branch reports true
       -> KILLED  same method

M3a is the gap I found in round 1: reporting true for something never held
survived both the default suite and -Pcontract against a real broker. It is now
closed in the adapter this daemon actually runs.

Both branches die, and they die through DIFFERENT assertions, which the worker
worked out and I confirmed by reading the code. held.get(target) is populated by
the deliver callback, not by own(), so after own() with nothing ever delivered
the map entry is still null — the "never held" assertion therefore exercises
`perTarget == null`, and only the double-ack assertion reaches `h == null`. Both
assertions are load-bearing; neither is redundant.

This PR was pushed before #451 landed, so the battery above tested the branch,
not the result. Merging main (bdcf285) in gave 0 conflicts and #448 adds no
PeerLauncher implementer, but a clean auto-merge is not a compiling merge, so I
built the merged tree 0829542:

  FULL BUILD               Tests run: 1578, Failures: 0  BUILD SUCCESS, 0 compile errors
  contract, real broker    Tests run: 9,    Failures: 0
  CompositePeerLauncherTest (#447's guarantee)  Tests run: 76, Failures: 0

The worker also declined to point the contract test at the shared local LavinMQ
broker, because this adapter never deletes queues and a run would leave orphaned
durable queues on the instance backing the live fleet. It started a disposable
rabbitmq:3.13-management container on a throwaway port instead, then removed it.
That was its own judgment and it was correct.

Held peer mail is deliberately NOT ackable: fleet_ack against a coord-id errors
and names fleet_poll{coordId}. See #437 for why refusing is the right answer
today, and why that is a policy choice rather than a structural one.
This commit was merged in pull request #448.
This commit is contained in:
2026-09-10 12:20:49 +02:00
8 changed files with 103 additions and 23 deletions
@@ -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")));
}
@@ -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) {
@@ -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;
}
}
}
@@ -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);
}
/**
@@ -46,6 +46,13 @@ public interface ReplyInbox {
/** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */
List<InboxMessage> 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);
}
@@ -352,7 +352,7 @@ class FleetMcpTest {
fail("a lead fleet_reply must not publish to the worker inbox");
}
@Override public List<InboxMessage> 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 <msgId>" 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
@@ -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 <msgId>" 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();
@@ -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