CB-590 follow-up: per-source reminder budgets for the push loop #90

Merged
ltms merged 2 commits from worker/cb590fix-185e9a-10 into main 2026-08-16 17:28:58 +02:00
2 changed files with 512 additions and 282 deletions
@@ -8,6 +8,8 @@ import dev.ltms.bridged.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@@ -16,35 +18,43 @@ 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
* herdr pane when a worker reply lands with no live {@code bridge_send} to resolve it.
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
* (CB-588).
*
* <p>The loop is triggered by {@link #onReplyQueued(String)} (called from
* {@link MessageService#reply} after the durable inbox publish). It checks four conditions
* at each tick via {@link #decide(String, int)}, then either injects a drain nudge,
* waits for the primary to become injectable, or stops reminding.
* <p><strong>CB-590: one schedule per lead.</strong> Both kinds of work are triggered through
* their own entry point — {@link #onReplyQueued(String)} and
* {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be
* nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}.
* Earlier this was two independent schedules (one keyed by worker target for replies, one keyed
* by lead for tickets) that could both decide to inject into the same pane in the same window —
* a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn
* twice. Collapsing to one schedule per lead makes that structurally impossible: at most one
* scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
*
* <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.
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
* via {@code STOP}.
*/
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. */
/** Coalesced form, several uncollected replies for the same lead. */
static final String REPLIES_NUDGE_FORMAT =
"%d workers returned replies — run bridge_poll(target=...) for each to collect them: %s";
/** 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. */
/** 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";
@@ -56,11 +66,11 @@ public final class ReplyPushLoop {
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/** 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. */
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
/** 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. */
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
@@ -89,178 +99,79 @@ public final class ReplyPushLoop {
}
}
// --- decision logic (package-private for unit-testing) -------------------------------------
// --- pending-work lookups (package-private for unit-testing) -------------------------------
/** The action the loop should take for a target at the given reminder count. */
/** The action the loop should take for a lead at the given reminder count. */
enum Action { INJECT, WAIT_BUSY, STOP }
/**
* Pure decision function: examine the current state and return what the loop should do.
*
* @param target the worker session (target terminal id)
* @param reminderCount how many nudges have been sent so far for this target
* @return the action the caller should take
* Reply targets still pending for {@code lead} — registered via {@link #onReplyQueued} and
* whose inbox still holds an unacked message. A target whose inbox has since drained (acked,
* or collected via a live {@code bridge_send} rendezvous instead) is dropped from
* {@link #pendingReplies} here rather than lingering forever; there is no explicit "reply
* collected" callback the way {@link #ticketCollected} exists for tickets, so the inbox itself
* is the only signal.
*/
Action decide(String target, int reminderCount) {
// CB-532: the destination is per-delegation — the lead that sent this worker its work, not
// "the primary". With two leads orchestrating one fleet the singular question has no right
// answer, and answering it anyway interrupted whichever lead happened to call bridge_send
// first with results it never asked for.
var nudgeTarget = primaryRegistry.nudgeTargetFor(target);
if (nudgeTarget.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, stopping reminder", target);
return Action.STOP;
}
if (inbox.peek(target).isEmpty()) {
log.debug("push: inbox empty for {}, stopping reminder", target);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: reminder cap ({}) reached for {}, stopping", maxReminders, target);
countNudge("exhausted");
return Action.STOP;
}
String leadTerminal = nudgeTarget.get();
AgentStatus status;
try {
status = agents.status(leadTerminal);
} catch (RuntimeException e) {
log.debug("push: status check failed for lead {}, will retry", leadTerminal, e);
return Action.WAIT_BUSY;
}
if (status.injectable()) {
return Action.INJECT;
}
log.debug("push: lead {} is {} (not injectable), waiting", leadTerminal, status);
return Action.WAIT_BUSY;
}
// --- public entrypoint ---------------------------------------------------------------------
/**
* Called when a reply is queued for {@code target}. Idempotent per target: a second call while
* a schedule is active is a no-op. The schedule nudges the primary, then schedules a follow-up
* check (reminder on backoff, or re-check on WAIT_BUSY), until the inbox is empty or the cap
* is reached.
*/
public void onReplyQueued(String target) {
if (activeTargets.putIfAbsent(target, Boolean.TRUE) != null) {
log.debug("push: already active for {}, ignoring duplicate trigger", target);
return; // already scheduled
}
log.debug("push: starting reminder loop for {}", target);
scheduleNext(target, 0);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String target, int reminderCount) {
var action = decide(target, reminderCount);
switch (action) {
case INJECT -> {
injectNudge(target, reminderCount);
scheduleNext(target, reminderCount + 1);
}
// Re-check after the configured backoff; the primary may become injectable soon.
case WAIT_BUSY -> scheduleNext(target, reminderCount);
case STOP -> {
activeTargets.remove(target);
log.debug("push: reminder loop ended for {}", target);
private Set<String> pendingReplyTargetsFor(String lead) {
Set<String> result = new HashSet<>();
for (var entry : pendingReplies.entrySet()) {
String target = entry.getKey();
String owningLead = entry.getValue();
if (!lead.equals(owningLead)) continue;
if (inbox.peek(target).isEmpty()) {
pendingReplies.remove(target, owningLead);
continue;
}
result.add(target);
}
return result;
}
/** Send the nudge and log the event. */
private void injectNudge(String target, int reminderCount) {
// Re-read rather than threading it down from decide(): the delegating lead can change
// between the decision and the injection, and the nudge should follow the current one.
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: lead for {} disappeared before the nudge could be sent", target);
return;
}
String leadTerminal = lead.get();
String nudge = NUDGE_FORMAT.formatted(target, target);
try {
agents.send(leadTerminal, nudge);
log.debug("push: nudge {}/{} sent to lead {} for target {}",
reminderCount + 1, maxReminders, leadTerminal, target);
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} for target {} (reminder {}/{}): {}",
leadTerminal, target, reminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String target, int nextReminderCount) {
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) {
private List<PendingTicket> pendingTicketsFor(String lead) {
return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList();
}
/** Ticket ids still pending for {@code lead} — a plain snapshot for race comparison. */
private Set<String> pendingTicketIdsFor(String lead) {
return pendingTicketsFor(lead).stream().map(PendingTicket::ticket)
.collect(Collectors.toUnmodifiableSet());
}
/**
* 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.
* Pure decision function: examine everything pending for {@code lead} — reply targets and
* tickets alike — and return what the loop should do.
*
* <p><strong>CB-590-fix: one schedule, two budgets.</strong> The single per-lead schedule
* (CB-590) still ticks once for both sources, but each source is capped independently —
* {@code replyReminderCount} against a reply target still pending, {@code ticketReminderCount}
* against a ticket still pending. A busy reply stream that exhausts its own cap must not stop
* the loop from nudging about a ticket that still has budget left, and vice versa: either
* source being eligible (has pending work AND is under its own cap) is enough for
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
* {@link Action#STOP}.
*
* @param lead the lead terminal to nudge
* @param replyReminderCount how many nudges have covered pending reply work for this lead
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
* @return the action the caller should take
*/
Action decideTickets(String lead, int reminderCount) {
if (pendingFor(lead).isEmpty()) {
log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork) {
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
if (!replyEligible && !ticketEligible) {
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
}
@@ -278,95 +189,194 @@ public final class ReplyPushLoop {
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);
}
}
// --- public entrypoints ----------------------------------------------------------------------
/** 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());
/**
* Called when a reply is queued for {@code target}. Resolves the lead delegating to
* {@code target} (CB-532) and coalesces onto that lead's single reminder schedule — starting
* one if none is active, joining an already-active one otherwise. A no-op if no lead is known
* to be waiting on {@code target}: there is nobody to nudge yet, and the durable inbox is the
* backstop until a lead is recorded.
*/
public void onReplyQueued(String target) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
return;
}
pendingReplies.put(target, lead.get());
startOrCoalesce(lead.get());
}
/**
* 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).
* 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>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>Resolves the delegating lead the same way {@link #onReplyQueued} does and coalesces onto
* the same per-lead schedule (CB-590) — several tickets, or a ticket and a reply, finishing
* for the same lead while its schedule is already active all ride the existing 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));
startOrCoalesce(lead.get());
}
/**
* 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. 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);
}
// --- the schedule ----------------------------------------------------------------------------
/** Start a reminder schedule for {@code lead}, or join the one already running. */
private void startOrCoalesce(String lead) {
if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead, 0, 0);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
var action = decide(lead, replyReminderCount, ticketReminderCount);
switch (action) {
case INJECT -> {
injectNudge(lead, replyReminderCount, ticketReminderCount);
// Only the source(s) actually eligible this tick spend a unit of their own budget —
// an exhausted source riding along in the combined message (still pending, still
// named) does not get charged again; its count stays put until it drains.
boolean replyEligible = !repliesBefore.isEmpty() && replyReminderCount < maxReminders;
boolean ticketEligible = !ticketsBefore.isEmpty() && ticketReminderCount < maxReminders;
scheduleNext(lead,
replyEligible ? replyReminderCount + 1 : replyReminderCount,
ticketEligible ? ticketReminderCount + 1 : ticketReminderCount);
}
// Re-check after the configured backoff; the lead may become injectable soon.
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
}
}
/**
* Release {@code lead}'s active-schedule slot, then restart it only if work landed that
* {@code repliesBefore} / {@code ticketsBefore} — the snapshots taken just before this tick's
* decision — did not already account for. {@link #onReplyQueued} / {@link #onTicketTerminal}
* read {@link #activeLeads} to decide whether to coalesce onto an existing schedule or start
* one, so work that lands between {@link #decide} 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, exactly what CB-588 (and
* now CB-590) exist to remove (originally found in review, gitea PR #73, for the ticket-only
* loop; carried forward here for the unified one).
*
* <p>Restarting on ANY non-empty pending set would be wrong: when STOP is reached because the
* reminder cap was hit rather than the backlog draining, the same never-collected work is
* expected to still be sitting there — that is the cap doing its job — and restarting would
* nudge about it forever, defeating the bound. Diffing the current pending sets against the
* "before" snapshots tells the two cases apart: an item present before this tick's decision is
* stale backlog, not a race; only an item absent from the "before" snapshot 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.
* genuine thread race.
*
* <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.
* <p>Terminates rather than spinning: this method restarts the schedule at most once per call,
* and a fresh {@link #onReplyQueued} / {@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 work in
* {@link #pendingReplies} / {@link #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) {
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
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);
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead, 0, 0);
return;
}
log.debug("push: ticket reminder loop ended for lead {}", lead);
log.debug("push: 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);
/** Send one combined nudge covering everything currently pending for {@code lead}. */
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
// collected (or another arrive), between the decision and the injection.
Set<String> replyTargets = pendingReplyTargetsFor(lead);
List<PendingTicket> tickets = pendingTicketsFor(lead);
if (replyTargets.isEmpty() && tickets.isEmpty()) {
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
return;
}
String nudge = formatTicketsNudge(pending);
String nudge = formatNudge(replyTargets, tickets);
try {
agents.send(lead, nudge);
log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)",
reminderCount + 1, maxReminders, lead, pending.size());
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
replyTargets.size(), tickets.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());
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 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);
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
backoffMs, TimeUnit.MILLISECONDS);
}
/** Render one or several pending tickets as a single nudge line. */
// --- nudge formatting ------------------------------------------------------------------------
/** Render everything pending for one lead as a single nudge line. */
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets) {
List<String> parts = new ArrayList<>();
if (!replyTargets.isEmpty()) {
parts.add(formatRepliesNudge(replyTargets));
}
if (!tickets.isEmpty()) {
parts.add(formatTicketsNudge(tickets));
}
return String.join(" | ", parts);
}
/** Render one or several pending reply targets. */
private static String formatRepliesNudge(Set<String> targets) {
if (targets.size() == 1) {
String target = targets.iterator().next();
return NUDGE_FORMAT.formatted(target, target);
}
String ids = String.join(", ", targets);
return REPLIES_NUDGE_FORMAT.formatted(targets.size(), ids);
}
/** Render one or several pending tickets. */
private static String formatTicketsNudge(List<PendingTicket> pending) {
if (pending.size() == 1) {
PendingTicket t = pending.get(0);
@@ -383,23 +393,22 @@ public final class ReplyPushLoop {
// --- 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} 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.
* Whether any reminder loop is currently active for some lead (CB-551). The idle-lead heartbeat
* uses this to stand aside: while the push loop is actively nudging a 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 #activeLeads},
* which now covers both reply-queued (CB-307) and ticket-terminal (CB-588) work (CB-590) —
* bounded by what has been triggered, not by any persistent state.
*/
public boolean isActive() {
return !activeTargets.isEmpty() || !activeLeads.isEmpty();
return !activeLeads.isEmpty();
}
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
activeLeads.clear();
pendingReplies.clear();
pendingTickets.clear();
}
@@ -20,6 +20,7 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.*;
@@ -27,6 +28,12 @@ import static org.junit.jupiter.api.Assertions.*;
* Unit tests for {@link ReplyPushLoop}: decision logic, nudge injection, idempotency,
* bounded reminders, and stop conditions.
*
* <p>CB-590 collapsed the CB-307 reply-nudge schedule and the CB-588 ticket-nudge schedule into
* one schedule per lead ({@link ReplyPushLoop#decide}), so most tests below register pending work
* through the public entry points ({@code onReplyQueued} / {@code onTicketTerminal}) before
* exercising {@code decide} directly, mirroring how the two entry points now share one decision
* function keyed by the lead terminal rather than by worker target.
*
* <p>Uses a {@link RecordingHerdrClient} that synchronizes access to its call list so the
* scheduler thread and test thread never have memory ordering issues. The {@code decide()}
* tests use a simple client with no concurrency concern.
@@ -58,38 +65,50 @@ class ReplyPushLoopTest {
// --- decide() logic ------------------------------------------------------------------------
@Test
void decideWithoutPrimaryIsStop() {
agents = agentWithStatus("idle");
var loop = new ReplyPushLoop(
new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0));
void onReplyQueuedWithNoKnownLeadNeverStartsASchedule() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
loop.onReplyQueued(WORKER); // no lead known -> never registered, never scheduled
assertFalse(loop.isActive(), "no lead known means nothing to nudge yet");
Thread.sleep(150);
assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting");
}
@Test
void decideWithEmptyInboxIsStop() {
void decideWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
void decideAtCapIsStop() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
var loop = loop(2, 100);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
}
@Test
void decideUnderCapWithInjectablePrimaryIsInject() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideUnderCapWithBlockedPrimaryIsInject() {
agents = agentWithStatus("blocked");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"BLOCKED is injectable");
}
@@ -97,7 +116,9 @@ class ReplyPushLoopTest {
void decideUnderCapWithDonePrimaryIsInject() {
agents = agentWithStatus("done");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"DONE is injectable");
}
@@ -105,23 +126,29 @@ class ReplyPushLoopTest {
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
agents = agentWithStatus("working");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
agents = agentWithStatus("unknown");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideStopsAfterInboxIsEmptied() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
inbox.ack(WORKER, "m1");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@@ -201,12 +228,19 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
}
// --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
@Test
void repliesNudgeFormatIsCorrect() {
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
assertTrue(multi.contains("2 workers"));
assertTrue(multi.contains("bridge_poll(target=...)"));
}
// --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
@Test
void decideTicketsWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
@@ -214,7 +248,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(2, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
}
@Test
@@ -222,7 +256,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -230,7 +264,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("working");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -238,7 +272,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false); // no lead known -> never registered as pending
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
@@ -255,7 +289,7 @@ class ReplyPushLoopTest {
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket nudge must not tell the lead to run the target-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
}
@Test
@@ -346,11 +380,11 @@ class ReplyPushLoopTest {
void aTicketStillPendingWhenTheLoopStopsIsNotStranded() {
// Regression for the race a reviewer found in gitea PR #73: onTicketTerminal's
// activeLeads.putIfAbsent can see the lead's slot as still occupied a moment before
// decideTickets' STOP releases it, so the ticket coalesces onto a schedule that is about to
// decide's STOP releases it, so the ticket coalesces onto a schedule that is about to
// die and nothing ever nudges about it. Forcing that exact thread interleaving is not
// reliable, so this drives stopOrRestartTicketLoop — the STOP path's own release-and-recheck —
// reliable, so this drives stopOrRestart — the STOP path's own release-and-recheck —
// directly, arranging the state it must not lose a ticket in: a ticket pending for the lead
// that was NOT part of the pre-decision snapshot (pendingBefore=empty), standing in for one
// that was NOT part of the pre-decision snapshot (ticketsBefore=empty), standing in for one
// that races in during the decision-to-release window.
var rec = recordingClient();
agents = new AgentControl(rec);
@@ -358,9 +392,9 @@ class ReplyPushLoopTest {
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
// Stand in for the scheduler thread reaching decideTickets==STOP for this lead — with nothing
// Stand in for the scheduler thread reaching decide==STOP for this lead — with nothing
// pending at decide time — while task-1 races in before the release below runs.
loop.stopOrRestartTicketLoop(PRIMARY, Set.of());
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
assertTrue(loop.isActive(), "a ticket that raced the loop's stop must reclaim the schedule "
+ "slot, not be stranded with no schedule left to ever nudge about it");
@@ -368,23 +402,67 @@ class ReplyPushLoopTest {
@Test
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
// The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) would be
// The other direction of the same fix: restarting on ANY non-empty pending set would be
// wrong. When STOP is reached because the reminder cap was hit, the same never-collected
// ticket is expected to still be there — that is the cap doing its job (acceptance criterion
// #7: nudges stay bounded). task-1 here was already accounted for at decide time (it is in
// pendingBefore), so it must not restart the loop just because it is still sitting there.
// #5: no spin / nudges stay bounded). task-1 here was already accounted for at decide time
// (it is in ticketsBefore), so it must not restart the loop just because it is still sitting
// there.
var rec = recordingClient();
agents = new AgentControl(rec);
ReplyPushLoop loop = loop(1, 100_000);
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
loop.stopOrRestartTicketLoop(PRIMARY, Set.of("task-1"));
loop.stopOrRestart(PRIMARY, Set.of(), Set.of("task-1"));
assertFalse(loop.isActive(), "a stale ticket already accounted for at decide time must not "
+ "restart the loop — that would defeat the reminder cap");
}
// --- mirror of the two stopOrRestart tests above, for the reply arm ------------------------
//
// Both tests above only ever passed Set.of() for repliesBefore, so racedIn's reply branch
// (`pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))`) was
// never exercised by anything other than an always-empty snapshot. The reviewer flagged this:
// racedIn is symmetric in the code, and only half of it was pinned by a test.
@Test
void aReplyStillPendingWhenTheLoopStopsIsNotStranded() {
// Mirrors aTicketStillPendingWhenTheLoopStopsIsNotStranded: a reply target that raced in
// during the decision-to-release window (absent from the "before" snapshot) must reclaim
// the schedule slot rather than being stranded with no schedule left to nudge about it.
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
ReplyPushLoop loop = loop(1, 100_000); // long backoff — no natural tick fires during this test
loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY}
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
assertTrue(loop.isActive(), "a reply that raced the loop's stop must reclaim the schedule "
+ "slot, not be stranded with no schedule left to ever nudge about it");
}
@Test
void aStaleUncollectedReplyAtCapDoesNotRestartTheLoop() {
// Mirrors aStaleUncollectedTicketAtCapDoesNotRestartTheLoop: a reply target already
// accounted for at decide time (present in repliesBefore) must not restart the loop —
// that is the reminder cap doing its job, not a race.
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
ReplyPushLoop loop = loop(1, 100_000);
loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY}
loop.stopOrRestart(PRIMARY, Set.of(WORKER), Set.of());
assertFalse(loop.isActive(), "a stale reply target already accounted for at decide time "
+ "must not restart the loop — that would defeat the reminder cap");
}
@Test
void ticketNudgeFormatIsCorrect() {
String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1");
@@ -402,7 +480,103 @@ class ReplyPushLoopTest {
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
"CB-588 must not change the CB-307 inbox nudge's call shape");
"CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
}
// --- CB-590: one schedule per lead — no overlap, no lost nudges -----------------------------
@Test
void replyAndTicketForTheSameLeadCoalesceIntoOneSendNeverOverlapping() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
Thread.sleep(300);
assertEquals(1, rec.sendCount(),
"a reply and a ticket for the same lead must coalesce onto ONE schedule — "
+ "two nudge injections into the same lead pane must never overlap");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
"the combined nudge must still mention the reply: " + nudge);
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
}
@Test
void aReplyQueuedWhileTheLeadIsBusyIsNotLostWhenATicketArrivesToo() throws Exception {
// The lead is busy for its first two status checks, then becomes injectable. A reply is
// queued while busy; a ticket for the same lead arrives before the lead frees up. Neither
// may be dropped — deferred is fine, lost is not (acceptance criterion #2).
var rec = new BusyThenIdleHerdrClient(2);
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 50); // cap=1: WAIT_BUSY doesn't count against it, so exactly one send once injectable
loop.onReplyQueued(WORKER); // schedule starts, first tick(s) WAIT_BUSY
loop.onTicketTerminal("task-1", WORKER, false); // coalesces onto the same waiting schedule
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"once the lead becomes injectable, the deferred work must still be nudged");
Thread.sleep(200);
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
}
// --- CB-590 follow-up: per-source reminder budgets — the regression this round exists for ---
@Test
void oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource() {
// The live trace this ticket was filed from: an undrained reply target got nudged up to
// its cap (5 reminders), then a ticket for the SAME lead went terminal shortly before the
// next scheduled tick — so it coalesced onto the still-active schedule (arriving BEFORE
// that tick's "before" snapshot, not during the decision-to-release race stopOrRestart
// guards). With CB-590's single shared reminder counter, that tick's decide() saw
// reminderCount already at the cap and returned STOP regardless of the ticket, and because
// the ticket was already present in that tick's "before" snapshot, stopOrRestart's
// racedIn check (proven correct on its own above) did not save it either — it is not a
// race, it looks like ordinary stale backlog. The ticket was then stranded: pending
// forever with no live schedule, never named in any nudge.
//
// Fixed by giving each source its own counter. Here the reply source is AT its cap (2/2)
// and the ticket source has NEVER been nudged (0/2) — decide() must still return INJECT,
// because the ticket is still eligible on its own budget.
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100_000); // huge backoff — this test drives decide()/isActive() directly
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 2, 0),
"the reply source is exhausted (2/2), but the ticket source has never been "
+ "nudged (0/2) — the lead must still be injected so the ticket is not "
+ "lost, exactly the CB-590 follow-up regression");
// isActive() (criterion #5): a real tick that takes the INJECT branch above never calls
// stopOrRestart, so the schedule started by onTicketTerminal above stays live — the
// ticket is not left stranded with isActive()==false while it is still pending.
assertTrue(loop.isActive(), "the schedule must stay active while the ticket source still "
+ "has budget left, even though the reply source sharing it is exhausted");
}
@Test
void bothSourcesExhaustedIsStillStop() {
// The flip side: per-source budgets must not turn into unbounded nudging. When BOTH
// sources are at their cap, decide() must still STOP — a per-source budget is still a
// budget.
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100_000);
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 2),
"both the reply and the ticket source are at their own cap — must still stop");
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@@ -430,8 +604,10 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
var loop = loop(2, 100, metrics);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
@@ -459,7 +635,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100_000, metrics);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the ticket reminder cap must count as exhausted");
@@ -547,4 +723,49 @@ class ReplyPushLoopTest {
private static RecordingHerdrClient recordingClient() {
return new RecordingHerdrClient();
}
/**
* Thread-safe fake that reports {@code working} (not injectable) for its first
* {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
* while the lead is busy is deferred, not dropped, once it becomes injectable.
*/
private static final class BusyThenIdleHerdrClient implements HerdrClient {
private final List<Map.Entry<String, Object>> calls =
Collections.synchronizedList(new ArrayList<>());
private final AtomicInteger statusChecks = new AtomicInteger();
private final int busyChecks;
volatile CountDownLatch sendLatch = new CountDownLatch(1);
BusyThenIdleHerdrClient(int busyChecks) {
this.busyChecks = busyChecks;
}
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
String status = statusChecks.getAndIncrement() < busyChecks ? "working" : "idle";
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", PRIMARY)
.put("agent_status", status));
}
if ("agent.prompt".equals(method)) {
calls.add(Map.entry(method, params));
sendLatch.countDown();
}
return MAPPER.createObjectNode();
}
long sendCount() {
return calls.size();
}
List<Map.Entry<String, Object>> sentParams() {
return List.copyOf(calls);
}
@Override
public void close() {
}
}
}