From 6ebad2a91f81dae6fcd1d50fdc45e819bae27323 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 16:41:23 +0200 Subject: [PATCH] CB-588 follow-up: reclaim a pruned ticket's pendingTickets entry too MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit tasks is the sole authority on whether a ticket exists, but pruneTerminalTickets dropped entries from it without telling ReplyPushLoop. ticketCollected(ticket) was only ever called from MessageService.poll's terminal branch, which pruneTerminalTickets short-circuits past once a ticket is gone (poll returns null at the top). A ticket the lead never polled — or one the reminder cap already gave up on — was pruned from tasks but never collected in ReplyPushLoop.pendingTickets, so it rode along on every later nudge to the same lead forever, naming a ticket bridge_poll could no longer find, and the map itself never shrank. pruneTerminalTickets now calls pushLoop.ticketCollected for every ticket it actually removes (guarded on pushLoop != null), so a pending nudge entry lives exactly as long as its ticket is pollable. Reused tasks as the only removal trigger rather than adding a second live query back into MessageService — no new source of truth. Added an injectable clock (LongSupplier nowNanos, defaulting to System::nanoTime) to MessageService, mirroring the SessionManager/ SessionReaper nowNanos seam, so a test can cross the 10-minute TICKET_TTL_NANOS deterministically instead of sleeping for real. Confirmed the new regression test fails against the prior pruneTerminalTickets (a stale ticket rides along on a later coalesced nudge) before restoring the fix. --- .../dev/ltms/bridged/msg/MessageService.java | 49 ++++++++++++--- .../ltms/bridged/msg/MessageServiceTest.java | 63 ++++++++++++++++++- 2 files changed, 103 insertions(+), 9 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 41d9170..eaa2bb3 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -20,6 +20,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.locks.ReentrantLock; +import java.util.function.LongSupplier; /** * The blocking delegation feature (CB-104): deliver {@code content} into a worker and block until @@ -52,8 +53,12 @@ public final class MessageService { */ private static final long ASYNC_TIMEOUT_MS = 30 * 60 * 1_000L; - /** How long a finished (terminal) ticket is retained for polling before it is pruned. */ - private static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L; + /** + * How long a finished (terminal) ticket is retained for polling before it is pruned. Package- + * private (not {@code private}) so a test can advance an injected clock past it deterministically + * instead of duplicating the magic number or sleeping for real. + */ + static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L; /** Outcome of a blocking send. */ public enum Outcome { @@ -161,12 +166,13 @@ public final class MessageService { private static final class Task { private final String target; private final CompletableFuture future = new CompletableFuture<>(); - private final long createdNanos = System.nanoTime(); + private final long createdNanos; private volatile Reply question; private volatile String turnId; - private Task(String target) { + private Task(String target, long createdNanos) { this.target = target; + this.createdNanos = createdNanos; } } @@ -176,6 +182,9 @@ public final class MessageService { private final ReplyInbox inbox; private final ReplyPushLoop pushLoop; private final Metrics metrics; // CB-502: nullable — no registry in unit tests + // CB-588: injectable so pruneTerminalTickets' 10-minute TICKET_TTL_NANOS can be exercised in a + // test without a real wait — same seam SessionManager already uses for its idle reaper (nowNanos). + private final LongSupplier nowNanos; private final ConcurrentHashMap sessionLocks = new ConcurrentHashMap<>(); private final ConcurrentHashMap tasks = new ConcurrentHashMap<>(); /** Async task that owns each exact forward rendezvous waiter. */ @@ -208,12 +217,19 @@ public final class MessageService { */ public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, ReplyPushLoop pushLoop, Metrics metrics) { + this(agents, injector, rendezvous, inbox, pushLoop, metrics, System::nanoTime); + } + + /** Test constructor with an injectable clock (CB-588: exercise the ticket-prune TTL without a real wait). */ + MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, + ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) { this.agents = agents; this.injector = injector; this.rendezvous = rendezvous; this.inbox = inbox; this.pushLoop = pushLoop; this.metrics = metrics; + this.nowNanos = nowNanos; } /** Create with an explicit {@link ReplyInbox} and no push loop. */ @@ -573,7 +589,7 @@ public final class MessageService { */ public String sendAsync(String target, String content, Runnable onAccepted) { String ticket = "task-" + ticketSeq.incrementAndGet(); - Task task = new Task(target); + Task task = new Task(target, nowNanos.getAsLong()); tasks.put(ticket, task); if (pushLoop != null) { // CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a @@ -659,10 +675,27 @@ public final class MessageService { } } - /** Drop finished tickets older than the TTL so the registry cannot grow without bound. */ + /** + * Drop finished tickets older than the TTL so {@link #tasks} cannot grow without bound. + * + *

{@code tasks} is the sole authority on whether a ticket still exists — {@link #poll} returns + * {@code null} the instant a ticket is gone from here, before it ever reaches the terminal branch + * that calls {@link ReplyPushLoop#ticketCollected}. Without telling the push loop about a prune + * too, its own {@code pendingTickets} entry would outlive the ticket it names: an unpolled ticket + * (or one the reminder cap already gave up on) is pruned here but never collected there, so it + * lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead + * — naming a ticket {@code bridge_poll} can no longer find (CB-588 follow-up). + */ private void pruneTerminalTickets() { - long cutoff = System.nanoTime() - TICKET_TTL_NANOS; - tasks.values().removeIf(t -> t.future.isDone() && t.createdNanos < cutoff); + long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS; + tasks.entrySet().removeIf(e -> { + Task t = e.getValue(); + boolean expired = t.future.isDone() && t.createdNanos < cutoff; + if (expired && pushLoop != null) { + pushLoop.ticketCollected(e.getKey()); + } + return expired; + }); } /** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */ diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index b257e50..67057f5 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -733,13 +733,18 @@ class MessageServiceTest { } private PushWiring wireWithPushLoop(int maxReminders, long backoffMs) { + return wireWithPushLoop(maxReminders, backoffMs, System::nanoTime); + } + + /** As above, with an injectable clock (CB-588 follow-up: exercise pruneTerminalTickets' TTL). */ + private PushWiring wireWithPushLoop(int maxReminders, long backoffMs, java.util.function.LongSupplier nowNanos) { PrimaryRegistry registry = new PrimaryRegistry(null); registry.recordDelegation(T, LEAD); FakeHerdr leadHerdr = new FakeHerdr(); AgentControl leadAgents = new AgentControl(leadHerdr); var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor(); ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs); - MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop); + MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos); return new PushWiring(service, leadHerdr, scheduler); } @@ -840,6 +845,62 @@ class MessageServiceTest { assertEquals("async result", view.reply()); } + /** + * CB-588 follow-up: {@code tasks} is the sole authority on whether a ticket exists, and + * {@code pruneTerminalTickets} drops entries from it once {@link MessageService#TICKET_TTL_NANOS} + * elapses. Before this test, that prune never told {@code ReplyPushLoop} — its own + * {@code pendingTickets} entry for a pruned, never-collected ticket had no remover at all, so it + * rode along on every later nudge to the same lead, naming a ticket {@code bridge_poll} could no + * longer find. Uses the injectable clock (mirroring {@code SessionManager}'s {@code nowNanos} seam + * for its idle reaper) to cross the 10-minute TTL without a real wait. + */ + @Test + void aPrunedTicketIsReclaimedFromThePushLoopNotLeakedForever() throws Exception { + java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L); + // maxReminders=1 + a short backoff: the stale ticket gets its one legitimate reminder, then + // decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix) + // pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of + // the bug report, reached deterministically rather than by timing it against a live tick. + try (var wiring = wireWithPushLoop(1, 50, clock::get)) { + String stale = wiring.service().sendAsync(T, "first task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "stale result")); + + // Let the reminder loop fire its one nudge and hit the cap (STOP removes it from + // activeLeads; pendingTickets is untouched either way — that asymmetry is the bug). + awaitNudge(wiring.leadHerdr()); + Thread.sleep(300); + assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale), + "sanity: the stale ticket's own reminder must have fired first"); + + // Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly. + clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1)); + + // A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes. + String fresh = wiring.service().sendAsync(T, "second task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "fresh result")); + + // The fresh ticket restarts the (now-dormant) reminder loop with its own nudge. + long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count(); + long deadline = System.currentTimeMillis() + 3000; + while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before + && System.currentTimeMillis() < deadline) { + Thread.sleep(10); + } + String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); + assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge); + assertFalse(latestNudge.contains(stale), + "a pruned ticket must never be named in a later nudge — it is gone and bridge_poll " + + "on it would return nothing: " + latestNudge); + + // And bridge_poll(ticket=stale) really does return nothing now — the nudge would have lied. + assertNull(wiring.service().poll(stale), "the pruned ticket must actually be gone, not just unmentioned"); + } + } + private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket, MessageService.Phase phase) throws Exception { long deadline = System.currentTimeMillis() + 3000;