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 7aab4b1..b7baab1 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -9,6 +9,7 @@ import dev.ltms.fleet.metrics.Metrics; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; import java.util.List; import java.util.UUID; import java.util.concurrent.CompletableFuture; @@ -167,6 +168,12 @@ public final class MessageService { private static final class Task { private final String ticket; private final String target; + /** + * When this task was created (#137 fix): the tiebreaker for which of several open tasks on + * one target gets a recovered reply in {@link #abandon} — the oldest, since it is the one + * that has been waiting longest. + */ + private final long createdNanos; private final CompletableFuture future = new CompletableFuture<>(); /** * When {@link #future} resolved, or {@code null} while it is still pending — the clock @@ -184,6 +191,7 @@ public final class MessageService { private Task(String ticket, String target, LongSupplier nowNanos) { this.ticket = ticket; this.target = target; + this.createdNanos = nowNanos.getAsLong(); future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong()); } } @@ -375,6 +383,19 @@ 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. {@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 + * it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly + * as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery + * deterministically instead. + * * @return always {@code true} — the reply resolved a live send, completed a parked ticket, or * was queued */ @@ -392,13 +413,21 @@ public final class MessageService { // FAILED with a misleading "session released before it replied" reason, even though the reply // had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket} // sees the real reply instead. - Task orphan = askAnsweredAsyncTask(session); - if (orphan != null && orphan.future.complete(new Reply(Outcome.REPLIED, content))) { - if (orphan.turnId != null) { - asyncTasksByTurn.remove(orphan.turnId, orphan); + List candidates = askAnsweredAsyncTasks(session); + if (candidates.size() == 1) { + Task orphan = candidates.get(0); + if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) { + if (orphan.turnId != null) { + asyncTasksByTurn.remove(orphan.turnId, orphan); + } + count(FleetMetrics.REPLIES, "path", "async-recovered"); + return true; // the ticket itself took it — no inbox stranding at all } - count(FleetMetrics.REPLIES, "path", "async-recovered"); - return true; // the ticket itself took it — no inbox stranding at all + } else if (candidates.size() > 1) { + List tickets = candidates.stream().map(t -> t.ticket).toList(); + log.warn("reply from {} matches {} open async tickets {} — cannot tell which one it " + + "answers, queuing to the inbox instead of guessing", session, candidates.size(), + tickets); } inbox.publish(session, UUID.randomUUID().toString(), content); // CB-640: record the stranding itself (not just the reply text) so fleet health can see a @@ -414,21 +443,38 @@ public final class MessageService { } /** - * The still-open async task on {@code target} whose {@code fleet_ask} was already answered — its - * {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer} — - * yet whose future is not resolved yet (#137). {@code null} if no such task exists, including the + * Every still-open async task on {@code target} whose {@code fleet_ask} was already answered — + * its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer} + * — yet whose future is not resolved yet (#137). Empty if no such task exists, including the * common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was * 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}). + * + *

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 Task askAnsweredAsyncTask(String target) { + private List askAnsweredAsyncTasks(String target) { + List candidates = new ArrayList<>(); for (Task task : tasks.values()) { if (target.equals(task.target) && task.question == null && task.turnId != null && !task.future.isDone()) { - return task; + candidates.add(task); } } - return null; + return candidates; } /** Record a counter sample when a registry is wired; a no-op in unit tests. */ @@ -473,16 +519,63 @@ public final class MessageService { * outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}. * *

#137 defence in depth. {@link #reply} already hands a worker's real - * {@code fleet_reply} straight to the async ticket it belongs to whenever one is still parked - * waiting for it (see {@link #askAnsweredAsyncTask}), so by the time a session is released its - * tasks are normally already resolved — this loop's {@code complete} calls are then harmless - * no-ops (a {@link CompletableFuture} can only resolve once). But should some other path someday - * strand a reply in the inbox without completing its ticket, checking + * {@code fleet_reply} straight to the async ticket it belongs to whenever exactly one is still + * parked waiting for it (see {@link #askAnsweredAsyncTasks}), so by the time a session is + * released its tasks are normally already resolved — this loop's {@code complete} calls are then + * harmless no-ops (a {@link CompletableFuture} can only resolve once). But should some other path + * someday strand a reply in the inbox without completing its ticket, checking * {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down * session whose worker in fact replied is still reported {@code REPLIED} with that reply's own * 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 — 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. + * + *

Do not "fix" that gap with a test that reaches past this class. A test can + * build both facts by calling {@link Rendezvous#close} itself on the waiter an accepted + * {@link #send} is still blocked on: the send keeps the lock, no waiter is registered any more, a + * second parked task stays open, and the next {@link #reply} then strands. That was checked, and + * such a test does go red against the pre-fix loop. But it only goes red because it broke the + * lock-and-waiter invariant above from outside — no caller of this class ever does that — so it + * pins a state production cannot reach, and would read to the next person as if it could. + * * @return true if a live waiter or an async task was failed (never true for one recovered as a * reply — see the note above) */ @@ -494,21 +587,39 @@ public final class MessageService { CompletableFuture waiter = rendezvous.currentWaiter(target); boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason); boolean asyncFailed = false; - Reply recovered = null; // lazily drained at most once, only if a task actually needs it + + List matching = new ArrayList<>(); for (Task task : tasks.values()) { - if (!target.equals(task.target) || task.question != null || task.future.isDone()) { - continue; + if (target.equals(task.target) && task.question == null && !task.future.isDone()) { + matching.add(task); } - if (hadStrandedReply && recovered == null) { - recovered = recoverStrandedReply(target); + } + Task recoveryTask = null; + if (hadStrandedReply && !matching.isEmpty()) { + recoveryTask = matching.get(0); + for (Task candidate : matching) { + if (candidate.createdNanos < recoveryTask.createdNanos) { + recoveryTask = candidate; + } } - Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason); + } + Reply recovered = recoveryTask != null ? recoverStrandedReply(target) : null; + for (Task task : matching) { + boolean isRecovery = task == recoveryTask && recovered != null; + Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason); if (task.future.complete(outcome)) { if (outcome.outcome() == Outcome.WORKER_FAILED) { asyncFailed = true; } else if (task.turnId != null) { asyncTasksByTurn.remove(task.turnId, task); } + } 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()); } } if (failed) { @@ -523,6 +634,10 @@ public final class MessageService { * stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already * drained it first). When more than one message is queued, only the newest is the worker's actual * final answer ({@link #drainReplies} returns them oldest-first). + * + *

This does drain (removes the messages from the inbox) before the caller knows whether the + * task it is recovering for will actually accept them — {@link #abandon} is the one that puts a + * reply back if its {@code complete} call turns out to lose the race. */ private Reply recoverStrandedReply(String target) { var messages = drainReplies(target); 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");