CB-588 follow-up: reclaim a pruned ticket's pendingTickets entry too
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m0s

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.
This commit is contained in:
Dai Ha
2026-08-15 16:41:23 +02:00
parent 3d10ed385c
commit 6ebad2a91f
2 changed files with 103 additions and 9 deletions
@@ -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<Reply> 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<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> 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.
*
* <p>{@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. */
@@ -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;