From 312c0584ce79bdabbe5461d475a79f9d4fc78871 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 27 Aug 2026 21:58:21 +0700 Subject: [PATCH] CB-640: add MessageService message-layer health evidence accessors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit hasQueuedDelivery/hasStrandedReply/hasOrphanedDelegation surface three of the message-layer facts FleetHealthMonitor needs but currently hardcodes to NOT_YET_OBSERVED. Additive only — no existing public method's signature or behavior changes. --- .../dev/ltms/fleet/msg/MessageService.java | 84 ++++++++++ .../ltms/fleet/msg/MessageServiceTest.java | 143 ++++++++++++++++++ 2 files changed, 227 insertions(+) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 1b024fc..4cb6b19 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -204,6 +204,25 @@ public final class MessageService { new ConcurrentHashMap<>(); /** Async tickets paused on a specific {@code fleet_ask} turn. */ private final ConcurrentHashMap asyncTasksByTurn = new ConcurrentHashMap<>(); + /** + * Targets whose most recent {@code fleet_reply} arrived with no send awaiting it (CB-640) — + * {@link Rendezvous#resolve} returned {@code false} and the reply was queued into the inbox + * instead (see {@link #reply}). The reply itself is not lost (it sits in the inbox for a + * later drain), but the stranding is a fact the health layer needs to see. Bounded by the + * target's own lifecycle rather than a TTL: an entry is cleared the next time this target's + * delivery is accepted ({@link #send}) or the target is torn down ({@link #abandon}), so the + * map holds at most one entry per session with an unresolved stranding right now. + */ + private final ConcurrentHashMap strandedReplies = new ConcurrentHashMap<>(); + /** + * Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — the + * message never reached the {@link Injector} delivery window before the caller's deadline, so + * it is still sitting in the injector's own per-target queue. Set where {@link #send} already + * computes {@code wasDelivered} for that outcome; no new queue is kept here, only the fact. + * 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 queuedDeliveries = new ConcurrentHashMap<>(); private final AtomicLong ticketSeq = new AtomicLong(); private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor( Thread.ofVirtual().name("bridge-async-", 0).factory()); @@ -271,6 +290,55 @@ public final class MessageService { return !inbox.peek(target).isEmpty(); } + /** + * Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out + * before the {@link Injector} ever delivered it — the caller saw + * {@link Outcome#TIMED_OUT_QUEUED} (see the {@code TimeoutException} branch of {@link #send}), + * and the message is still sitting in the injector's per-target queue waiting for the worker + * to go idle. 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); + } + + /** + * Read-only delegation fact for fleet views (CB-640): {@code target}'s last {@code fleet_reply} + * arrived while no send was waiting for it, so {@link Rendezvous#resolve} returned + * {@code false} and {@link #reply} fell back to queueing it in the inbox (see the CB-307 + * javadoc there and {@code docs/CB-307-Reliable-Delivery.md} §1). Cleared the next time this + * target's delivery is accepted or the target is abandoned — see {@link #strandedReplies}. + */ + public boolean hasStrandedReply(String target) { + return target != null && strandedReplies.containsKey(target); + } + + /** + * Read-only delegation fact for fleet views (CB-640): an async ticket is still + * {@link Phase#PENDING} against {@code target}, yet nothing is actually in flight for it — no + * open rendezvous waiter ({@link #hasAcceptedDelivery}) and no message still sitting in the + * injector's queue ({@link #hasQueuedDelivery}). A healthy PENDING ticket can briefly look this + * way while its virtual thread has not yet been scheduled or is blocked on the session lock + * behind another send to the same target, so this is a snapshot fact for the health classifier + * to weigh across ticks, not proof on its own that the ticket is stuck. It also genuinely + * persists — not just as a passing race — once an async {@code fleet_ask} lapses unanswered: + * {@link #ask} clears the ticket's question and returns it to {@code PENDING}, but {@link #send} + * already closed the forward waiter the instant the question surfaced, so the target has + * neither an accepted nor a queued delivery left to show for it. + */ + public boolean hasOrphanedDelegation(String target) { + if (target == null || hasAcceptedDelivery(target) || hasQueuedDelivery(target)) { + return false; + } + for (Task task : tasks.values()) { + if (target.equals(task.target) && task.question == null && !task.future.isDone()) { + return true; + } + } + return false; + } + /** * Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the * inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter @@ -288,6 +356,9 @@ public final class MessageService { return true; // a live send took it — unchanged fast path } inbox.publish(session, UUID.randomUUID().toString(), content); + // CB-640: record the stranding itself (not just the reply text) so fleet health can see a + // worker whose replies keep missing their waiter, not only the queue depth this leaves behind. + strandedReplies.put(session, Boolean.TRUE); // A rising inbox share is the signal CB-307 exists to make visible: the worker finished but // nobody was waiting, so delivery now depends on the push loop and a drain. count(FleetMetrics.REPLIES, "path", "inbox"); @@ -341,6 +412,9 @@ public final class MessageService { * @return true if a live waiter was failed */ public boolean abandon(String target, String reason) { + // CB-640: the session is gone — nothing will ever accept or deliver into it now. + strandedReplies.remove(target); + queuedDeliveries.remove(target); CompletableFuture waiter = rendezvous.currentWaiter(target); boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason); boolean asyncFailed = false; @@ -446,6 +520,11 @@ public final class MessageService { // enqueue failure is safely closed by the finally below: nothing is left queued, and the // failed send leaves no stale waiter behind. CompletableFuture reply = rendezvous.open(target); + // CB-640: this send now owns target's delivery, so any earlier stranded-reply or + // still-queued fact no longer describes the live state — clear both rather than let + // them outlive the send that supersedes them. + strandedReplies.remove(target); + queuedDeliveries.remove(target); try { if (task != null) { asyncTasksByWaiter.put(reply, task); @@ -464,6 +543,11 @@ public final class MessageService { } catch (TimeoutException e) { boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); log.debug("send to {} timed out (delivered={})", target, wasDelivered); + if (!wasDelivered) { + // CB-640: still sitting in the injector's queue, waiting for the member to + // go idle — record the fact for fleet health (see queuedDeliveries). + queuedDeliveries.put(target, Boolean.TRUE); + } return recorded(new Reply( wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null)); } catch (ExecutionException e) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 166a0c8..c0feb46 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1081,4 +1081,147 @@ class MessageServiceTest { assertEquals(phase, view.phase()); return view; } + + // --- CB-640: fleet health evidence accessors -------------------------------------------- + + @Test + void hasQueuedDeliveryIsFalseForAnUnknownTarget() { + assertFalse(messages.hasQueuedDelivery("nobody-ever-sent-here")); + } + + @Test + void hasQueuedDeliveryIsFalseBeforeAnyTimeout() { + assertFalse(messages.hasQueuedDelivery(T)); + } + + @Test + void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() { + // Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with + // TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above. + MessageService.Reply r = messages.send(T, "never delivered", 50); + assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome()); + + assertTrue(messages.hasQueuedDelivery(T), + "a TIMED_OUT_QUEUED send leaves the message still queued in the injector"); + } + + @Test + void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception { + assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome()); + assertTrue(messages.hasQueuedDelivery(T)); + + // A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared. + CompletableFuture second = sendAsync(); + awaitWaiting(); + assertFalse(messages.hasQueuedDelivery(T), + "a fresh accepted delivery supersedes the earlier queued fact"); + + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + assertTrue(rendezvous.resolve(T, "second done")); + second.get(5, TimeUnit.SECONDS); + } + + @Test + void hasQueuedDeliveryClearsOnAbandon() { + assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome()); + assertTrue(messages.hasQueuedDelivery(T)); + + messages.abandon(T, "session released"); + assertFalse(messages.hasQueuedDelivery(T), "a torn-down target has nothing left queued for it"); + } + + @Test + void hasStrandedReplyIsFalseForAnUnknownTarget() { + assertFalse(messages.hasStrandedReply("nobody-ever-sent-here")); + } + + @Test + void hasStrandedReplyIsFalseWhenTheReplyResolvedALiveSend() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + + assertTrue(messages.reply(T, "resolved-live")); + assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded"); + + MessageService.Reply r = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, r.outcome()); + } + + @Test + void hasStrandedReplyIsTrueWhenNoSendWasWaiting() { + // No send is open for T — the reply queues into the inbox and is recorded as stranded. + assertTrue(messages.reply(T, "nobody was waiting")); + assertTrue(messages.hasStrandedReply(T), + "a reply with no open send strands, even though it is safely queued in the inbox"); + } + + @Test + void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception { + assertTrue(messages.reply(T, "stray")); + assertTrue(messages.hasStrandedReply(T)); + + // The next accepted delivery for T clears the stale stranding fact — the one case the + // ticket calls out as the one that matters. + CompletableFuture send = sendAsync(); + awaitWaiting(); + assertFalse(messages.hasStrandedReply(T), + "a stranded reply must clear once the target's delivery is accepted again"); + + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + assertTrue(rendezvous.resolve(T, "done")); + send.get(5, TimeUnit.SECONDS); + } + + @Test + void hasStrandedReplyClearsOnAbandon() { + assertTrue(messages.reply(T, "stray")); + assertTrue(messages.hasStrandedReply(T)); + + messages.abandon(T, "session released"); + assertFalse(messages.hasStrandedReply(T), "a torn-down target has nothing left to strand"); + } + + @Test + void hasOrphanedDelegationIsFalseForAnUnknownTarget() { + assertFalse(messages.hasOrphanedDelegation("nobody-ever-sent-here")); + } + + @Test + void hasOrphanedDelegationIsFalseWhilePendingTicketsHaveAnAcceptedDelivery() throws Exception { + // Mirrors abandonFailsEveryPendingAsyncTicketForTheReleasedTarget above: "first" holds the + // session lock and its waiter is open, so the target genuinely has something in flight even + // though "second" and "third" are themselves parked (PENDING) behind the lock. + messages.sendAsync(T, "first task"); + awaitWaiting(); // first task owns the target lock and rendezvous waiter + String second = messages.sendAsync(T, "second task"); + assertEquals(MessageService.Phase.PENDING, messages.poll(second).phase()); + + assertFalse(messages.hasOrphanedDelegation(T), + "the target has an accepted delivery in flight (first), so nothing here is orphaned"); + + assertTrue(messages.abandon(T, "session released")); // release the lock and the parked tickets + } + + @Test + void hasOrphanedDelegationIsTrueOnceAnUnansweredAskLapsesBackToPending() throws Exception { + // Same setup as unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget above: + // once the ask lapses, the ticket goes back to PENDING but send() already closed the + // forward waiter the instant the question surfaced — nothing is left in flight for T. + String ticket = messages.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + assertEquals(MessageService.AskOutcome.TIMED_OUT, + messages.ask(T, "which config?", 200).outcome()); + assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase()); + + assertFalse(messages.hasAcceptedDelivery(T), "the forward waiter closed when the question surfaced"); + assertFalse(messages.hasQueuedDelivery(T), "this ticket never timed out as queued"); + assertTrue(messages.hasOrphanedDelegation(T), + "a PENDING ticket with no accepted or queued delivery for its target is orphaned"); + + assertTrue(messages.abandon(T, "session released")); // clean up the still-open ticket + } }