Compare commits

...

13 Commits

Author SHA1 Message Date
Dai Ha 23f299e105 Preserve strict AMQP exception handling
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m46s
2026-09-05 05:53:44 +07:00
Dai Ha d292522d00 Name AMQP connection failure logs
CI / contract (pull_request) Successful in 1m20s
CI / build (pull_request) Successful in 1m29s
2026-09-05 05:45:37 +07:00
Dai Ha b6b88c5f1c #334: pin the fresh-owner gate on ask()'s timeout teardown
CI / build (push) Successful in 2m7s
CI / contract (push) Successful in 34m37s
2026-09-04 17:12:49 +07:00
Dai Ha 86dddfe240 Merge #334: ask()'s timeout closes the turn before it forgets the task mapping 2026-09-04 17:08:40 +07:00
Dai Ha 0d5944af63 fleetd #334: close ask()'s turn before forgetting its Task, closing the last stranding window
CI / build (pull_request) Successful in 2m14s
CI / contract (pull_request) Successful in 19m14s
ask()'s TimeoutException catch used to run clearAsyncQuestion(turnId, true) -- forgetting the
Task's asyncTasksByTurn mapping -- before rendezvous.closeAsk(turnId) ran in the shared finally.
Between those two calls the ask was still "answerable" (askSession(turnId) non-null) but the Task
mapping was already gone, so a racing answer() call found task == null, skipped
finishAsyncTask, and stranded the async ticket at PENDING even though answer() itself reported a
result. #329 fixed one step of this same race; this closes the remaining one.

The fix reorders the fresh owner's teardown: closeAsk runs first, then markAskTimedOut and
clearAsyncQuestion. A racing answer() call now either sees the ask still open (and the Task
mapping guaranteed intact) or sees it already closed (STALE_TURN, before it ever reaches
asyncTasksByTurn). It also gates the whole block by ticket.fresh(), matching the invariant the
finally block already states ("only the fresh owner tears down the shared turn") -- a duplicate
coalesced ask() timing out no longer forgets bookkeeping the fresh owner still needs.

Adds aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket, which pins the exact window
with a new test-only hook (askTimeoutRaceHookForTest) and proves both invariants: a late answer()
racing the timeout sees STALE_TURN, and the async ticket still resolves DONE from the worker's
real reply. Reverting the reorder (verified locally, not committed) makes this test fail with
"expected STALE_TURN but was TIMED_OUT_WORKING".
2026-09-04 17:06:15 +07:00
Dai Ha 4a5030a5c6 #348: drop a chrome skip that cannot fire, and pin the live pattern shape
CI / contract (push) Successful in 2m10s
CI / build (push) Successful in 2m12s
2026-09-04 16:57:36 +07:00
Dai Ha c1c8794c48 Merge #348: a member's prose about a usage limit no longer quarantines a credential 2026-09-04 16:53:19 +07:00
Dai Ha 65a78932c1 Merge #335: a per-task cleanup throw in abandon() no longer strands the tasks behind it
CI / contract (push) Successful in 1m24s
CI / build (push) Successful in 2m8s
2026-09-04 16:44:04 +07:00
Dai Ha f429ca1a50 Avoid exhaustion cooldown for member prose 2026-09-04 16:43:16 +07:00
Dai Ha 73aab3f83e Merge #342: teardown resolves the pane's real tab instead of trusting the delegate's placement config
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m47s
2026-09-04 16:38:31 +07:00
Dai Ha 887aca0183 fleetd#335: abandon()'s per-task cleanup and sendAsync's terminal hook must not swallow throws
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 2m21s
Site 1 (abandon()'s matching loop, reachable): the recovery/put-back branch calls
inbox.publish, which AmqpReplyInbox implements as a real broker round trip that
throws IllegalStateException on an unroutable/unconfirmed/interrupted publish.
An uncaught throw there aborted the loop, stranding every task after it in
`matching` PENDING forever. Fixed by recording each task's own future.complete()
result before any cleanup runs, then wrapping the cleanup in try/catch so one
task's failure cannot stop its siblings from getting their outcome. Reaching the
throwing branch by real timing needs a race the file's own #137 follow-up already
found unreachable through the public API, so the reproducing test uses a
test-only hook (same technique as the existing fleetd #324/#329 hooks) to inject
the throw at that exact point.

Site 2 (sendAsync's task.future.whenComplete, reachable): the returned stage is
discarded, so an uncaught throw from pushLoop.onTicketTerminal vanished with no
log line. Reproduced for real: Fleetd's shutdown hook runs messages.close()
(stops the async executor from taking new work, but does not cancel a send
already in flight) before pushLoop.close() (shuts its scheduler down
immediately) — a ticket completing in that window makes onTicketTerminal's own
scheduler.schedule(...) throw a genuine RejectedExecutionException. Fixed with a
try/catch(Throwable) plus log.error inside the whenComplete action.

Site 3 (the two `finally { asyncTasksByWaiter.remove(reply); rendezvous.close(...);
}` blocks in send() and answer()): read Rendezvous.close/closeAsk and the
ConcurrentHashMap operations behind them — both are plain map ops on a non-null
key with no user-overridable code, so neither can throw. Left unchanged; not a
defect.

Mutation-proven: reverting either fix reproduces the failure it exists to catch
— removing site 1's try/catch aborts abandon() with the injected exception
(MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks
errors); removing site 2's try/catch leaves the RejectedExecutionException
unlogged (MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently
fails its log assertion). Full suite: mvn clean install, Tests run: 1365,
Failures: 0, Errors: 0, BUILD SUCCESS.
2026-09-04 16:38:14 +07:00
Dai Ha f379847942 Merge #345: the timeout path's use of Cancellation.DELIVERED is now pinned
CI / contract (push) Successful in 47s
CI / build (push) Successful in 2m8s
2026-09-04 16:32:45 +07:00
Dai Ha ea12107497 fleetd #345: test timeout cancellation race
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m3s
2026-09-04 16:30:46 +07:00
7 changed files with 714 additions and 54 deletions
@@ -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.
*
@@ -8,6 +8,7 @@ import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import com.rabbitmq.client.impl.DefaultExceptionHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -153,17 +154,22 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
public static AmqpReplyInbox open(String uri, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new AmqpReplyInbox(factory.newConnection("fleetd-reply-inbox"), prefetch);
return new AmqpReplyInbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.REPLY_INBOX), prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
}
}
static ConnectionFactory connectionFactory(String uri) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
factory.setExceptionHandler(new AmqpConnectionFailureLogger(AmqpConnectionFailureLogger.REPLY_INBOX, log));
return factory;
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
AmqpReplyInbox(Connection connection) {
this(connection, DEFAULT_PREFETCH);
@@ -571,3 +577,44 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
}
/**
* Keeps RabbitMQ's forgiving exception behaviour while adding the connection identity that its
* default logger drops. Package-private so both AMQP connections use the same two names.
*/
final class AmqpConnectionFailureLogger extends DefaultExceptionHandler {
static final String REPLY_INBOX = "fleetd-reply-inbox";
static final String LEAD_MAILBOX = "fleetd-lead-mailbox";
private final String connectionName;
private final Logger logger;
AmqpConnectionFailureLogger(String connectionName, Logger logger) {
this.connectionName = connectionName;
this.logger = logger;
}
String connectionName() {
return connectionName;
}
@Override
protected void log(String message, Throwable cause) {
if (isSocketClosedOrConnectionReset(cause)) {
logger.warn("AMQP connection {}: {} (Exception message: {})", connectionName, message, cause.getMessage());
} else {
logger.error("AMQP connection {}: {}", connectionName, message, cause);
}
}
private static boolean isSocketClosedOrConnectionReset(Throwable cause) {
// Deliberate copy of ForgivingExceptionHandler's private static helper; check it on amqp-client upgrades.
if (!(cause instanceof IOException)) {
return false;
}
return "Connection reset".equals(cause.getMessage())
|| "Socket closed".equals(cause.getMessage())
|| "Connection reset by peer".equals(cause.getMessage());
}
}
@@ -122,17 +122,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new LeadMailbox(factory.newConnection("fleetd-lead-mailbox"), selfCoordId, prefetch);
return new LeadMailbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.LEAD_MAILBOX), selfCoordId, prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
}
}
static ConnectionFactory connectionFactory(String uri) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
factory.setExceptionHandler(new AmqpConnectionFailureLogger(AmqpConnectionFailureLogger.LEAD_MAILBOX, log));
return factory;
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId) {
this(connection, selfCoordId, DEFAULT_PREFETCH);
@@ -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));
@@ -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 —
@@ -0,0 +1,134 @@
package dev.ltms.fleet.msg;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.impl.DefaultExceptionHandler;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.lang.reflect.Proxy;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class AmqpConnectionFailureLoggerTest {
@Test
void installedHandlersLogTheirOwnConnectionNamesAtErrorWithTheCause() throws Exception {
ConnectionFactory inboxFactory = AmqpReplyInbox.connectionFactory("amqp://127.0.0.1");
ConnectionFactory mailboxFactory = LeadMailbox.connectionFactory("amqp://127.0.0.1");
AmqpConnectionFailureLogger inboxHandler = installedStrictHandler(inboxFactory, "reply inbox");
AmqpConnectionFailureLogger mailboxHandler = installedStrictHandler(mailboxFactory, "lead mailbox");
assertEquals(AmqpConnectionFailureLogger.REPLY_INBOX, inboxHandler.connectionName());
assertEquals(AmqpConnectionFailureLogger.LEAD_MAILBOX, mailboxHandler.connectionName());
ListAppender<ILoggingEvent> inboxEvents = attach(AmqpReplyInbox.class);
ListAppender<ILoggingEvent> mailboxEvents = attach(LeadMailbox.class);
IllegalStateException inboxFailure = new IllegalStateException("inbox failure");
IllegalStateException mailboxFailure = new IllegalStateException("mailbox failure");
try {
inboxHandler.handleUnexpectedConnectionDriverException(null, inboxFailure);
mailboxHandler.handleConnectionRecoveryException(null, mailboxFailure);
assertError(inboxEvents, "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
inboxFailure, "inbox failure line");
assertError(mailboxEvents, "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
mailboxFailure, "mailbox recovery line");
} finally {
detach(AmqpReplyInbox.class, inboxEvents);
detach(LeadMailbox.class, mailboxEvents);
}
}
@Test
void connectionResetKeepsForgivingHandlerWarningSemantics() {
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
ListAppender<ILoggingEvent> events = attach(AmqpReplyInbox.class);
try {
handler.handleUnexpectedConnectionDriverException(null, new IOException("Connection reset"));
assertEquals(1, events.list.size(), "the handler must still log a reset");
ILoggingEvent event = events.list.getFirst();
assertEquals(Level.WARN, event.getLevel(), "ForgivingExceptionHandler logs connection resets at WARN");
assertEquals("AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred "
+ "(Exception message: Connection reset)", event.getFormattedMessage());
assertTrue(event.getThrowableProxy() == null, "ForgivingExceptionHandler does not attach a reset stack trace");
} finally {
detach(AmqpReplyInbox.class, events);
}
}
@Test
void connectionNamesStayDistinct() {
assertNotEquals(AmqpConnectionFailureLogger.REPLY_INBOX, AmqpConnectionFailureLogger.LEAD_MAILBOX,
"reply-inbox and lead-mailbox failures must be distinguishable");
}
@Test
void strictConsumerExceptionStillClosesItsChannel() {
AtomicInteger closes = new AtomicInteger();
Channel channel = (Channel) Proxy.newProxyInstance(getClass().getClassLoader(), new Class<?>[] {Channel.class},
(_, method, _) -> switch (method.getName()) {
case "close" -> {
closes.incrementAndGet();
yield null;
}
case "toString" -> "test-channel";
default -> throw new UnsupportedOperationException(method.getName());
});
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
handler.handleConsumerException(channel, new IllegalStateException("consumer failed"), null, "tag", "handleDelivery");
assertEquals(1, closes.get(), "DefaultExceptionHandler must close a channel after a consumer exception");
}
@Test
void handlerOnlyChangesDefaultHandlerLogging() {
assertEquals(DefaultExceptionHandler.class,
AmqpConnectionFailureLogger.class.getSuperclass());
assertFalse(java.util.Arrays.stream(AmqpConnectionFailureLogger.class.getDeclaredMethods())
.anyMatch(method -> method.getName().startsWith("handle")),
"all exception-handling methods must remain inherited from DefaultExceptionHandler");
}
private static ListAppender<ILoggingEvent> attach(Class<?> owner) {
Logger logger = (Logger) LoggerFactory.getLogger(owner);
logger.setLevel(Level.DEBUG);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(Class<?> owner, ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(owner)).detachAppender(appender);
}
private static AmqpConnectionFailureLogger installedStrictHandler(ConnectionFactory factory, String connection) {
assertInstanceOf(DefaultExceptionHandler.class, factory.getExceptionHandler(),
connection + " must keep DefaultExceptionHandler: replacing the strict handler with a forgiving one "
+ "changes when a channel is closed");
return assertInstanceOf(AmqpConnectionFailureLogger.class, factory.getExceptionHandler());
}
private static void assertError(ListAppender<ILoggingEvent> events, String message, Throwable cause, String name) {
assertEquals(1, events.list.size(), name);
ILoggingEvent event = events.list.getFirst();
assertEquals(Level.ERROR, event.getLevel(), name);
assertEquals(message, event.getFormattedMessage(), name);
assertEquals(cause.toString(), event.getThrowableProxy().getClassName() + ": "
+ event.getThrowableProxy().getMessage(), name);
}
}
@@ -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