From a4dbc8f8b7a33d53f85b204cd5a5e06fa9a98c83 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 18:23:54 +0700 Subject: [PATCH] fleetd #572: pin answer()'s session-lock release across all four exits MessageService.answer() releases its per-session lock in an outer finally (MessageService.java:1218) that mutation testing showed was covered but unasserted: removing that line left all 1750 existing tests green, because every existing test on this path checks answer()'s return value, never that the lock it took is actually reacquirable afterward. If it leaked, a session would be wedged forever with no exception and no log line. Adds four tests, one per exit of answer() (normal REPLIED reply, TIMED_OUT_WORKING, ExecutionException rethrow, InterruptedException rethrow), each proving the lock is reacquirable via a bounded (300ms) follow-up send on the same session rather than merely checking answer()'s own outcome. No production change. --- .../ltms/fleet/msg/MessageServiceTest.java | 178 ++++++++++++++++++ 1 file changed, 178 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 72d8b4f..4b83d04 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -22,6 +22,7 @@ import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -579,6 +580,183 @@ class MessageServiceTest { assertEquals("config.yaml", a.answer()); } + // --- fleetd #572: answer() must release the session lock on EVERY exit, not just return it ----- + // + // answer()'s session lock is released in an outer `finally` (MessageService.java:1217-1219) that + // wraps its whole body. Removing that one line survives the entire suite: coverage is not the + // gap, assertion is — every existing answer() test checks the RETURN VALUE, never that the lock + // it took is reacquirable afterward. If it is not, the session is wedged forever with no + // exception and no log line. Each test below drives answer() through one of its four exits + // (normal reply, TIMED_OUT_WORKING, ExecutionException, InterruptedException) and then proves + // reacquisition the only way that actually proves it: a bounded follow-up `send` on the SAME + // session must not come back BUSY. A `send` only reports BUSY when `tryLock` itself timed out + // (MessageService.java:926-928) — every other early return in `send` still passes through its own + // lock-acquired `try`, so a non-BUSY probe result is specifically evidence the lock was free. + + /** The REPLIED exit (the happy path) — the resumed worker's real {@code fleet_reply} arrives. */ + @Test + void answerReleasesTheSessionLockAfterANormalReply() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.Reply q = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, q.outcome()); + + CompletableFuture answer = + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); + ask.get(5, TimeUnit.SECONDS); // worker resumed with the answer + + awaitWaiting(); // the answering call has (re)opened its own forward waiter + assertTrue(rendezvous.resolve(T, "done"), "the worker's final reply resolves the answering send"); + MessageService.Reply done = answer.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, done.outcome()); + + MessageService.Reply probe = messages.send(T, "probe after normal reply", 300); + assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), + "the session lock must be released after a normal REPLIED answer(), or this bounded " + + "follow-up send would come back BUSY instead of timing out on its own work"); + } + + /** + * The {@code TIMED_OUT_WORKING} exit. Same setup as {@link #answerTimesOutWhenTheResumedWorkerNeverReplies} + * (which only checks {@code answer()}'s return value — exactly the assertion the fleetd #572 + * mutation survives), plus the reacquisition proof that test does not make. + * + *

{@code answer()} MUST run on its own thread here, not the test's main thread: {@code lock} + * is a {@link java.util.concurrent.locks.ReentrantLock}, so a probe {@code send} issued by the + * SAME thread that (under the mutation) still "holds" it would reenter for free and report a + * false pass — reentrancy, not release. Measured while writing this test: with {@code answer()} + * called inline, this test stayed green under the mutation while its three siblings correctly + * went red. + */ + @Test + void answerReleasesTheSessionLockAfterATimedOutWorkingReturn() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.Reply q = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, q.outcome()); + + // The primary answers, unblocking the worker; the worker never sends its follow-up + // fleet_reply, so the answering call rides out its short window as still-working. + CompletableFuture answer = + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200)); + MessageService.Reply answered = answer.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answered.outcome(), + "an answered worker that never replies times out as still working"); + ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer + + MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300); + assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), + "the session lock must be released after a TIMED_OUT_WORKING answer(), or this bounded " + + "follow-up send would come back BUSY instead of timing out on its own work"); + } + + /** + * The {@code ExecutionException} exit. Forces it directly on answer()'s own reopened forward + * waiter — {@code rendezvous.currentWaiter(T)} is the exact {@code CompletableFuture} its + * {@code reply.get(...)} is blocked on — rather than trying to make a real worker fail, since the + * failure mode under test is in answer()'s own wait, not in how it got triggered. + */ + @Test + void answerReleasesTheSessionLockWhenTheReplyFutureFailsExceptionally() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.Reply q = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, q.outcome()); + String turnId = q.turnId(); + + java.util.concurrent.atomic.AtomicReference caught = + new java.util.concurrent.atomic.AtomicReference<>(); + Thread answerer = new Thread(() -> { + try { + messages.answer(turnId, "config.yaml", 5000); + caught.set(new AssertionError("expected answer() to throw")); + } catch (Throwable t) { + caught.set(t); + } + }); + answerer.start(); + awaitWaiting(); // answer() re-opened its forward waiter and is about to block on it + + CompletableFuture waiter = rendezvous.currentWaiter(T); + assertNotNull(waiter, "answer() should have registered a forward waiter for the worker session"); + waiter.completeExceptionally(new RuntimeException("boom")); + + answerer.join(5000); + assertFalse(answerer.isAlive(), "answer() should have left by throwing once its reply future failed"); + Throwable thrown = caught.get(); + assertNotNull(thrown, "answer() should have thrown"); + assertTrue(thrown instanceof RuntimeException, "a RuntimeException cause is rethrown as-is: " + thrown); + assertEquals("boom", thrown.getMessage()); + + ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer + + MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300); + assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), + "the session lock must be released when answer()'s reply future fails exceptionally, " + + "or this bounded follow-up send would come back BUSY instead of timing out on its own work"); + } + + /** + * The {@code InterruptedException} exit. Same shape as the {@code ExecutionException} case above, + * but the thread blocked in {@code reply.get(...)} is interrupted instead of the future failing. + */ + @Test + void answerReleasesTheSessionLockWhenInterrupted() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); + MessageService.Reply q = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, q.outcome()); + String turnId = q.turnId(); + + java.util.concurrent.atomic.AtomicReference caught = + new java.util.concurrent.atomic.AtomicReference<>(); + Thread answerer = new Thread(() -> { + try { + messages.answer(turnId, "config.yaml", 5000); + caught.set(new AssertionError("expected answer() to throw")); + } catch (Throwable t) { + caught.set(t); + } + }); + answerer.start(); + awaitWaiting(); // answer() re-opened its forward waiter and is about to block on it + + answerer.interrupt(); + answerer.join(5000); + assertFalse(answerer.isAlive(), "answer() should have left by throwing once its thread was interrupted"); + Throwable thrown = caught.get(); + assertNotNull(thrown, "answer() should have thrown"); + assertTrue(thrown instanceof IllegalStateException, + "an interrupted wait is wrapped in IllegalStateException: " + thrown); + + ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer + + MessageService.Reply probe = messages.send(T, "probe after interruption", 300); + assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(), + "the session lock must be released when answer()'s wait is interrupted, or this bounded " + + "follow-up send would come back BUSY instead of timing out on its own work"); + } + @Test void pollReturnsNullForAnUnknownTicket() { assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown"); -- 2.52.0