Merge #561: the completion/session listener fan-out survives either half throwing, and both sites are pinned
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 50s
CI / build (push) Successful in 2m7s

fleetd #561. Fleetd composed two TurnListener halves as `completion.X(); sessions.X();`, so a throw
from the first half skipped the second. The anonymous class is now a package-private factory,
Fleetd.turnListener(completion, sessions), built on two helpers that always attempt both halves and
rethrow whatever escaped — a second failure attached with addSuppressed rather than dropped, so it
still reaches StatusPoller's catch (Throwable).

The wiring at Fleetd.java:498 calls that factory, so the seam under test is the real caller.

Lead verification, on a merged tree, re-running the checks rather than accepting the worker's:
exit 0, 1761 tests from Maven and from an independent sum over 131 surefire reports.

The interesting part is what the first round MISSED. Two helpers maintain one invariant — "the
second half always runs" — and the first round's five tests asserted it at only one site. Measured:

  bothMustRun                      reverted to the pre-fix bug -> 1 named failure   (pinned)
  bothMustRunKeepingSecondResult   the SAME bug                -> 1755/1755 GREEN   (unpinned)
  failure.addSuppressed(t) deleted                             -> 1755/1755 GREEN   (unpinned)

Both survivors are now killed by new tests, re-verified by the lead after the fix:
sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction and
bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions, each failing alone under its own mutation.

The rule this cost us, worth carrying: COUNT ASSERTIONS PER SITE, NOT PER INVARIANT. The total being
non-zero is what hides a zero at one site, and extracting a shared helper makes it worse rather than
better — it does not reduce the number of sites, only how many are visible. Credit to the fleet01
lead, who predicted this shape before an instance was found.

onDelivered stays deliberately unguarded. Its comment now gives the real reason —
CompletionResolver.captureBaseline already catches RuntimeException around its scrape and fails open,
so that half does not realistically throw — instead of the previous reason, which was true but about
registration rather than about this pair. A correct conclusion resting on a wrong premise reads
exactly like a verified one.
This commit was merged in pull request #570.
This commit is contained in:
2026-09-12 13:20:56 +02:00
2 changed files with 442 additions and 36 deletions
+156 -36
View File
@@ -79,6 +79,7 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BooleanSupplier;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Predicate;
@@ -491,42 +492,10 @@ public final class Fleetd {
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
completion.onTurnComplete(target);
sessions.onTurnComplete(target);
}
@Override
public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target);
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
completion.resolveBeforePostAction(target);
return sessions.onTurnCompleteWithPostAction(target);
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
public void onTurnFailed(String target) {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
}
};
// fleetd #561: extracted to a static factory (see turnListener below) — the anonymous class
// this replaced had two bare, unguarded statements per callback, and nothing enforced that
// the completion half went first beyond call order in the source.
TurnListener turnListener = turnListener(completion, sessions);
Predicate<String> deliverable = deliverableTo(presence, leads);
// fleetd #556: registration is wired directly to `completion`, not folded into the
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
@@ -1096,6 +1065,157 @@ public final class Fleetd {
.orElse(null);
}
/**
* fleetd #561: compose the production {@link TurnListener} from its two halves — the
* completion resolver (which resolves a blocked {@code fleet_send}'s waiter) and the session
* manager (which drives the member's lifecycle state) — hardened so a throw from either
* half's callback can never suppress the other half's callback for the same event.
*
* <p>Before this, each method below was two bare, unguarded statements: whichever ran first
* throwing meant the second one never ran at all, and nothing beyond call order in the
* source enforced "the completion half goes first". #556 fixed the identical shape for {@code
* onDelivered}'s registration by moving it off the fan-out entirely (see {@code
* TurnRegistrar}); this fixes the four remaining callbacks — {@code onTurnComplete}, {@code
* onTurnCompleteWithPostAction}, and both {@code onTurnFailed} overloads — by hardening the
* fan-out itself instead, since none of them can be pulled out of the listener the way
* registration was.
*
* <p>The invariant this composition guarantees: <b>a throwing session-listener must not
* prevent the completion resolver from being told the turn ended.</b> The completion half is
* always attempted first, and — as a bonus the resolver does not depend on — its own throw
* does not stop the session half from running either. Whatever escapes (from one half or
* both) is rethrown once both have been attempted, with a second failure recorded via {@link
* Throwable#addSuppressed} on the first, so it still reaches {@link StatusPoller}'s {@code
* catch (Throwable)} and logs at ERROR. Nothing here swallows a failure to make the two halves
* look safe.
*
* <p>{@code onTurnCompleteWithPostAction} is the one place order is not just fault-tolerance
* but a functional requirement: {@code completion.resolveBeforePostAction} must resolve the
* scrape before {@code sessions.onTurnCompleteWithPostAction}'s adapter housekeeping can erase
* the pane's rendered output (see {@link CompletionResolver#resolveBeforePostAction}). That is
* why this composition does not treat the pair symmetrically the way {@link #bothMustRun}
* does for the other three callbacks: the mirror case (completion half throws, session half's
* return value still observed) is not preserved here — once the completion half's failure
* escapes, the session half's return value is discarded, matching how the {@link Injector}
* already treats any throw from this callback as "the action did not start" (see the {@code
* started} default at its call site, {@code Injector.java} ~line 616).
*
* <p>Package-private so {@code FleetdTurnListenerCompositionTest} can build this listener
* directly from a real {@link CompletionResolver} and a fake {@link TurnListener} standing in
* for {@code sessions}, without booting the rest of {@code main}.
*/
static TurnListener turnListener(CompletionResolver completion, TurnListener sessions) {
return new TurnListener() {
@Override
public void onTurnComplete(String target) {
bothMustRun(() -> completion.onTurnComplete(target), () -> sessions.onTurnComplete(target));
}
@Override
public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target);
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
return bothMustRunKeepingSecondResult(() -> completion.resolveBeforePostAction(target),
() -> sessions.onTurnCompleteWithPostAction(target));
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
// fleetd #561: left unguarded on purpose, not because registration survives a
// throw elsewhere. completion.onDelivered runs CompletionResolver.captureBaseline,
// which already wraps its scrape read in its own catch (RuntimeException) and
// fails open (baseline = null) — so this call does not realistically throw, and
// there is nothing here for bothMustRun to protect.
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
public void onTurnFailed(String target) {
bothMustRun(() -> completion.onTurnFailed(target), () -> sessions.onTurnFailed(target));
}
@Override
public void onTurnFailed(String target, String reason) {
bothMustRun(() -> completion.onTurnFailed(target, reason), () -> sessions.onTurnFailed(target));
}
};
}
/**
* fleetd #561: run two listener-callback halves for one lifecycle event, guaranteeing the
* SECOND always runs even when the FIRST throws. Whatever is thrown is rethrown once both
* halves have been attempted — a second failure is attached to the first via {@link
* Throwable#addSuppressed} rather than dropped. Never swallows.
*/
private static void bothMustRun(Runnable completionHalf, Runnable sessionsHalf) {
Throwable failure = null;
try {
completionHalf.run();
} catch (Throwable t) {
failure = t;
}
try {
sessionsHalf.run();
} catch (Throwable t) {
if (failure == null) {
failure = t;
} else {
failure.addSuppressed(t);
}
}
if (failure != null) {
throwUnchecked(failure);
}
}
/**
* fleetd #561: like {@link #bothMustRun}, but for {@code onTurnCompleteWithPostAction}, whose
* session half returns the value the {@link Injector} needs. The second (session) half's
* result is what this method returns; if the first (completion) half throws, the second half
* still runs and its result is still computed here, but the throw is rethrown afterward
* regardless — so that result is discarded at the {@link Injector} call site exactly as it
* already is today when this callback throws (see {@link #turnListener}'s javadoc).
*/
private static boolean bothMustRunKeepingSecondResult(Runnable completionHalf,
BooleanSupplier sessionsHalf) {
Throwable failure = null;
try {
completionHalf.run();
} catch (Throwable t) {
failure = t;
}
boolean result = false;
try {
result = sessionsHalf.getAsBoolean();
} catch (Throwable t) {
if (failure == null) {
failure = t;
} else {
failure.addSuppressed(t);
}
}
if (failure != null) {
throwUnchecked(failure);
}
return result;
}
/**
* fleetd #561: rethrow a captured {@link Throwable} without a checked-exception wrapper. The
* two callback halves above never declare a checked exception (both existing production
* halves — {@code CompletionResolver} and {@code SessionManager} — only ever throw unchecked),
* so this only ever actually rethrows a {@link RuntimeException} or {@link Error}; the generic
* cast is the standard "sneaky throw" idiom, not a claim that a checked exception is expected.
*/
@SuppressWarnings("unchecked")
private static <T extends Throwable> void throwUnchecked(Throwable t) throws T {
throw (T) t;
}
/**
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
@@ -0,0 +1,286 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import org.junit.jupiter.api.Test;
import java.util.LinkedHashSet;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #561: {@code Fleetd.turnListener} composes the completion resolver and the session
* manager into one {@link TurnListener}. Each of the four callbacks below used to be two bare,
* unguarded statements (completion first, then sessions) — a throw from the session half used to
* skip nothing <em>after</em> it (there was nothing after it), but nothing enforced that the
* completion half had to come first either, beyond call order in the source. {@code onDelivered}
* had the identical shape and was fixed by #556 (moving its registration off the fan-out
* entirely); these four callbacks cannot be fixed that way, so the fan-out itself is hardened
* instead (see {@link Fleetd#turnListener} and {@code Fleetd.bothMustRun}/{@code
* bothMustRunKeepingSecondResult}).
*
* <p>The invariant under test: <b>a throwing session-listener must not prevent the completion
* resolver from being told the turn ended.</b> Each test below builds the exact production
* composition ({@link Fleetd#turnListener}) from a real {@link CompletionResolver} and a fake
* {@code sessions} half that throws, then asserts the completion half's effect (the captured
* {@link Rendezvous} waiter resolving) happened anyway — never by inspecting call order directly.
*/
class FleetdTurnListenerCompositionTest {
/** Records which callbacks ran and can be told to throw from a chosen one. */
private static final class RecordingSessions implements TurnListener {
final Set<String> called = new LinkedHashSet<>();
private final Set<String> throwing;
RecordingSessions(String... throwingMethods) {
this.throwing = Set.of(throwingMethods);
}
private void maybeThrow(String method) {
called.add(method);
if (throwing.contains(method)) {
throw new IllegalStateException("boom: sessions." + method);
}
}
@Override
public void onTurnComplete(String target) {
maybeThrow("onTurnComplete");
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
maybeThrow("onTurnCompleteWithPostAction");
return true;
}
@Override
public void onTurnFailed(String target) {
maybeThrow("onTurnFailed");
}
@Override
public void onTurnFailed(String target, String reason) {
maybeThrow("onTurnFailedWithReason");
}
}
private static CompletionResolver newResolver(FakeHerdr herdr, Rendezvous rendezvous) {
return new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
}
/**
* A resolver with a controllable clock, and the clock itself, so a test can register a turn
* "delivered" at time 0 and then jump the clock past {@link CompletionResolver#MIN_TURN_NANOS}
* before resolving it — otherwise {@code onTurnComplete}/{@code onTurnCompleteWithPostAction}
* resolve within microseconds of registering in-test, well inside the fleetd#164 floor, and
* get classified as a too-fast crash rather than a real completion. Mirrors the {@code
* LongSupplier} clock-injection pattern fleetd#164's own tests use.
*/
private static CompletionResolver newResolverPastTheFloor(FakeHerdr herdr, Rendezvous rendezvous,
AtomicLong clock) {
LongSupplier nowNanos = clock::get;
return new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), nowNanos);
}
@Test
void onTurnCompleteResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("⏺ real answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
AtomicLong clock = new AtomicLong(0);
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
var waiter = rendezvous.open("term_a");
completion.register("term_a", new TurnToken("term_a", waiter)); // delivered at clock=0
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnComplete("term_a"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnComplete", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must have run");
// onTurnComplete resolves off-thread (a virtual thread) — wait on the waiter itself,
// exactly like MessageServiceTest.completionFallbackIsNeverQueued does.
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind(),
"the completion resolver must still have resolved the send, despite the session "
+ "half throwing");
assertEquals("real answer", resolution.text());
}
@Test
void onTurnFailedResolvesTheWaiterAsFailedEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("⏺ crash context\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
var waiter = rendezvous.open("term_b");
completion.register("term_b", new TurnToken("term_b", waiter));
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnFailed("term_b"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
"the completion resolver must still have failed the send, despite the session "
+ "half throwing");
}
@Test
void onTurnFailedWithReasonResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
var waiter = rendezvous.open("term_c");
completion.register("term_c", new TurnToken("term_c", waiter));
// fleetd #561: the composed listener calls sessions.onTurnFailed(target) — the ONE-arg
// overload — for this reason-carrying event too (matching the pre-existing production
// behaviour: SessionManager never overrides the two-arg overload either), so the
// throwing key here is "onTurnFailed", not a distinct "...WithReason" one.
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnFailed("term_c", "worker unreachable"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
"the completion resolver must still have failed the send, despite the session "
+ "half throwing");
assertEquals("worker unreachable", resolution.text(),
"the explicit reason must still reach the resolved send");
}
@Test
void onTurnCompleteWithPostActionResolvesTheWaiterEvenWhenTheSessionHalfThrows() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ answer before reset\n❯ ");
Rendezvous rendezvous = new Rendezvous();
AtomicLong clock = new AtomicLong(0);
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
var waiter = rendezvous.open("term_d");
completion.register("term_d", new TurnToken("term_d", waiter)); // delivered at clock=0
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
RecordingSessions sessions = new RecordingSessions("onTurnCompleteWithPostAction");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnCompleteWithPostAction("term_d"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnCompleteWithPostAction", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"), "the session half must have run");
// resolveBeforePostAction is synchronous by design (it must run before the context reset
// can erase the pane) — the waiter is already resolved by the time the throw propagates.
assertTrue(waiter.isDone(), "resolveBeforePostAction is synchronous — the send must "
+ "already be resolved once the composed call returns (by throwing)");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("answer before reset", waiter.getNow(null).text());
}
/**
* The mirror case (comment 17041's acceptance item 2): the completion half throws, and the
* session half still ran. The production {@link CompletionResolver} is deliberately defensive
* (its scrape reads are wrapped in {@code catch (RuntimeException)}, by the same fail-open
* design {@link CompletionResolver#captureBaseline} documents) so it essentially never throws
* synchronously in normal operation — forcing it to do so needs a genuinely-real trigger, not
* a fabricated one. {@code ConcurrentHashMap.get(null)} is that trigger: passing a {@code null}
* target makes {@code onTurnComplete}'s {@code inFlight.get(target)} throw a
* {@link NullPointerException} before it ever starts its resolving thread — a real code path,
* not a contrived one.
*/
@Test
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronously() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
TurnListener composed = Fleetd.turnListener(completion, sessions);
assertThrows(NullPointerException.class, () -> composed.onTurnComplete(null),
"ConcurrentHashMap.get(null) inside CompletionResolver.onTurnComplete must still "
+ "escape the composed listener");
assertTrue(sessions.called.contains("onTurnComplete"),
"the session half must still have run even though the completion half threw first");
}
/**
* The mirror case above only exercises {@code onTurnComplete}/{@code bothMustRun}. {@code
* onTurnCompleteWithPostAction} is composed through the OTHER helper,
* {@code bothMustRunKeepingSecondResult}, and nothing previously asserted that its session half
* still runs when its completion half throws — that helper could be reverted to the pre-#561
* broken shape (run the session half only if the completion half did not throw) and the suite
* would still stay green. Uses the same real, non-fabricated trigger as the test above: a
* {@code null} target makes {@code resolveBeforePostAction}'s {@code inFlight.get(target)}
* throw a {@link NullPointerException} before {@code resolve} is ever entered.
*/
@Test
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
TurnListener composed = Fleetd.turnListener(completion, sessions);
assertThrows(NullPointerException.class, () -> composed.onTurnCompleteWithPostAction(null),
"ConcurrentHashMap.get(null) inside CompletionResolver.resolveBeforePostAction must "
+ "still escape the composed listener");
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"),
"the session half must still have run even though the completion half threw first");
}
/**
* {@code bothMustRun} must not just run both halves — it must not DROP a second failure when
* both halves throw. Forces the completion half to throw (the same real {@code
* inFlight.get(null)} NullPointerException trigger used above) while the session half throws a
* distinct {@link IllegalStateException}, and asserts the completion half's throwable is what
* escapes while the session half's throwable survives as a suppressed exception rather than
* being silently discarded.
*/
@Test
void bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
TurnListener composed = Fleetd.turnListener(completion, sessions);
NullPointerException thrown = assertThrows(NullPointerException.class,
() -> composed.onTurnComplete(null),
"the completion half's throw (NPE from inFlight.get(null)) must be what escapes");
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must still have run");
assertEquals(1, thrown.getSuppressed().length,
"the session half's distinct failure must be recorded as suppressed, not dropped");
assertEquals(IllegalStateException.class, thrown.getSuppressed()[0].getClass());
assertEquals("boom: sessions.onTurnComplete", thrown.getSuppressed()[0].getMessage());
}
}