From 97f6c33a45555c4313e1aa876c428e2667eefeeb Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 31 Aug 2026 22:46:23 +0700 Subject: [PATCH] #137: don't guess when a target has more than one open async task MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The CB-205-recovery fix in #205 assumed a target has at most one open async task, with no guard. Fix three consequences: - reply(): askAnsweredAsyncTask -> askAnsweredAsyncTasks (List). Exactly one candidate completes it (unchanged). Zero falls to the inbox (unchanged). More than one now ALSO falls to the inbox instead of picking an arbitrary ConcurrentHashMap iteration order, and logs a WARN naming the target and every candidate ticket. - abandon(): a stranded reply now settles at most one matching task — the oldest by Task#createdNanos (a new field, the tiebreaker). Every other matching task keeps WORKER_FAILED, same as today. - abandon(): if the chosen recovery task's complete() loses a race (another path resolved it first), the drained reply is republished to the inbox instead of being silently dropped. Single-task behavior is unchanged; only the ambiguous case changes. --- .../dev/ltms/fleet/msg/MessageService.java | 108 ++++++++++++++---- 1 file changed, 85 insertions(+), 23 deletions(-) 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);