diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 4bd2bfa..997c12f 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -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 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. + * + *

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. + * + *

The invariant this composition guarantees: a throwing session-listener must not + * prevent the completion resolver from being told the turn ended. 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. + * + *

{@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). + * + *

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 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 diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdTurnListenerCompositionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdTurnListenerCompositionTest.java new file mode 100644 index 0000000..207aee3 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdTurnListenerCompositionTest.java @@ -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 after 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}). + * + *

The invariant under test: a throwing session-listener must not prevent the completion + * resolver from being told the turn ended. 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 called = new LinkedHashSet<>(); + private final Set 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()); + } +}