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 19039ac..a2c1437 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -709,27 +709,47 @@ public final class MessageService { boolean isRecovery = task == recoveryTask && recovered != null; Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason); String turnId = task.turnId; - if (task.future.complete(outcome)) { - if (outcome.outcome() == Outcome.WORKER_FAILED) { - asyncFailed = true; + // fleetd #335: completing THIS task's future must not depend on any other task's + // cleanup succeeding — every task in `matching` is owed its own outcome regardless of + // what happens below, so decide and record that before doing anything that can throw. + boolean completedHere = task.future.complete(outcome); + if (completedHere && outcome.outcome() == Outcome.WORKER_FAILED) { + asyncFailed = true; + } + try { + if (abandonCleanupHookForTest != null) { + // Test-only (fleetd #335 site 1): see the field's own javadoc. + abandonCleanupHookForTest.run(); } - if (turnId != null) { - // #275: whether this task was swept out of ASKING or was already answered and - // only waiting on its resumed turn's real reply (#137), nothing will ever - // complete this turnId now — drop it from this class's own bookkeeping AND the - // reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn - // that is actually done, and a late answer() sees it as lapsed rather than - // resolving a question nothing is listening for any more. - asyncTasksByTurn.remove(turnId, task); - rendezvous.closeAsk(turnId); + if (completedHere) { + if (turnId != null) { + // #275: whether this task was swept out of ASKING or was already answered and + // only waiting on its resumed turn's real reply (#137), nothing will ever + // complete this turnId now — drop it from this class's own bookkeeping AND the + // reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn + // that is actually done, and a late answer() sees it as lapsed rather than + // resolving a question nothing is listening for any more. + asyncTasksByTurn.remove(turnId, task); + rendezvous.closeAsk(turnId); + } + } else if (isRecovery) { + // The recovered reply was already drained out of the inbox, but this task resolved + // through another path (e.g. a concurrent reply() or a second abandon() racing this + // one) between us choosing it and completing it here. Put the reply back rather than + // lose it silently — it may still belong to some other still-open task, or the next + // caller that drains this target's inbox. + inbox.publish(target, UUID.randomUUID().toString(), recovered.text()); } - } else if (isRecovery) { - // The recovered reply was already drained out of the inbox, but this task resolved - // through another path (e.g. a concurrent reply() or a second abandon() racing this - // one) between us choosing it and completing it here. Put the reply back rather than - // lose it silently — it may still belong to some other still-open task, or the next - // caller that drains this target's inbox. - inbox.publish(target, UUID.randomUUID().toString(), recovered.text()); + } catch (RuntimeException e) { + // fleetd #335: inbox.publish reaches a broker (AmqpReplyInbox throws + // IllegalStateException on an unroutable/unconfirmed/interrupted publish) and this + // loop has no other teardown path — a caller on the release path, or the health + // monitor's GONE/NEVER_READY sweep. Losing this exception uncaught would abort the + // loop and leave every task still to come in `matching` PENDING forever (fleetd + // #335 site 1). completedHere is already recorded above, so only this task's + // best-effort bookkeeping is lost — log it and let the loop reach the rest. + log.error("abandon: per-task cleanup failed for ticket {} (target {}, turnId {})", + task.ticket, target, turnId, e); } } if (failed) { @@ -1118,7 +1138,18 @@ public final class MessageService { // class javadoc on sendAsync/CB-107. task.future.whenComplete((reply, ex) -> { boolean failed = ex != null || reply == null || !reply.completed(); - pushLoop.onTicketTerminal(ticket, target, failed); + try { + pushLoop.onTicketTerminal(ticket, target, failed); + } catch (Throwable t) { + // fleetd #335 (site 2): this stage's own CompletableFuture is discarded, so an + // uncaught throw here (e.g. a RejectedExecutionException from + // ReplyPushLoop's scheduler, already shut down while this in-flight send's + // whenComplete fires during the daemon's own shutdown sequence — messages.close() + // only stops accepting NEW async work, it does not cancel a delivery already + // running) vanishes with no log line and no metric, and the push loop never learns + // the ticket went terminal — the exact thing this hook exists to tell it. + log.error("push loop failed to learn ticket {} (target {}) went terminal", ticket, target, t); + } }); } asyncExecutor.submit(() -> { @@ -1399,6 +1430,33 @@ public final class MessageService { clearAsyncQuestion(turnId, true); } + /** + * Null in production; test seam for fleetd #335 (site 1) — invoked from {@link #abandon(String, + * String, boolean)}'s per-task loop, once per task, right before that task's own cleanup + * (turnId bookkeeping, or the stranded-reply put-back) runs. A test installs this to inject a + * throw at that exact point deterministically. + * + *

The one production call there that can really throw is {@code inbox.publish} in the + * put-back branch — {@link AmqpReplyInbox#publish} reaches a broker and throws {@link + * IllegalStateException} on an unroutable, unconfirmed, or interrupted publish — but reaching + * that branch requires a second completion of the very task {@code abandon} is about to + * complete to win the race first (see the branch's own comment), and the #137 follow-up + * investigation above already found the combination this needs (a stranded reply coinciding + * with an open matching task) unreachable through the public API, not merely hard to time. + * This hook reproduces the resulting shape — a per-task cleanup throw — directly, the same + * technique {@link #finishAsyncTaskRaceHook} and {@link + * #afterFinishAsyncTaskCompleteHookForTest} already use for their own hard-to-time races. + */ + private volatile Runnable abandonCleanupHookForTest; + + /** + * Test-only (fleetd #335, site 1): install {@link #abandonCleanupHookForTest}. Package-private + * so the test, in the same package, can reach it without widening any production API. + */ + void setAbandonCleanupHookForTest(Runnable hook) { + this.abandonCleanupHookForTest = hook; + } + /** 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 f9f196b..b4bbea6 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -785,6 +785,49 @@ class MessageServiceTest { assertFailedTicket(third, "agent target term_a not found"); } + // --- fleetd #335 (site 1): a per-task cleanup failure inside the abandon() loop must not ----- + // strand the tasks that come after it. abandon()'s own comment on the loop documents the one + // real production call that can throw there (inbox.publish, in the stranded-reply put-back + // branch, reached when a concurrent reply() or a second abandon() races this one) — but + // reaching that branch requires the exact combination the #137 follow-up above already found + // unreachable through the public API. abandonCleanupHookForTest reproduces the resulting SHAPE + // (one task's cleanup throws) directly instead, the same technique this file already uses for + // fleetd #324/#329's own hard-to-time races. + @Test + void aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks() throws Exception { + ListAppender appender = attachMessageServiceLog(); + try { + String first = messages.sendAsync(T, "first task"); + awaitWaiting(); // first task owns the target lock and rendezvous waiter + String second = messages.sendAsync(T, "second task"); // parked on the same lock + String third = messages.sendAsync(T, "third task"); // parked too — the whole sweep must survive + + java.util.concurrent.atomic.AtomicInteger calls = new java.util.concurrent.atomic.AtomicInteger(); + messages.setAbandonCleanupHookForTest(() -> { + if (calls.getAndIncrement() == 0) { + throw new RuntimeException("PROBE-335-SITE1"); + } + }); + + assertTrue(messages.abandon(T, "agent target term_a not found")); + + // Every task in the loop still gets its own outcome — the one whose cleanup threw + // included — even though the loop had no way to know in advance which one that would be. + assertFailedTicket(first, "agent target term_a not found"); + assertFailedTicket(second, "agent target term_a not found"); + assertFailedTicket(third, "agent target term_a not found"); + + assertTrue(appender.list.stream().anyMatch(e -> + e.getLevel() == Level.ERROR + && e.getThrowableProxy() != null + && "PROBE-335-SITE1".equals(e.getThrowableProxy().getMessage())), + "a per-task cleanup failure must still reach the log, not vanish silently"); + } finally { + messages.setAbandonCleanupHookForTest(null); + detachMessageServiceLog(appender); + } + } + // --- #137 follow-up: abandon() must not guess when more than one task is open --------------- // // A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously @@ -1472,6 +1515,42 @@ class MessageServiceTest { } } + // --- fleetd #335 (site 2): task.future.whenComplete's own returned stage is discarded, so an -- + // uncaught throw from ReplyPushLoop.onTicketTerminal used to vanish with no log line and no + // metric. The real production trigger is the daemon's own shutdown sequence (Fleetd's shutdown + // hook): messages.close() only stops the async executor from taking NEW work — it does not + // cancel a send already in flight — while pushLoop.close() shuts its scheduler down immediately + // right after, so a ticket that completes in that narrow window has onTicketTerminal's own + // scheduler.schedule(...) throw a real RejectedExecutionException. Reproduced here by shutting + // the very same scheduler down before the ticket resolves — no test-only hook needed, this + // reachable path throws for real. + @Test + void aTicketTerminalPushFailureDoesNotVanishSilently() throws Exception { + ListAppender appender = attachMessageServiceLog(); + try (var wiring = wireWithPushLoop(1, 50)) { + String ticket = wiring.service().sendAsync(T, "long task"); + awaitWaiting(); + wiring.scheduler().shutdownNow(); // simulate pushLoop.close() racing an in-flight send + injectDelivery(); + assertTrue(rendezvous.resolve(T, "async result")); + + // The ticket's own outcome must be unaffected by the swallowed exception — finishAsyncTask + // completes task.future before whenComplete's action (and thus onTicketTerminal) ever runs. + MessageService.TaskView done = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE); + assertEquals("async result", done.reply()); + + assertTrue(appender.list.stream().anyMatch(e -> + e.getLevel() == Level.ERROR + && e.getFormattedMessage().contains(ticket) + && e.getThrowableProxy() != null + && "java.util.concurrent.RejectedExecutionException" + .equals(e.getThrowableProxy().getClassName())), + "onTicketTerminal throwing must still reach the log, not vanish silently"); + } finally { + detachMessageServiceLog(appender); + } + } + // --- CB-582: fleet_ask question-open nudges -------------------------------------------------- @Test