Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| dbf6fef0e9 | |||
| b6b88c5f1c | |||
| 86dddfe240 | |||
| 0d5944af63 | |||
| 4a5030a5c6 | |||
| c1c8794c48 | |||
| 65a78932c1 | |||
| f429ca1a50 | |||
| 73aab3f83e | |||
| 887aca0183 | |||
| f379847942 | |||
| ea12107497 |
@@ -380,7 +380,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
if (startsWithExhaustion(matchedLine, exhausted)) {
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
@@ -478,7 +480,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
+ "usable assistant block; no fleet_reply): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
if (startsWithExhaustion(matchedLine, exhausted)) {
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
@@ -613,6 +617,46 @@ public final class CompletionResolver implements TurnListener {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when nothing before the match on this pane line ends a sentence — that is, the match is
|
||||
* still inside the line's first sentence rather than inside prose a member wrote about it.
|
||||
* Used to decide whether an exhaustion match may quarantine a credential (fleetd #348).
|
||||
*
|
||||
* <p><strong>Why this is looser than {@link #startsWithBackendError}.</strong> An
|
||||
* {@code exhaustedPattern} is written per profile and may name only the decisive words of a
|
||||
* provider message — {@code "usage limit has been reached"} without its leading {@code "The"}.
|
||||
* A start-of-line check would then reject the genuine refusal. That is the false negative
|
||||
* fleetd #348's invariant 1 calls the worse direction: an unrecorded exhaustion leaves the
|
||||
* fleet spawning into a credential with no capacity, and a quarantine runs 1800s against the
|
||||
* backend-error cooldown's fixed 60s.
|
||||
*
|
||||
* <p>This rule accepts a superset of what a start-of-line check accepts: if the match begins
|
||||
* right after the chrome, there is nothing in front of it, so there is no sentence ending
|
||||
* either. So moving to it cannot add a false negative.
|
||||
*
|
||||
* <p><strong>No chrome skipping here, deliberately.</strong> The first version of this method
|
||||
* copied {@code startsWithBackendError}'s leading-chrome loop. Measured on merge: deleting that
|
||||
* loop left all 1369 tests green, and it must — the scan only looks for {@code . ! ?}, and no
|
||||
* terminal chrome character is one of those. A step that cannot change the result is worse than
|
||||
* no step, because the next reader takes it as evidence that chrome was handled.
|
||||
*
|
||||
* <p>It stays a heuristic. Prose whose <em>first</em> sentence carries the pattern still
|
||||
* notifies the sink, and a genuine refusal behind an earlier full stop (a hostname, a version
|
||||
* number) still does not. Both are known and neither is fixed here.
|
||||
*/
|
||||
private static boolean startsWithExhaustion(String line, Pattern pattern) {
|
||||
var matcher = pattern.matcher(line);
|
||||
if (!matcher.find()) {
|
||||
return false;
|
||||
}
|
||||
for (int prefix = 0; prefix < matcher.start(); prefix++) {
|
||||
if (".!?".indexOf(line.charAt(prefix)) >= 0) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when the error pattern begins the matched pane line, rather than appearing in prose.
|
||||
*
|
||||
|
||||
@@ -709,27 +709,47 @@ public final class MessageService {
|
||||
boolean isRecovery = task == recoveryTask && recovered != null;
|
||||
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
String turnId = task.turnId;
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
// fleetd #335: completing THIS task's future must not depend on any other task's
|
||||
// cleanup succeeding — every task in `matching` is owed its own outcome regardless of
|
||||
// what happens below, so decide and record that before doing anything that can throw.
|
||||
boolean completedHere = task.future.complete(outcome);
|
||||
if (completedHere && outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
}
|
||||
try {
|
||||
if (abandonCleanupHookForTest != null) {
|
||||
// Test-only (fleetd #335 site 1): see the field's own javadoc.
|
||||
abandonCleanupHookForTest.run();
|
||||
}
|
||||
if (turnId != null) {
|
||||
// #275: whether this task was swept out of ASKING or was already answered and
|
||||
// only waiting on its resumed turn's real reply (#137), nothing will ever
|
||||
// complete this turnId now — drop it from this class's own bookkeeping AND the
|
||||
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
|
||||
// that is actually done, and a late answer() sees it as lapsed rather than
|
||||
// resolving a question nothing is listening for any more.
|
||||
asyncTasksByTurn.remove(turnId, task);
|
||||
rendezvous.closeAsk(turnId);
|
||||
if (completedHere) {
|
||||
if (turnId != null) {
|
||||
// #275: whether this task was swept out of ASKING or was already answered and
|
||||
// only waiting on its resumed turn's real reply (#137), nothing will ever
|
||||
// complete this turnId now — drop it from this class's own bookkeeping AND the
|
||||
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
|
||||
// that is actually done, and a late answer() sees it as lapsed rather than
|
||||
// resolving a question nothing is listening for any more.
|
||||
asyncTasksByTurn.remove(turnId, task);
|
||||
rendezvous.closeAsk(turnId);
|
||||
}
|
||||
} 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());
|
||||
}
|
||||
} 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());
|
||||
} catch (RuntimeException e) {
|
||||
// fleetd #335: inbox.publish reaches a broker (AmqpReplyInbox throws
|
||||
// IllegalStateException on an unroutable/unconfirmed/interrupted publish) and this
|
||||
// loop has no other teardown path — a caller on the release path, or the health
|
||||
// monitor's GONE/NEVER_READY sweep. Losing this exception uncaught would abort the
|
||||
// loop and leave every task still to come in `matching` PENDING forever (fleetd
|
||||
// #335 site 1). completedHere is already recorded above, so only this task's
|
||||
// best-effort bookkeeping is lost — log it and let the loop reach the rest.
|
||||
log.error("abandon: per-task cleanup failed for ticket {} (target {}, turnId {})",
|
||||
task.ticket, target, turnId, e);
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -871,6 +891,10 @@ public final class MessageService {
|
||||
boolean wasDelivered = delivery.completion().isDone()
|
||||
&& !delivery.completion().isCompletedExceptionally();
|
||||
if (!wasDelivered) {
|
||||
if (timeoutCancellationRaceHookForTest != null) {
|
||||
// Test-only (fleetd #345): see the field's own javadoc.
|
||||
timeoutCancellationRaceHookForTest.run();
|
||||
}
|
||||
// The target monitor makes cancellation atomic with onStatus picking this
|
||||
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
|
||||
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
|
||||
@@ -941,13 +965,41 @@ 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);
|
||||
// Only the fresh owner tears down the shared turn (mirrors the finally block below).
|
||||
// A duplicate's own timeoutMillis says nothing about whether the SHARED ask is actually
|
||||
// done — it must leave the close/forget bookkeeping to the fresh owner, exactly as it
|
||||
// already leaves closeAsk to it.
|
||||
if (ticket.fresh()) {
|
||||
// fleetd #334: close the ask turn BEFORE forgetting this task's turnId mapping below.
|
||||
// Before this fix the order was reversed — the mapping was forgotten here first, and
|
||||
// rendezvous.closeAsk only ran afterward, in the shared finally. A primary's answer()
|
||||
// call racing this exact timeout could then find rendezvous.askSession(turnId) still
|
||||
// non-null (the ask still "answerable") after the Task mapping was already gone:
|
||||
// answer()'s own asyncTasksByTurn lookup returned null, its task != null guard skipped
|
||||
// the completion, and the async ticket sat at PENDING forever even though answer()
|
||||
// itself reported the worker's real reply. Closing here first removes that window:
|
||||
// any answer() call that still observes askSession(turnId) != null is necessarily
|
||||
// racing a point BEFORE the forgetting below runs (both happen on this one thread, in
|
||||
// this order, with nothing that yields in between), so the Task mapping is still there
|
||||
// for it to find; any call that observes askSession(turnId) == null now correctly
|
||||
// bails out STALE_TURN (see answer()'s own top check) before ever reaching
|
||||
// asyncTasksByTurn. rendezvous.closeAsk is idempotent — a no-op once the turn is
|
||||
// already removed, see its own javadoc — so the shared finally below re-running it
|
||||
// for this same fresh call is harmless.
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
if (askTimeoutRaceHookForTest != null) {
|
||||
// Test-only (fleetd #334): see the field's own javadoc.
|
||||
askTimeoutRaceHookForTest.run();
|
||||
}
|
||||
// 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) {
|
||||
Throwable cause = e.getCause();
|
||||
@@ -1046,19 +1098,28 @@ public final class MessageService {
|
||||
// the chained ask deliberately left open.
|
||||
//
|
||||
// A null task is NOT only "this was never an async ticket". That reading was in this
|
||||
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
|
||||
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
|
||||
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
|
||||
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
|
||||
// is still answerable but the map entry is already gone, so the lookup at :991
|
||||
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
|
||||
// firing only that first half before answer() runs printed
|
||||
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
|
||||
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
|
||||
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
|
||||
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
|
||||
// method's own waiter is live. So #329 narrows this window rather than closing it.
|
||||
// Open as fleetd #334 — do not read this guard as complete.
|
||||
// comment when #329 merged and it is wrong; it is still not the whole story after
|
||||
// #334. A genuine async ticket can still land here with task == null — a blocking
|
||||
// (wait:true) send's fleet_ask never has a Task at all, so that case is expected and
|
||||
// fine. What #334 fixed was a SECOND, unintended way to get here with task == null:
|
||||
// ask()'s timeout path used to run clearAsyncQuestion(turnId, true) — which drops the
|
||||
// asyncTasksByTurn entry — in its catch block, while rendezvous.closeAsk(turnId) ran
|
||||
// later, in its finally. Between those two the ask was still answerable but the map
|
||||
// entry was already gone, so the lookup at :1053 returned null and this ticket was
|
||||
// never completed. Measured on 2026-09-04: a probe firing only that first half before
|
||||
// answer() ran printed "answer=REPLIED phase=PENDING reply=null" — the same stranded
|
||||
// ticket #329 set out to fix, one step earlier in the same race; the probe used
|
||||
// forgetTurnForTest, which omits ask()'s markAskTimedOut, and that omission does not
|
||||
// change the outcome, because askTimedOut is read only by askAnsweredAsyncTasks, and
|
||||
// reply() never reaches it while this method's own waiter is live. #334's fix
|
||||
// reorders ask()'s timeout catch to run closeAsk before the forgetting (see the
|
||||
// fresh-owner block there), which removes this path entirely rather than narrowing
|
||||
// it further: once closeAsk has run, rendezvous.askSession(turnId) is null and
|
||||
// answer() returns STALE_TURN from its own top check, before it ever reaches this
|
||||
// lookup — see aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket in
|
||||
// MessageServiceTest, which pins the exact window with askTimeoutRaceHookForTest.
|
||||
// So by the time this line runs, task == null means only the ordinary blocking-send
|
||||
// case (or #329's own already-fixed race elsewhere) — not this one.
|
||||
if (result.outcome() != Outcome.QUESTION && task != null) {
|
||||
finishAsyncTask(task, result);
|
||||
}
|
||||
@@ -1118,7 +1179,18 @@ public final class MessageService {
|
||||
// class javadoc on sendAsync/CB-107.
|
||||
task.future.whenComplete((reply, ex) -> {
|
||||
boolean failed = ex != null || reply == null || !reply.completed();
|
||||
pushLoop.onTicketTerminal(ticket, target, failed);
|
||||
try {
|
||||
pushLoop.onTicketTerminal(ticket, target, failed);
|
||||
} catch (Throwable t) {
|
||||
// fleetd #335 (site 2): this stage's own CompletableFuture is discarded, so an
|
||||
// uncaught throw here (e.g. a RejectedExecutionException from
|
||||
// ReplyPushLoop's scheduler, already shut down while this in-flight send's
|
||||
// whenComplete fires during the daemon's own shutdown sequence — messages.close()
|
||||
// only stops accepting NEW async work, it does not cancel a delivery already
|
||||
// running) vanishes with no log line and no metric, and the push loop never learns
|
||||
// the ticket went terminal — the exact thing this hook exists to tell it.
|
||||
log.error("push loop failed to learn ticket {} (target {}) went terminal", ticket, target, t);
|
||||
}
|
||||
});
|
||||
}
|
||||
asyncExecutor.submit(() -> {
|
||||
@@ -1370,6 +1442,24 @@ public final class MessageService {
|
||||
this.afterFinishAsyncTaskCompleteHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after
|
||||
* {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel}
|
||||
* takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued
|
||||
* delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}.
|
||||
* This deterministically covers the caller's need to use that result rather than relying on a
|
||||
* timing-sensitive real race.
|
||||
*/
|
||||
private volatile Runnable timeoutCancellationRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private
|
||||
* so the test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setTimeoutCancellationRaceHookForTest(Runnable hook) {
|
||||
this.timeoutCancellationRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
|
||||
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
|
||||
@@ -1399,6 +1489,56 @@ public final class MessageService {
|
||||
clearAsyncQuestion(turnId, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #335 (site 1) — invoked from {@link #abandon(String,
|
||||
* String, boolean)}'s per-task loop, once per task, right before that task's own cleanup
|
||||
* (turnId bookkeeping, or the stranded-reply put-back) runs. A test installs this to inject a
|
||||
* throw at that exact point deterministically.
|
||||
*
|
||||
* <p>The one production call there that can really throw is {@code inbox.publish} in the
|
||||
* put-back branch — {@link AmqpReplyInbox#publish} reaches a broker and throws {@link
|
||||
* IllegalStateException} on an unroutable, unconfirmed, or interrupted publish — but reaching
|
||||
* that branch requires a second completion of the very task {@code abandon} is about to
|
||||
* complete to win the race first (see the branch's own comment), and the #137 follow-up
|
||||
* investigation above already found the combination this needs (a stranded reply coinciding
|
||||
* with an open matching task) unreachable through the public API, not merely hard to time.
|
||||
* This hook reproduces the resulting shape — a per-task cleanup throw — directly, the same
|
||||
* technique {@link #finishAsyncTaskRaceHook} and {@link
|
||||
* #afterFinishAsyncTaskCompleteHookForTest} already use for their own hard-to-time races.
|
||||
*/
|
||||
private volatile Runnable abandonCleanupHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #335, site 1): install {@link #abandonCleanupHookForTest}. Package-private
|
||||
* so the test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setAbandonCleanupHookForTest(Runnable hook) {
|
||||
this.abandonCleanupHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #334 — invoked from {@link #ask}'s {@code
|
||||
* TimeoutException} catch, only for the fresh owner, right after {@code rendezvous.closeAsk}
|
||||
* has run and before {@link #markAskTimedOut} / {@link #clearAsyncQuestion} forget this task's
|
||||
* turnId mapping. A test installs this to call {@link #answer} for the very same {@code turnId}
|
||||
* synchronously from inside that exact window, deterministically reproducing the race a real
|
||||
* concurrent {@code answer()} call could otherwise only win by timing luck: with the ask already
|
||||
* closed, that call must see {@code rendezvous.askSession(turnId) == null} and return {@link
|
||||
* Outcome#STALE_TURN} immediately, never reaching {@code asyncTasksByTurn} at all — proving the
|
||||
* window fleetd #334 describes (mapping forgotten while the ask was still "answerable") is
|
||||
* closed, rather than merely narrowed the way fleetd #329 narrowed the sibling race in {@link
|
||||
* #finishAsyncTask}.
|
||||
*/
|
||||
private volatile Runnable askTimeoutRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #334): install {@link #askTimeoutRaceHookForTest}. Package-private so the
|
||||
* test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setAskTimeoutRaceHookForTest(Runnable hook) {
|
||||
this.askTimeoutRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
||||
private boolean hasAsyncQuestion(String target) {
|
||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
||||
|
||||
+192
@@ -0,0 +1,192 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* Guards against a "defect factory" built into this file's own established pattern, found live
|
||||
* while building the (parked) idle-sleep-guard PR: every time a component is added to
|
||||
* {@link FleetConfig}, the record grows by one arg AND a new back-compat constructor is added at
|
||||
* the OLD arity, so existing callers keep compiling. That is correct and required — see the
|
||||
* constructor ladder just below the record header. But {@link #withDefaults()}'s own {@code return
|
||||
* new FleetConfig(...)} call sits in this same file, written at a literal argument count. The very
|
||||
* next time a component is added, the freshly-added back-compat constructor at the OLD arity
|
||||
* silently captures that stale call, because it is now a legal overload at that arg count too. It
|
||||
* compiles. Every other test passes, because nothing else exercises the new field. The new
|
||||
* component is defaulted away — {@code null}, or whatever that back-compat overload defaults it to
|
||||
* — on every {@link FleetConfig#load}. Measured, not theoretical: this exact sequence happened
|
||||
* live when the {@code idleSleepGuard} component was added on a sibling branch; it was caught only
|
||||
* because that branch's own new tests happened to assert on the new field's value.
|
||||
*
|
||||
* <p>This test proves the opposite property, and does it in a way that survives the next field
|
||||
* being added without being rewritten: reflectively enumerate {@link FleetConfig}'s own record
|
||||
* components (never a hardcoded count — the arity is exactly what changes over time), build one
|
||||
* config through the true canonical constructor with a real, distinctive, non-null value in EVERY
|
||||
* component (reusing the exact reflective-construction pattern
|
||||
* {@link ConfigRefTopLevelReportingCoverageTest} already established for this file:
|
||||
* {@code getDeclaredConstructor(exact record-component types)}, which resolves the canonical
|
||||
* constructor by its true shape, not by binding to whichever overload happens to match arg count —
|
||||
* the same way Jackson resolves it), call the real {@link FleetConfig#withDefaults()}, and assert
|
||||
* every one of those values survives unchanged.
|
||||
*
|
||||
* <p>Why this is a valid check for every component, not just some: {@link #withDefaults()}'s own
|
||||
* comments document that it only ever REPLACES a component when the incoming value is {@code null}
|
||||
* (or blank, for {@code placement}) — {@code broker}/{@code primary}/{@code leadHeartbeat}/
|
||||
* {@code configReload}/{@code coordinator}/{@code worktreeGroup}/{@code memberLoginShell} are left
|
||||
* as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
|
||||
* {@code fleet}/{@code quarantineCooldownSeconds}/{@code memberCredentials}/{@code placement} are
|
||||
* replaced only on null/blank input. A value that is never null or blank going in must therefore
|
||||
* never change coming out, for every current component. No exclusion is needed today.
|
||||
*
|
||||
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} exists anyway, kept deliberately empty and size-pinned
|
||||
* by {@link #exclusionListSizeIsPinned()}: a future component that {@code withDefaults()} is
|
||||
* <em>documented</em> to transform unconditionally (unlike every field today) would legitimately
|
||||
* need one. Pinning the size at 0 means growing that set to make a failure go away is itself a
|
||||
* visible diff to this test, not a silent one — a checker that can be silenced by adding to its
|
||||
* own escape hatch is not a checker.
|
||||
*/
|
||||
class FleetConfigWithDefaultsPreservesEveryComponentTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
|
||||
|
||||
/** See the class javadoc — deliberately empty today; grow it only with a matching justification. */
|
||||
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
|
||||
|
||||
/** One real, distinctive, non-null (non-blank where blankness would mean "unset") value per component. */
|
||||
private static Map<String, Object> baseValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("bind", new FleetConfig.Bind("127.0.0.1", 8765));
|
||||
v.put("herdrSocket", "~/.config/herdr/guard.sock");
|
||||
v.put("memberHerdrSocket", "~/.config/herdr/member-guard.sock");
|
||||
v.put("profiles", Map.of("sonnet", minimalProfile("sonnet")));
|
||||
v.put("guard", new FleetConfig.Guard(List.of("host-guard")));
|
||||
v.put("worktreeRoot", "/wt/guard");
|
||||
v.put("lifecycle", new FleetConfig.Lifecycle(300, 5, 30, true));
|
||||
v.put("spawnReadyTimeoutMs", 12_345);
|
||||
v.put("spawnReadyPollMs", 234);
|
||||
v.put("broker", new FleetConfig.Broker("amqp://guard", null, 7));
|
||||
v.put("primary", new FleetConfig.Primary("term-guard", 4, 4000));
|
||||
v.put("fleet", new FleetConfig.Fleet(
|
||||
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-guard", 1, null, 10,
|
||||
"claude", null, null, null)),
|
||||
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
|
||||
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(301, 61_000L, 4));
|
||||
v.put("health", new FleetConfig.Health(true, 31, 601, 61, null));
|
||||
v.put("placement", "round-robin");
|
||||
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
|
||||
v.put("configReload", new FleetConfig.ConfigReload(true, 11));
|
||||
v.put("quarantineCooldownSeconds", 1801);
|
||||
v.put("memberCredentials", new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
|
||||
List.of("git"), List.of("git", "ssh"), null));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3));
|
||||
v.put("worktreeGroup", "group-guard");
|
||||
v.put("memberLoginShell", "/bin/zsh");
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
/** A minimal, otherwise-null {@link FleetConfig.Profile} — just enough to name one in a map. */
|
||||
private static FleetConfig.Profile minimalProfile(String name) {
|
||||
return new FleetConfig.Profile(name, null, null, null, null, null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null, null, null, null,
|
||||
null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Guards {@link #baseValues()} itself against drifting from the record's real shape — the same
|
||||
* assurance {@link ConfigRefTopLevelReportingCoverageTest} already relies on. This is what makes
|
||||
* "no hardcoded arity" true in practice: forgetting to add a new component here fails this
|
||||
* assertion by name, rather than silently checking one component fewer than the record has.
|
||||
*/
|
||||
private static void assertNamesMatchComponents(Map<String, Object> values) {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
names.add(rc.getName());
|
||||
}
|
||||
assertEquals(names, new TreeSet<>(values.keySet()),
|
||||
"this test's value map has drifted from FleetConfig's actual top-level components — "
|
||||
+ "update baseValues() alongside the record");
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a {@link FleetConfig} through the TRUE canonical constructor — resolved by the record's
|
||||
* own component types, not by argument count — so this never accidentally exercises a
|
||||
* back-compat overload the way a literal {@code new FleetConfig(...)} call risks doing.
|
||||
*/
|
||||
private static FleetConfig configOf(Map<String, Object> values) throws ReflectiveOperationException {
|
||||
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
|
||||
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
|
||||
Constructor<FleetConfig> ctor = FleetConfig.class.getDeclaredConstructor(types);
|
||||
return ctor.newInstance(args);
|
||||
}
|
||||
|
||||
@Test
|
||||
void exclusionListSizeIsPinned() {
|
||||
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
|
||||
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
|
||||
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
|
||||
+ "growing exclusion list that silences failures on its own is not a guard");
|
||||
}
|
||||
|
||||
/**
|
||||
* The mutation this is built to catch: make {@code withDefaults()}'s final constructor call
|
||||
* literal at some arg count, add one more component to the record with a new back-compat
|
||||
* constructor at the old arity, and the stale call silently rebinds. Every component here is
|
||||
* real and non-null (non-blank for the one String — {@code placement} — where blank has
|
||||
* meaning), so none of it should be replaced by {@code withDefaults()}; any component that
|
||||
* comes back different was silently dropped.
|
||||
*/
|
||||
@Test
|
||||
void everyComponentGivenARealValueSurvivesWithDefaults() throws ReflectiveOperationException {
|
||||
Map<String, Object> base = baseValues();
|
||||
FleetConfig config = configOf(base);
|
||||
FleetConfig defaulted = config.withDefaults();
|
||||
|
||||
List<String> dropped = new ArrayList<>();
|
||||
int checked = 0;
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
String name = rc.getName();
|
||||
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
|
||||
continue;
|
||||
}
|
||||
checked++;
|
||||
Object expected = base.get(name);
|
||||
Object actual;
|
||||
try {
|
||||
actual = rc.getAccessor().invoke(defaulted);
|
||||
} catch (ReflectiveOperationException e) {
|
||||
throw new RuntimeException("failed to read FleetConfig." + name + "()", e);
|
||||
}
|
||||
if (!Objects.equals(expected, actual)) {
|
||||
dropped.add(String.format(Locale.ROOT,
|
||||
"%s: withDefaults() was given a real, non-null value (%s) for '%s' but "
|
||||
+ "returned %s — a component silently dropped by withDefaults(), the "
|
||||
+ "shape of the defect this test exists to catch (its final "
|
||||
+ "\"return new FleetConfig(...)\" call binding to a back-compat "
|
||||
+ "constructor instead of the true canonical one)",
|
||||
name, expected, name, actual));
|
||||
}
|
||||
}
|
||||
|
||||
System.out.printf(Locale.ROOT,
|
||||
"FleetConfig.withDefaults() component-survival coverage — %d components, %d checked, "
|
||||
+ "%d excluded, %d survived%n",
|
||||
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
|
||||
assertEquals(List.of(), dropped,
|
||||
"withDefaults() silently dropped these real, given components: " + dropped);
|
||||
}
|
||||
}
|
||||
@@ -484,6 +484,27 @@ class CompletionResolverTest {
|
||||
|
||||
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
|
||||
|
||||
@Test
|
||||
void aNormalMemberReportMentioningTheExhaustionPatternDoesNotNotifyTheSink() {
|
||||
String block = "⏺ I reviewed capacity handling. The usage limit has been reached means no more work can start.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
|
||||
"a matching report still fails the send as exhausted");
|
||||
assertTrue(waiter.getNow(null).text().contains("I reviewed capacity handling."),
|
||||
"the exhausted result keeps the whole matched pane line");
|
||||
assertTrue(notified.isEmpty(),
|
||||
"a normal report mentioning an exhaustion pattern must not quarantine a credential");
|
||||
}
|
||||
|
||||
@Test
|
||||
void classifiesAMatchingScrapeAsBackendExhaustedInsteadOfACompletedReply() {
|
||||
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
|
||||
@@ -534,6 +555,53 @@ class CompletionResolverTest {
|
||||
"the sink is told the matched reason: " + notified.get(0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRealExhaustionBehindTerminalChromeStillNotifiesTheSink() {
|
||||
String block = "⏺ │ The usage limit has been reached. Try again later.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
|
||||
"a real exhaustion must still fail the send as exhausted");
|
||||
assertEquals(1, notified.size(),
|
||||
"a real exhaustion behind terminal chrome must reach the sink");
|
||||
}
|
||||
|
||||
/**
|
||||
* The live fleet configures {@code exhaustedPattern: "The usage limit has been reached"} — with
|
||||
* the leading {@code "The"}. Every other test here uses a pattern without it, which is the shape
|
||||
* that made fleetd #348 need a looser rule than a start-of-line check. This pins the deployed
|
||||
* shape as well, so a later tightening of {@link CompletionResolver} cannot silently stop
|
||||
* recording the exhaustion this fleet actually reports.
|
||||
*
|
||||
* <p>What it does not prove: that this is the only pattern shape an operator will write.
|
||||
*/
|
||||
@Test
|
||||
void anExhaustionPatternCarryingItsLeadingWordsStillNotifiesTheSink() {
|
||||
String block = "⏺ │ The usage limit has been reached. Try again later.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("The usage limit has been reached");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
|
||||
"the live pattern shape must still fail the send as exhausted");
|
||||
assertEquals(1, notified.size(),
|
||||
"the live pattern shape must still reach the sink");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() {
|
||||
// The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed —
|
||||
|
||||
@@ -377,6 +377,126 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #334. {@code ask()}'s {@code TimeoutException} catch used to forget this task's
|
||||
* {@code turnId} mapping ({@code clearAsyncQuestion(turnId, true)}) BEFORE closing the ask
|
||||
* ({@code rendezvous.closeAsk}, in the shared {@code finally}). A primary's {@code answer()}
|
||||
* call racing that exact window found the ask still "answerable" ({@code
|
||||
* rendezvous.askSession(turnId)} still non-null) while the {@code Task} was already forgotten,
|
||||
* so its {@code task != null} guard skipped the completion and the async ticket sat at
|
||||
* {@code PENDING} forever even though {@code answer()} itself reported a result. The fix
|
||||
* (closing the ask first) makes this window impossible: a racing {@code answer()} call either
|
||||
* still finds the ask open (and the {@code Task} mapping guaranteed intact) or finds it already
|
||||
* closed (and bails {@code STALE_TURN} before ever touching the {@code Task}). This test pins
|
||||
* the exact window with {@code askTimeoutRaceHookForTest} and proves both invariants the ticket
|
||||
* named: (1) a late/racing {@code answer()} sees the ask as already lapsed ({@code STALE_TURN}),
|
||||
* never made answerable again, and (2) the async ticket still resolves {@code DONE} once the
|
||||
* worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}.
|
||||
*/
|
||||
/**
|
||||
* fleetd #334 gated the ask-timeout teardown on {@code ticket.fresh()}, matching the {@code
|
||||
* finally} block that already did. This pins that gate. A coalesced duplicate passes its own
|
||||
* {@code timeoutMillis}, which says nothing about whether the shared ask is done — so a
|
||||
* duplicate timing out first must leave the fresh owner's still-open ask answerable.
|
||||
*
|
||||
* <p>Measured on merge: without this test, removing the {@code ticket.fresh()} gate left all
|
||||
* 1371 tests green. The gate shipped with the reorder and nothing held it there.
|
||||
*
|
||||
* <p>What this does not prove: anything about the ordering inside the gate — that is
|
||||
* {@code aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket}'s job.
|
||||
*/
|
||||
@Test
|
||||
void aCoalescedDuplicateAskTimingOutLeavesTheFreshOwnersAskOpen() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
||||
|
||||
CompletableFuture<MessageService.AskResult> fresh =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
|
||||
MessageService.TaskView asking = null;
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
asking = messages.poll(ticket);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertNotNull(asking, "the fresh owner's question must surface before the duplicate asks");
|
||||
String turnId = asking.turnId();
|
||||
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
|
||||
|
||||
// A coalesced duplicate on the same session, with its own much shorter timeout.
|
||||
MessageService.AskResult duplicate = messages.ask(T, "which config file?", 100);
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, duplicate.outcome(),
|
||||
"the duplicate's own timeout elapses first");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answered =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500));
|
||||
|
||||
MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(),
|
||||
"a duplicate's timeout must not lapse the ask the fresh owner still holds");
|
||||
assertEquals("fleetd.yaml", a.answer());
|
||||
assertFalse(answered.get(5, TimeUnit.SECONDS).outcome() == MessageService.Outcome.STALE_TURN,
|
||||
"the answer must not be rejected as stale");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 200));
|
||||
|
||||
// Wait for the question to actually surface (poll sees ASKING) before racing the timeout.
|
||||
MessageService.TaskView asking = null;
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
asking = messages.poll(ticket);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertNotNull(asking, "the question must surface before the ask times out");
|
||||
String turnId = asking.turnId();
|
||||
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
|
||||
|
||||
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
|
||||
messages.setAskTimeoutRaceHookForTest(() ->
|
||||
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
|
||||
try {
|
||||
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
|
||||
|
||||
MessageService.Reply late = lateAnswer.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.STALE_TURN, late.outcome(),
|
||||
"a late answer racing the timeout teardown must see the ask as already lapsed");
|
||||
|
||||
// The worker resumes on its own (per the ask() contract) and eventually sends its real
|
||||
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
|
||||
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
|
||||
} finally {
|
||||
messages.setAskTimeoutRaceHookForTest(null);
|
||||
}
|
||||
|
||||
MessageService.TaskView done = null;
|
||||
deadline = System.currentTimeMillis() + 2000;
|
||||
while ((done == null || done.phase() == MessageService.Phase.PENDING)
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
done = messages.poll(ticket);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertNotNull(done);
|
||||
assertEquals(MessageService.Phase.DONE, done.phase(), "the async ticket must not be stranded PENDING");
|
||||
assertEquals("real result", done.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void answeringAnUnknownTurnIsStale() {
|
||||
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
|
||||
@@ -410,6 +530,29 @@ class MessageServiceTest {
|
||||
"a delivered send whose worker never replies times out as still working");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send}
|
||||
* first observes its completion as incomplete, but before {@code cancel} takes the target monitor.
|
||||
* The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report
|
||||
* {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed.
|
||||
*
|
||||
* <p>What this does not prove: that this precise interleaving happens by itself under production
|
||||
* timing. The test forces it through a test-only hook; it proves the timeout caller handles the
|
||||
* injector result when the interleaving occurs.
|
||||
*/
|
||||
@Test
|
||||
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
|
||||
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
|
||||
try {
|
||||
MessageService.Reply reply = messages.send(T, "race delivery", 50);
|
||||
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
|
||||
"cancel reporting DELIVERED means the worker received the timed-out message");
|
||||
} finally {
|
||||
messages.setTimeoutCancellationRaceHookForTest(null);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
@@ -785,6 +928,49 @@ class MessageServiceTest {
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
}
|
||||
|
||||
// --- fleetd #335 (site 1): a per-task cleanup failure inside the abandon() loop must not -----
|
||||
// strand the tasks that come after it. abandon()'s own comment on the loop documents the one
|
||||
// real production call that can throw there (inbox.publish, in the stranded-reply put-back
|
||||
// branch, reached when a concurrent reply() or a second abandon() races this one) — but
|
||||
// reaching that branch requires the exact combination the #137 follow-up above already found
|
||||
// unreachable through the public API. abandonCleanupHookForTest reproduces the resulting SHAPE
|
||||
// (one task's cleanup throws) directly instead, the same technique this file already uses for
|
||||
// fleetd #324/#329's own hard-to-time races.
|
||||
@Test
|
||||
void aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks() throws Exception {
|
||||
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
|
||||
try {
|
||||
String first = messages.sendAsync(T, "first task");
|
||||
awaitWaiting(); // first task owns the target lock and rendezvous waiter
|
||||
String second = messages.sendAsync(T, "second task"); // parked on the same lock
|
||||
String third = messages.sendAsync(T, "third task"); // parked too — the whole sweep must survive
|
||||
|
||||
java.util.concurrent.atomic.AtomicInteger calls = new java.util.concurrent.atomic.AtomicInteger();
|
||||
messages.setAbandonCleanupHookForTest(() -> {
|
||||
if (calls.getAndIncrement() == 0) {
|
||||
throw new RuntimeException("PROBE-335-SITE1");
|
||||
}
|
||||
});
|
||||
|
||||
assertTrue(messages.abandon(T, "agent target term_a not found"));
|
||||
|
||||
// Every task in the loop still gets its own outcome — the one whose cleanup threw
|
||||
// included — even though the loop had no way to know in advance which one that would be.
|
||||
assertFailedTicket(first, "agent target term_a not found");
|
||||
assertFailedTicket(second, "agent target term_a not found");
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.ERROR
|
||||
&& e.getThrowableProxy() != null
|
||||
&& "PROBE-335-SITE1".equals(e.getThrowableProxy().getMessage())),
|
||||
"a per-task cleanup failure must still reach the log, not vanish silently");
|
||||
} finally {
|
||||
messages.setAbandonCleanupHookForTest(null);
|
||||
detachMessageServiceLog(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- #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
|
||||
@@ -1472,6 +1658,42 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
// --- fleetd #335 (site 2): task.future.whenComplete's own returned stage is discarded, so an --
|
||||
// uncaught throw from ReplyPushLoop.onTicketTerminal used to vanish with no log line and no
|
||||
// metric. The real production trigger is the daemon's own shutdown sequence (Fleetd's shutdown
|
||||
// hook): messages.close() only stops the async executor from taking NEW work — it does not
|
||||
// cancel a send already in flight — while pushLoop.close() shuts its scheduler down immediately
|
||||
// right after, so a ticket that completes in that narrow window has onTicketTerminal's own
|
||||
// scheduler.schedule(...) throw a real RejectedExecutionException. Reproduced here by shutting
|
||||
// the very same scheduler down before the ticket resolves — no test-only hook needed, this
|
||||
// reachable path throws for real.
|
||||
@Test
|
||||
void aTicketTerminalPushFailureDoesNotVanishSilently() throws Exception {
|
||||
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
|
||||
try (var wiring = wireWithPushLoop(1, 50)) {
|
||||
String ticket = wiring.service().sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
wiring.scheduler().shutdownNow(); // simulate pushLoop.close() racing an in-flight send
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "async result"));
|
||||
|
||||
// The ticket's own outcome must be unaffected by the swallowed exception — finishAsyncTask
|
||||
// completes task.future before whenComplete's action (and thus onTicketTerminal) ever runs.
|
||||
MessageService.TaskView done = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE);
|
||||
assertEquals("async result", done.reply());
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.ERROR
|
||||
&& e.getFormattedMessage().contains(ticket)
|
||||
&& e.getThrowableProxy() != null
|
||||
&& "java.util.concurrent.RejectedExecutionException"
|
||||
.equals(e.getThrowableProxy().getClassName())),
|
||||
"onTicketTerminal throwing must still reach the log, not vanish silently");
|
||||
} finally {
|
||||
detachMessageServiceLog(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user