fleetd #546: widen Injector's delivery catch to Throwable, stop re-delivery on Error
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:
@@ -306,7 +306,7 @@ public final class Injector {
|
|||||||
if (t == null) return;
|
if (t == null) return;
|
||||||
|
|
||||||
Pending sent = null;
|
Pending sent = null;
|
||||||
RuntimeException sendError = null;
|
Throwable sendError = null;
|
||||||
boolean turnCompleted = false;
|
boolean turnCompleted = false;
|
||||||
boolean turnFailed = false;
|
boolean turnFailed = false;
|
||||||
boolean resubmit = false;
|
boolean resubmit = false;
|
||||||
@@ -388,9 +388,14 @@ public final class Injector {
|
|||||||
t.turnObserved = false;
|
t.turnObserved = false;
|
||||||
t.injectableSincePickup = 0;
|
t.injectableSincePickup = 0;
|
||||||
sent = p;
|
sent = p;
|
||||||
} catch (RuntimeException e) {
|
} catch (Throwable e) {
|
||||||
// Delivery failed at herdr; drop the poisoned message and surface it
|
// 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();
|
t.queue.poll();
|
||||||
p.state = Pending.State.NOT_DELIVERED;
|
p.state = Pending.State.NOT_DELIVERED;
|
||||||
sent = p;
|
sent = p;
|
||||||
|
|||||||
@@ -2,9 +2,11 @@ package dev.ltms.fleet.inject;
|
|||||||
|
|
||||||
import ch.qos.logback.classic.Level;
|
import ch.qos.logback.classic.Level;
|
||||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||||
|
import com.fasterxml.jackson.databind.JsonNode;
|
||||||
import dev.ltms.fleet.herdr.AgentControl;
|
import dev.ltms.fleet.herdr.AgentControl;
|
||||||
import dev.ltms.fleet.herdr.AgentStatus;
|
import dev.ltms.fleet.herdr.AgentStatus;
|
||||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||||
|
import dev.ltms.fleet.herdr.HerdrClient;
|
||||||
import dev.ltms.fleet.herdr.HerdrException;
|
import dev.ltms.fleet.herdr.HerdrException;
|
||||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||||
import dev.ltms.fleet.testing.CapturedLog;
|
import dev.ltms.fleet.testing.CapturedLog;
|
||||||
@@ -664,4 +666,111 @@ class InjectorTest {
|
|||||||
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
||||||
assertInstanceOf(HerdrException.class, ex.getCause());
|
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");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user