From 78ca24dc3feecd2279bc848221ac29737406bd9e Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 16 Aug 2026 16:56:04 +0200 Subject: [PATCH 1/2] CB-590: collapse the CB-307 and CB-588 nudge schedules into one per lead MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both reply-queued and ticket-terminal nudges could independently decide to inject into the same lead pane in the same window, since they ran as two separate schedules keyed differently (worker target vs. lead) that never checked each other. Replace both with a single per-lead schedule (activeLeads) that drains pending reply targets and pending tickets together, sends at most one combined nudge per tick, and shares one reminder cap across both sources — so two injections into the same pane can no longer overlap, and work queued while the lead is busy is never lost, only deferred. --- .../dev/ltms/bridged/msg/ReplyPushLoop.java | 480 +++++++++--------- .../ltms/bridged/msg/ReplyPushLoopTest.java | 195 +++++-- 2 files changed, 394 insertions(+), 281 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index 77aefb5..f368f1e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -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,41 @@ 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). * - *

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. + *

CB-590: one schedule per lead. 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. * - *

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. - * - *

CB-588 ticket nudges. {@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. + *

Each tick examines everything 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)}). 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 the shared reminder cap + * ({@link #maxReminders}) is reached, 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 +64,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 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 pendingReplies = new ConcurrentHashMap<>(); + /** Tickets that have gone terminal but not yet been polled, keyed by ticket. */ private final ConcurrentHashMap 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 activeLeads = new ConcurrentHashMap<>(); public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, @@ -89,178 +97,67 @@ 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 pendingReplyTargetsFor(String lead) { + Set 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). - * - *

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 pendingFor(String lead) { + private List 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 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. + * + * @param lead the lead terminal to nudge + * @param reminderCount how many nudges have been sent so far for this lead (shared across + * both reply and ticket work — CB-590 collapses the reminder cap onto + * one counter per lead, so alternating sources cannot outrun the bound) + * @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 reminderCount) { + boolean anyPending = !pendingReplyTargetsFor(lead).isEmpty() || !pendingTicketIdsFor(lead).isEmpty(); + if (!anyPending) { + 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); + log.debug("push: reminder cap ({}) reached for lead {}, stopping", maxReminders, lead); countNudge("exhausted"); return Action.STOP; } @@ -278,95 +175,185 @@ 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 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 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). * - *

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. + *

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); + } + + /** Execute one loop tick — called on the scheduler thread. */ + private void tick(String lead, int reminderCount) { + Set repliesBefore = pendingReplyTargetsFor(lead); + Set ticketsBefore = pendingTicketIdsFor(lead); + var action = decide(lead, reminderCount); + switch (action) { + case INJECT -> { + injectNudge(lead, reminderCount); + scheduleNext(lead, reminderCount + 1); + } + // Re-check after the configured backoff; the lead may become injectable soon. + case WAIT_BUSY -> scheduleNext(lead, reminderCount); + 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). + * + *

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. * *

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. * - *

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. + *

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 pendingBefore) { + void stopOrRestart(String lead, Set repliesBefore, Set 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); 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 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 reminderCount) { + // 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 replyTargets = pendingReplyTargetsFor(lead); + List 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 target(s), {} ticket(s))", + reminderCount + 1, maxReminders, lead, 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 {} (reminder {}/{}): {}", + lead, 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); + /** Schedule the next tick on the scheduler thread pool. */ + private void scheduleNext(String lead, int nextReminderCount) { + scheduler.schedule(() -> tick(lead, nextReminderCount), 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 replyTargets, List tickets) { + List 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 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 pending) { if (pending.size() == 1) { PendingTicket t = pending.get(0); @@ -383,23 +370,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(); } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java index 3cfdbf5..974245d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -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. * + *

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. + * *

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)); } @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)); } @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)); } @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), "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), "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)); } @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)); } @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)); inbox.ack(WORKER, "m1"); - assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 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)); } @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, 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)); } @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)); } @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)); } // --- 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,18 +402,19 @@ 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"); @@ -402,7 +437,52 @@ 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); } // --- metrics (CB-512) ---------------------------------------------------------------------- @@ -430,8 +510,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)); assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"), "hitting the reminder cap must count as exhausted"); @@ -459,7 +541,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, 2)); assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"), "hitting the ticket reminder cap must count as exhausted"); @@ -547,4 +629,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> 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> sentParams() { + return List.copyOf(calls); + } + + @Override + public void close() { + } + } } -- 2.52.0 From 88b9503c3beab253c394784ca9d63bce3659a430 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 16 Aug 2026 17:26:07 +0200 Subject: [PATCH 2/2] CB-590 follow-up: give each nudge source its own reminder budget MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit decide(lead, reminderCount) shared one counter across the reply and ticket sources after PR #84 collapsed both onto a single per-lead schedule. A reply stream that used up the whole budget could then make decide() STOP even for a ticket that had never been nudged and had coalesced onto the same still-active schedule — stranding it with no live schedule left, since stopOrRestart's racedIn check does not save work that was already present in the "before" snapshot. decide() now tracks a per-source count (replyReminderCount, ticketReminderCount) and returns INJECT while either source is still under its own cap, STOP only when both are exhausted. Still exactly one schedule per lead; stopOrRestart's snapshot-diff logic is untouched. --- .../dev/ltms/bridged/msg/ReplyPushLoop.java | 77 +++++++---- .../ltms/bridged/msg/ReplyPushLoopTest.java | 126 +++++++++++++++--- 2 files changed, 160 insertions(+), 43 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index f368f1e..6e5bcdf 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -37,10 +37,12 @@ import java.util.stream.Collectors; *

Each tick examines everything 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)}). 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 the shared reminder cap - * ({@link #maxReminders}) is reached, whichever the durable inbox / pending-ticket set doesn't - * already answer via {@code STOP}. + * ({@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 { @@ -144,20 +146,32 @@ public final class ReplyPushLoop { * Pure decision function: examine everything pending for {@code lead} — reply targets and * tickets alike — and return what the loop should do. * - * @param lead the lead terminal to nudge - * @param reminderCount how many nudges have been sent so far for this lead (shared across - * both reply and ticket work — CB-590 collapses the reminder cap onto - * one counter per lead, so alternating sources cannot outrun the bound) + *

CB-590-fix: one schedule, two budgets. 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 decide(String lead, int reminderCount) { - boolean anyPending = !pendingReplyTargetsFor(lead).isEmpty() || !pendingTicketIdsFor(lead).isEmpty(); - if (!anyPending) { + 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: 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; } @@ -241,21 +255,28 @@ public final class ReplyPushLoop { return; } log.debug("push: starting reminder loop for lead {}", lead); - scheduleNext(lead, 0); + scheduleNext(lead, 0, 0); } /** Execute one loop tick — called on the scheduler thread. */ - private void tick(String lead, int reminderCount) { + private void tick(String lead, int replyReminderCount, int ticketReminderCount) { Set repliesBefore = pendingReplyTargetsFor(lead); Set ticketsBefore = pendingTicketIdsFor(lead); - var action = decide(lead, reminderCount); + var action = decide(lead, replyReminderCount, ticketReminderCount); switch (action) { case INJECT -> { - injectNudge(lead, reminderCount); - scheduleNext(lead, reminderCount + 1); + 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, reminderCount); + case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount); case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore); } } @@ -296,14 +317,14 @@ public final class ReplyPushLoop { || 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); + scheduleNext(lead, 0, 0); return; } log.debug("push: reminder loop ended for lead {}", lead); } /** Send one combined nudge covering everything currently pending for {@code lead}. */ - private void injectNudge(String lead, int reminderCount) { + 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 replyTargets = pendingReplyTargetsFor(lead); @@ -315,18 +336,20 @@ public final class ReplyPushLoop { String nudge = formatNudge(replyTargets, tickets); try { agents.send(lead, nudge); - log.debug("push: nudge {}/{} sent to lead {} ({} reply target(s), {} ticket(s))", - reminderCount + 1, maxReminders, lead, replyTargets.size(), tickets.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 {} (reminder {}/{}): {}", - lead, 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 tick on the scheduler thread pool. */ - private void scheduleNext(String lead, int nextReminderCount) { - scheduler.schedule(() -> tick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS); + private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) { + scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount), + backoffMs, TimeUnit.MILLISECONDS); } // --- nudge formatting ------------------------------------------------------------------------ diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java index 974245d..ad58a15 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -81,7 +81,7 @@ class ReplyPushLoopTest { @Test void decideWithNothingPendingIsStop() { agents = agentWithStatus("idle"); - assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0)); } @Test @@ -90,7 +90,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(2, 100); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0)); } @Test @@ -99,7 +99,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0)); } @Test @@ -108,7 +108,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0), + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0), "BLOCKED is injectable"); } @@ -118,7 +118,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0), + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0), "DONE is injectable"); } @@ -128,7 +128,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0)); } @Test @@ -137,7 +137,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0)); } @Test @@ -146,9 +146,9 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0)); inbox.ack(WORKER, "m1"); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0)); } // --- onReplyQueued integration ------------------------------------------------------------- @@ -240,7 +240,7 @@ class ReplyPushLoopTest { @Test void decideTicketsWithNothingPendingIsStop() { agents = agentWithStatus("idle"); - assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0)); } @Test @@ -248,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.decide(PRIMARY, 2)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2)); } @Test @@ -256,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.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0)); } @Test @@ -264,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.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0)); } @Test @@ -272,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.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0)); } // --- CB-588: onTicketTerminal integration --------------------------------------------------- @@ -420,6 +420,49 @@ class ReplyPushLoopTest { + "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"); @@ -485,6 +528,57 @@ class ReplyPushLoopTest { 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) ---------------------------------------------------------------------- @Test @@ -513,7 +607,7 @@ class ReplyPushLoopTest { var loop = loop(2, 100, metrics); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 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"); @@ -541,7 +635,7 @@ class ReplyPushLoopTest { var loop = loop(2, 100_000, metrics); loop.onTicketTerminal("task-1", WORKER, false); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(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"); -- 2.52.0