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 55f44fa..d944780 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -466,8 +466,18 @@ public final class MessageService { if (candidates.size() == 1) { Task orphan = candidates.get(0); if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) { - if (orphan.turnId != null) { - asyncTasksByTurn.remove(orphan.turnId, orphan); + // fleetd #329 (F3): read orphan.turnId once. It used to be read twice, under no + // lock — once for this null check, once as the removal key — the same double-read + // shape fleetd #324 fixed in finishAsyncTask. clearAsyncQuestion's unlocked + // forgetTurn=true path (ask()'s timeout cleanup) can null the field between the two + // reads; capturing it once removes the torn read here too. + String turnId = orphan.turnId; + if (turnId != null) { + if (replyOrphanTurnIdRaceHookForTest != null) { + // Test-only (fleetd #329, F3): see the field's own javadoc. + replyOrphanTurnIdRaceHookForTest.run(); + } + asyncTasksByTurn.remove(turnId, orphan); } count(FleetMetrics.REPLIES, "path", "async-recovered"); return true; // the ticket itself took it — no inbox stranding at all @@ -1002,13 +1012,36 @@ public final class MessageService { // Measured when #282 was merged: this guard is DEFENCE IN DEPTH, not the thing // that makes the chained ask work. ask() calls markAsyncQuestion (:860) before // resolveQuestion (:861), so by the time this thread wakes, the task has already - // moved to the new turnId and finishAsyncTask(oldTurnId, ...) finds nothing. Removing - // this guard alone leaves the test green. Keep it anyway: it mirrors sendAsync's - // sibling guard, and that sibling's own comment (:1017) warns the two orderings are - // not something to rely on. Do NOT delete it as dead code without re-checking that - // ordering, and do not treat it as the sole protection either. - if (result.outcome() != Outcome.QUESTION) { - finishAsyncTask(turnId, result); + // moved to the new turnId — this outcome check, not a Task lookup, is what tells the + // two cases apart (see fleetd #329 below). Removing this guard alone leaves the test + // green. Keep it anyway: it mirrors sendAsync's sibling guard, and that sibling's own + // comment (:1017) warns the two orderings are not something to rely on. Do NOT delete + // it as dead code without re-checking that ordering, and do not treat it as the sole + // protection either. + // + // fleetd #329 (F1): complete the SAME Task object this method already looked up at + // :991, rather than re-resolving it from turnId a second time. The old + // finishAsyncTask(turnId, result) did its own asyncTasksByTurn.get(turnId) here, and + // that second lookup races ask()'s own timeout path: ask()'s ticket.answer().get(...) + // can time out (or lose that exact race) at essentially the same instant this method's + // rendezvous.answerAsk(turnId, ...) above already succeeded, running + // markAskTimedOut + clearAsyncQuestion(turnId, true) with no lock at all — which + // forgets turnId (removes it from asyncTasksByTurn, nulls Task.turnId) before this + // thread ever gets here. The worker's real reply then arrived, this method's own wait + // woke up with it, and the by-turnId lookup found nothing: the async ticket's future + // was never completed, so fleet_poll{ticket} stayed PENDING forever even though + // answer() itself correctly returned REPLIED. Reusing the reference captured at :991 + // — before that race window opens — sidesteps the second lookup entirely: it is the + // identical Task whatever asyncTasksByTurn or Task.turnId say by the time we reach + // this point, so finishAsyncTask(Task, Reply) can still complete its future and detach + // it using whatever turnId it now reads. The #282 chained-ask case is unaffected + // because it is still gated purely by result.outcome() == QUESTION above, which does + // not depend on this lookup — widening what "task" means here cannot complete a ticket + // the chained ask deliberately left open. When task is null (this was never an async + // ticket — a blocking fleet_ask's answer() call has no Task at all), there is nothing + // to complete, matching the old lookup-miss behaviour. + if (result.outcome() != Outcome.QUESTION && task != null) { + finishAsyncTask(task, result); } return result; } catch (TimeoutException e) { @@ -1059,7 +1092,7 @@ public final class MessageService { // this fires exactly once, from whichever path completes it: finishAsyncTask(task, result) // below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106 // completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same - // finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION + // finishAsyncTask reached via answer()'s own finishAsyncTask(task, result) once a QUESTION // is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516 // abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the // one an async ticket always takes) never told the push loop anything happened — see the @@ -1079,7 +1112,22 @@ public final class MessageService { finishAsyncTask(task, result); } } catch (Throwable t) { - task.future.completeExceptionally(t); + // fleetd #329 (F2): finishAsyncTask above already completes task.future — on its very + // first line — before doing anything else, so anything that throws afterward (inside + // finishAsyncTask's own cleanup, or from a future addition to this try block) lands + // here with the future already resolved. completeExceptionally on an already-completed + // future is a silent no-op: it returns false and does nothing, so without the check + // below the exception simply vanished — no log, no metric, nothing. Measured (see the + // ticket): temporarily reintroducing the fleetd #324 NPE reproduced 19 real exceptions + // on the ordinary path, with 74/74 tests staying green and not one log line produced. + // Do not stop completing the future first — finishAsyncTask completing before it + // cleans up is what makes a late failure harmless to the ticket's own result — only + // add the missing visibility for the case where that step, or whatever ran after it, + // has already lost the race to report through the future. + if (!task.future.completeExceptionally(t)) { + log.error("async send {} -> {} threw after its ticket was already resolved", + ticket, target, t); + } } }); pruneTerminalTickets(); @@ -1247,6 +1295,10 @@ public final class MessageService { */ private void finishAsyncTask(Task task, Reply result) { task.future.complete(result); + if (afterFinishAsyncTaskCompleteHookForTest != null) { + // Test-only (fleetd #329, F2): see the field's own javadoc. + afterFinishAsyncTaskCompleteHookForTest.run(); + } String turnId = task.turnId; if (turnId != null) { if (finishAsyncTaskRaceHook != null) { @@ -1276,6 +1328,49 @@ public final class MessageService { this.finishAsyncTaskRaceHook = hook; } + /** + * Null in production; test seam for fleetd #329 (F2) — invoked from {@link #finishAsyncTask(Task, + * Reply)} unconditionally, immediately after {@code task.future.complete(result)} runs (before + * {@link Task#turnId} is even read, so it fires regardless of whether this task was ever asked). + * A test installs this to force an exception into exactly the shape fleetd #329 identified: + * something throws inside {@link #sendAsync}'s executor task after the async ticket's future is + * already resolved, so the surrounding {@code catch (Throwable t)} can only report failure through + * {@code completeExceptionally} — a silent no-op on an already-completed future. Engineering a + * real exception to land in that exact post-completion window is what the ticket itself had to do + * by temporarily deleting a production guard (fleetd #324's {@code turnId != null} check); this + * hook drives the identical shape deterministically instead. + */ + private volatile Runnable afterFinishAsyncTaskCompleteHookForTest; + + /** + * Test-only (fleetd #329, F2): install {@link #afterFinishAsyncTaskCompleteHookForTest}. + * Package-private so the test, in the same package, can reach it without widening any production + * API. + */ + void setAfterFinishAsyncTaskCompleteHookForTest(Runnable hook) { + this.afterFinishAsyncTaskCompleteHookForTest = hook; + } + + /** + * Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after + * the single local read of {@code orphan.turnId} passes its null-check and before that (now-local) + * value is used as the {@code asyncTasksByTurn} removal key. Mirrors {@link + * #finishAsyncTaskRaceHook} exactly, for the structurally identical double-read fleetd #324 fixed + * in {@link #finishAsyncTask}: a test installs this to force, deterministically, {@code + * clearAsyncQuestion}'s unlocked {@code forgetTurn=true} path nulling {@link Task#turnId} in that + * exact window, and to confirm the single-read fix tolerates it (the captured local is used + * unconditionally, so a hook that nulls the field afterward cannot affect this call). + */ + private volatile Runnable replyOrphanTurnIdRaceHookForTest; + + /** + * Test-only (fleetd #329, F3): install {@link #replyOrphanTurnIdRaceHookForTest}. Package-private + * so the test, in the same package, can reach it without widening any production API. + */ + void setReplyOrphanTurnIdRaceHookForTest(Runnable hook) { + this.replyOrphanTurnIdRaceHookForTest = hook; + } + /** * Test-only (fleetd #324): run the exact production cleanup {@link #ask}'s own timeout path runs * unlocked — {@link #clearAsyncQuestion(String, boolean)} with {@code forgetTurn=true} — so a test @@ -1285,14 +1380,6 @@ public final class MessageService { clearAsyncQuestion(turnId, true); } - /** Complete the async ticket correlated to a specific answered turn. */ - private void finishAsyncTask(String turnId, Reply result) { - Task task = asyncTasksByTurn.get(turnId); - if (task != null) { - finishAsyncTask(task, result); - } - } - /** A new send must not open a waiter while an async ticket owns this worker's paused turn. */ private boolean hasAsyncQuestion(String target) { return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target)); 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 27834dc..9a6cbda 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1,5 +1,9 @@ package dev.ltms.fleet.msg; +import ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentStatus; import dev.ltms.fleet.herdr.FakeHerdr; @@ -11,6 +15,7 @@ import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.inject.Injector; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; @@ -1041,6 +1046,105 @@ class MessageServiceTest { "the reply completed its own ticket directly and never touched the inbox"); } + /** + * fleetd #329 (F1). {@code answer()} completes the async ticket by looking {@code turnId} up in + * {@code asyncTasksByTurn} a SECOND time (the first is at :991, purely to re-register the + * {@code asyncTasksByWaiter} entry for #282's chained-ask case). That second lookup races + * {@code ask()}'s own unlocked timeout cleanup ({@code clearAsyncQuestion(turnId, true)}, run from + * {@code markAskTimedOut} + forgetting): {@code ask()}'s {@code ticket.answer().get(timeoutMillis)} + * can time out at essentially the same instant {@code answer()}'s {@code rendezvous.answerAsk} + * call above already succeeded and unblocked the worker. fleetd #324's single-read fix does not + * help here — that fixed a torn read of one already-held {@link MessageService} internal + * {@code Task}; this is a second, independent map lookup by a caller that no longer holds the + * {@code Task} it already found once. + * + *

This test does not wait for that real race to land on its own schedule — it drives the exact + * sequence the ticket describes (worker asks, primary answers, worker's real reply arrives) and + * fires the identical production cleanup {@code ask()}'s timeout path runs + * ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd #324 + * already merged) at the point between "the worker's own {@code ask()} call has unblocked" and + * "the worker's real {@code fleet_reply} arrives" — the exact window fleetd #329 names. + * + *

What this proves: given that exact interleaving, the async ticket must still resolve to the + * worker's real reply, not stay stuck {@link MessageService.Phase#PENDING} forever. What it does + * not prove: that the interleaving itself is reachable in production on its own timing — that is + * established by reading the code (see the ticket's "path in"), not by this test, for the same + * reason fleetd #324's own race test says so. + */ + @Test + void aReplyRacingAsksTimeoutCleanupStillCompletesTheAsyncTicket() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); + MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + String turnId = asking.turnId(); + + CompletableFuture answer = + CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(), + "the worker's own ask() call must have already unblocked with the primary's answer " + + "before we force the race below"); + + // ask()'s own unlocked timeout cleanup can forget this exact turnId at essentially the same + // instant answer() has already unblocked the worker and is now waiting on the resumed turn's + // real reply — reproduce that interleaving directly instead of trying to win a real race. + messages.forgetTurnForTest(turnId); + + assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); + + assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(), + "the primary's own answer() call must still see the worker's real reply"); + MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); + assertEquals("PR opened: https://example/pulls/42", done.reply(), + "fleet_poll{ticket} must return the worker's real reply, not stay PENDING forever " + + "just because ask()'s timeout cleanup forgot this turnId first"); + } + + /** + * fleetd #329 (F3). {@link MessageService#reply} reads the same orphan {@code Task}'s {@code + * turnId} twice on this path (once to check it is non-null, once as the {@code + * asyncTasksByTurn.remove} key) — the exact double-read shape fleetd #324 fixed in {@code + * finishAsyncTask}. This test drives the same real sequence as {@code + * aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket} (worker asks, primary answers, the + * primary's own bounded wait for the resumed turn expires) to reach a task with {@code turnId} + * genuinely stamped and no live rendezvous waiter open — the state {@code askAnsweredAsyncTasks} + * matches here — then fires the identical production cleanup {@code ask()}'s own timeout path + * runs ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd + * #324 merged) at the point between the check and the removal use, via a dedicated test hook + * mirroring {@code finishAsyncTaskRaceHook}. + */ + @Test + void replyToAnOrphanedTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); + MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + String turnId = asking.turnId(); + + MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(), + "the primary's own bounded wait must give up first, leaving turnId stamped with no " + + "live waiter — the state reply()'s F3 code path matches"); + + messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId)); + try { + assertTrue(messages.reply(T, "PR opened: https://example/pulls/42")); + + MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); + assertEquals("PR opened: https://example/pulls/42", done.reply(), + "the ticket must still resolve to the worker's real reply despite the forced race"); + } finally { + messages.setReplyOrphanTurnIdRaceHookForTest(null); + } + } + /** * fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion} * becomes false the instant it lapses — proven above), so a second, independent delegation can @@ -1151,6 +1255,63 @@ class MessageServiceTest { assertEquals("done", done.reply()); } + // --- fleetd #329 (F2): an exception after the ticket's future completes must reach a log ----- + // + // sendAsync's own executor task ends with catch (Throwable t) { task.future.completeExceptionally(t); } + // — but finishAsyncTask completes that same future on its first line, so anything that throws + // afterward hits an already-completed future: completeExceptionally returns false and does + // nothing, and (before this fix) nothing logged it either. Measured on the ticket: temporarily + // reintroducing the fleetd #324 NPE reproduced 19 real exceptions on the ordinary path with + // 74/74 tests staying green and zero log lines. There is no reachable production call site where + // finishAsyncTask throws after completing the future (fleetd #324 already closed the one that + // used to), so this test drives the shape directly via a dedicated test-only hook + // (afterFinishAsyncTaskCompleteHookForTest) rather than trying to engineer a real exception into + // that narrow window — the same technique fleetd #324's own race test and this ticket's F1/F3 + // tests use for their own races. + + private static ListAppender attachMessageServiceLog() { + Logger logger = (Logger) LoggerFactory.getLogger(MessageService.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + return appender; + } + + private static void detachMessageServiceLog(ListAppender appender) { + ((Logger) LoggerFactory.getLogger(MessageService.class)).detachAppender(appender); + } + + @Test + void anExceptionAfterTheTicketFutureCompletesStillReachesTheLog() throws Exception { + ListAppender appender = attachMessageServiceLog(); + try { + messages.setAfterFinishAsyncTaskCompleteHookForTest(() -> { + throw new RuntimeException("PROBE-329-F2"); + }); + + String ticket = messages.sendAsync(T, "do the task"); + awaitWaiting(); + assertTrue(rendezvous.resolve(T, "done")); + + // finishAsyncTask's own first line already completed the future with the real result + // BEFORE the hook threw — F2 is a visibility gap, not a correctness gap for the ticket + // itself, so the ticket's own outcome must be unaffected by the swallowed exception. + MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); + assertEquals("done", done.reply()); + + assertTrue(appender.list.stream().anyMatch(e -> + e.getLevel() == Level.ERROR + && e.getFormattedMessage().contains(ticket) + && e.getThrowableProxy() != null + && "PROBE-329-F2".equals(e.getThrowableProxy().getMessage())), + "an exception thrown after the ticket's future already completed must still " + + "reach the log, not vanish silently"); + } finally { + messages.setAfterFinishAsyncTaskCompleteHookForTest(null); + detachMessageServiceLog(appender); + } + } + // --- CB-582: fleet_status pendingAsk() ------------------------------------------------------ @Test