Tearing a worker down left the send that was waiting on it stranded. The
rendezvous waiter stayed open, so a blocking bridge_send kept blocking and an
async one kept reporting PENDING until ASYNC_TIMEOUT_MS — thirty minutes —
even though the worker provably no longer existed and the delegation could
never complete.
Observed repeatedly this session: a task sitting at
{"phase":"pending","detail":"worker unknown"} long after I had personally
deleted the pane. The maddening part is that poll() already HAD the evidence —
it calls liveStatus() to build that detail string, gets back "unknown", and
reports PENDING anyway. It also never reached /metrics under any outcome, so a
stalled delegation was invisible to both the task view and the dashboard. The
only way I ever diagnosed one was reading the worker's pane over the herdr
socket by hand.
Fixed at the choke point rather than by guessing from status strings:
SessionManager.release() is the single path every teardown funnels through
(REST stop, MCP stop, idle-TTL reaper, recycle, shutdown drain), so it now
notifies a release listener with the terminalId, and Bridged.main wires that to
MessageService.abandon(). Abandon fails the open waiter with a real reason.
Deliberately NOT done by inferring "gone" from liveStatus(): that method
collapses a vanished pane, a wedged worker and a herdr hiccup into the same
"unknown" string, so acting on it would fail live delegations during a
transient blip. An explicit lifecycle signal cannot be ambiguous.
Notified before launcher.stop() so a blocked caller fails fast, and wrapped so
a listener failure can never prevent the teardown it is reacting to.
Resolving as a failure rather than letting it time out also means the outcome
is counted — a torn-down delegation now shows up as
sends_total{outcome="failed"} instead of nothing at all.
353 tests (was 346): abandon fails a waiting send / is a no-op with no waiter /
never clobbers a send the worker already answered, an abandoned async task
polls FAILED rather than PENDING, release notifies with the right terminal,
releasing an unknown pane notifies nobody, and a throwing listener does not
block teardown.
Verified live on the running daemon, reproducing the original scenario:
async send -> {"phase":"pending","detail":"worker working"}
DELETE the worker -> 204
poll -> {"phase":"failed","detail":"the worker session was
released before it replied"}
/metrics -> bridged_sends_total{outcome="failed"} 1
Previously that poll returned PENDING for thirty minutes and the counter stayed
empty.
NOT addressed here, and worth deciding separately: a worker that simply never
replies (rather than being released) still rides out the full 30-minute
ASYNC_TIMEOUT_MS. That is a policy question about long-running delegations, not
a correctness bug.
This commit is contained in:
@@ -209,6 +209,12 @@ public final class Bridged {
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
|
||||
pushLoop, metrics);
|
||||
|
||||
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
|
||||
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
|
||||
// reached /metrics — the delegation was unresolvable and nothing said so.
|
||||
sessions.onRelease(terminal ->
|
||||
messages.abandon(terminal, "the worker session was released before it replied"));
|
||||
|
||||
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(
|
||||
|
||||
@@ -255,6 +255,33 @@ public final class MessageService {
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Abandon any send still waiting on {@code target} because its session has gone away (CB-516).
|
||||
*
|
||||
* <p>Without this, tearing a worker down left its rendezvous waiter open: a blocking
|
||||
* {@code bridge_send} kept blocking, and an async one kept reporting {@code PENDING} until
|
||||
* {@link #ASYNC_TIMEOUT_MS} — thirty minutes — even though the worker provably no longer
|
||||
* existed and the delegation could never complete. Worse, {@code poll} already had the evidence
|
||||
* (it calls {@code liveStatus} to build its detail string and gets back {@code "unknown"}) and
|
||||
* reported {@code PENDING} anyway.
|
||||
*
|
||||
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
|
||||
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
||||
*
|
||||
* @return true if a live waiter was failed
|
||||
*/
|
||||
public boolean abandon(String target, String reason) {
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
if (waiter == null || waiter.isDone()) {
|
||||
return false; // nobody is blocked on this worker — nothing to abandon
|
||||
}
|
||||
boolean failed = rendezvous.resolveFailure(waiter, reason);
|
||||
if (failed) {
|
||||
log.debug("abandoned send to {}: {}", target, reason);
|
||||
}
|
||||
return failed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
|
||||
@@ -17,6 +17,7 @@ import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
/**
|
||||
@@ -47,6 +48,9 @@ public final class SessionManager implements TurnListener {
|
||||
private final LongSupplier nowNanos;
|
||||
private final int contextCap;
|
||||
|
||||
/** CB-516: notified with a terminalId on every release; no-op until wired. */
|
||||
private volatile Consumer<String> releaseListener = _ -> { };
|
||||
|
||||
/** Backward-compatible constructor: shared-tree sessions, production git seam. */
|
||||
public SessionManager(PeerLauncher launcher) {
|
||||
this(launcher, new GitWorktrees(), System::nanoTime, 0);
|
||||
@@ -136,6 +140,10 @@ public final class SessionManager implements TurnListener {
|
||||
if (removed != null) {
|
||||
log.debug("releasing session pane={} terminal={} state={}",
|
||||
removed.paneId(), removed.terminalId(), removed.state());
|
||||
// CB-516: a send still waiting on this worker can never be answered now. Tell the
|
||||
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
|
||||
// reason instead of sitting on a rendezvous nothing will ever resolve.
|
||||
notifyReleased(removed.terminalId());
|
||||
}
|
||||
launcher.stop(paneId);
|
||||
if (removed != null && removed.worktree() != null) {
|
||||
@@ -143,6 +151,32 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Register a callback invoked with a session's {@code terminalId} whenever it is released
|
||||
* (CB-516). Every teardown path funnels through {@link #release}, so one hook covers the REST
|
||||
* and MCP stop tools, the idle-TTL reaper, {@code recycle}, and shutdown drain alike.
|
||||
*
|
||||
* <p>Set rather than injected because {@code MessageService} — the intended listener — is
|
||||
* constructed after this manager (it needs the injector and rendezvous, which need the session
|
||||
* presence view this manager exposes). Wiring it at construction would require breaking that
|
||||
* cycle for one callback.
|
||||
*/
|
||||
public void onRelease(Consumer<String> listener) {
|
||||
this.releaseListener = (listener == null) ? _ -> { } : listener;
|
||||
}
|
||||
|
||||
/** A listener failure must never prevent the teardown it is reacting to. */
|
||||
private void notifyReleased(String terminalId) {
|
||||
if (terminalId == null) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
releaseListener.accept(terminalId);
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
|
||||
}
|
||||
}
|
||||
|
||||
private WorkerSession acquireWithWorktree(String profile, String requestedCwd, String callerCwd,
|
||||
String ownerTerminal, WorktreeRequest wt) {
|
||||
String resolvedProfile = (profile == null || profile.isBlank())
|
||||
|
||||
@@ -401,4 +401,67 @@ class MessageServiceTest {
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver the task
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up
|
||||
}
|
||||
|
||||
// --- CB-516: a released session must not leave a send hanging ------------------------------
|
||||
|
||||
/**
|
||||
* The bug this fixes: tearing a worker down left its rendezvous waiter open, so a blocking send
|
||||
* kept blocking and an async one kept reporting PENDING until the 30-minute async timeout —
|
||||
* even though the worker provably no longer existed.
|
||||
*/
|
||||
@Test
|
||||
void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
|
||||
awaitWaiting();
|
||||
|
||||
assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned");
|
||||
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.WORKER_FAILED, r.outcome(),
|
||||
"an abandoned send fails rather than riding out its timeout");
|
||||
assertEquals("session released", r.text(), "the caller is told why");
|
||||
}
|
||||
|
||||
@Test
|
||||
void abandonIsANoOpWhenNobodyIsWaiting() {
|
||||
assertFalse(messages.abandon(T, "session released"),
|
||||
"no open send ⇒ nothing to abandon");
|
||||
}
|
||||
|
||||
@Test
|
||||
void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "the real answer"));
|
||||
|
||||
assertFalse(messages.abandon(T, "session released"),
|
||||
"a send already answered by the worker must not be clobbered");
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals("the real answer", r.text());
|
||||
}
|
||||
|
||||
/** The async path is the one that hung: poll must report FAILED, not PENDING forever. */
|
||||
@Test
|
||||
void anAbandonedAsyncTaskPollsAsFailedNotPending() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) break;
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertNotNull(view);
|
||||
assertEquals(MessageService.Phase.FAILED, view.phase(),
|
||||
"a delegation whose worker is gone must not keep reporting PENDING");
|
||||
assertTrue(view.detail() != null && view.detail().contains("released"),
|
||||
"and the detail says why, rather than 'worker unknown'");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -362,4 +362,46 @@ class SessionManagerTest {
|
||||
assertTrue(sessions.roster().isEmpty(),
|
||||
"no session is registered when spawn times out (roster empty)");
|
||||
}
|
||||
|
||||
// --- CB-516: release must notify, so a blocked send can be failed --------------------------
|
||||
|
||||
@Test
|
||||
void releaseNotifiesTheListenerWithTheReleasedTerminal() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(released::add);
|
||||
|
||||
WorkerSession s = sessions.acquire("ltms-local", null, "/caller", null);
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(java.util.List.of(s.terminalId()), released,
|
||||
"every teardown path funnels through release, so one hook must see the terminal");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releasingAnUnknownPaneNotifiesNobody() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(released::add);
|
||||
|
||||
sessions.release("w9:p404"); // idempotent teardown of something already gone
|
||||
|
||||
assertTrue(released.isEmpty(), "no session removed ⇒ no send was waiting on it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aThrowingReleaseListenerDoesNotBlockTheTeardown() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
sessions.onRelease(_ -> {
|
||||
throw new IllegalStateException("listener blew up");
|
||||
});
|
||||
|
||||
WorkerSession s = sessions.acquire("ltms-local", null, "/caller", null);
|
||||
assertDoesNotThrow(() -> sessions.release(s.paneId()),
|
||||
"a listener failure must never prevent the teardown it is reacting to");
|
||||
assertTrue(sessions.get(s.paneId()).isEmpty(), "and the session is still deregistered");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user