Merge worker/571-attempted-outcome-5739f7-2
This commit is contained in:
@@ -805,6 +805,12 @@ public final class FleetMcp {
|
||||
+ "answered (turnId stale)");
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
|
||||
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
|
||||
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
|
||||
// so the message may already be sitting in the pane. Do not invite a blind retry the way
|
||||
// the case above does; a resend on this route can double-deliver the same brief.
|
||||
case TIMED_OUT_UNCONFIRMED -> text("[no reply within " + timeout + "ms — delivery unconfirmed; "
|
||||
+ "the message may already have reached the worker, so a retry risks sending it "
|
||||
+ "twice — poll status before resending]");
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -93,25 +93,30 @@ public final class MessageService {
|
||||
/** Timed out after the message was delivered — the worker is still working. */
|
||||
TIMED_OUT_WORKING,
|
||||
/**
|
||||
* Timed out with no confirmed delivery. Despite the name, this does not mean the message
|
||||
* is sitting in a queue. {@link #send} reaches this outcome through {@link Injector#cancel},
|
||||
* whose result tells three routes apart:
|
||||
* {@link Injector.Cancellation#CANCELLED} means the message was still queued and this call
|
||||
* removed it, so the target saw nothing and it will not arrive later;
|
||||
* {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was
|
||||
* cleared because the target never became ready or was abandoned, or the injector's call to
|
||||
* the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr
|
||||
* error that this codebase already treats as a confirmed absence — so this route too
|
||||
* establishes that the target saw nothing and it will not arrive later; but
|
||||
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) means that call was made and its
|
||||
* outcome is unknown. {@code agent.prompt} pastes <em>and submits</em> in one call, so on
|
||||
* this route the target may hold a complete, already-submitted turn and be working on it
|
||||
* right now — {@link Outcome#TIMED_OUT_WORKING}'s meaning, reported here as
|
||||
* {@code TIMED_OUT_QUEUED} only because this caller never observed the pickup. Only
|
||||
* {@code CANCELLED} and {@code NOT_DELIVERED} establish that the target saw nothing;
|
||||
* {@code ATTEMPTED} does not.
|
||||
* Timed out with no confirmed delivery, and the target saw nothing — the message will not
|
||||
* arrive later, so a caller may resend. {@link #send} reaches this outcome through {@link
|
||||
* Injector#cancel} reporting one of two routes: {@link Injector.Cancellation#CANCELLED}
|
||||
* means the message was still queued and this call removed it; {@link
|
||||
* Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was cleared
|
||||
* because the target never became ready or was abandoned, or the injector's call to the
|
||||
* target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr
|
||||
* error that this codebase already treats as a confirmed absence. A third route,
|
||||
* {@link Injector.Cancellation#ATTEMPTED}, used to be folded into this same outcome
|
||||
* (fleetd #571) — it no longer is; see {@link #TIMED_OUT_UNCONFIRMED}.
|
||||
*/
|
||||
TIMED_OUT_QUEUED,
|
||||
/**
|
||||
* Timed out with delivery unknown. {@link #send} reaches this outcome when {@link
|
||||
* Injector#cancel} reports {@link Injector.Cancellation#ATTEMPTED} (fleetd #551): the call
|
||||
* to the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) was made, but
|
||||
* this caller never observed whether it reached the pane. {@code agent.prompt} pastes
|
||||
* <em>and submits</em> in one call, so the target may already hold a complete, submitted
|
||||
* turn and be working on it right now — the same reality as {@link #TIMED_OUT_WORKING},
|
||||
* just not confirmed. The message may or may not have arrived. Treat this as neither a
|
||||
* confirmed delivery nor a confirmed absence: a caller that resends on this outcome risks a
|
||||
* double delivery — the same brief typed into the pane twice (fleetd #571).
|
||||
*/
|
||||
TIMED_OUT_UNCONFIRMED,
|
||||
/** Another send to this session was in flight for the whole window. */
|
||||
BUSY,
|
||||
/**
|
||||
@@ -317,20 +322,21 @@ public final class MessageService {
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> strandedReplies = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — {@link
|
||||
* #send} called {@link Injector#cancel} and got back something other than {@code DELIVERED}.
|
||||
* That covers three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message
|
||||
* was still queued and {@code cancel} removed it right there; {@link
|
||||
* Injector.Cancellation#NOT_DELIVERED} — nothing was ever sent, because the target never became
|
||||
* ready, was torn down, or the call to its terminal failed with a herdr error this codebase
|
||||
* already treats as a confirmed absence; or {@link Injector.Cancellation#ATTEMPTED} (fleetd
|
||||
* #551) — the call to the target's terminal was made and its outcome is unknown, so the target
|
||||
* may already hold a complete, submitted turn. Only the first two mean the message will not
|
||||
* arrive later and the target saw nothing; on the third it may already have arrived in full.
|
||||
* Set where {@link #send} already computes {@code wasDelivered} for that outcome; no queue is
|
||||
* kept here, only the fact that the send ended with no confirmed delivery. Cleared the same way
|
||||
* as {@link #strandedReplies}: the next accepted delivery for the target ({@link #send} opening
|
||||
* a fresh waiter) or a teardown ({@link #abandon}).
|
||||
* Targets whose last send timed out with no confirmed delivery (CB-640) — {@link #send} called
|
||||
* {@link Injector#cancel} and got back something other than {@code DELIVERED}. That covers
|
||||
* three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message was still
|
||||
* queued and {@code cancel} removed it right there; {@link Injector.Cancellation#NOT_DELIVERED}
|
||||
* — nothing was ever sent, because the target never became ready, was torn down, or the call to
|
||||
* its terminal failed with a herdr error this codebase already treats as a confirmed absence; or
|
||||
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) — the call to the target's terminal was
|
||||
* made and its outcome is unknown, so the target may already hold a complete, submitted turn.
|
||||
* Only the first two mean the message will not arrive later and the target saw nothing; on the
|
||||
* third it may already have arrived in full — and the caller sees a different outcome for it
|
||||
* ({@link Outcome#TIMED_OUT_UNCONFIRMED}, fleetd #571) than for the first two ({@link
|
||||
* Outcome#TIMED_OUT_QUEUED}). Set where {@link #send} already computes {@code wasDelivered} for
|
||||
* that outcome; no queue is kept here, only the fact that the send ended with no confirmed
|
||||
* delivery. Cleared the same way as {@link #strandedReplies}: the next accepted delivery for the
|
||||
* target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
|
||||
private final AtomicLong ticketSeq = new AtomicLong();
|
||||
@@ -417,21 +423,22 @@ public final class MessageService {
|
||||
|
||||
/**
|
||||
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out
|
||||
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} (see the
|
||||
* {@code TimeoutException} branch of {@link #send}). Despite the method's name, this is not
|
||||
* proof that a message is sitting in a queue: {@link Injector#cancel} reports this outcome
|
||||
* through three routes. {@link Injector.Cancellation#CANCELLED} means the message was still
|
||||
* queued and got removed right there. {@link Injector.Cancellation#NOT_DELIVERED} means
|
||||
* nothing was ever sent — the target never became ready, was torn down, or the call to its
|
||||
* terminal failed with a herdr error this codebase already treats as a confirmed absence.
|
||||
* Only these two routes mean the message will not arrive later. {@link
|
||||
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} or {@link
|
||||
* Outcome#TIMED_OUT_UNCONFIRMED} (fleetd #571; see the {@code TimeoutException} branch of
|
||||
* {@link #send}). Despite the method's name, this is not proof that a message is sitting in a
|
||||
* queue: {@link Injector#cancel} reports this outcome through three routes. {@link
|
||||
* Injector.Cancellation#CANCELLED} means the message was still queued and got removed right
|
||||
* there. {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the target
|
||||
* never became ready, was torn down, or the call to its terminal failed with a herdr error this
|
||||
* codebase already treats as a confirmed absence. Only these two routes mean the message will
|
||||
* not arrive later, and both report {@code TIMED_OUT_QUEUED}. {@link
|
||||
* Injector.Cancellation#ATTEMPTED} (fleetd #551) means the call to the target's terminal was
|
||||
* made and its outcome is unknown: {@code agent.prompt} pastes <em>and submits</em> in one
|
||||
* call, so on this route the target may already hold a complete, submitted turn and be
|
||||
* working on it right now — it does NOT follow that the target saw nothing. Distinct from
|
||||
* {@link Outcome#TIMED_OUT_WORKING}, where delivery already happened and only the reply is
|
||||
* outstanding. Cleared the next time this target's delivery is accepted or the target is
|
||||
* abandoned — see {@link #queuedDeliveries}.
|
||||
* working on it right now — it does NOT follow that the target saw nothing, and this route
|
||||
* reports {@code TIMED_OUT_UNCONFIRMED} instead. Distinct from {@link Outcome#TIMED_OUT_WORKING},
|
||||
* where delivery already happened and only the reply is outstanding. Cleared the next time this
|
||||
* target's delivery is accepted or the target is abandoned — see {@link #queuedDeliveries}.
|
||||
*/
|
||||
public boolean hasQueuedDelivery(String target) {
|
||||
return target != null && queuedDeliveries.containsKey(target);
|
||||
@@ -646,7 +653,7 @@ public final class MessageService {
|
||||
return switch (o) {
|
||||
case REPLIED -> "replied";
|
||||
case COMPLETED_UNREPLIED -> "completion_fallback";
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
case STALE_TURN, QUESTION -> null; // not a completed delegation
|
||||
@@ -975,6 +982,7 @@ public final class MessageService {
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivery.completion().isDone()
|
||||
&& !delivery.completion().isCompletedExceptionally();
|
||||
Injector.Cancellation cancellation = null;
|
||||
if (!wasDelivered) {
|
||||
if (timeoutCancellationRaceHookForTest != null) {
|
||||
// Test-only (fleetd #345): see the field's own javadoc.
|
||||
@@ -982,20 +990,27 @@ public final class MessageService {
|
||||
}
|
||||
// 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;
|
||||
cancellation = injector.cancel(delivery);
|
||||
wasDelivered = cancellation == Injector.Cancellation.DELIVERED;
|
||||
}
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
if (!wasDelivered) {
|
||||
Outcome outcome;
|
||||
if (wasDelivered) {
|
||||
outcome = Outcome.TIMED_OUT_WORKING;
|
||||
} else if (cancellation == Injector.Cancellation.ATTEMPTED) {
|
||||
// fleetd #571: the call to the target's terminal was made and its outcome is
|
||||
// unknown — the message may already have arrived in full, so this must not
|
||||
// be reported as TIMED_OUT_QUEUED, which promises it never will.
|
||||
outcome = Outcome.TIMED_OUT_UNCONFIRMED;
|
||||
} else {
|
||||
// CB-640: record that delivery is not confirmed, for fleet health (see
|
||||
// queuedDeliveries). Whatever injector.cancel() reported above — this call
|
||||
// removed a still-queued Pending (CANCELLED), an earlier attempt already
|
||||
// failed with a confirmed absence (NOT_DELIVERED), or an earlier attempt was
|
||||
// made and its outcome is unknown (ATTEMPTED, fleetd #551 — the message may
|
||||
// already have arrived in full) — the send ends with no confirmed delivery.
|
||||
// queuedDeliveries). cancellation is CANCELLED (this call removed a
|
||||
// still-queued Pending) or NOT_DELIVERED (an earlier attempt already failed
|
||||
// with a confirmed absence) — both mean the target saw nothing.
|
||||
queuedDeliveries.put(target, Boolean.TRUE);
|
||||
outcome = Outcome.TIMED_OUT_QUEUED;
|
||||
}
|
||||
return recorded(new Reply(
|
||||
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
|
||||
return recorded(new Reply(outcome, null));
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
||||
|
||||
@@ -640,15 +640,31 @@ public final class FleetApp {
|
||||
}
|
||||
default -> ctx.status(202).json(Map.of(
|
||||
"sessionId", id,
|
||||
// fleetd #571 (ticket comment 17126): no `default` here on purpose. This switch
|
||||
// is an expression, so the compiler already demands every Outcome constant have
|
||||
// an arm — adding an 11th constant to Outcome is a compile error here, not a
|
||||
// silent fall-through. That is exactly the bug this ticket exists to fix:
|
||||
// `default -> "done"` used to sit here and would have told a REST caller the
|
||||
// delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery
|
||||
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never
|
||||
// actually reach this inner switch — the outer switch above always dispatches
|
||||
// them first — but they still need an arm to keep this switch exhaustive.
|
||||
"status", switch (reply.outcome()) {
|
||||
case TIMED_OUT_WORKING -> "working";
|
||||
case TIMED_OUT_QUEUED -> "queued";
|
||||
// Delivery here is unknown, not merely still queued — see
|
||||
// Outcome#TIMED_OUT_UNCONFIRMED's own javadoc.
|
||||
case TIMED_OUT_UNCONFIRMED -> "unconfirmed";
|
||||
case BUSY -> "busy";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
default -> "done"; // unreachable (terminal outcomes handled above)
|
||||
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
|
||||
},
|
||||
"detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|
||||
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
|
||||
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
|
||||
+ "message may already have reached the worker, so a resend "
|
||||
+ "risks sending it twice; poll status first"
|
||||
: (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|
||||
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
|
||||
&& reply.text() != null
|
||||
? reply.text()
|
||||
|
||||
@@ -7,6 +7,7 @@ import dev.ltms.fleet.auth.Role;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
@@ -57,9 +58,10 @@ class FleetMcpTest {
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private final Injector injector = new Injector(agents);
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous, inbox);
|
||||
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
@@ -298,6 +300,35 @@ class FleetMcpTest {
|
||||
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is
|
||||
* the one message whose whole job is to stop a caller retrying a delivery that may already have
|
||||
* arrived. Pin that its wording is actually distinct from the queued/working arm's retry
|
||||
* invitation — a mutation that swapped this arm's text for that one still passed every other
|
||||
* test in this suite, because nothing asserted the specific wording.
|
||||
*/
|
||||
@Test
|
||||
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt -> ATTEMPTED
|
||||
|
||||
McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS);
|
||||
String text = textOf(res);
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(text.contains("delivery unconfirmed"), "got: " + text);
|
||||
assertFalse(text.contains("retry or poll status"),
|
||||
"an unconfirmed delivery must not carry the queued/working arm's retry invitation — "
|
||||
+ "a resend here can double-deliver the same brief: got " + text);
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsMissingArgs() {
|
||||
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
|
||||
|
||||
@@ -601,6 +601,30 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571 (the acceptance test the ticket was filed for). The worker is idle so the injector
|
||||
* attempts delivery, but the {@code agent.prompt} call itself fails with a herdr error that is
|
||||
* not a confirmed absence (not a {@code *_not_found} code) — {@link Injector} marks the Pending
|
||||
* {@code ATTEMPTED} (fleetd #551), meaning the call was made and whether it reached the pane is
|
||||
* unknown. Before this fix, {@code send}'s {@code TimeoutException} branch collapsed
|
||||
* {@code ATTEMPTED} into {@code TIMED_OUT_QUEUED} — a promise that the message will never arrive,
|
||||
* which may already be false: {@code agent.prompt} pastes and submits in one call.
|
||||
*/
|
||||
@Test
|
||||
void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception {
|
||||
herdr.agentSendFailsWith("send_failed");
|
||||
CompletableFuture<MessageService.Reply> send =
|
||||
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED
|
||||
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_UNCONFIRMED, r.outcome(),
|
||||
"an ATTEMPTED delivery must not collapse into TIMED_OUT_QUEUED — the message may "
|
||||
+ "already have arrived in full, and TIMED_OUT_QUEUED promises it never will");
|
||||
assertNull(r.text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
|
||||
@@ -623,6 +623,28 @@ class FleetAppTest {
|
||||
assertTrue(herdr.called("agent.prompt"), "message was injected");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #571: the worker is idle, so the poller attempts delivery, but the {@code agent.prompt}
|
||||
* call itself fails with a herdr error that is not a confirmed absence — {@link
|
||||
* dev.ltms.fleet.inject.Injector} marks this {@code ATTEMPTED}, meaning the call was made and
|
||||
* whether it reached the pane is unknown. {@code writeReply}'s default arm must map this to its
|
||||
* own {@code "unconfirmed"} status, not silently fall through to {@code "done"} (which would
|
||||
* claim the delegation completed) nor collapse into {@code "queued"} (which would claim the
|
||||
* message will never arrive, when it may already be sitting in the pane).
|
||||
*/
|
||||
@Test
|
||||
void messageTimesOutUnconfirmedWhenDeliveryAttemptFails() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("send_failed");
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
|
||||
assertEquals(202, res.statusCode());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("unconfirmed", body.get("status").asText(),
|
||||
"an ATTEMPTED delivery must report its own status, not \"queued\" or \"done\"");
|
||||
assertTrue(herdr.called("agent.prompt"), "delivery must have been attempted");
|
||||
}
|
||||
|
||||
@Test
|
||||
void messageRejectsBlankContent() throws Exception {
|
||||
int port = startHealthy();
|
||||
|
||||
Reference in New Issue
Block a user