|
|
|
@@ -187,6 +187,24 @@ public final class MessageService {
|
|
|
|
|
private volatile Long completedNanos;
|
|
|
|
|
private volatile Reply question;
|
|
|
|
|
private volatile String turnId;
|
|
|
|
|
/**
|
|
|
|
|
* Set when this task's {@code fleet_ask} lapsed with no answer (fleetd #307):
|
|
|
|
|
* {@link #clearAsyncQuestion} then forgets {@link #turnId} (nulls it and drops the task from
|
|
|
|
|
* {@code asyncTasksByTurn}) so {@link #hasAsyncQuestion} stops reporting the target BUSY — a
|
|
|
|
|
* later {@code fleet_send} to it must be accepted, not refused. But the worker's turn is
|
|
|
|
|
* still genuinely live: it resumed on its own and will eventually call its real
|
|
|
|
|
* {@code fleet_reply}. Losing {@link #turnId} loses {@link #askAnsweredAsyncTasks}' only
|
|
|
|
|
* signal that such a reply belongs to this task, so that reply used to fall straight to the
|
|
|
|
|
* inbox and strand — {@code fleet_poll} stayed {@code PENDING} forever, later force-failed by
|
|
|
|
|
* {@link #abandon} with the misleading "session released before it replied". This flag is a
|
|
|
|
|
* second, independent signal that survives the forgetting: {@link #askAnsweredAsyncTasks}
|
|
|
|
|
* accepts it in place of a live {@link #turnId}, without ever re-adding the task to
|
|
|
|
|
* {@code asyncTasksByTurn} (so the BUSY release is untouched). Cleared implicitly once
|
|
|
|
|
* {@link #future} resolves — every match in {@link #askAnsweredAsyncTasks} already requires
|
|
|
|
|
* {@code !future.isDone()}, so a task that recovered (or was later failed by
|
|
|
|
|
* {@link #abandon}) can never match again regardless of this flag's value.
|
|
|
|
|
*/
|
|
|
|
|
private volatile boolean askTimedOut;
|
|
|
|
|
|
|
|
|
|
private Task(String ticket, String target, LongSupplier nowNanos) {
|
|
|
|
|
this.ticket = ticket;
|
|
|
|
@@ -395,18 +413,16 @@ public final class MessageService {
|
|
|
|
|
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
|
|
|
|
* are interactive and must never be queued.
|
|
|
|
|
*
|
|
|
|
|
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@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
|
|
|
|
|
* <em>wrong</em> 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.
|
|
|
|
|
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks} can
|
|
|
|
|
* return more than one entry — a reachable state, not a hypothetical one (see its own javadoc:
|
|
|
|
|
* an {@code fleet_ask} that lapsed with no answer, fleetd #307, frees the target for a completely fresh
|
|
|
|
|
* delegation, which can itself go on to ask-and-lapse before the first worker's real reply
|
|
|
|
|
* arrives). Returning whichever candidate a {@code ConcurrentHashMap} iteration reaches first
|
|
|
|
|
* would let a genuine reply complete the <em>wrong</em> 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.
|
|
|
|
|
*
|
|
|
|
|
* <p><strong>{@code content} is required (fleetd #302).</strong> Both doors that reach this
|
|
|
|
|
* method must reject a missing/blank reply the same way, so the check lives here rather than in
|
|
|
|
@@ -434,15 +450,18 @@ public final class MessageService {
|
|
|
|
|
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
|
|
|
|
return true; // a live send took it — unchanged fast path
|
|
|
|
|
}
|
|
|
|
|
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
|
|
|
|
|
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
|
|
|
|
|
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
|
|
|
|
|
// waiter long before the worker — now actually resuming real work — finishes and replies. That
|
|
|
|
|
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
|
|
|
|
|
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
|
|
|
|
|
// 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.
|
|
|
|
|
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
|
|
|
|
|
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
|
|
|
|
|
// - answer()'s own bounded wait (the primary's fleet_send{turnId} call, capped well under a
|
|
|
|
|
// minute) can time out and close its waiter long before the worker — now actually resuming
|
|
|
|
|
// real work — finishes and replies.
|
|
|
|
|
// - ask()'s own wait for the primary can time out first, with the worker resuming on its own
|
|
|
|
|
// and finishing unanswered.
|
|
|
|
|
// Either way that reply used to have nowhere to land but the session inbox, leaving the async
|
|
|
|
|
// ticket's future unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's
|
|
|
|
|
// abandon() forced it 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.
|
|
|
|
|
List<Task> candidates = askAnsweredAsyncTasks(session);
|
|
|
|
|
if (candidates.size() == 1) {
|
|
|
|
|
Task orphan = candidates.get(0);
|
|
|
|
@@ -473,34 +492,42 @@ public final class MessageService {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* 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}).
|
|
|
|
|
* Every still-open async task on {@code target} whose worker is genuinely expected to send a
|
|
|
|
|
* real {@code fleet_reply} next with nothing left registered to catch it: either its
|
|
|
|
|
* {@code fleet_ask} was already answered — {@link Task#turnId} is stamped but {@link
|
|
|
|
|
* Task#question} was cleared by {@link #answer} — or its {@code fleet_ask} lapsed unanswered and
|
|
|
|
|
* {@link Task#askTimedOut} marks that (fleetd #307; {@link Task#turnId} is {@code null} by then, forgotten
|
|
|
|
|
* so the target is not left BUSY — see {@link Task#askTimedOut}'s own javadoc). Either way the
|
|
|
|
|
* task's future is not resolved yet. 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 both {@code turnId == null} and {@code askTimedOut == false}, so it can never match here and
|
|
|
|
|
* only ever completes through the ordinary rendezvous fast path in {@link #reply}).
|
|
|
|
|
*
|
|
|
|
|
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@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).
|
|
|
|
|
* <p><strong>Can return more than one entry — reachable, not just defence in depth.</strong>
|
|
|
|
|
* {@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}. While a task's {@code turnId} stays stamped — {@link #answer} leaves
|
|
|
|
|
* it in place ({@code clearAsyncQuestion(turnId, false)}) until {@link #finishAsyncTask} removes
|
|
|
|
|
* the stamp and completes the future in the same call — no second task on the same target can
|
|
|
|
|
* reach an eligible state, because {@link #send} would refuse it as BUSY first. That single-task
|
|
|
|
|
* guarantee holds only for the {@code turnId}-stamped half of this method's match: an
|
|
|
|
|
* {@link Task#askTimedOut} task is, by construction, no longer stamped in {@code asyncTasksByTurn}
|
|
|
|
|
* (that is the whole point of forgetting {@code turnId} in {@link #clearAsyncQuestion}), so the
|
|
|
|
|
* target is free the moment one ask lapses. A fresh, independent {@code sendAsync} to the same
|
|
|
|
|
* target can then be dispatched, itself pause on {@code fleet_ask}, and itself time out — landing
|
|
|
|
|
* a second {@code askTimedOut} task on the very target the first one is still waiting to answer
|
|
|
|
|
* for. Two (or more) genuinely open tasks on one target is therefore a real, reachable state
|
|
|
|
|
* today, not a hypothetical: {@link #reply} treats it as unresolvable and falls back to the
|
|
|
|
|
* inbox rather than guess which task a reply belongs to (guessing wrong would hand the lead a
|
|
|
|
|
* plausible-looking answer to a delegation the worker never touched — worse than a failure,
|
|
|
|
|
* because the lead acts on it); {@link #abandon} instead picks the oldest deterministically (its
|
|
|
|
|
* own {@code matching} list has a different, wider match — see its javadoc).
|
|
|
|
|
*/
|
|
|
|
|
private List<Task> askAnsweredAsyncTasks(String target) {
|
|
|
|
|
List<Task> candidates = new ArrayList<>();
|
|
|
|
|
for (Task task : tasks.values()) {
|
|
|
|
|
if (target.equals(task.target) && task.question == null && task.turnId != null
|
|
|
|
|
&& !task.future.isDone()) {
|
|
|
|
|
if (target.equals(task.target) && task.question == null && !task.future.isDone()
|
|
|
|
|
&& (task.turnId != null || task.askTimedOut)) {
|
|
|
|
|
candidates.add(task);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@@ -898,6 +925,12 @@ public final class MessageService {
|
|
|
|
|
return new AskResult(AskOutcome.ANSWERED, answer);
|
|
|
|
|
} catch (TimeoutException e) {
|
|
|
|
|
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
|
|
|
|
|
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
|
|
|
|
|
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
|
|
|
|
|
// what keeps the target from staying BUSY forever), but it would otherwise also erase
|
|
|
|
|
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
|
|
|
|
|
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
|
|
|
|
|
markAskTimedOut(ticket.turnId());
|
|
|
|
|
clearAsyncQuestion(ticket.turnId(), true);
|
|
|
|
|
return new AskResult(AskOutcome.TIMED_OUT, null);
|
|
|
|
|
} catch (ExecutionException e) {
|
|
|
|
@@ -1162,6 +1195,23 @@ public final class MessageService {
|
|
|
|
|
return task;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Mark {@code turnId}'s task as having a {@code fleet_ask} that lapsed with no answer (fleetd #307), so
|
|
|
|
|
* {@link #askAnsweredAsyncTasks} still recognizes the worker's eventual real {@code fleet_reply}
|
|
|
|
|
* as belonging to it after {@link #clearAsyncQuestion}'s {@code forgetTurn=true} erases
|
|
|
|
|
* {@link Task#turnId} — see {@link Task#askTimedOut}. Must be called before that forgetting, while
|
|
|
|
|
* {@code turnId} can still resolve the task in {@code asyncTasksByTurn}; a lookup afterward would
|
|
|
|
|
* find nothing. Only when it matches the task's current turn — same guard as
|
|
|
|
|
* {@link #clearAsyncQuestion} — so a chained second {@code fleet_ask} (#282) that already moved
|
|
|
|
|
* the task to a fresh {@code turnId} cannot mark it for a turn that is no longer its own.
|
|
|
|
|
*/
|
|
|
|
|
private void markAskTimedOut(String turnId) {
|
|
|
|
|
Task task = asyncTasksByTurn.get(turnId);
|
|
|
|
|
if (task != null && turnId.equals(task.turnId)) {
|
|
|
|
|
task.askTimedOut = true;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
|
|
|
|
|
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
|
|
|
|
|
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
|
|
|
|
|