fleetd #556: make turn registration structural, independent of any TurnListener
CI / shell-tests (pull_request) Successful in 4s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m29s

The Injector owns the invariant "every delivered turn has a registered
waiter," but before this the only thing that satisfied it was
CompletionResolver.captureBaseline, called from inside a TurnListener
callback wired in Fleetd.java. Any TurnListener that throws (from
onDelivered or elsewhere) could break the invariant with no way for the
Injector to detect it. #553 only made the one reachable listener behave
via a try/finally backstop; it did not remove this structural dependency.

Add a narrow TurnRegistrar functional interface, decoupled from
TurnListener, whose only job is registering a delivered turn's waiter.
CompletionResolver now implements it via a new register() method
(extracted from captureBaseline's registration half; captureBaseline
keeps its own full body unchanged, so existing direct callers/tests are
untouched). Injector gets an explicit registrar field/constructor family
(auto-derived from the TurnListener via instanceof where that still
works, explicit where Fleetd's anonymous fan-out listener can't
implement two interfaces at once) and calls registrar.register(...)
directly and unconditionally in both the ordinary delivery path and the
#553 finally backstop, before turnListener.onDelivered(...) — so
registration no longer depends on that notification callback succeeding.
Fleetd.java wires completion::register explicitly as the registrar,
bypassing the turnListener fan-out for registration purposes.

CB-116 ordering (onTurnComplete reads the PREVIOUS turn's inFlight entry
before the new turn's registrar.register() runs) and the two-arg
inFlight.remove(target, turn) vs one-arg distinction on the completion
path are both preserved unchanged.

#561's order-dependent test (asserting Fleetd.java's
completion.onDelivered -> sessions.onDelivered call order) does not
exist anywhere in this repo at this branch point — nothing to delete.

Adds 5 tests: a TurnListener that throws from every callback still
leaves the delivered turn registered and resolvable; the new turn's
registration still runs after the previous turn's completion is read
(CB-116 guard, pinned as a call-order assertion); and three tests naming
the two-arg-remove invariant directly (a superseded turn's terminal
handling must not evict its successor's registration) across resolve()'s
plain-completion branch, its echoed-noReportMessage sub-path, and
fail().
This commit is contained in:
Dai Ha
2026-09-12 16:50:09 +07:00
parent 4f9aba40e7
commit a46e4058ac
6 changed files with 381 additions and 37 deletions
@@ -528,8 +528,11 @@ public final class Fleetd {
}
};
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
// listener) throwing, regardless of call order. See TurnRegistrar's javadoc.
Injector injector = new Injector(router, turnListener, deliverable,
presence::forget);
presence::forget, completion::register);
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start();
@@ -50,7 +50,7 @@ import java.util.regex.Pattern;
* so the scrape's herdr round-trip never stalls the status poller. The captured waiter is read on the
* poller thread (before any next-turn delivery can overwrite it) and passed into the virtual thread.
*/
public final class CompletionResolver implements TurnListener {
public final class CompletionResolver implements TurnListener, TurnRegistrar {
private static final Logger log = LoggerFactory.getLogger(CompletionResolver.class);
@@ -235,6 +235,26 @@ public final class CompletionResolver implements TurnListener {
this.worktreeBranches = Objects.requireNonNull(worktreeBranches, "worktreeBranches");
}
/**
* fleetd #556: {@link TurnRegistrar}'s structural half of what {@link #captureBaseline} used to
* do alone — record the waiter, no I/O. Called directly and unconditionally by the {@link
* Injector} as part of delivering a turn, before {@link #onDelivered} ever runs, so this
* registration cannot be skipped by a {@link TurnListener} throwing (from {@code onDelivered} or
* any other callback). {@link #captureBaseline} still performs this exact check-and-put itself
* as well — harmless and idempotent when it runs right after this — so a caller that only wires
* the {@link TurnListener} path (e.g. an existing test that never mentions {@link TurnRegistrar})
* keeps working unchanged.
*/
@Override
public void register(String target, TurnToken token) {
CompletableFuture<Rendezvous.Resolution> waiter = token.waiter();
if (waiter == null) {
inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later
return;
}
inFlight.put(target, new InFlight(waiter, null, nowNanos.getAsLong(), token.injectedText()));
}
@Override
public void onDelivered(String target, TurnToken token) {
// Capture the exact waiter this turn belongs to (CB-116) and snapshot the pane's pre-turn
@@ -244,7 +264,14 @@ public final class CompletionResolver implements TurnListener {
captureBaseline(target, token);
}
/** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */
/**
* Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link
* #onDelivered}). fleetd #556: this still does its own full waiter check and {@code inFlight.put}
* — the same registration {@link #register} performs — so it keeps working standalone (as every
* test calling it directly already does) even though the {@link Injector} now also calls {@link
* #register} on its own, earlier and unconditionally. The two writes are idempotent with each
* other; only the second (this one) carries a real baseline, since only this one pays for a scrape.
*/
void captureBaseline(String target, TurnToken token) {
CompletableFuture<Rendezvous.Resolution> waiter = token.waiter();
if (waiter == null) {
@@ -11,6 +11,7 @@ import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
@@ -93,6 +94,11 @@ public final class Injector {
private final AgentControl agents;
private final HerdrRouter router;
private final TurnListener turnListener;
/**
* fleetd #556: the Injector's own registration invariant, called directly and unconditionally —
* never folded into {@link #turnListener}'s fan-out. See {@link TurnRegistrar}.
*/
private final TurnRegistrar registrar;
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
/**
@@ -134,15 +140,43 @@ public final class Injector {
*/
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget) {
this(agents, turnListener, ready, forget, System::currentTimeMillis);
// fleetd #556: no explicit registrar named here, so fall back to turnListener itself when it
// happens to also implement TurnRegistrar (true for CompletionResolver, the only production
// TurnListener that needs CB-106 registration) — every existing caller of this overload keeps
// working unchanged. A turnListener that does NOT implement TurnRegistrar (a fan-out object,
// or a bare test lambda) gets TurnRegistrar.NOOP here, same as before this ticket.
this(agents, turnListener, ready, forget,
turnListener instanceof TurnRegistrar r ? r : TurnRegistrar.NOOP);
}
/** Full constructor — for tests: an injectable wall-clock supplier (fleetd #501). */
/**
* fleetd #556: explicit-registrar overload. Use this whenever {@code turnListener} does not
* itself implement {@link TurnRegistrar} — e.g. a fan-out object that also notifies unrelated
* observers — so registration is wired directly to the one component that owns it, rather than
* relying on {@code turnListener} happening to implement both interfaces.
*/
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget, TurnRegistrar registrar) {
this(agents, turnListener, ready, forget, registrar, System::currentTimeMillis);
}
/**
* Package-private constructor — for tests: an injectable wall-clock supplier (fleetd #501),
* with no explicit registrar (fleetd #556 fallback, same as the public 4-arg overload above).
*/
Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget, LongSupplier nowMillis) {
this(agents, turnListener, ready, forget,
turnListener instanceof TurnRegistrar r ? r : TurnRegistrar.NOOP, nowMillis);
}
/** Full constructor — for tests: an explicit registrar and an injectable wall-clock supplier. */
Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget, TurnRegistrar registrar, LongSupplier nowMillis) {
this.agents = agents;
this.router = null;
this.turnListener = turnListener;
this.registrar = Objects.requireNonNull(registrar, "registrar");
this.ready = ready;
this.forget = forget;
this.nowMillis = nowMillis;
@@ -150,15 +184,28 @@ public final class Injector {
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget) {
this(router, turnListener, ready, forget, System::currentTimeMillis);
// fleetd #556: see the AgentControl overload above for the fallback rationale.
this(router, turnListener, ready, forget,
turnListener instanceof TurnRegistrar r ? r : TurnRegistrar.NOOP);
}
/**
* fleetd #556: explicit-registrar overload — production wiring ({@code Fleetd}) uses this, since
* its {@code turnListener} is a fan-out object that also notifies {@code SessionManager} and does
* not itself implement {@link TurnRegistrar}.
*/
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget, TurnRegistrar registrar) {
this(router, turnListener, ready, forget, registrar, System::currentTimeMillis);
}
/** Full constructor — for tests: an injectable wall-clock supplier (fleetd #501). */
Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget, LongSupplier nowMillis) {
Consumer<String> forget, TurnRegistrar registrar, LongSupplier nowMillis) {
this.agents = null;
this.router = router;
this.turnListener = turnListener;
this.registrar = Objects.requireNonNull(registrar, "registrar");
this.ready = ready;
this.forget = forget;
this.nowMillis = nowMillis;
@@ -506,20 +553,25 @@ public final class Injector {
//
// 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.
// on for the worker's ANSWER). fleetd #556: the waiter is registered by `registrar.register`
// now — called directly and unconditionally, structural rather than delegated to a listener
// callback — never by turnListener.onDelivered() alone (before this ticket, that WAS the
// only path, and any TurnListener wired ahead of the one that mattered could throw and skip
// it; #553 could only make the reachable listener behave, not remove that dependency).
// Completing `sent.delivered()` without also calling `registrar.register` leaves the waiter
// unregistered forever — a hang with a success receipt, which is worse than the plain hang
// #553 was 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 —
// register, then onDelivered (notification only, may throw), 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
// was withdrawn — `registrar.register` WRITES CompletionResolver's inFlight record for this
// turn, onTurnComplete READS it for the PREVIOUS turn, and running register 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.
// not a replacement. fleetd #556 keeps this ordering unchanged — it only splits what used
// to be one listener call (register + notify) into two calls in the same position.
boolean sentHandled = false;
try {
if (resubmit) {
@@ -585,13 +637,13 @@ public final class Injector {
}
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
// having completed it. If register()/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 re-run this block's job 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;
@@ -599,32 +651,41 @@ public final class Injector {
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).
// fleetd #556: register the waiter FIRST and unconditionally — this is the
// Injector's own invariant, kept structural rather than delegated. Even if
// turnListener.onDelivered below throws (a bug in an unrelated observer, e.g.
// SessionManager's bookkeeping, or a fan-out wired in a different order), the
// waiter is already registered and resolvable — a throwing TurnListener can no
// longer take this invariant down with it.
registrar.register(target, sent.token());
// Notification only, from here down: 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). Allowed to throw and to fail (the herdr
// scrape inside already fails open with baseline == null on its own).
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.
// Backstop (fleetd #553, extended #556): 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 registrar.register()
// 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. This block 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
// fleetd #553 (lead review, comment 16916): `sentHandled` guards ONLY the register/notify
// 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" —
// carry both meanings — "the turn was handled" 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
// this target has anything to finish; `!sentHandled` gates only the register/notify
// 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.
//
@@ -635,6 +696,11 @@ public final class Injector {
if (sent != null) {
try {
if (!sentHandled && sendError == null) {
// fleetd #556: register first here too, same as the ordinary path above — this
// recovery path is reached only when an EARLIER block (resubmit/turnCompleted/
// turnFailed) threw before the ordinary path ever ran, so nothing has
// registered this turn yet.
registrar.register(target, sent.token());
turnListener.onDelivered(target, sent.token());
}
} catch (Throwable recoveryError) {
@@ -0,0 +1,29 @@
package dev.ltms.fleet.inject;
import dev.ltms.fleet.msg.TurnToken;
/**
* fleetd #556: the {@link Injector}'s own invariant — every delivered turn has a registered
* waiter — kept structural rather than delegated to a {@link TurnListener} callback that can
* throw. The {@link Injector} calls {@link #register} directly and unconditionally as part of
* delivering a turn, before it ever calls {@code turnListener.onDelivered}. A {@link TurnListener}
* that throws from every callback (a bug in an unrelated observer, e.g. session bookkeeping, or a
* fan-out object wired in the wrong order) can therefore never leave a delivered turn unregistered
* — the shape #553 could not close, since that fix could only make the reachable listener behave,
* never remove the Injector's dependency on a listener behaving at all.
*
* <p>Deliberately narrow: registration only, no I/O, nothing that scrapes a pane or can reasonably
* fail. The CB-115 staleness baseline (a herdr round-trip, allowed to fail) stays a
* {@link TurnListener#onDelivered} concern — fired by the {@link Injector} only after this call has
* already run, so its own failure cannot un-register anything.
*/
@FunctionalInterface
public interface TurnRegistrar {
/** Record {@code token}'s waiter as the turn currently in flight for {@code target}. */
void register(String target, TurnToken token);
/** No-op registrar for callers that don't need CB-106 rendezvous tracking. */
TurnRegistrar NOOP = (_, _) -> {
};
}
@@ -19,6 +19,7 @@ import java.util.regex.Pattern;
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.assertTrue;
/** Unit behaviour of the CB-106 completion resolver in isolation from the injector. */
@@ -519,6 +520,111 @@ class CompletionResolverTest {
assertTrue(rendezvous.isWaiting("term_a"), "turn N+1 is still awaiting its own resolution");
}
// --- fleetd #556 (comment 16908): name the invariant the two-arg remove(target, turn) idiom
// exists to maintain — a superseded turn's terminal handling must not evict its successor's
// registration — and test THAT, not the idiom's literal spelling. Three entry points reach a
// conditional remove: resolve()'s plain completion path, resolve()'s echoed-brief/noReportMessage
// sub-path, and fail(). This is a non-goal-to-break for fleetd #556's redesign: flattening any of
// these removes to the one-arg form must turn one of these three red. -------------------------
@Test
void aSupersededTurnsPlainCompletionMustNotEvictItsSuccessorsRegistration() {
// Turn A is overtaken by turn B on the same target (a rapid back-to-back send): B's own
// delivery overwrites the map entry A's delivery installed. A's late completion fallback
// then finally fires and reaches resolve()'s ordinary (non-echoed) completion branch — the
// two-arg remove(target, turnA) there must be a no-op once the map holds B, not a blind
// evict.
FakeHerdr herdr = new FakeHerdr().readText("⏺ A's pre-turn pane\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null); // baseline null: fails open, never suppressed
rendezvous.close("term_a", waiterA); // A's own send's finally, freeing the session for B
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB)); // the injector registering B
herdr.readText("⏺ A's late, genuine answer\n❯ "); // what the pane shows when A's fallback fires
resolver.resolve("term_a", turnA);
assertTrue(waiterA.isDone(), "A's own resolve must still complete A's waiter normally");
assertEquals(Rendezvous.Kind.COMPLETION, waiterA.getNow(null).kind());
CompletionResolver.InFlight afterA = resolver.inFlight("term_a");
assertNotNull(afterA, "B's registration must still be present — a one-arg remove(target) "
+ "here would evict it even though the map no longer holds turnA");
assertEquals(waiterB, afterA.waiter(), "the surviving entry must still be B's, untouched by A's "
+ "terminal handling");
assertTrue(rendezvous.resolve("term_a", "B replied"),
"B must still resolve normally through the ordinary reply path");
assertEquals(Rendezvous.Kind.REPLY, waiterB.getNow(null).kind());
}
@Test
void aSupersededTurnsEchoedNoReportSubPathMustNotEvictItsSuccessorsRegistration() {
// Same invariant, driven down resolve()'s OTHER branch: a pane that merely echoes the
// injected brief back (no real report) resolves through noReportMessage(), not the plain
// tail — comment 16908 calls this "a sub-path of the first, not a separate method", reaching
// the same terminal remove() call site from different logic above it.
String echoedBrief = "y".repeat(450); // >= CompletionResolver.ECHO_MIN_CHARS normalised chars
FakeHerdr herdr = new FakeHerdr().readText("⏺ " + echoedBrief + "\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
// Back-dated deliveredAtNanos (like the 2-arg InFlight convenience ctor) so MIN_TURN_NANOS
// never fires; injectedText == the pane's own content so echoesInjectedBrief() is true.
var turnA = new CompletionResolver.InFlight(waiterA, null, Long.MIN_VALUE / 2, echoedBrief);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertTrue(waiterA.isDone(), "A's own resolve must still complete A's waiter via the echoed "
+ "no-report branch");
assertTrue(waiterA.getNow(null).text().contains(CompletionResolver.NO_REPORT_PREFIX),
"sanity: A must actually have gone down the noReportMessage sub-path, not the plain one");
CompletionResolver.InFlight afterA = resolver.inFlight("term_a");
assertNotNull(afterA, "B's registration must still be present after A's echoed-brief resolve");
assertEquals(waiterB, afterA.waiter());
assertTrue(rendezvous.resolve("term_a", "B replied"),
"B must still resolve normally through the ordinary reply path");
assertEquals(Rendezvous.Kind.REPLY, waiterB.getNow(null).kind());
}
@Test
void aSupersededTurnsFailMustNotEvictItsSuccessorsRegistration() {
// The fail() entry point (CB-109 wedge / CB-110 drop): turn A is failed while turn B has
// already superseded it in the map. fail()'s own two-arg remove(target, turnA) must be a
// no-op once the map holds B.
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.fail("term_a", turnA, "A wedged in an unknown state");
assertTrue(waiterA.isDone(), "A's own fail must still complete A's waiter");
assertEquals(Rendezvous.Kind.FAILED, waiterA.getNow(null).kind());
CompletionResolver.InFlight afterA = resolver.inFlight("term_a");
assertNotNull(afterA, "B's registration must still be present — a one-arg remove(target) in "
+ "fail() would evict it even though the map no longer holds turnA");
assertEquals(waiterB, afterA.waiter());
assertTrue(rendezvous.resolve("term_a", "B replied"),
"B must still resolve normally through the ordinary reply path");
assertEquals(Rendezvous.Kind.REPLY, waiterB.getNow(null).kind());
}
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
@Test
@@ -1197,4 +1197,117 @@ class InjectorTest {
+ "every path out of onStatus, including the one where onDelivered itself is "
+ "what threw. Leaving it pending here is a hang, not a fix");
}
// --- fleetd #556: registration is the Injector's own invariant now, structural rather than
// delegated to a TurnListener callback that can throw. #553 could only make the ONE reachable
// listener behave; this closes the shape itself. ---
@Test
void aTurnListenerThatThrowsFromEveryCallbackStillLeavesTheDeliveredTurnRegisteredAndResolvable() {
// fleetd #556 acceptance criterion 1 (the ticket's own, restated by comment 16966 as the
// replacement for #561's order-dependent test): a TurnListener that throws from EVERY
// callback it implements — nothing about it is safe to lean on — must still leave a
// delivered turn registered, with its waiter still resolvable through the ordinary path.
// Before this ticket, registration lived INSIDE onDelivered, one of the callbacks that
// throws here — this test is red against that shape, because the only listener callback
// this specific delivery scenario exercises (a first delivery, no previous turn to
// complete) is exactly the one that throws.
RuntimeException boom = new RuntimeException("fleetd #556: throws from every callback");
TurnListener throwsEverywhere = new TurnListener() {
@Override
public void onTurnComplete(String target) {
throw boom;
}
@Override
public boolean hasPostTurnAction(String target) {
throw boom;
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
throw boom;
}
@Override
public void onTurnFailed(String target) {
throw boom;
}
@Override
public void onTurnFailed(String target, String reason) {
throw boom;
}
@Override
public void onDelivered(String target, TurnToken token) {
throw boom;
}
};
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
Injector inj = new Injector(new AgentControl(herdr), throwsEverywhere, _ -> true, _ -> {
}, completion::register);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
inj.enqueue(T, "brief", new TurnToken(T, waiter));
RuntimeException thrown = assertThrows(RuntimeException.class,
() -> inj.onStatus(T, AgentStatus.IDLE), // delivers "brief"; onDelivered throws
"the throwing listener's own bug must still escape onStatus — a listener bug stays "
+ "loud (fleetd #553's own guarantee, unchanged by this ticket)");
assertSame(boom, thrown);
CompletionResolver.InFlight inFlight = completion.inFlight(T);
assertNotNull(inFlight,
"fleetd #556: the delivered turn must be registered even though the ONLY "
+ "TurnListener wired throws from every single callback, including "
+ "onDelivered — registration must not depend on that call succeeding");
assertSame(waiter, inFlight.waiter(),
"the registered entry must carry the exact waiter this turn's send opened");
assertFalse(waiter.isDone(), "sanity: nothing has resolved this waiter yet");
assertTrue(rendezvous.resolve(T, "hello"),
"the waiter registrar.register captured must be the real, live rendezvous waiter — "
+ "still resolvable through the ordinary fleet_reply path, unaffected by the "
+ "throwing listener");
assertEquals(Rendezvous.Kind.REPLY, waiter.getNow(null).kind());
}
@Test
void theNewTurnsRegistrationRunsAfterThePreviousTurnsCompletionIsRead() {
// fleetd #556 acceptance criterion 2 (CB-116 guard, ticket body + comment 16966): the
// ordering constraint from fleetd #553 must survive this redesign — onTurnComplete reads
// CompletionResolver's inFlight entry for the PREVIOUS turn, so the new turn's
// registrar.register() call must not run before that read, or a redesign trades this
// ticket's defect for the CB-116 cross-turn stale reply. Pinned here as a call-ORDER
// assertion so a future redesign that hoists registrar.register() earlier (e.g. "for
// simplicity", ahead of the turnCompleted block) goes red immediately, before it can ever
// reach a real cross-turn scrape race.
List<String> events = new ArrayList<>();
TurnListener listener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
events.add("onTurnComplete:" + target);
}
};
TurnRegistrar registrar = (target, token) -> events.add("register:" + target);
Injector inj = new Injector(new AgentControl(herdr), listener, _ -> true, _ -> {
}, registrar);
inj.enqueue(T, "first", TestTurnTokens.inert(T));
inj.enqueue(T, "second", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // delivers "first" -> register:term_a (not asserted below)
inj.onStatus(T, AgentStatus.WORKING); // picked up
events.clear();
inj.onStatus(T, AgentStatus.IDLE); // "first" completes AND "second" delivers, same cycle
assertEquals(List.of("onTurnComplete:" + T, "register:" + T), events,
"onTurnComplete for the previous turn (\"first\") must run before register for the "
+ "new turn (\"second\") within the same onStatus cycle — reversing this order "
+ "is the CB-116 cross-turn stale reply fleetd #553 fixed, and this redesign "
+ "must not reopen it");
}
}