diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index ba069ef..522a8e7 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -383,9 +383,12 @@ public final class MessageService { * {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions * are interactive and must never be queued. * - *

Ambiguous match also falls to the inbox. A target may have more than one - * open async task waiting on an answered turn — a lead can, and does, issue a second - * {@code fleet_send{wait:false}} at a member that is still busy. Returning whichever candidate a + *

Ambiguous match also falls to the inbox. {@link #askAnsweredAsyncTasks} + * cannot actually return more than one entry today (see its own javadoc for why — in short, + * {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as + * long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from + * two other facts, not one this method enforces, so this branch stays in as defence in depth + * rather than being removed as dead code: if it ever weakens, returning whichever candidate a * {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the * wrong ticket — silently handing the lead something that reads like a correct answer to * a delegation the worker never touched, which is worse than a failure because the lead acts on @@ -447,10 +450,21 @@ public final class MessageService { * never asked has {@code turnId == null}, so it can never match here and only ever completes * through the ordinary rendezvous fast path in {@link #reply}). * - *

Usually holds at most one entry, but not always: a target has no built-in limit of one open - * async task (see the {@link #reply} javadoc), so more than one task can independently reach this - * exact state at once. The caller decides what "more than one" means — {@link #reply} treats it - * as unresolvable and falls back to the inbox; {@link #abandon} picks the oldest deterministically. + *

Returns at most one entry today — verified, not assumed. {@link #send} + * refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that + * check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} — + * not only while its question is still open. {@link #answer} deliberately leaves that stamp in + * place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually + * resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that + * task's future in the same call. So a second task can never reach "{@code turnId} stamped, future + * still open" — the exact pair this method matches on — while a first one already holds it: by + * the time the stamp is gone, so is the eligibility. This is an emergent property of those two + * facts holding together, not something this method (or its callers) enforces on its own — flip + * {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being + * true, with nothing left to fail loudly. The callers below still handle "more than one" as + * defence in depth against exactly that, not because they exercise it today: {@link #reply} + * treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest + * deterministically (its own {@code matching} list has no such guarantee — see its javadoc). */ private List askAnsweredAsyncTasks(String target) { List candidates = new ArrayList<>(); @@ -515,13 +529,44 @@ public final class MessageService { * text, never the misleading "the worker session was released before it replied" (which also * means the snapshot/worktree recovery hint that follows it never prints once a reply exists). * - *

At most one task gets the recovered reply. A stranded reply is one worker - * answer, so it can settle at most one open task on this target — never every open task, and - * never a guess. When more than one task is still open here (the same "more than one open async - * task per target" situation {@link #reply} defers on), the recovered reply goes to the - * oldest (lowest {@link Task#createdNanos}) — it has been waiting longest, so it is the - * one most likely to be what the reply actually answers. Every other open task keeps the ordinary - * {@code WORKER_FAILED} path it would take without a stranded reply at all. + *

At most one task gets the recovered reply — and here, unlike {@link #reply}'s + * {@link #askAnsweredAsyncTasks}, {@code matching.size() >= 2} alone is reachable today. + * This method's {@code matching} filter has no {@code turnId != null} requirement, so it matches + * any plain (never-asked) open task too — and {@link #sendAsync} does not limit a target to one + * of those: a second {@code fleet_send{wait:false}} at a target that is still busy returns its own + * ticket immediately and simply parks its {@link #send} behind the target's session lock for up + * to {@link #ASYNC_TIMEOUT_MS}, exactly as {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} + * already proves. Before this fix, the loop below drained the strand once and then reused that + * same {@code Reply} for every task it walked past — so two open tasks really did both + * complete {@code REPLIED} with the same text (see the pre-fix loop in commit 97f6c33's parent). + * A stranded reply is one worker answer, so it can settle at most one open task on this target — + * never every open task, and never a guess. When more than one task is still open here, the + * recovered reply goes to the oldest (lowest {@link Task#createdNanos}) — it has been + * waiting longest, so it is the one most likely to be what the reply actually answers. Every + * other open task keeps the ordinary {@code WORKER_FAILED} path it would take without a stranded + * reply at all. + * + *

{@code matching.size() >= 2} together with {@code hadStrandedReply} is a different + * question, and today it is defence in depth rather than a path this codebase's public API can + * drive. This class has exactly two sites that ever acquire a target's entry in + * {@code sessionLocks} — {@link #send} and {@link #answer} — and both open a {@link Rendezvous} + * waiter for that same target as the very first thing they do after acquiring the lock, then hold + * lock and waiter together for the rest of their critical section ({@link #send} also clears + * {@link #strandedReplies} right there, the instant it opens its waiter — before it ever enqueues + * delivery). So "the session lock is held" and "a live waiter is open for it" are the same fact + * throughout this class, and {@link #reply}'s fast path always resolves a currently-open waiter + * directly rather than stranding. The two facts this method wants therefore cannot be produced + * side by side: while the lock is held, a real reply resolves the open waiter directly and never + * reaches {@link #strandedReplies}; the instant the lock is free, any parked matching task's own + * {@link #send} that is scheduled next wins it and, by opening its waiter, clears the strand again + * before this method ever runs. There is no way to hold that lock open-but-unaccepted from outside + * {@link #send}/{@link #answer} to freeze a window in between. Constructing both facts at once + * through {@code sendAsync}/{@code reply}/{@code ask}/{@code answer} would need a race against + * virtual-thread scheduling, not a deterministic sequence — so the oldest-wins code below stays as + * defence in depth against a regression to that mechanism (e.g. clearing {@link #strandedReplies} + * on a narrower condition than "any acceptance"), not because today's test suite exercises the + * conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what + * {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers. * * @return true if a live waiter or an async task was failed (never true for one recovered as a * reply — see the note above) diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index bf27236..95f1e92 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -688,6 +688,32 @@ class MessageServiceTest { assertFailedTicket(third, "agent target term_a not found"); } + // --- #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 + // open matching tasks was attempted here and removed after investigation showed the combination + // is not reachable through the public API today, not merely hard to time right: + // + // This class has exactly two call sites that ever hold a target's entry in the session-lock map + // (send() and answer()), and both open a Rendezvous waiter for that same target as the first thing + // they do after acquiring the lock, holding lock and waiter together for their whole critical + // section. So "the lock is held" and "a live waiter is open" are the same fact throughout this + // class. reply()'s fast path always resolves a currently-open waiter directly instead of + // stranding — so a strand can only be created while NO task is accepted (lock free), and the + // instant the lock is next taken (by any parked matching task's own send(), the moment it is + // scheduled), that acceptance clears strandedReplies again (see send()'s CB-640 comment) before + // abandon() can ever observe both facts together. Confirmed empirically too: an earlier version of + // this test stranded a reply, then created an "accepted" task (awaitWaiting()) followed by a + // "parked" one — and the accepted task's own acceptance silently cleared the strand it was + // supposed to be racing against, so the parked task came back WORKER_FAILED instead of DONE, not + // because the fix was missing but because the test's premise could not be constructed. + // + // The reachable half — matching.size() >= 2 alone, no strand — is exactly what + // abandonFailsEveryPendingAsyncTicketForTheReleasedTarget already covers (all fail, none guess). + // The oldest-wins code in abandon() stays as defence in depth (see its own javadoc) against a + // regression that would make the conjunction reachable, e.g. clearing strandedReplies on a + // narrower condition than "any acceptance" — not because this suite exercises it today. + @Test void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception { String ticket = messages.sendAsync(T, "task that asks");