From b6b88c5f1c69cf4f8d1228c67810eafbe2871e34 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 17:12:49 +0700 Subject: [PATCH] #334: pin the fresh-owner gate on ask()'s timeout teardown --- .../ltms/fleet/msg/MessageServiceTest.java | 50 +++++++++++++++++++ 1 file changed, 50 insertions(+) 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 d965d10..c3dca80 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -393,6 +393,56 @@ class MessageServiceTest { * never made answerable again, and (2) the async ticket still resolves {@code DONE} once the * worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}. */ + /** + * fleetd #334 gated the ask-timeout teardown on {@code ticket.fresh()}, matching the {@code + * finally} block that already did. This pins that gate. A coalesced duplicate passes its own + * {@code timeoutMillis}, which says nothing about whether the shared ask is done — so a + * duplicate timing out first must leave the fresh owner's still-open ask answerable. + * + *

Measured on merge: without this test, removing the {@code ticket.fresh()} gate left all + * 1371 tests green. The gate shipped with the reorder and nothing held it there. + * + *

What this does not prove: anything about the ordering inside the gate — that is + * {@code aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket}'s job. + */ + @Test + void aCoalescedDuplicateAskTimingOutLeavesTheFreshOwnersAskOpen() throws Exception { + String ticket = messages.sendAsync(T, "long task"); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); // deliver + injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask + + CompletableFuture fresh = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + + MessageService.TaskView asking = null; + long deadline = System.currentTimeMillis() + 2000; + while ((asking == null || asking.phase() != MessageService.Phase.ASKING) + && System.currentTimeMillis() < deadline) { + asking = messages.poll(ticket); + //noinspection BusyWait + Thread.sleep(5); + } + assertNotNull(asking, "the fresh owner's question must surface before the duplicate asks"); + String turnId = asking.turnId(); + assertNotNull(turnId, "an ASKING view carries the turnId to answer on"); + + // A coalesced duplicate on the same session, with its own much shorter timeout. + MessageService.AskResult duplicate = messages.ask(T, "which config file?", 100); + assertEquals(MessageService.AskOutcome.TIMED_OUT, duplicate.outcome(), + "the duplicate's own timeout elapses first"); + + CompletableFuture answered = + CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500)); + + MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(), + "a duplicate's timeout must not lapse the ask the fresh owner still holds"); + assertEquals("fleetd.yaml", a.answer()); + assertFalse(answered.get(5, TimeUnit.SECONDS).outcome() == MessageService.Outcome.STALE_TURN, + "the answer must not be rejected as stale"); + } + @Test void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception { String ticket = messages.sendAsync(T, "long task");