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..ba069ef 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,16 @@ 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 + * {@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 +410,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 +440,27 @@ 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}). + * + *

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. */ - 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 +505,24 @@ 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. 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. + * * @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 +534,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 +581,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);