fleetd #572: pin answer()'s session-lock release across all four exits #574
@@ -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");
|
||||
|
||||
Reference in New Issue
Block a user