Compare commits

...

3 Commits

Author SHA1 Message Date
Dai Ha f429ca1a50 Avoid exhaustion cooldown for member prose 2026-09-04 16:43:16 +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
4 changed files with 119 additions and 11 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;
}
@@ -614,19 +618,38 @@ public final class CompletionResolver implements TurnListener {
}
/**
* True when the error pattern begins the matched pane line, rather than appearing in prose.
* True when the pattern appears before the first sentence ending in the matched pane line, rather
* than after a member has started prose about it.
*
* <p>Leading terminal chrome is skipped first — box-drawing characters, bullets, gutter bars and
* spaces. #339 introduced this check with a bare {@code lookingAt}, and that rejected a genuine
* error line rendered as {@code "| 503 Service Unavailable: ..."}: the send still failed, but the
* credential outage was never recorded. That is the false negative #339's own invariant 3 called
* worse than the false positive it set out to fix — measured with a throwaway probe on the
* raw-scrape path, which is exactly the path whose comment says to expect leading chrome.
* spaces. Exhaustion patterns often name only the decisive words in a provider message, such as
* {@code "usage limit has been reached"}; they do not include its leading {@code "The"}. A bare
* {@code lookingAt} would therefore reject that genuine refusal, including one behind terminal
* chrome.
*
* <p>Skipping only a leading run of non-letter, non-digit characters keeps the fix's intent. A
* member's prose ({@code "I checked the retry path. An API Error: makes it back off."}) still
* does not match, because there the pattern sits after words, not after chrome.
* <p>A member's prose about a refusal normally follows a completed sentence. That sentence ending
* is enough to keep it from reaching a cooldown sink while the send still fails with the full pane
* tail. This is deliberately less strict than the backend-error check because exhausted patterns
* may start after a provider's leading words.
*/
private static boolean startsWithExhaustion(String line, Pattern pattern) {
int i = 0;
while (i < line.length() && !Character.isLetterOrDigit(line.charAt(i))) {
i++;
}
var matcher = pattern.matcher(line);
if (!matcher.find()) {
return false;
}
for (int prefix = i; prefix < matcher.start(); prefix++) {
if (".!?".indexOf(line.charAt(prefix)) >= 0) {
return false;
}
}
return true;
}
/** True when an error pattern begins the matched pane line after optional terminal chrome. */
private static boolean startsWithBackendError(String line, Pattern pattern) {
int i = 0;
while (i < line.length() && !Character.isLetterOrDigit(line.charAt(i))) {
@@ -871,6 +871,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;
@@ -1370,6 +1374,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)
@@ -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,25 @@ 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");
}
@Test
void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() {
// The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed —
@@ -410,6 +410,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();