CB-640: add MessageService message-layer health evidence accessors
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m11s

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.
This commit is contained in:
Dai Ha
2026-08-27 21:58:21 +07:00
parent 97d9cebc59
commit 312c0584ce
2 changed files with 227 additions and 0 deletions
@@ -204,6 +204,25 @@ public final class MessageService {
new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code fleet_ask} turn. */
private final ConcurrentHashMap<String, Task> 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<String, Boolean> 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<String, Boolean> 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<Rendezvous.Resolution> 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<Rendezvous.Resolution> 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) {
@@ -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<MessageService.Reply> 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<MessageService.Reply> 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<MessageService.Reply> 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
}
}