fleetd #780: tell the truth about a queued-but-undelivered send
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 2m10s

A send to a pane that never frees up was accepted, then polled as
"pending — worker working" forever, with no WARN and no timeout. That
text is only ever correct for a message the injector already handed
off; a message still sitting in its queue now reports as queued, with
the target's live status and how long it has been waiting.

Injector.enqueue() now stamps Pending.enqueuedAtMillis so
queuedWaitMillis(target) can tell "never attempted" apart from
"delivered, now being worked on". MessageService.pendingDetail() uses
it to pick the poll wording. FleetMcp.sendAsync() adds a best-effort
warning to the accept text when the target is not injectable at send
time, per the brief's stated preference for a warning over a hard
refusal (a short busy spell is normal).

Injector gets a new bound, QUEUE_WAIT_GRACE_POLLS (4800 polls, ~20min
at the 250ms poll interval), mirroring READINESS_GRACE_POLLS: a
message that sits at the head of the queue with no delivery attempt
for that long fails via onTurnFailed, the same path a readiness
timeout uses, without touching presence — the target is busy, not
gone.

Fixed InjectorTest.readinessGraceExpiryLogsTheMeasuredElapsedTimeNotArithmeticOnConstants's
stub clock, which now sees one extra, legitimate nowMillis read from
enqueue()'s new stamp.

Tests: aQueuedButNeverInjectedMessageDoesNotPollAsWorking +
aDeliveredMessageStillPollsAsWorkerWorking (poll wording, with
positive control); failsAQueuedMessageWhoseTargetNeverFreesUp +
aTargetThatFreesUpBeforeTheQueueWaitGraceIsDeliveredNormally (timeout,
with positive control).
This commit is contained in:
Dai Ha
2026-10-05 19:54:37 +02:00
parent c043d149cf
commit 1cc0322782
5 changed files with 197 additions and 13 deletions
@@ -94,6 +94,17 @@ public final class Injector {
*/
private static final int READINESS_GRACE_POLLS = 240;
/**
* How many consecutive polls a message may sit queued with no delivery attempt at all before
* it is failed and the queue cleared — covers every reason the head of the queue is never
* reached, including a target that stays busy ({@code working}) or unclassifiable
* ({@code unknown}) for the whole window. Set well above an ordinary turn so a worker
* genuinely mid-task is never cut off, and below a caller's own overall timeout so a target
* that never frees up fails with this specific reason instead of riding out that longer wait
* silently.
*/
private static final int QUEUE_WAIT_GRACE_POLLS = 4800;
/**
* The single source for the injector poll cadence — how often the {@link StatusPoller} drives
* {@link #onStatus} at. {@code Fleetd} passes this to every {@link StatusPoller} it constructs,
@@ -296,13 +307,16 @@ public final class Injector {
final String text;
final TurnToken token;
final CompletableFuture<Void> delivered;
final long enqueuedAtMillis;
volatile State state = State.QUEUED; // written under the owning Target monitor
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered) {
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered,
long enqueuedAtMillis) {
this.target = target;
this.text = text;
this.token = token;
this.delivered = delivered;
this.enqueuedAtMillis = enqueuedAtMillis;
}
String text() {
@@ -329,6 +343,7 @@ public final class Injector {
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
long notReadySinceMillis; // wall-clock time of the FIRST non-ready sample in the current notReadySincePoll streak (fleetd #501); reset alongside it
int queueWaitSincePoll; // consecutive polls the queue has held an undelivered message with no attempt made
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
boolean awaitingPostTurnPickup;
boolean postTurnObserved;
@@ -349,7 +364,7 @@ public final class Injector {
*/
public Delivery enqueue(String target, String text, TurnToken token) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(target, text, token, delivered);
Pending p = new Pending(target, text, token, delivered, nowMillis.getAsLong());
targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
@@ -415,6 +430,7 @@ public final class Injector {
boolean resubmit = false;
boolean startPostTurn = false;
List<Pending> notReady = null; // queued messages failed because the worker never became ready
List<Pending> queueStalled = null; // queued messages failed because the queue never drained
synchronized (t) {
if (status == AgentStatus.WORKING) {
if (t.awaitingPostTurnPickup) {
@@ -606,6 +622,28 @@ public final class Injector {
}
}
// A message still queued and never attempted this poll is bounded on its own, whatever
// the reason the head of the queue was never reached — a target stuck WORKING or
// UNKNOWN for the whole window hits this even though neither branch above ever looks at
// the queue. Any poll that did attempt the head (`sent != null`, success or failure
// alike) counts as progress and resets the streak, even if messages remain behind it.
if (!t.queue.isEmpty() && sent == null) {
if (++t.queueWaitSincePoll >= QUEUE_WAIT_GRACE_POLLS) {
queueStalled = new ArrayList<>(t.queue);
for (Pending pending : queueStalled) {
pending.state = Pending.State.NOT_DELIVERED;
}
log.warn("queue for {} never drained after {} polls (limit={} polls/{}s): "
+ "failing {} queued message(s) that were never attempted",
target, t.queueWaitSincePoll, QUEUE_WAIT_GRACE_POLLS,
QUEUE_WAIT_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, queueStalled.size());
t.queue.clear();
t.queueWaitSincePoll = 0;
}
} else {
t.queueWaitSincePoll = 0;
}
// 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)) {
@@ -676,6 +714,19 @@ public final class Injector {
forget.accept(target);
turnListener.onTurnFailed(target);
}
if (queueStalled != null) {
// The target is not gone — it may still be genuinely busy — so this does not call
// forget.accept: that would clear presence/readiness state for a worker that is
// simply taking a long turn. It still resolves the awaiting send's own waiter via
// onTurnFailed (mirroring notReady above), so a caller learns this specific message
// never reached the pane instead of riding out its own much longer timeout.
RuntimeException cause = new IllegalStateException(
target + " never freed up to receive this message within the queue wait grace");
for (Pending p : queueStalled) {
p.delivered().completeExceptionally(cause);
}
turnListener.onTurnFailed(target, cause.getMessage());
}
if (turnCompleted) {
if (startPostTurn) {
// fleetd #553: the listener call is wrapped so `t.postTurnPending` (set true inside
@@ -812,6 +863,23 @@ public final class Injector {
.collect(Collectors.toSet());
}
/**
* How long the oldest still-queued, never-attempted message for {@code target} has been
* waiting, or {@code null} when nothing is queued (including when the head has already been
* attempted or delivered). A caller uses this to tell a message that genuinely never reached
* the pane apart from one that was delivered and is now simply being worked on.
*/
public Long queuedWaitMillis(String target) {
Target t = targets.get(target);
if (t == null) {
return null;
}
synchronized (t) {
Pending head = t.queue.peek();
return head != null ? nowMillis.getAsLong() - head.enqueuedAtMillis : null;
}
}
/**
* Forget a target whose worker is gone, failing every still-queued message so awaiting callers
* unblock instead of hanging forever. If a message had already been <em>delivered</em> but its
@@ -8,6 +8,7 @@ import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.inject.MemberPresence;
@@ -1106,7 +1107,29 @@ public final class FleetMcp {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted, creator);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket
+ notInjectableWarning(messages, sessionId));
}
/**
* A synchronous status check at accept time: when {@code sessionId} is not currently
* idle/blocked/done, this message is queued rather than reaching the pane right away. Empty
* when the target is injectable or its status could not be read — a best-effort warning, not
* a reason to withhold the accept receipt.
*/
private static String notInjectableWarning(MessageService messages, String sessionId) {
AgentStatus status;
try {
status = messages.status(sessionId);
} catch (RuntimeException e) {
return "";
}
if (status.injectable()) {
return "";
}
return "\n\nWarning: " + sessionId + " is currently " + status.name().toLowerCase()
+ ", not idle/blocked/done — this message is queued, not yet delivered, and will "
+ "wait until the target frees up. Poll fleet_poll to see when it lands.";
}
/** A configured profile is never a send target; other unknown values may be herdr-owned panes. */
@@ -1426,7 +1426,7 @@ public final class MessageService {
return new TaskView(ticket, Phase.ASKING, question.text(), null,
"worker is waiting for your answer", question.turnId());
}
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
return new TaskView(ticket, Phase.PENDING, null, null, pendingDetail(task.target), null);
}
// CB-588: the ticket is terminal and being handed to the caller right here — tell the push
// loop it is collected so a later tick's nudge never names a ticket the lead already has.
@@ -1499,6 +1499,21 @@ public final class MessageService {
return task != null && task.completedNanos != null;
}
/**
* Detail text for a {@link Phase#PENDING} poll of a plain (non-asking) delegation.
* Distinguishes a message still sitting in the injector's queue, never delivered, from one
* that already reached the pane and is simply being worked on — so a caller cannot read
* "worker working" as "received" when it was not.
*/
private String pendingDetail(String target) {
Long queuedMillis = injector.queuedWaitMillis(target);
if (queuedMillis != null) {
return "queued, not yet delivered (target is " + liveStatus(target) + "; queued "
+ (queuedMillis / 1000) + "s)";
}
return "worker " + liveStatus(target);
}
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
private String liveStatus(String target) {
try {
@@ -559,20 +559,21 @@ class InjectorTest {
void readinessGraceExpiryLogsTheMeasuredElapsedTimeNotArithmeticOnConstants() {
// fleetd #501, defect 2: the old line computed "({}s)" as READINESS_GRACE_POLLS *
// POLL_INTERVAL_MILLIS / 1000 — arithmetic on two constants, never a measurement, and wrong
// in the direction that says everything ran on schedule. This stub clock returns two FIXED
// values (1_000ms at the first non-ready sample, 318_412ms at the poll that trips the grace)
// whose difference — 317_412ms — does NOT equal 240 * POLL_INTERVAL_MILLIS (=60_000ms).
// Asserting on that literal, non-derived number is what makes this test able to fail if the
// production code goes back to printing the constant-arithmetic value instead of the
// injected clock's measurement.
long[] readings = {1_000L, 318_412L};
// in the direction that says everything ran on schedule. This stub clock returns three FIXED
// values (an unused enqueue-time stamp, 1_000ms at the first non-ready sample, 318_412ms at
// the poll that trips the grace) whose last two differ — 317_412ms — which does NOT equal
// 240 * POLL_INTERVAL_MILLIS (=60_000ms). Asserting on that literal, non-derived number is
// what makes this test able to fail if the production code goes back to printing the
// constant-arithmetic value instead of the injected clock's measurement.
long[] readings = {0L, 1_000L, 318_412L};
AtomicInteger call = new AtomicInteger(0);
LongSupplier stubClock = () -> {
int i = call.getAndIncrement();
if (i >= readings.length) {
throw new AssertionError("nowMillis read more times than this fixture expects (" + i
+ "); the readiness-not-ready branch should read the clock exactly twice — "
+ "once to stamp the first non-ready sample, once at grace expiry");
+ "); the readiness-not-ready branch should read the clock exactly three times —"
+ " once to stamp the queued message's own enqueue time, once to stamp the "
+ "first non-ready sample, once at grace expiry");
}
return readings[i];
};
@@ -614,6 +615,46 @@ class InjectorTest {
assertEquals(List.of("task"), sent(), "a worker that connects within the grace is delivered to");
}
// ~20 minutes of a continuously busy target at the 250ms prod poll interval; enough to trip
// the queue-wait grace.
private static final int QUEUE_WAIT_SAMPLES = 4800;
@Test
void failsAQueuedMessageWhoseTargetNeverFreesUp() {
// A target that stays WORKING the whole time never reaches the branch that looks at the
// queue at all, so nothing else bounds this. The message must fail rather than wait
// forever, and the caller's future unblocks through the same turn-failure path a readiness
// timeout uses — but presence must not be touched, since this target is merely busy, not
// gone.
Captor cap = new Captor();
List<String> forgotten = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> true, forgotten::add);
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
for (int i = 0; i < QUEUE_WAIT_SAMPLES; i++) inj.onStatus(T, AgentStatus.WORKING);
assertEquals(List.of(), sent(), "a target that never frees up is never delivered to");
assertTrue(f.isCompletedExceptionally(), "the caller's future fails instead of hanging forever");
assertEquals(List.of(T), cap.failed, "the awaiting send resolves through the turn-failure path");
assertEquals(List.of(), cap.completed, "a never-freed target is a failure, not a completion");
assertEquals(List.of(), forgotten, "the target is busy, not gone — presence must not be cleared");
assertTrue(inj.activeTargets().isEmpty(), "the target is reclaimed, not polled forever");
}
@Test
void aTargetThatFreesUpBeforeTheQueueWaitGraceIsDeliveredNormally() {
// Positive control: a target that is merely busy for a while, then frees up before the
// grace elapses, is still delivered normally — the long-task case this grace must not break.
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP);
inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < 100; i++) inj.onStatus(T, AgentStatus.WORKING); // busy, well under the grace
assertEquals(List.of(), sent());
inj.onStatus(T, AgentStatus.IDLE); // frees up
assertEquals(List.of("task"), sent(), "a target that frees up within the grace is delivered to");
}
@Test
void dropClearsWorkerPresence() {
// CB-114 (finding #1): a vanished worker's readiness must be forgotten so a stale entry cannot
@@ -1706,6 +1706,43 @@ class MessageServiceTest {
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* A message still sitting in the injector's queue — never attempted, let alone delivered —
* must not poll as the worker working on it. {@code herdr.agentStatus("working")} supplies the
* target's real, independent live status (busy with something else entirely), while no
* {@code injector.onStatus} call is ever made, so the message is never even attempted.
*/
@Test
void aQueuedButNeverInjectedMessageDoesNotPollAsWorking() throws Exception {
herdr.agentStatus("working");
String ticket = messages.sendAsync(T, "task");
awaitWaiting();
MessageService.TaskView view = messages.poll(ticket);
assertEquals(MessageService.Phase.PENDING, view.phase());
assertFalse(view.detail().contains("worker working"),
"a message never injected must not read as the worker working on it: " + view.detail());
assertTrue(view.detail().contains("queued") && view.detail().contains("not yet delivered"),
"must report the message as queued, not delivered: " + view.detail());
}
/**
* Positive control for the test above: once the message is actually delivered, the poll's
* detail is the plain live-status text again.
*/
@Test
void aDeliveredMessageStillPollsAsWorkerWorking() throws Exception {
herdr.agentStatus("working");
String ticket = messages.sendAsync(T, "task");
awaitWaiting();
injectDelivery();
MessageService.TaskView view = messages.poll(ticket);
assertEquals(MessageService.Phase.PENDING, view.phase());
assertEquals("worker working", view.detail(),
"once actually delivered, the detail reports the worker's live status directly");
}
/**
* fleetd #329 (F1). {@code answer()} completes the async ticket by looking {@code turnId} up in
* {@code asyncTasksByTurn} a SECOND time (the first is at :991, purely to re-register the