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 16af1ef..e79c67c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -805,6 +805,12 @@ public final class FleetMcp { + "answered (turnId stale)"); case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker " + r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]"); + // fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call, + // so the message may already be sitting in the pane. Do not invite a blind retry the way + // the case above does; a resend on this route can double-deliver the same brief. + case TIMED_OUT_UNCONFIRMED -> text("[no reply within " + timeout + "ms — delivery unconfirmed; " + + "the message may already have reached the worker, so a retry risks sending it " + + "twice — poll status before resending]"); }; } 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 3345935..3f35fb8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -93,25 +93,30 @@ public final class MessageService { /** Timed out after the message was delivered — the worker is still working. */ TIMED_OUT_WORKING, /** - * Timed out with no confirmed delivery. Despite the name, this does not mean the message - * is sitting in a queue. {@link #send} reaches this outcome through {@link Injector#cancel}, - * whose result tells three routes apart: - * {@link Injector.Cancellation#CANCELLED} means the message was still queued and this call - * removed it, so the target saw nothing and it will not arrive later; - * {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was - * cleared because the target never became ready or was abandoned, or the injector's call to - * the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr - * error that this codebase already treats as a confirmed absence — so this route too - * establishes that the target saw nothing and it will not arrive later; but - * {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) means that call was made and its - * outcome is unknown. {@code agent.prompt} pastes and submits in one call, so on - * this route the target may hold a complete, already-submitted turn and be working on it - * right now — {@link Outcome#TIMED_OUT_WORKING}'s meaning, reported here as - * {@code TIMED_OUT_QUEUED} only because this caller never observed the pickup. Only - * {@code CANCELLED} and {@code NOT_DELIVERED} establish that the target saw nothing; - * {@code ATTEMPTED} does not. + * Timed out with no confirmed delivery, and the target saw nothing — the message will not + * arrive later, so a caller may resend. {@link #send} reaches this outcome through {@link + * Injector#cancel} reporting one of two routes: {@link Injector.Cancellation#CANCELLED} + * means the message was still queued and this call removed it; {@link + * Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was cleared + * because the target never became ready or was abandoned, or the injector's call to the + * target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr + * error that this codebase already treats as a confirmed absence. A third route, + * {@link Injector.Cancellation#ATTEMPTED}, used to be folded into this same outcome + * (fleetd #571) — it no longer is; see {@link #TIMED_OUT_UNCONFIRMED}. */ TIMED_OUT_QUEUED, + /** + * Timed out with delivery unknown. {@link #send} reaches this outcome when {@link + * Injector#cancel} reports {@link Injector.Cancellation#ATTEMPTED} (fleetd #551): the call + * to the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) was made, but + * this caller never observed whether it reached the pane. {@code agent.prompt} pastes + * and submits in one call, so the target may already hold a complete, submitted + * turn and be working on it right now — the same reality as {@link #TIMED_OUT_WORKING}, + * just not confirmed. The message may or may not have arrived. Treat this as neither a + * confirmed delivery nor a confirmed absence: a caller that resends on this outcome risks a + * double delivery — the same brief typed into the pane twice (fleetd #571). + */ + TIMED_OUT_UNCONFIRMED, /** Another send to this session was in flight for the whole window. */ BUSY, /** @@ -317,20 +322,21 @@ public final class MessageService { */ private final ConcurrentHashMap strandedReplies = new ConcurrentHashMap<>(); /** - * Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — {@link - * #send} called {@link Injector#cancel} and got back something other than {@code DELIVERED}. - * That covers three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message - * was still queued and {@code cancel} removed it right there; {@link - * Injector.Cancellation#NOT_DELIVERED} — nothing was ever sent, because the target never became - * ready, was torn down, or the call to its terminal failed with a herdr error this codebase - * already treats as a confirmed absence; or {@link Injector.Cancellation#ATTEMPTED} (fleetd - * #551) — the call to the target's terminal was made and its outcome is unknown, so the target - * may already hold a complete, submitted turn. Only the first two mean the message will not - * arrive later and the target saw nothing; on the third it may already have arrived in full. - * Set where {@link #send} already computes {@code wasDelivered} for that outcome; no queue is - * kept here, only the fact that the send ended with no confirmed delivery. Cleared the same way - * as {@link #strandedReplies}: the next accepted delivery for the target ({@link #send} opening - * a fresh waiter) or a teardown ({@link #abandon}). + * Targets whose last send timed out with no confirmed delivery (CB-640) — {@link #send} called + * {@link Injector#cancel} and got back something other than {@code DELIVERED}. That covers + * three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message was still + * queued and {@code cancel} removed it right there; {@link Injector.Cancellation#NOT_DELIVERED} + * — nothing was ever sent, because the target never became ready, was torn down, or the call to + * its terminal failed with a herdr error this codebase already treats as a confirmed absence; or + * {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) — the call to the target's terminal was + * made and its outcome is unknown, so the target may already hold a complete, submitted turn. + * Only the first two mean the message will not arrive later and the target saw nothing; on the + * third it may already have arrived in full — and the caller sees a different outcome for it + * ({@link Outcome#TIMED_OUT_UNCONFIRMED}, fleetd #571) than for the first two ({@link + * Outcome#TIMED_OUT_QUEUED}). Set where {@link #send} already computes {@code wasDelivered} for + * that outcome; no queue is kept here, only the fact that the send ended with no confirmed + * delivery. Cleared the same way as {@link #strandedReplies}: the next accepted delivery for the + * target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}). */ private final ConcurrentHashMap queuedDeliveries = new ConcurrentHashMap<>(); private final AtomicLong ticketSeq = new AtomicLong(); @@ -417,21 +423,22 @@ public final class MessageService { /** * Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out - * with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} (see the - * {@code TimeoutException} branch of {@link #send}). Despite the method's name, this is not - * proof that a message is sitting in a queue: {@link Injector#cancel} reports this outcome - * through three routes. {@link Injector.Cancellation#CANCELLED} means the message was still - * queued and got removed right there. {@link Injector.Cancellation#NOT_DELIVERED} means - * nothing was ever sent — the target never became ready, was torn down, or the call to its - * terminal failed with a herdr error this codebase already treats as a confirmed absence. - * Only these two routes mean the message will not arrive later. {@link + * with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} or {@link + * Outcome#TIMED_OUT_UNCONFIRMED} (fleetd #571; see the {@code TimeoutException} branch of + * {@link #send}). Despite the method's name, this is not proof that a message is sitting in a + * queue: {@link Injector#cancel} reports this outcome through three routes. {@link + * Injector.Cancellation#CANCELLED} means the message was still queued and got removed right + * there. {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the target + * never became ready, was torn down, or the call to its terminal failed with a herdr error this + * codebase already treats as a confirmed absence. Only these two routes mean the message will + * not arrive later, and both report {@code TIMED_OUT_QUEUED}. {@link * Injector.Cancellation#ATTEMPTED} (fleetd #551) means the call to the target's terminal was * made and its outcome is unknown: {@code agent.prompt} pastes and submits in one * call, so on this route the target may already hold a complete, submitted turn and be - * working on it right now — it does NOT follow that the target saw nothing. Distinct from - * {@link Outcome#TIMED_OUT_WORKING}, where delivery already happened and only the reply is - * outstanding. Cleared the next time this target's delivery is accepted or the target is - * abandoned — see {@link #queuedDeliveries}. + * working on it right now — it does NOT follow that the target saw nothing, and this route + * reports {@code TIMED_OUT_UNCONFIRMED} instead. Distinct from {@link Outcome#TIMED_OUT_WORKING}, + * where delivery already happened and only the reply is outstanding. Cleared the next time this + * target's delivery is accepted or the target is abandoned — see {@link #queuedDeliveries}. */ public boolean hasQueuedDelivery(String target) { return target != null && queuedDeliveries.containsKey(target); @@ -646,7 +653,7 @@ public final class MessageService { return switch (o) { case REPLIED -> "replied"; case COMPLETED_UNREPLIED -> "completion_fallback"; - case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout"; + case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout"; case WORKER_FAILED -> "failed"; case BACKEND_EXHAUSTED -> "backend_exhausted"; case STALE_TURN, QUESTION -> null; // not a completed delegation @@ -975,6 +982,7 @@ public final class MessageService { } catch (TimeoutException e) { boolean wasDelivered = delivery.completion().isDone() && !delivery.completion().isCompletedExceptionally(); + Injector.Cancellation cancellation = null; if (!wasDelivered) { if (timeoutCancellationRaceHookForTest != null) { // Test-only (fleetd #345): see the field's own javadoc. @@ -982,20 +990,27 @@ public final class MessageService { } // The target monitor makes cancellation atomic with onStatus picking this // Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed. - wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED; + cancellation = injector.cancel(delivery); + wasDelivered = cancellation == Injector.Cancellation.DELIVERED; } log.debug("send to {} timed out (delivered={})", target, wasDelivered); - if (!wasDelivered) { + Outcome outcome; + if (wasDelivered) { + outcome = Outcome.TIMED_OUT_WORKING; + } else if (cancellation == Injector.Cancellation.ATTEMPTED) { + // fleetd #571: the call to the target's terminal was made and its outcome is + // unknown — the message may already have arrived in full, so this must not + // be reported as TIMED_OUT_QUEUED, which promises it never will. + outcome = Outcome.TIMED_OUT_UNCONFIRMED; + } else { // CB-640: record that delivery is not confirmed, for fleet health (see - // queuedDeliveries). Whatever injector.cancel() reported above — this call - // removed a still-queued Pending (CANCELLED), an earlier attempt already - // failed with a confirmed absence (NOT_DELIVERED), or an earlier attempt was - // made and its outcome is unknown (ATTEMPTED, fleetd #551 — the message may - // already have arrived in full) — the send ends with no confirmed delivery. + // queuedDeliveries). cancellation is CANCELLED (this call removed a + // still-queued Pending) or NOT_DELIVERED (an earlier attempt already failed + // with a confirmed absence) — both mean the target saw nothing. queuedDeliveries.put(target, Boolean.TRUE); + outcome = Outcome.TIMED_OUT_QUEUED; } - return recorded(new Reply( - wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null)); + return recorded(new Reply(outcome, null)); } catch (ExecutionException e) { Throwable cause = e.getCause(); throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause); diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 87acb9f..e451b31 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -640,15 +640,31 @@ public final class FleetApp { } default -> ctx.status(202).json(Map.of( "sessionId", id, + // fleetd #571 (ticket comment 17126): no `default` here on purpose. This switch + // is an expression, so the compiler already demands every Outcome constant have + // an arm — adding an 11th constant to Outcome is a compile error here, not a + // silent fall-through. That is exactly the bug this ticket exists to fix: + // `default -> "done"` used to sit here and would have told a REST caller the + // delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery + // is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never + // actually reach this inner switch — the outer switch above always dispatches + // them first — but they still need an arm to keep this switch exhaustive. "status", switch (reply.outcome()) { case TIMED_OUT_WORKING -> "working"; case TIMED_OUT_QUEUED -> "queued"; + // Delivery here is unknown, not merely still queued — see + // Outcome#TIMED_OUT_UNCONFIRMED's own javadoc. + case TIMED_OUT_UNCONFIRMED -> "unconfirmed"; case BUSY -> "busy"; case WORKER_FAILED -> "failed"; case BACKEND_EXHAUSTED -> "backend_exhausted"; - default -> "done"; // unreachable (terminal outcomes handled above) + case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable }, - "detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED + "detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED + ? "no reply within " + timeout + "ms; delivery is unconfirmed — the " + + "message may already have reached the worker, so a resend " + + "risks sending it twice; poll status first" + : (reply.outcome() == MessageService.Outcome.WORKER_FAILED || reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED) && reply.text() != null ? reply.text() 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 0664118..94e74ac 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -7,6 +7,7 @@ import dev.ltms.fleet.auth.Role; import dev.ltms.fleet.config.FleetConfig; import dev.ltms.fleet.guard.SubscriptionGuard; import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.AgentStatus; import dev.ltms.fleet.herdr.FakeHerdr; import dev.ltms.fleet.herdr.PaneLocator; import dev.ltms.fleet.herdr.WorkspaceControl; @@ -57,9 +58,10 @@ class FleetMcpTest { private final FakeHerdr herdr = new FakeHerdr(); private final AgentControl agents = new AgentControl(herdr); + private final Injector injector = new Injector(agents); private final Rendezvous rendezvous = new Rendezvous(); private final InMemoryReplyInbox inbox = new InMemoryReplyInbox(); - private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous, inbox); + private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox); @BeforeEach void setUp() { @@ -298,6 +300,35 @@ class FleetMcpTest { assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res)); } + /** + * fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is + * the one message whose whole job is to stop a caller retrying a delivery that may already have + * arrived. Pin that its wording is actually distinct from the queued/working arm's retry + * invitation — a mutation that swapped this arm's text for that one still passed every other + * test in this suite, because nothing asserted the specific wording. + */ + @Test + void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception { + herdr.agentSendFailsWith("send_failed"); + CompletableFuture send = CompletableFuture.supplyAsync( + () -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of())); + long deadline = System.currentTimeMillis() + 2000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + //noinspection BusyWait + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter"); + injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt -> ATTEMPTED + + McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS); + String text = textOf(res); + assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error"); + assertTrue(text.contains("delivery unconfirmed"), "got: " + text); + assertFalse(text.contains("retry or poll status"), + "an unconfirmed delivery must not carry the queued/working arm's retry invitation — " + + "a resend here can double-deliver the same brief: got " + text); + } + @Test void sendRejectsMissingArgs() { assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError()); diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 97f2380..0155b35 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -601,6 +601,30 @@ class MessageServiceTest { } } + /** + * fleetd #571 (the acceptance test the ticket was filed for). The worker is idle so the injector + * attempts delivery, but the {@code agent.prompt} call itself fails with a herdr error that is + * not a confirmed absence (not a {@code *_not_found} code) — {@link Injector} marks the Pending + * {@code ATTEMPTED} (fleetd #551), meaning the call was made and whether it reached the pane is + * unknown. Before this fix, {@code send}'s {@code TimeoutException} branch collapsed + * {@code ATTEMPTED} into {@code TIMED_OUT_QUEUED} — a promise that the message will never arrive, + * which may already be false: {@code agent.prompt} pastes and submits in one call. + */ + @Test + void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception { + herdr.agentSendFailsWith("send_failed"); + CompletableFuture send = + CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150)); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED + + MessageService.Reply r = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.TIMED_OUT_UNCONFIRMED, r.outcome(), + "an ATTEMPTED delivery must not collapse into TIMED_OUT_QUEUED — the message may " + + "already have arrived in full, and TIMED_OUT_QUEUED promises it never will"); + assertNull(r.text()); + } + @Test void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception { CompletableFuture send = sendAsync(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java index 7d014c1..10f23d9 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java @@ -623,6 +623,28 @@ class FleetAppTest { assertTrue(herdr.called("agent.prompt"), "message was injected"); } + /** + * fleetd #571: the worker is idle, so the poller attempts delivery, but the {@code agent.prompt} + * call itself fails with a herdr error that is not a confirmed absence — {@link + * dev.ltms.fleet.inject.Injector} marks this {@code ATTEMPTED}, meaning the call was made and + * whether it reached the pane is unknown. {@code writeReply}'s default arm must map this to its + * own {@code "unconfirmed"} status, not silently fall through to {@code "done"} (which would + * claim the delegation completed) nor collapse into {@code "queued"} (which would claim the + * message will never arrive, when it may already be sitting in the pane). + */ + @Test + void messageTimesOutUnconfirmedWhenDeliveryAttemptFails() throws Exception { + FakeHerdr herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("send_failed"); + int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")); + + HttpResponse res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}"); + assertEquals(202, res.statusCode()); + JsonNode body = mapper.readTree(res.body()); + assertEquals("unconfirmed", body.get("status").asText(), + "an ATTEMPTED delivery must report its own status, not \"queued\" or \"done\""); + assertTrue(herdr.called("agent.prompt"), "delivery must have been attempted"); + } + @Test void messageRejectsBlankContent() throws Exception { int port = startHealthy();