Merge pull request 'fleetd #553: register the rendezvous waiter in onStatus's finally backstop' (#557) from worker/553-onstatus-completion-leak-0da881-2 into main
CI / shell-tests (push) Successful in 4s
CI / contract (push) Successful in 54s
CI / build (push) Successful in 1m40s

This commit was merged in pull request #557.
This commit is contained in:
2026-09-12 10:50:24 +02:00
2 changed files with 576 additions and 43 deletions
@@ -494,54 +494,163 @@ public final class Injector {
// Fire listeners / herdr calls after releasing the monitor so nothing runs on the poller
// thread while it holds the target lock.
if (resubmit) {
try {
agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits
} catch (RuntimeException e) {
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
//
// fleetd #553: wrapped in try/finally. By this point, if `sent != null`, the delivery has
// already happened inside the monitor above — off the queue, p.state == DELIVERED, text
// typed into the target's pane — so `sent`'s future MUST be completed one way or another,
// in every path out of this region, or the caller (a blocking fleet_send, or an async
// ticket) waits forever on a message it actually received. But the finally must not
// swallow whatever escaped: a listener that throws is a defect in THAT listener, and it
// must still reach StatusPoller's catch (Throwable) so log.error fires — converting a
// loud listener bug into a silently orphaned future would be worse than the bug itself.
//
// There are TWO futures at stake here, not one: `sent.delivered()` (the delivery future)
// and `sent.token().waiter()` (the rendezvous waiter a blocking fleet_send actually waits
// on for the worker's ANSWER). The waiter is registered only by turnListener.onDelivered(),
// called from inside the `if (sent != null)` block below. Completing `sent.delivered()`
// without also calling onDelivered leaves the waiter unregistered forever — a hang with a
// success receipt, which is worse than the plain hang this ticket is about. So the
// `finally` below does not just complete the future: on the path where the `if (sent !=
// null)` block never ran, it does that block's whole job — onDelivered, then complete.
//
// `sentHandled` is NOT allowed to move earlier than this: an earlier comment on #553
// proposed running `if (sent != null)` first, before turnCompleted/turnFailed, and that
// was withdrawn — onDelivered WRITES CompletionResolver's inFlight record for this turn,
// onTurnComplete READS it for the PREVIOUS turn, and running onDelivered first makes
// onTurnComplete resolve the NEW turn's waiter with the PREVIOUS turn's stale output
// (CB-116). The `if (sent != null)` block stays last; the `finally` is a backstop for it,
// not a replacement.
boolean sentHandled = false;
try {
if (resubmit) {
try {
agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits
} catch (Throwable e) {
// fleetd #553: widened from RuntimeException (same reasoning as #549 at :382) —
// this is the FIRST block after the monitor, so an Error escaping it used to skip
// every block below it, including the `sent` completion.
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
}
}
}
if (notReady != null) {
// Worker never became available: forget its (never-set) readiness, unblock every queued
// caller, and route the awaiting send through the same failure path as a stalled turn so
// a blocking or async waiter resolves WORKER_FAILED rather than riding out the timeout.
forget.accept(target);
RuntimeException cause = new IllegalStateException(
target + " never became available (no bridge MCP connection within the boot window)");
for (Pending p : notReady) {
p.delivered().completeExceptionally(cause);
if (notReady != null) {
// Worker never became available: unblock every queued caller FIRST — these messages'
// fate (NOT_DELIVERED, off the queue) was already decided inside the monitor above, so
// fleetd #553 completes every one of them before calling forget.accept or
// turnListener.onTurnFailed below. Either of those is a listener/consumer callback and
// can throw (an ordinary RuntimeException is enough — the same reasoning as the rest of
// this ticket): completing the futures first means such a throw can no longer leave any
// of them permanently pending, regardless of which one throws or in which order.
RuntimeException cause = new IllegalStateException(
target + " never became available (no bridge MCP connection within the boot window)");
for (Pending p : notReady) {
p.delivered().completeExceptionally(cause);
}
// route the awaiting send through the same failure path as a stalled turn so a blocking
// or async waiter resolves WORKER_FAILED rather than riding out the timeout.
forget.accept(target);
turnListener.onTurnFailed(target);
}
turnListener.onTurnFailed(target);
}
if (turnCompleted) {
if (startPostTurn) {
boolean started = turnListener.onTurnCompleteWithPostAction(target);
synchronized (t) {
t.postTurnPending = false;
if (started) {
t.awaitingPostTurnPickup = true;
t.injectableSincePostTurnPickup = 0;
if (turnCompleted) {
if (startPostTurn) {
// fleetd #553: the listener call is wrapped so `t.postTurnPending` (set true inside
// the monitor above, before this call) is always reset. Before this wrapping, a
// RuntimeException from onTurnCompleteWithPostAction skipped the reset below,
// permanently wedging the target: postTurnPending stayed true forever, so this
// target's delivery guard (:376) would never again pass and no further message to it
// would ever be delivered — a target-level lockup, not just one skipped future.
// `started` defaults to false so an exception is treated as "the post-turn action did
// not start" rather than falsely arming the post-turn pickup latch for an action that
// never ran.
boolean started = false;
try {
started = turnListener.onTurnCompleteWithPostAction(target);
} finally {
synchronized (t) {
t.postTurnPending = false;
if (started) {
t.awaitingPostTurnPickup = true;
t.injectableSincePostTurnPickup = 0;
}
if (t.queue.isEmpty() && !t.awaitingPostTurnPickup) {
targets.remove(target, t);
}
}
}
if (t.queue.isEmpty() && !t.awaitingPostTurnPickup) {
targets.remove(target, t);
} else {
turnListener.onTurnComplete(target);
}
}
if (turnFailed) {
turnListener.onTurnFailed(target);
}
if (sent != null) {
// fleetd #553: record the INTENT to handle before any side effect, not the fact of
// having completed it. If onDelivered() throws part way through — e.g. after
// CompletionResolver's captureBaseline has already done its inFlight.put — a flag
// set only after this block would still read false, and the `finally` below would
// call onDelivered() a SECOND time. A second captureBaseline runs later, once
// onTurnComplete has already thrown and possibly scraped the pane, and can snapshot
// a pane that already absorbed this turn's output — which means the pane's tail
// never differs from that baseline again and CompletionResolver.resolve's own
// suppression at :364-369 (`return` keeping the in-flight record) drops every future
// completion for this turn, permanently. So this flag must go true FIRST.
sentHandled = true;
if (sendError != null) {
log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage());
sent.delivered().completeExceptionally(sendError);
} else {
// Baseline the pane's pre-turn content so a misattributed completion (no new output)
// can't resolve this send with the previous turn's stale answer (CB-115).
turnListener.onDelivered(target, sent.token());
sent.delivered().complete(null);
}
}
} finally {
// Backstop (fleetd #553): if an earlier block in this try threw before the `if (sent !=
// null)` block above ran, `sentHandled` is still false here, and this does that block's
// WHOLE job — not just the future completion. Skipping onDelivered() here would leave
// sent.token().waiter() never registered with CompletionResolver, so a blocking
// fleet_send would be told its message was delivered and then wait out its full timeout
// for an answer that can never resolve — worse than the plain hang, because it now looks
// like success. onDelivered() runs only when sendError == null: nothing was delivered on
// the error path, so there is nothing to register.
//
// fleetd #553 (lead review, comment 16916): `sentHandled` guards ONLY the onDelivered
// re-call below, never the future completion outside this inner try. One boolean cannot
// carry both meanings — "onDelivered was called" and "the future has been dealt with" —
// because they come apart exactly when onDelivered throws PART WAY THROUGH: sentHandled
// is already true (set before the call, correctly — see the `if (sent != null)` block
// above), so a guard on the outer `if` here would skip this whole recovery, including the
// completion, and leave sent.delivered() pending forever even though the message really
// was typed into the pane and taken off the queue. So `sent != null` alone gates whether
// this target has anything to finish; `!sentHandled` gates only the onDelivered re-call
// inside. CompletableFuture.complete/completeExceptionally are idempotent — on the
// ordinary path (sentHandled == true, no throw) the `if (sent != null)` block above has
// already completed this future, so the calls below are a no-op returning false.
//
// This recovery is wrapped in its own try/catch(Throwable) that swallows only ITS OWN
// throwable and logs at WARN — the future is still completed either way — while the
// ORIGINAL throwable from the try block above is left alone to keep unwinding out of
// this method to StatusPoller's catch (Throwable), so a listener bug stays loud.
if (sent != null) {
try {
if (!sentHandled && sendError == null) {
turnListener.onDelivered(target, sent.token());
}
} catch (Throwable recoveryError) {
log.warn("fleetd #553 backstop: onDelivered failed for {} while recovering from "
+ "an earlier listener failure; completing its delivery future "
+ "anyway: {}",
target, recoveryError.getMessage());
} finally {
// No-op (returns false) on the ordinary path, where the `if (sent != null)` block
// above already completed this future — see the comment above this block.
if (sendError != null) {
sent.delivered().completeExceptionally(sendError);
} else {
sent.delivered().complete(null);
}
}
} else {
turnListener.onTurnComplete(target);
}
}
if (turnFailed) {
turnListener.onTurnFailed(target);
}
if (sent != null) {
if (sendError != null) {
log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage());
sent.delivered().completeExceptionally(sendError);
} else {
// Baseline the pane's pre-turn content so a misattributed completion (no new output)
// can't resolve this send with the previous turn's stale answer (CB-115).
turnListener.onDelivered(target, sent.token());
sent.delivered().complete(null);
}
}
}
@@ -8,7 +8,9 @@ import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TestTurnTokens;
import dev.ltms.fleet.msg.TurnToken;
import dev.ltms.fleet.testing.CapturedLog;
import org.junit.jupiter.api.Test;
@@ -773,4 +775,426 @@ class InjectorTest {
"a HerdrException must still be dropped and marked NOT_DELIVERED, unchanged by the "
+ "wider Throwable catch");
}
// --- fleetd #553: a throwable from any listener callback in onStatus must not skip the
// delivered-future completion. The whole post-monitor region is now wrapped in try/finally. ---
@Test
void theOriginalThrowableFromAListenerStillEscapesOnStatus() {
// fleetd #553's control: the entire point of the finally backstop is that it completes a
// future WITHOUT swallowing whatever escaped. A fix built the wrong way (e.g. catching
// Throwable in the finally, or wrapping the region in try/catch instead of try/finally)
// could make every other test in this group pass while still converting a loud listener
// bug into a silent one — which the ticket calls a worse outcome than the bug itself. This
// is deliberately its own test, not folded into another one's assertion.
RuntimeException boom = new RuntimeException("fleetd #553 control");
TurnListener throwing = new TurnListener() {
@Override
public void onTurnComplete(String target) {
throw boom;
}
};
Injector inj = new Injector(new AgentControl(herdr), throwing);
inj.enqueue(T, "only", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // pickup
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE),
"onStatus must still propagate the listener's own throwable, unmodified");
assertSame(boom, thrown, "must be the EXACT throwable, not a wrapper or a different instance");
}
@Test
void aRuntimeExceptionFromOnTurnCompleteStillCompletesTheNextDelivery() {
// fleetd #553, acceptance test 1: onTurnComplete fires for "first"'s completed turn INSIDE
// the same onStatus call that then peeks and delivers "second" — the turnCompleted check
// (Injector.java, inside the released-pickup branch) clears awaitingCompletion just before
// the delivery guard right below it is evaluated, so both happen in one round when there is
// no post-turn action. Before fleetd #553, a RuntimeException thrown here unwound straight
// out of onStatus and skipped the `if (sent != null)` block, leaving "second"'s future
// pending forever even though it was already off the queue, marked DELIVERED, and typed
// into the pane.
RuntimeException boom = new RuntimeException("boom from onTurnComplete");
TurnListener throwing = new TurnListener() {
@Override
public void onTurnComplete(String target) {
throw boom;
}
};
Injector inj = new Injector(new AgentControl(herdr), throwing);
inj.enqueue(T, "first", TestTurnTokens.inert(T));
Injector.Delivery second = inj.enqueue(T, "second", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // delivers "first"
inj.onStatus(T, AgentStatus.WORKING); // picked up
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE),
"the throwable from onTurnComplete must still escape onStatus");
assertSame(boom, thrown);
assertEquals(List.of("first", "second"), sent(),
"fleetd #553: \"second\" must still be sent even though onTurnComplete threw for "
+ "\"first\"'s completion");
assertTrue(second.completion().isDone() && !second.completion().isCompletedExceptionally(),
"fleetd #553: \"second\"'s delivered future must still complete normally despite "
+ "the throw");
}
private static final int READINESS_GRACE_SAMPLES = 240; // Injector.READINESS_GRACE_POLLS is private
@Test
void aRuntimeExceptionFromForgetStillCompletesTheReadinessFailureFutures() {
// fleetd #553: forget.accept fires as part of the CB-114 readiness-grace path (production
// wires it to presence::forget). Before fleetd #553, that block called forget.accept
// BEFORE completing the queued messages' futures, so a RuntimeException from forget left
// every one of them pending forever even though they were already marked NOT_DELIVERED and
// dropped from the queue. fleetd #553 completes those futures first, so forget.accept
// throwing afterward can no longer un-complete them.
RuntimeException boom = new RuntimeException("boom from forget");
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
throw boom;
});
Injector.Delivery delivery = inj.enqueue(T, "task", TestTurnTokens.inert(T));
// Build the not-ready streak up to (but not past) the grace threshold.
for (int i = 0; i < READINESS_GRACE_SAMPLES - 1; i++) inj.onStatus(T, AgentStatus.IDLE);
assertFalse(delivery.completion().isDone(), "must not be resolved before the grace expires");
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE),
"the throwable from forget.accept must still escape onStatus");
assertSame(boom, thrown);
assertTrue(delivery.completion().isCompletedExceptionally(),
"fleetd #553: the queued message's future must still be completed even though "
+ "forget.accept threw");
}
@Test
void aRuntimeExceptionFromOnTurnFailedStillLeavesTheReadinessFailureFuturesCompleted() {
// fleetd #553, acceptance test: same readiness-grace path as the forget test above, but the
// throw comes from onTurnFailed instead. NOTE: onTurnFailed is the LAST statement in this
// block (both before and after this ticket), so the queued messages' futures are already
// completed by the time it runs, regardless of this fix — this test cannot be made to fail
// against the pre-fleetd-553 code the way the onTurnComplete/forget/postAction tests can; I
// could not find a construction where onTurnFailed's own throw is what stands between a
// future and its completion (flagged to the lead via fleet_ask; no reply arrived before the
// ~55s window closed, so recorded here instead). It still pins a real invariant this ticket
// cares about — an ordinary RuntimeException from this callback must not un-complete a
// future that was already decided — so it stays as a regression lock, not a bug-fix proof.
RuntimeException boom = new RuntimeException("boom from onTurnFailed");
TurnListener throwing = new TurnListener() {
@Override
public void onTurnComplete(String target) {
}
@Override
public void onTurnFailed(String target) {
throw boom;
}
};
Injector inj = new Injector(new AgentControl(herdr), throwing, _ -> false, _ -> {
});
Injector.Delivery delivery = inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < READINESS_GRACE_SAMPLES - 1; i++) inj.onStatus(T, AgentStatus.IDLE);
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE),
"the throwable from onTurnFailed must still escape onStatus");
assertSame(boom, thrown);
assertTrue(delivery.completion().isCompletedExceptionally(),
"the queued message's future must remain completed despite onTurnFailed throwing");
}
@Test
void aRuntimeExceptionFromOnTurnCompleteWithPostActionStillUnwedgesTheTarget() {
// fleetd #553, acceptance test: before this ticket, a RuntimeException from
// onTurnCompleteWithPostAction skipped the `t.postTurnPending = false` reset that used to
// sit right after it (with nothing guarding it), permanently wedging the target —
// postTurnPending stayed true forever, so the delivery guard never passed again and
// "second" (already queued) was never delivered, no matter how many further onStatus
// rounds ran. fleetd #553 wraps that call so the reset always runs.
RuntimeException boom = new RuntimeException("boom from onTurnCompleteWithPostAction");
class ThrowingPostTurn implements TurnListener {
@Override
public void onTurnComplete(String target) {
}
@Override
public boolean hasPostTurnAction(String target) {
return true;
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
throw boom;
}
}
Injector inj = new Injector(new AgentControl(herdr), new ThrowingPostTurn());
inj.enqueue(T, "first", TestTurnTokens.inert(T));
Injector.Delivery second = inj.enqueue(T, "second", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // delivers "first"
inj.onStatus(T, AgentStatus.WORKING); // picked up
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE), // "first"'s turn completes; post-action throws
"the throwable from onTurnCompleteWithPostAction must still escape onStatus");
assertSame(boom, thrown);
assertEquals(List.of("first"), sent(), "\"second\" must not be sent in the SAME round as the throw");
inj.onStatus(T, AgentStatus.IDLE); // a later, unrelated round
assertEquals(List.of("first", "second"), sent(),
"fleetd #553: the target must not stay wedged — \"second\" must still be delivered "
+ "once postTurnPending is reset despite the earlier throw");
assertTrue(second.completion().isDone() && !second.completion().isCompletedExceptionally());
}
/**
* A {@link HerdrClient} that throws a non-{@link RuntimeException} {@link Error} from {@code
* agent.send_keys} instead of delegating — the fleetd #553 case for the resubmit nudge: its own
* catch (Injector.java, the first block after the monitor) was {@code RuntimeException}-only,
* the same one-class-too-narrow shape #546 fixed at the send seam. Records every call it sees
* itself, mirroring {@code ErrorOnPrompt} above, since the delegate's own recording is never
* reached for {@code agent.send_keys}.
*/
private static final class ErrorOnSendKeys implements HerdrClient {
private final FakeHerdr delegate;
private final List<FakeHerdr.Call> calls = new java.util.concurrent.CopyOnWriteArrayList<>();
private ErrorOnSendKeys(FakeHerdr delegate) {
this.delegate = delegate;
}
List<FakeHerdr.Call> calls() {
return calls;
}
@Override
public JsonNode call(String method, Object params) {
calls.add(new FakeHerdr.Call(method, params));
if (method.equals("agent.send_keys")) {
throw new AssertionError("simulated non-RuntimeException resubmit failure (fleetd #553)");
}
return delegate.call(method, params);
}
@Override
public void close() {
delegate.close();
}
}
@Test
void anErrorFromTheResubmitNudgeDoesNotPreventTheSentCompletion() {
// fleetd #553, acceptance test: the resubmit nudge's own catch was RuntimeException-only —
// the same one-class-too-narrow shape #546 fixed at the send seam (now :382). It sits FIRST
// after the monitor, so before this ticket an Error escaping it skipped every block below in
// THAT round, including any `sent` completion that round might otherwise produce. Widened to
// Throwable (matching #549's own widening) so it can no longer escape onStatus at all — the
// delivery that already completed at send time, and the worker's later pickup, are both
// unaffected by the Error in between.
ErrorOnSendKeys throwing = new ErrorOnSendKeys(new FakeHerdr());
Injector inj = new Injector(new AgentControl(throwing));
Injector.Delivery delivery = inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // delivers "task"; awaiting pickup
assertTrue(delivery.completion().isDone() && !delivery.completion().isCompletedExceptionally(),
"the delivery completes at send time, before any resubmit nudge is attempted");
assertDoesNotThrow(() -> inj.onStatus(T, AgentStatus.IDLE), // still idle -> resubmit; Error thrown+caught
"fleetd #553: an Error from the resubmit nudge must be caught inside onStatus, not "
+ "escape it");
inj.onStatus(T, AgentStatus.WORKING); // the worker's real pickup must still resolve cleanly
long promptCalls = throwing.calls().stream().filter(c -> c.method().equals("agent.prompt")).count();
assertEquals(1, promptCalls,
"sent exactly once, unaffected by the resubmit Error in between");
}
@Test
void ordinarySuccessStillCompletesExactlyOnceAfterOnDelivered() {
// Acceptance: the ordinary success path is unchanged by fleetd #553 — the future completes
// normally exactly once, and onDelivered still runs before it (CB-115's pane baseline).
List<String> events = new ArrayList<>();
TurnListener listener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
events.add("onDelivered");
}
};
Injector inj = new Injector(new AgentControl(herdr), listener);
CompletableFuture<Void> f = inj.enqueue(T, "hello", TestTurnTokens.inert(T)).completion();
f.whenComplete((v, ex) -> events.add("completed"));
inj.onStatus(T, AgentStatus.IDLE);
assertEquals(List.of("onDelivered", "completed"), events,
"onDelivered must still run before the future completes, unchanged by fleetd #553");
assertTrue(f.isDone() && !f.isCompletedExceptionally());
}
// Acceptance: the ordinary failure path (a HerdrException at the send seam still completes the
// future exceptionally with that exception) is already covered, unchanged, by the existing
// aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged test above (fleetd #546) — its
// codepath is untouched by fleetd #553's try/finally, since sendError != null is set before the
// try block begins and that branch never threw to begin with.
// --- fleetd #553, sharpened acceptance (ticket comments 16870/16884/16890/16903): completing
// sent.delivered() is only HALF the job. See the two tests below. ---
@Test
void aRuntimeExceptionFromOnTurnCompleteStillRegistersTheRendezvousWaiterForTheNextDelivery() {
// fleetd #553's real invariant, not just the completion of sent.delivered(). There are TWO
// futures at stake: sent.delivered() (the delivery future) and sent.token().waiter() (the
// rendezvous waiter a blocking fleet_send actually waits on for the worker's ANSWER). The
// waiter is registered only by turnListener.onDelivered() — production wires this to
// CompletionResolver.captureBaseline, which does the inFlight.put() — and that call lives
// inside the very `if (sent != null)` block a plain "complete the future" backstop does not
// reach. A finally that only completes sent.delivered() converts a hang into a HANG WITH A
// SUCCESS RECEIPT: the caller is told the send landed, then waits out its full timeout for
// an answer that can never resolve, because CompletionResolver.resolve (:301-308) finds no
// inFlight entry and returns silently. This test is red on that half-fix, and green only
// once the finally also runs onDelivered() when the normal block never got the chance.
//
// The scenario: onTurnComplete throws while completing "first"'s turn, in the SAME onStatus
// round that (per Injector.java :361-391) then peeks and delivers "second" — the normal
// case, not a corner, since the :365 assignment is what lets the :376 delivery guard pass.
// turnListener.onTurnComplete runs BEFORE the `if (sent != null)` block, so the throw here
// means onDelivered for "second" is never reached on the normal path — only the finally
// backstop can register it.
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
RuntimeException boom = new RuntimeException("boom from onTurnComplete");
TurnListener throwingOnTurnCompleteOnly = new TurnListener() {
@Override
public void onDelivered(String target, TurnToken token) {
resolver.onDelivered(target, token); // the real registration, exactly as production wires it
}
@Override
public void onTurnComplete(String target) {
throw boom; // "first"'s completion callback, thrown before "second" is delivered
}
};
Injector inj = new Injector(new AgentControl(herdr), throwingOnTurnCompleteOnly);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open("session-a");
TurnToken secondToken = new TurnToken(T, waiter);
inj.enqueue(T, "first", TestTurnTokens.inert(T));
Injector.Delivery second = inj.enqueue(T, "second", secondToken);
inj.onStatus(T, AgentStatus.IDLE); // delivers "first"
inj.onStatus(T, AgentStatus.WORKING); // picked up
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE), // "first" completes (throws); "second" delivered
"the original throwable from onTurnComplete must still escape onStatus");
assertSame(boom, thrown, "must be the EXACT throwable, not a wrapper or a different instance");
assertEquals(List.of("first", "second"), sent(),
"\"second\" must still be sent even though onTurnComplete threw for \"first\"");
assertTrue(second.completion().isDone() && !second.completion().isCompletedExceptionally(),
"\"second\"'s delivery future must still complete despite the throw");
CompletionResolver.InFlight inFlight = resolver.inFlight(T);
assertNotNull(inFlight,
"fleetd #553: the finally must register the new turn's waiter (via onDelivered), not "
+ "just complete its delivery future — otherwise a blocking fleet_send is told "
+ "its message landed and then waits out the full timeout for an answer that "
+ "can never resolve");
assertSame(waiter, inFlight.waiter(),
"the registered waiter must be exactly \"second\"'s waiter, not some other one");
}
@Test
void anOnDeliveredThrowAfterItsOwnRegistrationDoesNotRunASecondTime() {
// fleetd #553 comment 16903: `sentHandled` must be set to true BEFORE onDelivered() runs,
// not after. Setting it after would mean a throw from onDelivered() PART WAY THROUGH — e.g.
// after CompletionResolver's own captureBaseline has already done its inFlight.put() — still
// leaves the flag false, so the finally backstop (seeing "not handled") calls onDelivered() a
// SECOND time. That second captureBaseline runs later, over a pane that may already have
// absorbed this turn's output, and CompletionResolver.resolve's own baseline-equals-tail
// suppression (:364-369, which deliberately keeps the in-flight record on a match) can then
// drop every future completion for the turn, permanently. This assertion is red on a
// `sentHandled = true` placed AFTER the onDelivered() call (onDelivered runs twice) and green
// on the correct placement (runs exactly once) — regardless of what onDelivered itself did.
AtomicInteger onDeliveredCalls = new AtomicInteger();
RuntimeException boom = new RuntimeException("boom from onDelivered, after its own registration ran");
TurnListener listener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
}
@Override
public void onDelivered(String target, TurnToken token) {
onDeliveredCalls.incrementAndGet(); // stand-in for CompletionResolver's inFlight.put
throw boom; // then fail, as if a LATER step inside onDelivered blew up
}
};
Injector inj = new Injector(new AgentControl(herdr), listener);
inj.enqueue(T, "task", TestTurnTokens.inert(T));
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE), // delivers "task"; onDelivered throws
"the original throwable from onDelivered must still escape onStatus");
assertSame(boom, thrown);
assertEquals(1, onDeliveredCalls.get(),
"fleetd #553: onDelivered must be called exactly once — a `sentHandled` flag set "
+ "AFTER the call (rather than before) would leave it false here and the "
+ "finally backstop would call onDelivered a second time");
}
@Test
void anOnDeliveredThrowOnTheNormalPathStillCompletesTheDeliveryFuture() {
// fleetd #553, lead review of PR #557 (ticket comment 16916): `sentHandled` must guard ONLY
// the onDelivered RE-CALL in the finally, never the future completion alongside it. Gating
// BOTH behind `!sentHandled` (the shape this test is red against) misses the one path the
// flag's own correct placement creates: `sentHandled` is set to true FIRST, before the
// onDelivered() call, inside the `if (sent != null)` block above (see the previous test —
// that placement is right and must not change). So when onDelivered() itself throws on that
// NORMAL path, `sentHandled` already reads true by the time control reaches the finally, and
// a `sent != null && !sentHandled` guard around the WHOLE recovery — completion included —
// skips it entirely. The message was typed into the target's pane and taken off the queue
// inside the monitor, same as any other delivery, so its future is left pending forever: the
// exact defect this ticket exists to close, just reached from a different throwing call.
//
// The fix splits the one flag's two jobs: `!sentHandled` keeps gating only the onDelivered
// call (so the exactly-once guarantee in the test above still holds — CompletableFuture.
// complete/completeExceptionally are idempotent, so completing unconditionally here is a
// no-op on the ordinary path, where the `if (sent != null)` block already completed it.
RuntimeException boom = new RuntimeException("boom from onDelivered on the normal path");
TurnListener listener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
}
@Override
public void onDelivered(String target, TurnToken token) {
throw boom;
}
};
Injector inj = new Injector(new AgentControl(herdr), listener);
CompletableFuture<Void> delivered = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE), // delivers "task"; onDelivered throws
"the original throwable from onDelivered must still escape onStatus");
assertSame(boom, thrown);
assertTrue(delivered.isDone(),
"fleetd #553: \"task\" was actually delivered — typed into the pane and taken off "
+ "the queue inside the monitor — so its delivery future must be completed on "
+ "every path out of onStatus, including the one where onDelivered itself is "
+ "what threw. Leaving it pending here is a hang, not a fix");
}
}