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..4bba924 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -871,6 +871,10 @@ public final class MessageService { boolean wasDelivered = delivery.completion().isDone() && !delivery.completion().isCompletedExceptionally(); if (!wasDelivered) { + if (timeoutCancellationRaceHookForTest != null) { + // Test-only (fleetd #345): see the field's own javadoc. + timeoutCancellationRaceHookForTest.run(); + } // 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; @@ -1370,6 +1374,24 @@ public final class MessageService { this.afterFinishAsyncTaskCompleteHookForTest = hook; } + /** + * Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after + * {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel} + * takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued + * delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}. + * This deterministically covers the caller's need to use that result rather than relying on a + * timing-sensitive real race. + */ + private volatile Runnable timeoutCancellationRaceHookForTest; + + /** + * Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private + * so the test, in the same package, can reach it without widening any production API. + */ + void setTimeoutCancellationRaceHookForTest(Runnable hook) { + this.timeoutCancellationRaceHookForTest = 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) 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..86eea10 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -410,6 +410,29 @@ class MessageServiceTest { "a delivered send whose worker never replies times out as still working"); } + /** + * fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send} + * first observes its completion as incomplete, but before {@code cancel} takes the target monitor. + * The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report + * {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed. + * + *
What this does not prove: that this precise interleaving happens by itself under production
+ * timing. The test forces it through a test-only hook; it proves the timeout caller handles the
+ * injector result when the interleaving occurs.
+ */
+ @Test
+ void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
+ messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
+ try {
+ MessageService.Reply reply = messages.send(T, "race delivery", 50);
+
+ assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
+ "cancel reporting DELIVERED means the worker received the timed-out message");
+ } finally {
+ messages.setTimeoutCancellationRaceHookForTest(null);
+ }
+ }
+
@Test
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
CompletableFuture