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 4bba924..0a73221 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) { @@ -1122,7 +1142,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(() -> { @@ -1421,6 +1452,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 86eea10..001870d 100644
--- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java
+++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java
@@ -808,6 +808,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