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");