Merge CB-588: nudge the lead when an async ticket goes terminal

An async delegation ticket (bridge_send wait:false) resolves on MessageService.reply's
rendezvous fast path, which returns before onReplyQueued. So CB-307's push loop only ever
heard about the durable-inbox case, and the mode CLAUDE.md tells leads to prefer never
nudged anyone. Closes gitea #72.

Adds a second, independent reminder schedule keyed by the lead terminal, so several
tickets finishing together coalesce into one nudge. The CB-307 path is untouched.

Three defects were found in review and fixed before merge:
 * a pendingTickets entry outlived the ticket it named. poll() returns null once
   pruneTerminalTickets drops a ticket, so ticketCollected was never reached and the
   entry leaked for the daemon's life, riding along on every later nudge and sending
   the lead after a ticket bridge_poll can no longer find.
 * a lost nudge: a ticket landing between decideTickets returning STOP and
   activeLeads.remove coalesced onto a schedule that was about to die. That is the
   exact failure this ticket exists to remove, reintroduced in a narrow window.
 * the success direction was unpinned in tests, and the comment listing the paths that
   complete the future was short by several.

The obvious fix for the second one was wrong: restarting on any pending ticket defeats
the reminder cap, because a never-collected ticket at cap is expected to still be there.
The fix diffs against a snapshot taken before the decision, so only a ticket that truly
arrived during the window restarts the schedule.

Verified on my own unpiped build: 802 tests, 0 failures, BUILD SUCCESS.
Two reviewers on the diff; the loop-gating finding they raised is split out as CB-590.
This commit is contained in:
Dai Ha
2026-08-15 17:06:12 +02:00
4 changed files with 715 additions and 12 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. */
@@ -191,7 +200,9 @@ public final class MessageService {
* Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}.
*
* @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply
* branch ({@link #reply}) so it can nudge the primary to drain the inbox
* branch ({@link #reply}) so it can nudge the primary to drain the inbox, and
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
* terminal phase, and whenever {@link #poll} hands a terminal ticket to its caller
*/
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
ReplyInbox inbox, ReplyPushLoop pushLoop) {
@@ -206,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. */
@@ -571,8 +589,24 @@ 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
// worker paused in bridge_ask leaves it running, per finishAsyncTask's own contract — so
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
// below on any non-QUESTION outcome of send() — a worker's bridge_reply, the CB-106
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
// abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the
// one an async ticket always takes) never told the push loop anything happened — see the
// class javadoc on sendAsync/CB-107.
task.future.whenComplete((reply, ex) -> {
boolean failed = ex != null || reply == null || !reply.completed();
pushLoop.onTicketTerminal(ticket, target, failed);
});
}
asyncExecutor.submit(() -> {
try {
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
@@ -610,6 +644,11 @@ public final class MessageService {
}
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
}
// CB-588: the ticket is terminal and being handed to the caller right here — tell the push
// loop it is collected so a later tick's nudge never names a ticket the lead already has.
if (pushLoop != null) {
pushLoop.ticketCollected(ticket);
}
Reply r;
try {
r = f.getNow(null);
@@ -640,10 +679,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. */
@@ -8,9 +8,12 @@ import dev.ltms.bridged.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/**
* Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own
@@ -23,11 +26,27 @@ import java.util.concurrent.TimeUnit;
*
* <p>Bounded: at most {@link #maxReminders} nudges per target, with a configurable backoff
* between them. The reply is never lost — the durable inbox is the backstop.
*
* <p><strong>CB-588 ticket nudges.</strong> {@link #onTicketTerminal(String, String, boolean)} is a
* second, independent entry point for an async delegation ticket ({@code bridge_send(wait:false)})
* reaching a terminal phase. That path takes the rendezvous fast path in {@link MessageService#reply}
* and never reaches {@link #onReplyQueued}, so without this the lead's own charter — prefer
* {@code wait:false} for anything non-trivial — was exactly the mode this loop failed to cover. It
* reuses the same status gating, bounded/backoff reminders, and metrics, keyed by the nudge-receiving
* lead terminal rather than the worker target so several tickets finishing together coalesce into one
* nudge. The two entry points do not interact: {@link #onReplyQueued} / {@link #decide} / their nudge
* text and bound are unchanged.
*/
public final class ReplyPushLoop {
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
/** CB-588: singular form, one uncollected ticket. */
static final String TICKET_NUDGE_FORMAT =
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
/** CB-588: coalesced form, several uncollected tickets for the same lead. */
static final String TICKETS_NUDGE_FORMAT =
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
private final PrimaryRegistry primaryRegistry;
private final AgentControl agents;
@@ -39,6 +58,10 @@ public final class ReplyPushLoop {
/** Track targets that have an active schedule. */
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
/** CB-588: tickets that have gone terminal but not yet been polled, keyed by ticket. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
/** CB-588: leads with an active ticket-reminder schedule. */
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
@@ -173,23 +196,211 @@ public final class ReplyPushLoop {
scheduler.schedule(() -> tick(target, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
}
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
private record PendingTicket(String ticket, String lead, boolean failed) {
}
/**
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
* before {@link #onReplyQueued} is ever called, so without this entry point a ticket finishing
* that way never nudged anyone (CB-588 / gitea #72).
*
* <p>Idempotent per lead: several tickets going terminal for the same lead while its schedule is
* already active coalesce onto that schedule's next tick rather than firing a nudge each.
*
* @param ticket the ticket to nudge about
* @param target the worker session the ticket was sent to — resolves which lead delegated it
* @param failed whether the ticket ended in a failure phase rather than {@code DONE}
*/
public void onTicketTerminal(String ticket, String target, boolean failed) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
ticket, target);
return;
}
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
if (activeLeads.putIfAbsent(lead.get(), Boolean.TRUE) != null) {
log.debug("push: ticket reminder loop already active for lead {}, {} coalesced in",
lead.get(), ticket);
return;
}
log.debug("push: starting ticket reminder loop for lead {}", lead.get());
scheduleTicketTick(lead.get(), 0);
}
/**
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
* lead already has (CB-588 acceptance #5). A ticket that was never pending (unknown ticket, or
* one nudged with no push loop configured) is a no-op.
*/
public void ticketCollected(String ticket) {
pendingTickets.remove(ticket);
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
private List<PendingTicket> pendingFor(String lead) {
return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList();
}
/**
* Pure decision function for ticket nudges, mirroring {@link #decide(String, int)} but keyed by
* the nudge-receiving lead terminal rather than the worker session — several tickets from
* different workers delegated by the same lead coalesce onto it.
*/
Action decideTickets(String lead, int reminderCount) {
if (pendingFor(lead).isEmpty()) {
log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
}
AgentStatus status;
try {
status = agents.status(lead);
} catch (RuntimeException e) {
log.debug("push: status check failed for lead {}, will retry", lead, e);
return Action.WAIT_BUSY;
}
if (status.injectable()) {
return Action.INJECT;
}
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
return Action.WAIT_BUSY;
}
/** Execute one ticket-loop tick — called on the scheduler thread. */
private void ticketTick(String lead, int reminderCount) {
Set<String> pendingBefore = pendingIdsFor(lead);
var action = decideTickets(lead, reminderCount);
switch (action) {
case INJECT -> {
injectTicketNudge(lead, reminderCount);
scheduleTicketTick(lead, reminderCount + 1);
}
case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount);
case STOP -> stopOrRestartTicketLoop(lead, pendingBefore);
}
}
/** Ticket IDs pending for {@code lead} right now, as a plain snapshot for race comparison. */
private Set<String> pendingIdsFor(String lead) {
return pendingFor(lead).stream().map(PendingTicket::ticket).collect(Collectors.toUnmodifiableSet());
}
/**
* Release {@code lead}'s active-schedule slot, then restart it only if a ticket landed that
* {@code pendingBefore} — the snapshot taken just before this tick's decision — did not already
* account for. {@code onTicketTerminal} reads {@code activeLeads} to decide whether to coalesce
* onto an existing schedule or start one, so a ticket that lands between {@code decideTickets}
* returning {@link Action#STOP} and this removal running sees the (soon-to-be-stale) slot as
* occupied, coalesces onto a schedule that is about to die, and gets no nudge scheduled at all —
* a lost nudge, the exact failure CB-588 exists to remove (found in review, gitea PR #73).
*
* <p>Restarting on ANY non-empty {@code pendingFor(lead)} would be wrong: when STOP is reached
* because the reminder cap was hit rather than the backlog draining, the same never-collected
* ticket is expected to still be sitting there — that is the cap doing its job — and restarting
* would nudge about it forever, defeating the bound (the original CB-307 bounded-reminder
* guarantee, carried into CB-588 by acceptance criterion #7 — this exact regression showed up as
* two existing tests failing once a naive "any pending ticket restarts" version of this fix went
* in: {@code successfulTicketNudgeIncrementsDelivered} and {@code ticketNudgesSendUpToCapThenStop}).
* Diffing the current pending set against {@code pendingBefore} tells the two cases apart: a
* ticket present before this tick's decision is stale backlog, not a race; only a ticket absent
* from {@code pendingBefore} can only have arrived during the decision-to-release window, which is
* exactly the race this method closes.
*
* <p>Package-private so a test can drive the interleaving directly rather than trying to force a
* genuine thread race: pass the exact {@code pendingBefore} snapshot a race requires (or does
* not) and call this to prove the recheck responds correctly either way.
*
* <p>Terminates rather than spinning: this method restarts the schedule at most once per call, and
* a fresh {@link #onTicketTerminal} racing the recheck below still terminates in one of two ways —
* either it observes the slot already vacated (by the {@code activeLeads.remove} above, which
* happens-before this recheck in program order) and claims it itself, or it lands first and this
* recheck then observes its ticket in {@code pendingTickets} and reclaims the slot instead. Exactly
* one side always wins; neither can miss the other, so this never loops on its own account.
*/
void stopOrRestartTicketLoop(String lead, Set<String> pendingBefore) {
activeLeads.remove(lead);
boolean ticketRacedIn = pendingFor(lead).stream().anyMatch(t -> !pendingBefore.contains(t.ticket()));
if (ticketRacedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: a ticket for lead {} raced the reminder loop's stop — restarting", lead);
scheduleTicketTick(lead, 0);
return;
}
log.debug("push: ticket reminder loop ended for lead {}", lead);
}
/** Send the coalesced ticket nudge and log the event. */
private void injectTicketNudge(String lead, int reminderCount) {
// Re-read rather than threading it down from decideTickets(): a ticket can be collected (or
// another can arrive) between the decision and the injection.
List<PendingTicket> pending = pendingFor(lead);
if (pending.isEmpty()) {
log.debug("push: pending tickets for lead {} drained before the nudge could be sent", lead);
return;
}
String nudge = formatTicketsNudge(pending);
try {
agents.send(lead, nudge);
log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)",
reminderCount + 1, maxReminders, lead, pending.size());
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} for {} ticket(s) (reminder {}/{}): {}",
lead, pending.size(), reminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next ticket-loop tick on the scheduler thread pool. */
private void scheduleTicketTick(String lead, int nextReminderCount) {
scheduler.schedule(() -> ticketTick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
}
/** Render one or several pending tickets as a single nudge line. */
private static String formatTicketsNudge(List<PendingTicket> pending) {
if (pending.size() == 1) {
PendingTicket t = pending.get(0);
return TICKET_NUDGE_FORMAT.formatted(t.ticket(), t.failed() ? " (FAILED)" : "", t.ticket());
}
long failedCount = pending.stream().filter(PendingTicket::failed).count();
String ids = pending.stream()
.map(t -> t.failed() ? t.ticket() + " (FAILED)" : t.ticket())
.collect(Collectors.joining(", "));
String failedNote = failedCount > 0 ? " (%d failed)".formatted(failedCount) : "";
return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids);
}
// --- lifecycle -----------------------------------------------------------------------------
/**
* Whether any reminder loop is currently active for some target (CB-551). The idle-lead heartbeat
* uses this to stand aside: while the push loop is actively nudging the lead, a concurrent
* heartbeat injection would start a second competing turn in the same pane — racing loops multiply
* turns and context burn. "Active" means a schedule exists in {@link #activeTargets}; the set is
* bounded by what has been triggered, not by any persistent state.
* turns and context burn. "Active" means a schedule exists in {@link #activeTargets} or
* {@link #activeLeads} (CB-588 ticket nudges are a second source of pane injections the heartbeat
* must equally stand aside for); the sets are bounded by what has been triggered, not by any
* persistent state.
*/
public boolean isActive() {
return !activeTargets.isEmpty();
return !activeTargets.isEmpty() || !activeLeads.isEmpty();
}
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
activeLeads.clear();
pendingTickets.clear();
}
/** @see #stop() */