Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 887aca0183 |
@@ -937,20 +937,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
*/
|
||||
@Override
|
||||
public void stop(String idOrPane) {
|
||||
// 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.
|
||||
// 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.
|
||||
String paneId = paneByAgentId.remove(idOrPane);
|
||||
if (paneId == null) {
|
||||
paneId = idOrPane; // raw-pane fallback (reap, gate timeout, pane-addressed callers)
|
||||
}
|
||||
WorkspaceControl.PaneLocation loc = spaces.locatePane(paneId);
|
||||
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
|
||||
try {
|
||||
agents.close(paneId);
|
||||
} catch (HerdrException e) {
|
||||
@@ -991,6 +985,11 @@ 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) {
|
||||
@@ -1118,7 +1138,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(() -> {
|
||||
@@ -1399,6 +1430,33 @@ 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;
|
||||
}
|
||||
|
||||
/** 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,41 +299,6 @@ 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
|
||||
|
||||
@@ -785,6 +785,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
|
||||
@@ -1472,6 +1515,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,21 +711,13 @@ class FleetAppTest {
|
||||
|
||||
@Test
|
||||
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
||||
// 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);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
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"));
|
||||
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");
|
||||
assertFalse(herdr.called("tab.close"), "pane placement owns no tab to close");
|
||||
assertFalse(herdr.called("pane.get"), "no tab resolution in pane placement");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user