Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1cc0322782 |
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user