fleetd #572: pin answer()'s session-lock release across all four exits #574

Merged
ltms merged 1 commits from worker/572-answer-lock-release-46a9ae-5 into main 2026-09-12 13:30:23 +02:00
@@ -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<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> 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<MessageService.Reply> 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.
*
* <p>{@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<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> 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<MessageService.Reply> 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<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> 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<Throwable> 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<Rendezvous.Resolution> 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<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> 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<Throwable> 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");