Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ea12107497 |
@@ -937,20 +937,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
*/
|
*/
|
||||||
@Override
|
@Override
|
||||||
public void stop(String idOrPane) {
|
public void stop(String idOrPane) {
|
||||||
// Teardown knows only the pane, not which profile spawned it — and, when this call is
|
// Teardown knows only the pane, not which profile spawned it. Attempt tab cleanup when any
|
||||||
// routed here through CompositePeerLauncher's single-daemon stop() shortcut (fleetd #342),
|
// profile uses tab placement (so the bridge may have created a dedicated peer tab); the
|
||||||
// not even which adapter's config actually governed the spawn: the shortcut can hand the
|
// single-occupant check below is what actually protects the user's shared tabs.
|
||||||
// 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);
|
String paneId = paneByAgentId.remove(idOrPane);
|
||||||
if (paneId == null) {
|
if (paneId == null) {
|
||||||
paneId = idOrPane; // raw-pane fallback (reap, gate timeout, pane-addressed callers)
|
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 {
|
try {
|
||||||
agents.close(paneId);
|
agents.close(paneId);
|
||||||
} catch (HerdrException e) {
|
} 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). */
|
/** True when a herdr error means the target is already gone (safe to treat as done). */
|
||||||
private static boolean isAlreadyGone(HerdrException e) {
|
private static boolean isAlreadyGone(HerdrException e) {
|
||||||
return e.code() != null && e.code().endsWith("_not_found");
|
return e.code() != null && e.code().endsWith("_not_found");
|
||||||
|
|||||||
@@ -871,6 +871,10 @@ public final class MessageService {
|
|||||||
boolean wasDelivered = delivery.completion().isDone()
|
boolean wasDelivered = delivery.completion().isDone()
|
||||||
&& !delivery.completion().isCompletedExceptionally();
|
&& !delivery.completion().isCompletedExceptionally();
|
||||||
if (!wasDelivered) {
|
if (!wasDelivered) {
|
||||||
|
if (timeoutCancellationRaceHookForTest != null) {
|
||||||
|
// Test-only (fleetd #345): see the field's own javadoc.
|
||||||
|
timeoutCancellationRaceHookForTest.run();
|
||||||
|
}
|
||||||
// The target monitor makes cancellation atomic with onStatus picking this
|
// The target monitor makes cancellation atomic with onStatus picking this
|
||||||
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
|
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
|
||||||
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
|
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
|
||||||
@@ -1370,6 +1374,24 @@ public final class MessageService {
|
|||||||
this.afterFinishAsyncTaskCompleteHookForTest = hook;
|
this.afterFinishAsyncTaskCompleteHookForTest = hook;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after
|
||||||
|
* {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel}
|
||||||
|
* takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued
|
||||||
|
* delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}.
|
||||||
|
* This deterministically covers the caller's need to use that result rather than relying on a
|
||||||
|
* timing-sensitive real race.
|
||||||
|
*/
|
||||||
|
private volatile Runnable timeoutCancellationRaceHookForTest;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private
|
||||||
|
* so the test, in the same package, can reach it without widening any production API.
|
||||||
|
*/
|
||||||
|
void setTimeoutCancellationRaceHookForTest(Runnable hook) {
|
||||||
|
this.timeoutCancellationRaceHookForTest = hook;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
|
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
|
||||||
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
|
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
|
||||||
|
|||||||
@@ -299,41 +299,6 @@ class CompositePeerLauncherTest {
|
|||||||
"two adapter kinds sharing one daemon keep the fallback route");
|
"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
|
@Test
|
||||||
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
|
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
|
||||||
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
|
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
|
||||||
|
|||||||
@@ -410,6 +410,29 @@ class MessageServiceTest {
|
|||||||
"a delivered send whose worker never replies times out as still working");
|
"a delivered send whose worker never replies times out as still working");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send}
|
||||||
|
* first observes its completion as incomplete, but before {@code cancel} takes the target monitor.
|
||||||
|
* The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report
|
||||||
|
* {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed.
|
||||||
|
*
|
||||||
|
* <p>What this does not prove: that this precise interleaving happens by itself under production
|
||||||
|
* timing. The test forces it through a test-only hook; it proves the timeout caller handles the
|
||||||
|
* injector result when the interleaving occurs.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
|
||||||
|
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
|
||||||
|
try {
|
||||||
|
MessageService.Reply reply = messages.send(T, "race delivery", 50);
|
||||||
|
|
||||||
|
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
|
||||||
|
"cancel reporting DELIVERED means the worker received the timed-out message");
|
||||||
|
} finally {
|
||||||
|
messages.setTimeoutCancellationRaceHookForTest(null);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
||||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||||
|
|||||||
@@ -711,21 +711,13 @@ class FleetAppTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
||||||
// fleetd #342: tab-cleanup resolution is no longer skipped based on a profile's declared
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
// 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");
|
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "pane");
|
||||||
|
|
||||||
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
|
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
|
||||||
assertTrue(herdr.called("pane.close"));
|
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 owns no tab to close");
|
||||||
assertFalse(herdr.called("tab.close"), "pane placement's shared tab must never be closed");
|
assertFalse(herdr.called("pane.get"), "no tab resolution in pane placement");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
Reference in New Issue
Block a user