From 147f50c19e81a486b76244e73e257650873683f5 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 15:56:10 +0700 Subject: [PATCH] #338: cancel timed-out queued deliveries --- .../java/dev/ltms/fleet/inject/Injector.java | 101 ++++++++++++++++-- .../dev/ltms/fleet/msg/MessageService.java | 14 ++- .../dev/ltms/fleet/inject/InjectorTest.java | 42 ++++++-- .../fleet/inject/StatusPollerRoutingTest.java | 4 +- .../ltms/fleet/msg/MessageServiceTest.java | 7 +- 5 files changed, 146 insertions(+), 22 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java index 638d0c9..08c154a 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java @@ -145,8 +145,58 @@ public final class Injector { return router != null ? router.agentsFor(target) : agents; } + /** The result of trying to remove an undelivered message from the injector. */ + public enum Cancellation { + CANCELLED, + DELIVERED, + NOT_DELIVERED + } + + /** + * An identity handle for one queued delivery. It is the only value accepted by + * {@link #cancel(Delivery)}, so a caller cannot cancel a different message with the same target + * or text. + */ + public static final class Delivery { + private final Pending pending; + + private Delivery(Pending pending) { + this.pending = pending; + } + + public CompletableFuture completion() { + return pending.delivered; + } + } + /** A pending message and the future that completes when it has been delivered. */ - private record Pending(String text, TurnToken token, CompletableFuture delivered) { + private static final class Pending { + enum State { QUEUED, DELIVERED, NOT_DELIVERED, CANCELLED } + + final String target; + final String text; + final TurnToken token; + final CompletableFuture delivered; + volatile State state = State.QUEUED; // written under the owning Target monitor + + Pending(String target, String text, TurnToken token, CompletableFuture delivered) { + this.target = target; + this.text = text; + this.token = token; + this.delivered = delivered; + } + + String text() { + return text; + } + + TurnToken token() { + return token; + } + + CompletableFuture delivered() { + return delivered; + } } /** Per-worker delivery state, guarded by its own monitor (single writer per worker). */ @@ -177,15 +227,47 @@ public final class Injector { *

Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the * target" and "queue the message" and orphan it in a target it just removed. */ - public CompletableFuture enqueue(String target, String text, TurnToken token) { + public Delivery enqueue(String target, String text, TurnToken token) { CompletableFuture delivered = new CompletableFuture<>(); - Pending p = new Pending(text, token, delivered); + Pending p = new Pending(target, text, token, delivered); targets.compute(target, (_, existing) -> { Target t = (existing != null) ? existing : new Target(); t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop return t; }); - return delivered; + return new Delivery(p); + } + + /** + * Cancel this exact queued delivery. The target monitor serializes this operation with + * {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED} + * rather than claiming the message remained queued. + */ + public Cancellation cancel(Delivery delivery) { + Pending p = delivery.pending; + Target t = targets.get(p.target); + if (t == null) { + return cancellationOf(p); + } + synchronized (t) { + if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) { + return cancellationOf(p); + } + p.state = Pending.State.CANCELLED; + if (isQuiescent(t)) { + targets.remove(p.target, t); + } + return Cancellation.CANCELLED; + } + } + + private static Cancellation cancellationOf(Pending p) { + return p.state == Pending.State.DELIVERED ? Cancellation.DELIVERED : Cancellation.NOT_DELIVERED; + } + + private static boolean isQuiescent(Target t) { + return t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion + && !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved; } /** @@ -274,6 +356,7 @@ public final class Injector { try { agentsFor(target).send(target, p.text()); t.queue.poll(); + p.state = Pending.State.DELIVERED; t.awaitingPickup = true; t.awaitingCompletion = true; t.turnObserved = false; @@ -283,6 +366,7 @@ public final class Injector { // Delivery failed at herdr; drop the poisoned message and surface it // rather than blocking the queue behind it. t.queue.poll(); + p.state = Pending.State.NOT_DELIVERED; sent = p; sendError = e; } @@ -293,6 +377,9 @@ public final class Injector { // fail every queued message and release the target (CB-114) instead of // polling it indefinitely with the caller's future never completing. notReady = new ArrayList<>(t.queue); + for (Pending pending : notReady) { + pending.state = Pending.State.NOT_DELIVERED; + } log.warn("readiness grace for {} expired after {} polls ({}s): target never " + "became deliverable, so failing {} queued message(s) that never " + "reached its pane", @@ -340,8 +427,7 @@ public final class Injector { // 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 (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion - && !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved) { + if (isQuiescent(t)) { targets.remove(target, t); } } @@ -432,6 +518,9 @@ public final class Injector { boolean hadDeliveredTurn; synchronized (t) { pending = new ArrayList<>(t.queue); + for (Pending p : pending) { + p.state = Pending.State.NOT_DELIVERED; + } t.queue.clear(); hadDeliveredTurn = t.awaitingCompletion; t.awaitingCompletion = false; 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 557a864..19039ac 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -863,16 +863,22 @@ public final class MessageService { if (onAccepted != null) { onAccepted.run(); } - CompletableFuture delivered = injector.enqueue(target, content, token); + Injector.Delivery delivery = injector.enqueue(target, content, token); try { Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId())); } catch (TimeoutException e) { - boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); + boolean wasDelivered = delivery.completion().isDone() + && !delivery.completion().isCompletedExceptionally(); + if (!wasDelivered) { + // 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; + } 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). + // CB-640: record that delivery did not happen for fleet health (see + // queuedDeliveries). The exact Pending was cancelled, so it cannot arrive later. queuedDeliveries.put(target, Boolean.TRUE); } return recorded(new Reply( diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java index 72e0c90..77d4056 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java @@ -48,7 +48,7 @@ class InjectorTest { @Test void deliversWhenIdle() { - CompletableFuture f = injector.enqueue(T, "hello", TestTurnTokens.inert(T)); + CompletableFuture f = injector.enqueue(T, "hello", TestTurnTokens.inert(T)).completion(); assertFalse(f.isDone(), "not delivered until an injectable status arrives"); injector.onStatus(T, AgentStatus.IDLE); assertTrue(f.isDone()); @@ -170,6 +170,32 @@ class InjectorTest { assertEquals(List.of("a", "b", "c"), sent()); } + @Test + void cancellingTheMiddleDeliveryKeepsTheFollowingDeliveryReachable() { + Injector.Delivery first = injector.enqueue(T, "same text", TestTurnTokens.inert(T)); + Injector.Delivery cancelled = injector.enqueue(T, "same text", TestTurnTokens.inert(T)); + injector.enqueue(T, "after cancelled", TestTurnTokens.inert(T)); + + assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(cancelled)); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + injector.onStatus(T, AgentStatus.IDLE); + + assertEquals(List.of("same text", "after cancelled"), sent(), + "cancellation must match the exact Delivery and preserve the remaining FIFO queue"); + assertTrue(first.completion().isDone()); + } + + @Test + void cancellationReportsDeliveredWhenPickupWonTheRace() { + Injector.Delivery delivery = injector.enqueue(T, "already sent", TestTurnTokens.inert(T)); + injector.onStatus(T, AgentStatus.IDLE); + + assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(delivery), + "a cancellation after pickup must not claim that the text stayed queued"); + assertEquals(List.of("already sent"), sent()); + } + @Test void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() { assertTrue(injector.activeTargets().isEmpty()); @@ -378,7 +404,7 @@ class InjectorTest { void sendFailureDropsMessageAndFailsItsFuture() { FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed"); Injector inj = new Injector(new AgentControl(failing)); - CompletableFuture f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)); + CompletableFuture f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion(); inj.onStatus(T, AgentStatus.IDLE); assertTrue(f.isCompletedExceptionally()); @@ -387,7 +413,7 @@ class InjectorTest { @Test void dropFailsPendingWaiters() { - CompletableFuture f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T)); + CompletableFuture f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T)).completion(); injector.drop(T, new HerdrException("worker gone", "pane_not_found", null)); assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes"); } @@ -396,8 +422,8 @@ class InjectorTest { void dropPassesTheRealCauseForQueuedAndDeliveredWork() { Captor cap = new Captor(); Injector inj = new Injector(new AgentControl(herdr), cap); - CompletableFuture delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T)); - CompletableFuture queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T)); + CompletableFuture delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T)).completion(); + CompletableFuture queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T)).completion(); inj.onStatus(T, AgentStatus.IDLE); // deliver the first message inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued @@ -449,7 +475,7 @@ class InjectorTest { Captor cap = new Captor(); List forgotten = new ArrayList<>(); Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add); - CompletableFuture f = inj.enqueue(T, "task", TestTurnTokens.inert(T)); + CompletableFuture f = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion(); for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE); @@ -563,7 +589,7 @@ class InjectorTest { StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10); poller.start(); try { - CompletableFuture delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T)); + CompletableFuture delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T)).completion(); delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send } finally { poller.stop(); @@ -582,7 +608,7 @@ class InjectorTest { void deliveredFutureCarriesSendFailure() { FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed"); Injector inj = new Injector(new AgentControl(failing)); - CompletableFuture f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)); + CompletableFuture f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion(); inj.onStatus(T, AgentStatus.IDLE); ExecutionException ex = assertThrows(ExecutionException.class, f::get); assertInstanceOf(HerdrException.class, ex.getCause()); diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java index d748534..739a86d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java @@ -41,7 +41,7 @@ class StatusPollerRoutingTest { poller.start(); try { CompletableFuture delivered = - injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)); + injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion(); // Must resolve quickly: refining against the WRONG daemon (member) never classifies // out of UNKNOWN, so this would time out under the bug. delivered.get(2, TimeUnit.SECONDS); @@ -65,7 +65,7 @@ class StatusPollerRoutingTest { poller.start(); try { CompletableFuture delivered = - injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)); + injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion(); assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS), "a lead target must never be refined from the member daemon's pane content"); } finally { 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 9a6cbda..f9f196b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -258,7 +258,7 @@ class MessageServiceTest { injector.onStatus(T, AgentStatus.IDLE); // first delivery injector.onStatus(T, AgentStatus.WORKING); // first turn in flight - CompletableFuture queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T)); + CompletableFuture queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T)).completion(); CompletableFuture waiter = rendezvous.currentWaiter(T); injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null)); @@ -1768,7 +1768,10 @@ class MessageServiceTest { 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"); + "a TIMED_OUT_QUEUED send still records the undelivered delivery for fleet health"); + injector.onStatus(T, AgentStatus.IDLE); + assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.prompt")), + "a TIMED_OUT_QUEUED send must be cancelled, not delivered when the worker later goes idle"); } @Test