fleetd #546: widen Injector's delivery catch to Throwable, stop re-delivery on Error
CI / contract (pull_request) Successful in 1m29s
CI / build (pull_request) Successful in 1m49s

Injector.java:391 caught only RuntimeException around the herdr send seam. PR #543
(fleetd #538) widened StatusPoller's per-target catch to Throwable so the polling
loop now survives an Error there, which means it comes back round — and Injector's
narrower catch let the poisoned message stay QUEUED (the loop peeks, not polls),
so the next round re-sent the same text into the member's pane.

Widen the catch to Throwable, matching #543 one layer down. sendError's declared
type widens from RuntimeException to Throwable to keep compiling; its only consumer
(CompletableFuture.completeExceptionally(Throwable)) already accepts that type, so
no other caller-visible behavior changes. The ordinary HerdrException/RuntimeException
path is unchanged.

Adds three tests: an Error at the send seam is dropped and marked NOT_DELIVERED, a
second onStatus round does not re-send it, and a HerdrException control proves the
ordinary path is untouched.
This commit is contained in:
Dai Ha
2026-09-12 14:26:27 +07:00
parent 7611b69667
commit 87871eaefb
2 changed files with 117 additions and 3 deletions
@@ -306,7 +306,7 @@ public final class Injector {
if (t == null) return;
Pending sent = null;
RuntimeException sendError = null;
Throwable sendError = null;
boolean turnCompleted = false;
boolean turnFailed = false;
boolean resubmit = false;
@@ -388,9 +388,14 @@ public final class Injector {
t.turnObserved = false;
t.injectableSincePickup = 0;
sent = p;
} catch (RuntimeException e) {
} catch (Throwable e) {
// Delivery failed at herdr; drop the poisoned message and surface it
// rather than blocking the queue behind it.
// rather than blocking the queue behind it. Catches Throwable, not just
// RuntimeException: fleetd #546 — an Error escaping this send (e.g. a
// NoClassDefFoundError, see #413) would otherwise leave the entry QUEUED
// at the head of t.queue. Line :378 peeks rather than polls, so the next
// onStatus round would re-enter this try and send the same text again,
// typing the same brief into the member's pane a second time.
t.queue.poll();
p.state = Pending.State.NOT_DELIVERED;
sent = p;
@@ -2,9 +2,11 @@ package dev.ltms.fleet.inject;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import com.fasterxml.jackson.databind.JsonNode;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.msg.TestTurnTokens;
import dev.ltms.fleet.testing.CapturedLog;
@@ -664,4 +666,111 @@ class InjectorTest {
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
assertInstanceOf(HerdrException.class, ex.getCause());
}
/**
* A {@link HerdrClient} that throws a non-{@link RuntimeException} {@link Error} from {@code
* agent.prompt} instead of delegating — the fleetd #546 case: a stray non-RuntimeException
* throwable (e.g. a {@code NoClassDefFoundError}, fleetd #413) escaping the send seam at
* {@code Injector.java:383}. Records every call it sees itself, including the ones it throws
* for, since the delegate's own recording is never reached for {@code agent.prompt} — so a test
* can assert on exactly what this fake actually received.
*/
private static final class ErrorOnPrompt implements HerdrClient {
private final FakeHerdr delegate;
private final List<FakeHerdr.Call> calls = new java.util.concurrent.CopyOnWriteArrayList<>();
private ErrorOnPrompt(FakeHerdr delegate) {
this.delegate = delegate;
}
List<FakeHerdr.Call> calls() {
return calls;
}
@Override
public JsonNode call(String method, Object params) {
calls.add(new FakeHerdr.Call(method, params));
if (method.equals("agent.prompt")) {
throw new AssertionError("simulated non-RuntimeException send failure (fleetd #546)");
}
return delegate.call(method, params);
}
@Override
public void close() {
delegate.close();
}
}
/**
* Mirrors StatusPoller's own per-target {@code catch (Throwable)} (fleetd #538 / PR #543): the
* production polling loop already swallows whatever escapes one target's round and comes back
* for the next one. A unit test that calls {@code onStatus} directly (bypassing StatusPoller)
* needs the same survival so it can observe what a SECOND round does, regardless of whether
* fleetd #546's fix is present.
*/
private static void pollOnceSurviving(Injector inj, String target, AgentStatus status) {
try {
inj.onStatus(target, status);
} catch (Throwable ignored) {
// matches StatusPoller.loop's own catch (Throwable) added by PR #543
}
}
@Test
void anErrorFromSendRemovesTheMessageAndMarksItNotDelivered() {
// fleetd #546, acceptance test 1: an Error (not a RuntimeException) escaping the send seam
// at Injector.java:383 must still be caught, the message dropped from the queue, and its
// state set to NOT_DELIVERED — never left QUEUED at the head.
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
Injector inj = new Injector(new AgentControl(throwing));
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
assertDoesNotThrow(() -> inj.onStatus(T, AgentStatus.IDLE),
"fleetd #546: an Error from the send seam must be caught inside onStatus, not "
+ "escape it");
assertTrue(delivery.completion().isCompletedExceptionally(),
"the delivery's future must surface the send failure");
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
"the message must be dropped and marked NOT_DELIVERED, not left QUEUED at the "
+ "head of the queue");
}
@Test
void anErrorFromSendDoesNotRedeliverOnASecondRound() {
// fleetd #546, acceptance test 2 — the test this ticket exists for. Before the fix, an
// Error at the send seam left the message QUEUED (Injector.java:378 peeks, not polls), so a
// second onStatus round re-entered the same try and sent the same text again: the member's
// pane got the same brief typed into it twice.
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
Injector inj = new Injector(new AgentControl(throwing));
inj.enqueue(T, "brief", TestTurnTokens.inert(T));
pollOnceSurviving(inj, T, AgentStatus.IDLE);
pollOnceSurviving(inj, T, AgentStatus.IDLE);
long promptCalls = throwing.calls().stream().filter(c -> c.method().equals("agent.prompt")).count();
assertEquals(1, promptCalls, "fleetd #546: the poisoned text must be sent exactly once — a "
+ "second onStatus round must not re-enter send for the same message");
}
@Test
void aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged() {
// fleetd #546, acceptance test 3 (control): HerdrException extends RuntimeException, so it
// was already caught before this ticket's widening. This pins that the ordinary path is
// unchanged — still dropped, still NOT_DELIVERED, still surfaced to the caller.
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
Injector.Delivery delivery = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
assertTrue(delivery.completion().isCompletedExceptionally(),
"a HerdrException at the send seam must still surface to the caller, unchanged by "
+ "fleetd #546's widening");
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
"a HerdrException must still be dropped and marked NOT_DELIVERED, unchanged by the "
+ "wider Throwable catch");
}
}