Compare commits

...

6 Commits

Author SHA1 Message Date
Dai Ha 0d5944af63 fleetd #334: close ask()'s turn before forgetting its Task, closing the last stranding window
CI / build (pull_request) Successful in 2m14s
CI / contract (pull_request) Successful in 19m14s
ask()'s TimeoutException catch used to run clearAsyncQuestion(turnId, true) -- forgetting the
Task's asyncTasksByTurn mapping -- before rendezvous.closeAsk(turnId) ran in the shared finally.
Between those two calls the ask was still "answerable" (askSession(turnId) non-null) but the Task
mapping was already gone, so a racing answer() call found task == null, skipped
finishAsyncTask, and stranded the async ticket at PENDING even though answer() itself reported a
result. #329 fixed one step of this same race; this closes the remaining one.

The fix reorders the fresh owner's teardown: closeAsk runs first, then markAskTimedOut and
clearAsyncQuestion. A racing answer() call now either sees the ask still open (and the Task
mapping guaranteed intact) or sees it already closed (STALE_TURN, before it ever reaches
asyncTasksByTurn). It also gates the whole block by ticket.fresh(), matching the invariant the
finally block already states ("only the fresh owner tears down the shared turn") -- a duplicate
coalesced ask() timing out no longer forgets bookkeeping the fresh owner still needs.

Adds aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket, which pins the exact window
with a new test-only hook (askTimeoutRaceHookForTest) and proves both invariants: a late answer()
racing the timeout sees STALE_TURN, and the async ticket still resolves DONE from the worker's
real reply. Reverting the reorder (verified locally, not committed) makes this test fail with
"expected STALE_TURN but was TIMED_OUT_WORKING".
2026-09-04 17:06:15 +07:00
Dai Ha 65a78932c1 Merge #335: a per-task cleanup throw in abandon() no longer strands the tasks behind it
CI / contract (push) Successful in 1m24s
CI / build (push) Successful in 2m8s
2026-09-04 16:44:04 +07:00
Dai Ha 73aab3f83e Merge #342: teardown resolves the pane's real tab instead of trusting the delegate's placement config
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m47s
2026-09-04 16:38:31 +07:00
Dai Ha 887aca0183 fleetd#335: abandon()'s per-task cleanup and sendAsync's terminal hook must not swallow throws
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 2m21s
Site 1 (abandon()'s matching loop, reachable): the recovery/put-back branch calls
inbox.publish, which AmqpReplyInbox implements as a real broker round trip that
throws IllegalStateException on an unroutable/unconfirmed/interrupted publish.
An uncaught throw there aborted the loop, stranding every task after it in
`matching` PENDING forever. Fixed by recording each task's own future.complete()
result before any cleanup runs, then wrapping the cleanup in try/catch so one
task's failure cannot stop its siblings from getting their outcome. Reaching the
throwing branch by real timing needs a race the file's own #137 follow-up already
found unreachable through the public API, so the reproducing test uses a
test-only hook (same technique as the existing fleetd #324/#329 hooks) to inject
the throw at that exact point.

Site 2 (sendAsync's task.future.whenComplete, reachable): the returned stage is
discarded, so an uncaught throw from pushLoop.onTicketTerminal vanished with no
log line. Reproduced for real: Fleetd's shutdown hook runs messages.close()
(stops the async executor from taking new work, but does not cancel a send
already in flight) before pushLoop.close() (shuts its scheduler down
immediately) — a ticket completing in that window makes onTicketTerminal's own
scheduler.schedule(...) throw a genuine RejectedExecutionException. Fixed with a
try/catch(Throwable) plus log.error inside the whenComplete action.

Site 3 (the two `finally { asyncTasksByWaiter.remove(reply); rendezvous.close(...);
}` blocks in send() and answer()): read Rendezvous.close/closeAsk and the
ConcurrentHashMap operations behind them — both are plain map ops on a non-null
key with no user-overridable code, so neither can throw. Left unchanged; not a
defect.

Mutation-proven: reverting either fix reproduces the failure it exists to catch
— removing site 1's try/catch aborts abandon() with the injected exception
(MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks
errors); removing site 2's try/catch leaves the RejectedExecutionException
unlogged (MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently
fails its log assertion). Full suite: mvn clean install, Tests run: 1365,
Failures: 0, Errors: 0, BUILD SUCCESS.
2026-09-04 16:38:14 +07:00
Dai Ha 3fd23ecafa fleetd #342: base tab-cleanup teardown on the pane's real placement, not a delegate's static config
CI / contract (pull_request) Successful in 1m22s
CI / build (pull_request) Failing after 1m38s
HerdrPeerLauncher.stop() used to gate spaces.locatePane() on usesTabPlacement(),
which reads the delegate's OWN configured profiles. When CompositePeerLauncher's
single-daemon stop() shortcut hands a pane to a delegate that never spawned it
(spawnedBy empty after a daemon restart, herdrDaemonCount()==1), that delegate's
placement config says nothing true about how the pane was actually placed, and a
dedicated tab could be skipped and leaked.

Resolve the tab unconditionally instead — WorkspaceControl#locatePane already
tolerates a missing pane by returning null — and let the existing single-occupant
check (tabPaneCount()==1) be the only thing that decides whether to close it, same
as it already protects a shared tab regardless of declared placement.

Adds a mixed-placement CompositePeerLauncherTest (every existing stop-fallback test
configured both adapters as tab placement, so the mis-routing never showed) and
updates FleetAppTest#stopWorkerInPanePlacementClosesOnlyThePane, whose old
assertion (no pane.get on pane placement) documented exactly the skip this fix
removes.
2026-09-04 16:36:11 +07:00
Dai Ha f379847942 Merge #345: the timeout path's use of Cancellation.DELIVERED is now pinned
CI / contract (push) Successful in 47s
CI / build (push) Successful in 2m8s
2026-09-04 16:32:45 +07:00
5 changed files with 363 additions and 52 deletions
@@ -937,14 +937,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
*/
@Override
public void stop(String idOrPane) {
// Teardown knows only the pane, not which profile spawned it. Attempt tab cleanup when any
// profile uses tab placement (so the bridge may have created a dedicated peer tab); the
// single-occupant check below is what actually protects the user's shared tabs.
// Teardown knows only the pane, not which profile spawned it — and, when this call is
// routed here through CompositePeerLauncher's single-daemon stop() shortcut (fleetd #342),
// not even which adapter's config actually governed the spawn: the shortcut can hand the
// pane to a delegate that never spawned it, whose own profiles say nothing about how THIS
// pane was placed. So the decision to look for a tab to clean up is made from the pane's
// actual state, not from this delegate's static profile config: resolve the tab
// unconditionally and let {@link WorkspaceControl#locatePane} tolerate "not found" (it
// returns null rather than throwing); the single-occupant check below is what actually
// protects the user's shared tabs, exactly as it always has.
String paneId = paneByAgentId.remove(idOrPane);
if (paneId == null) {
paneId = idOrPane; // raw-pane fallback (reap, gate timeout, pane-addressed callers)
}
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
WorkspaceControl.PaneLocation loc = spaces.locatePane(paneId);
try {
agents.close(paneId);
} catch (HerdrException e) {
@@ -985,11 +991,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
}
/** Whether any configured profile places peers in their own tab (so tabs may need cleanup). */
private boolean usesTabPlacement() {
return profiles.values().stream().anyMatch(FleetConfig.Profile::tabPlacement);
}
/** True when a herdr error means the target is already gone (safe to treat as done). */
private static boolean isAlreadyGone(HerdrException e) {
return e.code() != null && e.code().endsWith("_not_found");
@@ -709,27 +709,47 @@ public final class MessageService {
boolean isRecovery = task == recoveryTask && recovered != null;
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
String turnId = task.turnId;
if (task.future.complete(outcome)) {
if (outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
// fleetd #335: completing THIS task's future must not depend on any other task's
// cleanup succeeding — every task in `matching` is owed its own outcome regardless of
// what happens below, so decide and record that before doing anything that can throw.
boolean completedHere = task.future.complete(outcome);
if (completedHere && outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
}
try {
if (abandonCleanupHookForTest != null) {
// Test-only (fleetd #335 site 1): see the field's own javadoc.
abandonCleanupHookForTest.run();
}
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
if (completedHere) {
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
} catch (RuntimeException e) {
// fleetd #335: inbox.publish reaches a broker (AmqpReplyInbox throws
// IllegalStateException on an unroutable/unconfirmed/interrupted publish) and this
// loop has no other teardown path — a caller on the release path, or the health
// monitor's GONE/NEVER_READY sweep. Losing this exception uncaught would abort the
// loop and leave every task still to come in `matching` PENDING forever (fleetd
// #335 site 1). completedHere is already recorded above, so only this task's
// best-effort bookkeeping is lost — log it and let the loop reach the rest.
log.error("abandon: per-task cleanup failed for ticket {} (target {}, turnId {})",
task.ticket, target, turnId, e);
}
}
if (failed) {
@@ -945,13 +965,41 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
// what keeps the target from staying BUSY forever), but it would otherwise also erase
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
// Only the fresh owner tears down the shared turn (mirrors the finally block below).
// A duplicate's own timeoutMillis says nothing about whether the SHARED ask is actually
// done — it must leave the close/forget bookkeeping to the fresh owner, exactly as it
// already leaves closeAsk to it.
if (ticket.fresh()) {
// fleetd #334: close the ask turn BEFORE forgetting this task's turnId mapping below.
// Before this fix the order was reversed — the mapping was forgotten here first, and
// rendezvous.closeAsk only ran afterward, in the shared finally. A primary's answer()
// call racing this exact timeout could then find rendezvous.askSession(turnId) still
// non-null (the ask still "answerable") after the Task mapping was already gone:
// answer()'s own asyncTasksByTurn lookup returned null, its task != null guard skipped
// the completion, and the async ticket sat at PENDING forever even though answer()
// itself reported the worker's real reply. Closing here first removes that window:
// any answer() call that still observes askSession(turnId) != null is necessarily
// racing a point BEFORE the forgetting below runs (both happen on this one thread, in
// this order, with nothing that yields in between), so the Task mapping is still there
// for it to find; any call that observes askSession(turnId) == null now correctly
// bails out STALE_TURN (see answer()'s own top check) before ever reaching
// asyncTasksByTurn. rendezvous.closeAsk is idempotent — a no-op once the turn is
// already removed, see its own javadoc — so the shared finally below re-running it
// for this same fresh call is harmless.
rendezvous.closeAsk(ticket.turnId());
if (askTimeoutRaceHookForTest != null) {
// Test-only (fleetd #334): see the field's own javadoc.
askTimeoutRaceHookForTest.run();
}
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it
// out of asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and
// stays (it is what keeps the target from staying BUSY forever), but it would
// otherwise also erase askAnsweredAsyncTasks' only signal that the worker's eventual
// real fleet_reply still belongs to this task, stranding it in the inbox with a false
// "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
}
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -1050,19 +1098,28 @@ public final class MessageService {
// the chained ask deliberately left open.
//
// A null task is NOT only "this was never an async ticket". That reading was in this
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
// is still answerable but the map entry is already gone, so the lookup at :991
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
// firing only that first half before answer() runs printed
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
// method's own waiter is live. So #329 narrows this window rather than closing it.
// Open as fleetd #334 — do not read this guard as complete.
// comment when #329 merged and it is wrong; it is still not the whole story after
// #334. A genuine async ticket can still land here with task == null — a blocking
// (wait:true) send's fleet_ask never has a Task at all, so that case is expected and
// fine. What #334 fixed was a SECOND, unintended way to get here with task == null:
// ask()'s timeout path used to run clearAsyncQuestion(turnId, true) — which drops the
// asyncTasksByTurn entry — in its catch block, while rendezvous.closeAsk(turnId) ran
// later, in its finally. Between those two the ask was still answerable but the map
// entry was already gone, so the lookup at :1053 returned null and this ticket was
// never completed. Measured on 2026-09-04: a probe firing only that first half before
// answer() ran printed "answer=REPLIED phase=PENDING reply=null" — the same stranded
// ticket #329 set out to fix, one step earlier in the same race; the probe used
// forgetTurnForTest, which omits ask()'s markAskTimedOut, and that omission does not
// change the outcome, because askTimedOut is read only by askAnsweredAsyncTasks, and
// reply() never reaches it while this method's own waiter is live. #334's fix
// reorders ask()'s timeout catch to run closeAsk before the forgetting (see the
// fresh-owner block there), which removes this path entirely rather than narrowing
// it further: once closeAsk has run, rendezvous.askSession(turnId) is null and
// answer() returns STALE_TURN from its own top check, before it ever reaches this
// lookup — see aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket in
// MessageServiceTest, which pins the exact window with askTimeoutRaceHookForTest.
// So by the time this line runs, task == null means only the ordinary blocking-send
// case (or #329's own already-fixed race elsewhere) — not this one.
if (result.outcome() != Outcome.QUESTION && task != null) {
finishAsyncTask(task, result);
}
@@ -1122,7 +1179,18 @@ public final class MessageService {
// class javadoc on sendAsync/CB-107.
task.future.whenComplete((reply, ex) -> {
boolean failed = ex != null || reply == null || !reply.completed();
pushLoop.onTicketTerminal(ticket, target, failed);
try {
pushLoop.onTicketTerminal(ticket, target, failed);
} catch (Throwable t) {
// fleetd #335 (site 2): this stage's own CompletableFuture is discarded, so an
// uncaught throw here (e.g. a RejectedExecutionException from
// ReplyPushLoop's scheduler, already shut down while this in-flight send's
// whenComplete fires during the daemon's own shutdown sequence — messages.close()
// only stops accepting NEW async work, it does not cancel a delivery already
// running) vanishes with no log line and no metric, and the push loop never learns
// the ticket went terminal — the exact thing this hook exists to tell it.
log.error("push loop failed to learn ticket {} (target {}) went terminal", ticket, target, t);
}
});
}
asyncExecutor.submit(() -> {
@@ -1421,6 +1489,56 @@ public final class MessageService {
clearAsyncQuestion(turnId, true);
}
/**
* Null in production; test seam for fleetd #335 (site 1) — invoked from {@link #abandon(String,
* String, boolean)}'s per-task loop, once per task, right before that task's own cleanup
* (turnId bookkeeping, or the stranded-reply put-back) runs. A test installs this to inject a
* throw at that exact point deterministically.
*
* <p>The one production call there that can really throw is {@code inbox.publish} in the
* put-back branch — {@link AmqpReplyInbox#publish} reaches a broker and throws {@link
* IllegalStateException} on an unroutable, unconfirmed, or interrupted publish — but reaching
* that branch requires a second completion of the very task {@code abandon} is about to
* complete to win the race first (see the branch's own comment), and the #137 follow-up
* investigation above already found the combination this needs (a stranded reply coinciding
* with an open matching task) unreachable through the public API, not merely hard to time.
* This hook reproduces the resulting shape — a per-task cleanup throw — directly, the same
* technique {@link #finishAsyncTaskRaceHook} and {@link
* #afterFinishAsyncTaskCompleteHookForTest} already use for their own hard-to-time races.
*/
private volatile Runnable abandonCleanupHookForTest;
/**
* Test-only (fleetd #335, site 1): install {@link #abandonCleanupHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setAbandonCleanupHookForTest(Runnable hook) {
this.abandonCleanupHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #334 — invoked from {@link #ask}'s {@code
* TimeoutException} catch, only for the fresh owner, right after {@code rendezvous.closeAsk}
* has run and before {@link #markAskTimedOut} / {@link #clearAsyncQuestion} forget this task's
* turnId mapping. A test installs this to call {@link #answer} for the very same {@code turnId}
* synchronously from inside that exact window, deterministically reproducing the race a real
* concurrent {@code answer()} call could otherwise only win by timing luck: with the ask already
* closed, that call must see {@code rendezvous.askSession(turnId) == null} and return {@link
* Outcome#STALE_TURN} immediately, never reaching {@code asyncTasksByTurn} at all — proving the
* window fleetd #334 describes (mapping forgotten while the ask was still "answerable") is
* closed, rather than merely narrowed the way fleetd #329 narrowed the sibling race in {@link
* #finishAsyncTask}.
*/
private volatile Runnable askTimeoutRaceHookForTest;
/**
* Test-only (fleetd #334): install {@link #askTimeoutRaceHookForTest}. Package-private so the
* test, in the same package, can reach it without widening any production API.
*/
void setAskTimeoutRaceHookForTest(Runnable hook) {
this.askTimeoutRaceHookForTest = hook;
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -299,6 +299,41 @@ class CompositePeerLauncherTest {
"two adapter kinds sharing one daemon keep the fallback route");
}
@Test
void stopThroughTheSingleDaemonShortcutStillClosesTheTabWhenTheFallbackDelegateUsesPanePlacement() {
// fleetd #342: the single-daemon shortcut (spawnedBy empty, herdrDaemonCount()==1) always
// routes stop() through delegates.getFirst() — here the claude adapter, configured for
// PANE placement (its own profiles never create a dedicated tab). The pane being torn
// down here actually belongs to the opencode adapter's TAB placement, sharing the same
// herdr daemon — mixing placements is the point: every existing stop-fallback test in this
// class configures BOTH adapters as "tab", so usesTabPlacement() was true either way and
// the mis-routing never showed.
//
// Before the fix, HerdrPeerLauncher#stop gated tab resolution on usesTabPlacement() of the
// delegate it happened to be called through, so the wrongly-routed (pane-placement) claude
// adapter never even looked for a tab to close, and the now-empty tab leaked with nothing
// to reap it. The fix (fleetd #342) resolves the pane's real tab unconditionally, so the
// decision follows the pane's actual placement rather than the fallback delegate's static
// config.
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile claudePane = new FleetConfig.Profile("claude", "http://gx00.gw:8000",
"coder", null, "FLEETD_WORKER_TOKEN", List.of("claude"), "pane", "fleetd-workers",
"w #{n}", null, null, null);
ClaudeCodeLauncher claude = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of("claude", claudePane), "claude", _ -> null);
PeerLauncher composite = new CompositePeerLauncher(List.of(claude, opencodeAdapter(herdr)), "claude");
// "w9:pW" was never spawned through this composite instance, so spawnedBy has no entry for
// it (the same in-memory-cache-miss shape a daemon restart leaves behind) — stop() falls
// through to the single-daemon shortcut and hands it to delegates.getFirst() (claude).
composite.stop("w9:pW");
assertTrue(herdr.called("pane.close"), "the pane itself is still closed on every routed path");
assertTrue(herdr.called("tab.close"),
"the pane's real (sole-occupant) tab must be closed even though the fallback routed "
+ "through a delegate configured for pane placement");
}
@Test
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
@@ -377,6 +377,76 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
}
/**
* fleetd #334. {@code ask()}'s {@code TimeoutException} catch used to forget this task's
* {@code turnId} mapping ({@code clearAsyncQuestion(turnId, true)}) BEFORE closing the ask
* ({@code rendezvous.closeAsk}, in the shared {@code finally}). A primary's {@code answer()}
* call racing that exact window found the ask still "answerable" ({@code
* rendezvous.askSession(turnId)} still non-null) while the {@code Task} was already forgotten,
* so its {@code task != null} guard skipped the completion and the async ticket sat at
* {@code PENDING} forever even though {@code answer()} itself reported a result. The fix
* (closing the ask first) makes this window impossible: a racing {@code answer()} call either
* still finds the ask open (and the {@code Task} mapping guaranteed intact) or finds it already
* closed (and bails {@code STALE_TURN} before ever touching the {@code Task}). This test pins
* the exact window with {@code askTimeoutRaceHookForTest} and proves both invariants the ticket
* named: (1) a late/racing {@code answer()} sees the ask as already lapsed ({@code STALE_TURN}),
* never made answerable again, and (2) the async ticket still resolves {@code DONE} once the
* worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}.
*/
@Test
void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 200));
// Wait for the question to actually surface (poll sees ASKING) before racing the timeout.
MessageService.TaskView asking = null;
long deadline = System.currentTimeMillis() + 2000;
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
&& System.currentTimeMillis() < deadline) {
asking = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(asking, "the question must surface before the ask times out");
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
messages.setAskTimeoutRaceHookForTest(() ->
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
try {
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
MessageService.Reply late = lateAnswer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.STALE_TURN, late.outcome(),
"a late answer racing the timeout teardown must see the ask as already lapsed");
// The worker resumes on its own (per the ask() contract) and eventually sends its real
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
} finally {
messages.setAskTimeoutRaceHookForTest(null);
}
MessageService.TaskView done = null;
deadline = System.currentTimeMillis() + 2000;
while ((done == null || done.phase() == MessageService.Phase.PENDING)
&& System.currentTimeMillis() < deadline) {
done = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(done);
assertEquals(MessageService.Phase.DONE, done.phase(), "the async ticket must not be stranded PENDING");
assertEquals("real result", done.reply());
}
@Test
void answeringAnUnknownTurnIsStale() {
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
@@ -808,6 +878,49 @@ class MessageServiceTest {
assertFailedTicket(third, "agent target term_a not found");
}
// --- fleetd #335 (site 1): a per-task cleanup failure inside the abandon() loop must not -----
// strand the tasks that come after it. abandon()'s own comment on the loop documents the one
// real production call that can throw there (inbox.publish, in the stranded-reply put-back
// branch, reached when a concurrent reply() or a second abandon() races this one) — but
// reaching that branch requires the exact combination the #137 follow-up above already found
// unreachable through the public API. abandonCleanupHookForTest reproduces the resulting SHAPE
// (one task's cleanup throws) directly instead, the same technique this file already uses for
// fleetd #324/#329's own hard-to-time races.
@Test
void aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try {
String first = messages.sendAsync(T, "first task");
awaitWaiting(); // first task owns the target lock and rendezvous waiter
String second = messages.sendAsync(T, "second task"); // parked on the same lock
String third = messages.sendAsync(T, "third task"); // parked too — the whole sweep must survive
java.util.concurrent.atomic.AtomicInteger calls = new java.util.concurrent.atomic.AtomicInteger();
messages.setAbandonCleanupHookForTest(() -> {
if (calls.getAndIncrement() == 0) {
throw new RuntimeException("PROBE-335-SITE1");
}
});
assertTrue(messages.abandon(T, "agent target term_a not found"));
// Every task in the loop still gets its own outcome — the one whose cleanup threw
// included — even though the loop had no way to know in advance which one that would be.
assertFailedTicket(first, "agent target term_a not found");
assertFailedTicket(second, "agent target term_a not found");
assertFailedTicket(third, "agent target term_a not found");
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getThrowableProxy() != null
&& "PROBE-335-SITE1".equals(e.getThrowableProxy().getMessage())),
"a per-task cleanup failure must still reach the log, not vanish silently");
} finally {
messages.setAbandonCleanupHookForTest(null);
detachMessageServiceLog(appender);
}
}
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
//
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
@@ -1495,6 +1608,42 @@ class MessageServiceTest {
}
}
// --- fleetd #335 (site 2): task.future.whenComplete's own returned stage is discarded, so an --
// uncaught throw from ReplyPushLoop.onTicketTerminal used to vanish with no log line and no
// metric. The real production trigger is the daemon's own shutdown sequence (Fleetd's shutdown
// hook): messages.close() only stops the async executor from taking NEW work — it does not
// cancel a send already in flight — while pushLoop.close() shuts its scheduler down immediately
// right after, so a ticket that completes in that narrow window has onTicketTerminal's own
// scheduler.schedule(...) throw a real RejectedExecutionException. Reproduced here by shutting
// the very same scheduler down before the ticket resolves — no test-only hook needed, this
// reachable path throws for real.
@Test
void aTicketTerminalPushFailureDoesNotVanishSilently() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try (var wiring = wireWithPushLoop(1, 50)) {
String ticket = wiring.service().sendAsync(T, "long task");
awaitWaiting();
wiring.scheduler().shutdownNow(); // simulate pushLoop.close() racing an in-flight send
injectDelivery();
assertTrue(rendezvous.resolve(T, "async result"));
// The ticket's own outcome must be unaffected by the swallowed exception — finishAsyncTask
// completes task.future before whenComplete's action (and thus onTicketTerminal) ever runs.
MessageService.TaskView done = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE);
assertEquals("async result", done.reply());
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getFormattedMessage().contains(ticket)
&& e.getThrowableProxy() != null
&& "java.util.concurrent.RejectedExecutionException"
.equals(e.getThrowableProxy().getClassName())),
"onTicketTerminal throwing must still reach the log, not vanish silently");
} finally {
detachMessageServiceLog(appender);
}
}
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
@Test
@@ -711,13 +711,21 @@ class FleetAppTest {
@Test
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
FakeHerdr herdr = new FakeHerdr();
// fleetd #342: tab-cleanup resolution is no longer skipped based on a profile's declared
// placement — a stop() routed through the wrong delegate (e.g. CompositePeerLauncher's
// single-daemon shortcut after a daemon restart) could carry a placement config that says
// nothing true about how THIS pane was actually placed. So the pane's real tab is now
// always resolved, and the single-occupant check below is what protects a pane-placement
// peer's shared tab, exactly as it always protected a tab-placement one. Model that
// realistically: the peer's pane was split into an existing tab that already held another
// occupant, so the tab must never be closed.
FakeHerdr herdr = new FakeHerdr().withWorkerTabPaneCount(2);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "pane");
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("pane.close"));
assertFalse(herdr.called("tab.close"), "pane placement owns no tab to close");
assertFalse(herdr.called("pane.get"), "no tab resolution in pane placement");
assertTrue(herdr.called("pane.get"), "tab resolution now always runs, regardless of placement");
assertFalse(herdr.called("tab.close"), "pane placement's shared tab must never be closed");
}
@Test