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() { + } + } }