Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 57b8c0b56d fleetd #339: guard backend error sink
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 1m50s
2026-09-04 15:57:18 +07:00
8 changed files with 70 additions and 164 deletions
@@ -1,10 +1,12 @@
package dev.ltms.fleet.inject;
/**
* Notified when {@link CompletionResolver} actually delivers a typed backend-error classification
* to a waiting send (fleetd#201 / #227) — never on a race that lost. {@link CompletionResolver}
* calls this only after {@code Rendezvous.resolveFailure} returns {@code true} for that exact
* waiter, mirroring the win-only race rule {@link ExhaustionSink} already uses.
* Notified when {@link CompletionResolver} has a backend-error match at the start of a pane line,
* or both a match and its too-fast crash signature, for a waiting send (fleetd#201 / #227). A text
* match inside ordinary pane prose can be a member's report about an error, so it fails the send
* without notifying this sink.
* {@link CompletionResolver} calls this only after {@code Rendezvous.resolveFailure} returns
* {@code true} for that exact waiter, mirroring the win-only race rule {@link ExhaustionSink} uses.
*
* <p>The public send result is unchanged by this classification — it is still a failed send
* ({@code Rendezvous.Kind#FAILED}); this sink is the internal seam a later stage (fleetd#201 Unit
@@ -387,9 +387,10 @@ public final class CompletionResolver implements TurnListener {
// fleetd#164 (part 2) / fleetd#201: a scrape that read cleanly and produced content still
// isn't a real reply when that content is the backend's own rejection (e.g. an HTTP 400
// before the worker did any work). Classify it as a failure naming the member, rather than
// handing the caller a scrape that reads like a completed answer, and — only on the
// resolution that actually wins the race, mirroring the exhaustion sink above — notify the
// typed backend-error sink so a later stage can act on repeated failures.
// handing the caller a scrape that reads like a completed answer. A text match alone is not
// enough to notify the typed backend-error sink: this assistant block can be a member's
// normal prose about an error. A line that starts with the error match is stronger evidence;
// the too-fast path below also has its crash signature before it records a credential failure.
String backendError = firstMatchingLine(assistantBlock, backendErrorPatternOrFallback(target));
if (backendError != null) {
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
@@ -401,9 +402,9 @@ public final class CompletionResolver implements TurnListener {
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race — a late
// duplicate must never double-count one backend failure.
backendErrorSink.onBackendError(target, backendError, reason);
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
backendErrorSink.onBackendError(target, backendError, reason);
}
}
return;
}
@@ -492,8 +493,9 @@ public final class CompletionResolver implements TurnListener {
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback from the raw scrape: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race.
backendErrorSink.onBackendError(target, backendError, reason);
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
backendErrorSink.onBackendError(target, backendError, reason);
}
}
return true;
}
@@ -540,10 +542,10 @@ public final class CompletionResolver implements TurnListener {
* {@link #MIN_TURN_NANOS} — a crash signature (e.g. a backend HTTP 400 before the worker did
* anything) that a bare {@code BUSY -> DONE} transition cannot be told apart from a genuinely
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
* apply, against whatever is on screen right now: a match is a typed failure that notifies
* {@link #backendErrorSink} (only on the resolution that wins the race); a non-match stays the
* original generic too-fast failure, naming the member and both timings, with whatever the pane
* shows appended so the caller sees the cause, not just "it failed".
* apply, against whatever is on screen right now: a match together with the too-fast crash
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
* non-match stays the original generic too-fast failure, naming the member and both timings,
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
*/
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
long elapsedNanos) {
@@ -611,6 +613,11 @@ public final class CompletionResolver implements TurnListener {
return null;
}
/** True when the error pattern begins the matched pane line, rather than appearing in prose. */
private static boolean startsWithBackendError(String line, Pattern pattern) {
return pattern.matcher(line).lookingAt();
}
/**
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
@@ -145,58 +145,8 @@ public final class Injector {
return router != null ? router.agentsFor(target) : agents;
}
/** The result of trying to remove an undelivered message from the injector. */
public enum Cancellation {
CANCELLED,
DELIVERED,
NOT_DELIVERED
}
/**
* An identity handle for one queued delivery. It is the only value accepted by
* {@link #cancel(Delivery)}, so a caller cannot cancel a different message with the same target
* or text.
*/
public static final class Delivery {
private final Pending pending;
private Delivery(Pending pending) {
this.pending = pending;
}
public CompletableFuture<Void> completion() {
return pending.delivered;
}
}
/** A pending message and the future that completes when it has been delivered. */
private static final class Pending {
enum State { QUEUED, DELIVERED, NOT_DELIVERED, CANCELLED }
final String target;
final String text;
final TurnToken token;
final CompletableFuture<Void> delivered;
volatile State state = State.QUEUED; // written under the owning Target monitor
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered) {
this.target = target;
this.text = text;
this.token = token;
this.delivered = delivered;
}
String text() {
return text;
}
TurnToken token() {
return token;
}
CompletableFuture<Void> delivered() {
return delivered;
}
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
}
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
@@ -227,47 +177,15 @@ public final class Injector {
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the
* target" and "queue the message" and orphan it in a target it just removed.
*/
public Delivery enqueue(String target, String text, TurnToken token) {
public CompletableFuture<Void> enqueue(String target, String text, TurnToken token) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(target, text, token, delivered);
Pending p = new Pending(text, token, delivered);
targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
return t;
});
return new Delivery(p);
}
/**
* Cancel this exact queued delivery. The target monitor serializes this operation with
* {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED}
* rather than claiming the message remained queued.
*/
public Cancellation cancel(Delivery delivery) {
Pending p = delivery.pending;
Target t = targets.get(p.target);
if (t == null) {
return cancellationOf(p);
}
synchronized (t) {
if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) {
return cancellationOf(p);
}
p.state = Pending.State.CANCELLED;
if (isQuiescent(t)) {
targets.remove(p.target, t);
}
return Cancellation.CANCELLED;
}
}
private static Cancellation cancellationOf(Pending p) {
return p.state == Pending.State.DELIVERED ? Cancellation.DELIVERED : Cancellation.NOT_DELIVERED;
}
private static boolean isQuiescent(Target t) {
return t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion
&& !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved;
return delivered;
}
/**
@@ -356,7 +274,6 @@ public final class Injector {
try {
agentsFor(target).send(target, p.text());
t.queue.poll();
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
t.turnObserved = false;
@@ -366,7 +283,6 @@ public final class Injector {
// Delivery failed at herdr; drop the poisoned message and surface it
// rather than blocking the queue behind it.
t.queue.poll();
p.state = Pending.State.NOT_DELIVERED;
sent = p;
sendError = e;
}
@@ -377,9 +293,6 @@ public final class Injector {
// fail every queued message and release the target (CB-114) instead of
// polling it indefinitely with the caller's future never completing.
notReady = new ArrayList<>(t.queue);
for (Pending pending : notReady) {
pending.state = Pending.State.NOT_DELIVERED;
}
log.warn("readiness grace for {} expired after {} polls ({}s): target never "
+ "became deliverable, so failing {} queued message(s) that never "
+ "reached its pane",
@@ -427,7 +340,8 @@ public final class Injector {
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
// completion awaited), so the map cannot grow without bound across short-lived workers.
if (isQuiescent(t)) {
if (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion
&& !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved) {
targets.remove(target, t);
}
}
@@ -518,9 +432,6 @@ public final class Injector {
boolean hadDeliveredTurn;
synchronized (t) {
pending = new ArrayList<>(t.queue);
for (Pending p : pending) {
p.state = Pending.State.NOT_DELIVERED;
}
t.queue.clear();
hadDeliveredTurn = t.awaitingCompletion;
t.awaitingCompletion = false;
@@ -863,22 +863,16 @@ public final class MessageService {
if (onAccepted != null) {
onAccepted.run();
}
Injector.Delivery delivery = injector.enqueue(target, content, token);
CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
} catch (TimeoutException e) {
boolean wasDelivered = delivery.completion().isDone()
&& !delivery.completion().isCompletedExceptionally();
if (!wasDelivered) {
// 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;
}
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
if (!wasDelivered) {
// CB-640: record that delivery did not happen for fleet health (see
// queuedDeliveries). The exact Pending was cancelled, so it cannot arrive later.
// CB-640: still sitting in the injector's queue, waiting for the member to
// go idle — record the fact for fleet health (see queuedDeliveries).
queuedDeliveries.put(target, Boolean.TRUE);
}
return recorded(new Reply(
@@ -804,7 +804,28 @@ class CompletionResolverTest {
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
@Test
void aConfiguredBackendErrorPatternClassifiesAMatchAsAFailureAndNotifiesTheSinkOnce() {
void aNormalMemberReportMentioningTheFallbackErrorPatternFailsButDoesNotNotifyTheSink() {
String block = "⏺ I checked the retry path. An API Error: makes it back off.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), BackendErrorPatternLookup.legacy(), sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a scrape mentioning the pattern must still fail the send");
assertTrue(waiter.getNow(null).text().contains("I checked the retry path. An API Error: makes it back off."),
"the failure must keep the whole pane tail");
assertTrue(notified.isEmpty(),
"a normal report mentioning the fallback pattern must not record a credential failure");
}
@Test
void aConfiguredBackendErrorPatternAtTheStartOfALineClassifiesAMatchAndNotifiesTheSinkOnce() {
String block = "⏺ 503 Service Unavailable: upstream credential rejected\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
@@ -924,7 +945,7 @@ class CompletionResolverTest {
}
@Test
void aConfiguredPatternAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
void aConfiguredPatternAtTheStartOfALineAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
// No ⏺ marker and leading TUI chrome ⇒ lastAssistantBlock() yields "", so classification must
// fall back to the raw scrape (fleetd#211) — and it must use the configured pattern too.
String block = """
@@ -48,7 +48,7 @@ class InjectorTest {
@Test
void deliversWhenIdle() {
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
assertFalse(f.isDone(), "not delivered until an injectable status arrives");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(f.isDone());
@@ -170,32 +170,6 @@ class InjectorTest {
assertEquals(List.of("a", "b", "c"), sent());
}
@Test
void cancellingTheMiddleDeliveryKeepsTheFollowingDeliveryReachable() {
Injector.Delivery first = injector.enqueue(T, "same text", TestTurnTokens.inert(T));
Injector.Delivery cancelled = injector.enqueue(T, "same text", TestTurnTokens.inert(T));
injector.enqueue(T, "after cancelled", TestTurnTokens.inert(T));
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(cancelled));
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
injector.onStatus(T, AgentStatus.IDLE);
assertEquals(List.of("same text", "after cancelled"), sent(),
"cancellation must match the exact Delivery and preserve the remaining FIFO queue");
assertTrue(first.completion().isDone());
}
@Test
void cancellationReportsDeliveredWhenPickupWonTheRace() {
Injector.Delivery delivery = injector.enqueue(T, "already sent", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.IDLE);
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(delivery),
"a cancellation after pickup must not claim that the text stayed queued");
assertEquals(List.of("already sent"), sent());
}
@Test
void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() {
assertTrue(injector.activeTargets().isEmpty());
@@ -404,7 +378,7 @@ class InjectorTest {
void sendFailureDropsMessageAndFailsItsFuture() {
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
assertTrue(f.isCompletedExceptionally());
@@ -413,7 +387,7 @@ class InjectorTest {
@Test
void dropFailsPendingWaiters() {
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T));
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
}
@@ -422,8 +396,8 @@ class InjectorTest {
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T));
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
@@ -475,7 +449,7 @@ class InjectorTest {
Captor cap = new Captor();
List<String> forgotten = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add);
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
@@ -589,7 +563,7 @@ class InjectorTest {
StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10);
poller.start();
try {
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T));
delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send
} finally {
poller.stop();
@@ -608,7 +582,7 @@ class InjectorTest {
void deliveredFutureCarriesSendFailure() {
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
assertInstanceOf(HerdrException.class, ex.getCause());
@@ -41,7 +41,7 @@ class StatusPollerRoutingTest {
poller.start();
try {
CompletableFuture<Void> delivered =
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion();
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
// Must resolve quickly: refining against the WRONG daemon (member) never classifies
// out of UNKNOWN, so this would time out under the bug.
delivered.get(2, TimeUnit.SECONDS);
@@ -65,7 +65,7 @@ class StatusPollerRoutingTest {
poller.start();
try {
CompletableFuture<Void> delivered =
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion();
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS),
"a lead target must never be refined from the member daemon's pane content");
} finally {
@@ -258,7 +258,7 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.IDLE); // first delivery
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T)).completion();
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T));
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
@@ -1768,10 +1768,7 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
assertTrue(messages.hasQueuedDelivery(T),
"a TIMED_OUT_QUEUED send still records the undelivered delivery for fleet health");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.prompt")),
"a TIMED_OUT_QUEUED send must be cancelled, not delivered when the worker later goes idle");
"a TIMED_OUT_QUEUED send leaves the message still queued in the injector");
}
@Test