diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index eaa2bb3..fe1a0b9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -594,10 +594,14 @@ public final class MessageService { if (pushLoop != null) { // CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a // worker paused in bridge_ask leaves it running, per finishAsyncTask's own contract — so - // this fires exactly once, from whichever path completes it: a worker's bridge_reply, the - // CB-106 completion fallback, a CB-109 wedge, or a CB-516 abandon() on teardown. Without - // this, MessageService.reply's rendezvous fast path (the one an async ticket always takes) - // never told the push loop anything happened — see the class javadoc on sendAsync/CB-107. + // this fires exactly once, from whichever path completes it: finishAsyncTask(task, result) + // below on any non-QUESTION outcome of send() — a worker's bridge_reply, the CB-106 + // completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same + // finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION + // is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516 + // abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the + // one an async ticket always takes) never told the push loop anything happened — see the + // class javadoc on sendAsync/CB-107. task.future.whenComplete((reply, ex) -> { boolean failed = ex != null || reply == null || !reply.completed(); pushLoop.onTicketTerminal(ticket, target, failed); 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 6a6e103..77aefb5 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -9,6 +9,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.List; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -279,6 +280,7 @@ public final class ReplyPushLoop { /** 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 -> { @@ -286,13 +288,58 @@ public final class ReplyPushLoop { scheduleTicketTick(lead, reminderCount + 1); } case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount); - case STOP -> { - activeLeads.remove(lead); - log.debug("push: ticket reminder loop ended for lead {}", lead); - } + case STOP -> stopOrRestartTicketLoop(lead, pendingBefore); } } + /** 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()); + } + + /** + * 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). + * + *

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

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

Terminates rather than spinning: this method restarts the schedule at most once per call, and + * a fresh {@link #onTicketTerminal} racing the recheck below still terminates in one of two ways — + * either it observes the slot already vacated (by the {@code activeLeads.remove} above, which + * happens-before this recheck in program order) and claims it itself, or it lands first and this + * recheck then observes its ticket in {@code pendingTickets} and reclaims the slot instead. Exactly + * one side always wins; neither can miss the other, so this never loops on its own account. + */ + void stopOrRestartTicketLoop(String lead, Set pendingBefore) { + activeLeads.remove(lead); + boolean ticketRacedIn = pendingFor(lead).stream().anyMatch(t -> !pendingBefore.contains(t.ticket())); + if (ticketRacedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) { + log.debug("push: a ticket for lead {} raced the reminder loop's stop — restarting", lead); + scheduleTicketTick(lead, 0); + return; + } + log.debug("push: ticket reminder loop ended for lead {}", lead); + } + /** Send the coalesced ticket nudge and log the event. */ private void injectTicketNudge(String lead, int reminderCount) { // Re-read rather than threading it down from decideTickets(): a ticket can be collected (or diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 67057f5..07a0a33 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -769,6 +769,8 @@ class MessageServiceTest { assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge); assertTrue(nudge.contains("bridge_poll(ticket="), "the nudge should name the exact ticket-collecting call: " + nudge); + assertFalse(nudge.toUpperCase().contains("FAILED"), + "a successfully-replied ticket's nudge must not say it failed: " + nudge); } } 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 ca7da42..3cfdbf5 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -15,6 +15,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; @@ -341,6 +342,49 @@ class ReplyPushLoopTest { assertFalse(loop.isActive(), "stopping clears the active ticket schedule too"); } + @Test + 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 + // 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 — + // 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 races in during the decision-to-release window. + var rec = recordingClient(); + agents = new AgentControl(rec); + ReplyPushLoop loop = loop(1, 100_000); // long backoff — no natural tick fires during this test + + 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 + // pending at decide time — while task-1 races in before the release below runs. + loop.stopOrRestartTicketLoop(PRIMARY, 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"); + } + + @Test + void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() { + // The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) 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. + 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")); + + assertFalse(loop.isActive(), "a stale ticket 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");