From 3d10ed385c9c17a4ca2c5a86b0aeff528ad9e59d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 16:10:22 +0200 Subject: [PATCH 1/3] CB-588: nudge the lead when an async ticket reaches a terminal phase MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An async bridge_send(wait:false) registers a rendezvous waiter, so its reply always takes MessageService.reply's fast path and returns before ReplyPushLoop.onReplyQueued is ever called — the exact mode the charter tells leads to prefer never nudged. Add ReplyPushLoop.onTicketTerminal(ticket, target, failed), a second entry point reusing the loop's status gating, bounded/backoff reminders and metrics, keyed by the nudge-receiving lead so several tickets finishing together coalesce into one nudge naming the count. MessageService wires it via task.future.whenComplete in sendAsync (covers reply, completion fallback, wedge, and abandon() alike) and calls the new ticketCollected(ticket) from poll() once a terminal view is handed back, so an already-collected ticket is never nudged again. The CB-307 inbox path (onReplyQueued/decide/NUDGE_FORMAT) is untouched. --- .../dev/ltms/bridged/msg/MessageService.java | 21 +- .../dev/ltms/bridged/msg/ReplyPushLoop.java | 170 +++++++++++++++- .../ltms/bridged/msg/MessageServiceTest.java | 141 +++++++++++++ .../ltms/bridged/msg/ReplyPushLoopTest.java | 188 ++++++++++++++++++ 4 files changed, 516 insertions(+), 4 deletions(-) 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 c35076d..41d9170 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -191,7 +191,9 @@ public final class MessageService { * Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}. * * @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply - * branch ({@link #reply}) so it can nudge the primary to drain the inbox + * branch ({@link #reply}) so it can nudge the primary to drain the inbox, and + * (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a + * terminal phase, and whenever {@link #poll} hands a terminal ticket to its caller */ public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, ReplyPushLoop pushLoop) { @@ -573,6 +575,18 @@ public final class MessageService { String ticket = "task-" + ticketSeq.incrementAndGet(); Task task = new Task(target); tasks.put(ticket, task); + if (pushLoop != null) { + // CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a + // worker paused in bridge_ask leaves it running, per finishAsyncTask's own contract — so + // this fires exactly once, from whichever path completes it: 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. + task.future.whenComplete((reply, ex) -> { + boolean failed = ex != null || reply == null || !reply.completed(); + pushLoop.onTicketTerminal(ticket, target, failed); + }); + } asyncExecutor.submit(() -> { try { Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task); @@ -610,6 +624,11 @@ public final class MessageService { } return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null); } + // CB-588: the ticket is terminal and being handed to the caller right here — tell the push + // loop it is collected so a later tick's nudge never names a ticket the lead already has. + if (pushLoop != null) { + pushLoop.ticketCollected(ticket); + } Reply r; try { r = f.getNow(null); 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 890a0c0..6a6e103 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -8,9 +8,11 @@ import dev.ltms.bridged.metrics.Metrics; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.List; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; /** * Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own @@ -23,11 +25,27 @@ import java.util.concurrent.TimeUnit; * *

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. */ public final class ReplyPushLoop { private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class); static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it"; + /** CB-588: singular form, one uncollected ticket. */ + static final String TICKET_NUDGE_FORMAT = + "Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it"; + /** CB-588: coalesced form, several uncollected tickets for the same lead. */ + static final String TICKETS_NUDGE_FORMAT = + "%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s"; private final PrimaryRegistry primaryRegistry; private final AgentControl agents; @@ -39,6 +57,10 @@ public final class ReplyPushLoop { /** 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. */ + private final ConcurrentHashMap pendingTickets = new ConcurrentHashMap<>(); + /** CB-588: leads with an active ticket-reminder schedule. */ + private final ConcurrentHashMap activeLeads = new ConcurrentHashMap<>(); public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, ScheduledExecutorService scheduler, @@ -173,23 +195,165 @@ public final class ReplyPushLoop { scheduler.schedule(() -> tick(target, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS); } + // --- CB-588: async ticket terminal nudges --------------------------------------------------- + + /** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */ + private record PendingTicket(String ticket, String lead, boolean failed) { + } + + /** + * Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a + * terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the + * durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a + * fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns + * before {@link #onReplyQueued} is ever called, so without this entry point a ticket finishing + * that way never nudged anyone (CB-588 / gitea #72). + * + *

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) { + return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList(); + } + + /** + * Pure decision function for ticket nudges, mirroring {@link #decide(String, int)} but keyed by + * the nudge-receiving lead terminal rather than the worker session — several tickets from + * different workers delegated by the same lead coalesce onto it. + */ + Action decideTickets(String lead, int reminderCount) { + if (pendingFor(lead).isEmpty()) { + log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead); + return Action.STOP; + } + if (reminderCount >= maxReminders) { + log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead); + countNudge("exhausted"); + return Action.STOP; + } + AgentStatus status; + try { + status = agents.status(lead); + } catch (RuntimeException e) { + log.debug("push: status check failed for lead {}, will retry", lead, e); + return Action.WAIT_BUSY; + } + if (status.injectable()) { + return Action.INJECT; + } + log.debug("push: lead {} is {} (not injectable), waiting", lead, status); + return Action.WAIT_BUSY; + } + + /** Execute one ticket-loop tick — called on the scheduler thread. */ + private void ticketTick(String lead, int reminderCount) { + var action = decideTickets(lead, reminderCount); + switch (action) { + case INJECT -> { + injectTicketNudge(lead, reminderCount); + 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); + } + } + } + + /** 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); + return; + } + String nudge = formatTicketsNudge(pending); + try { + agents.send(lead, nudge); + log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)", + reminderCount + 1, maxReminders, lead, pending.size()); + countNudge("delivered"); + } catch (RuntimeException e) { + log.warn("push: failed to nudge lead {} for {} ticket(s) (reminder {}/{}): {}", + lead, pending.size(), reminderCount + 1, maxReminders, e.toString()); + } + } + + /** Schedule the next ticket-loop tick on the scheduler thread pool. */ + private void scheduleTicketTick(String lead, int nextReminderCount) { + scheduler.schedule(() -> ticketTick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS); + } + + /** Render one or several pending tickets as a single nudge line. */ + private static String formatTicketsNudge(List pending) { + if (pending.size() == 1) { + PendingTicket t = pending.get(0); + return TICKET_NUDGE_FORMAT.formatted(t.ticket(), t.failed() ? " (FAILED)" : "", t.ticket()); + } + long failedCount = pending.stream().filter(PendingTicket::failed).count(); + String ids = pending.stream() + .map(t -> t.failed() ? t.ticket() + " (FAILED)" : t.ticket()) + .collect(Collectors.joining(", ")); + String failedNote = failedCount > 0 ? " (%d failed)".formatted(failedCount) : ""; + return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids); + } + // --- lifecycle ----------------------------------------------------------------------------- /** * Whether any reminder loop is currently active for some target (CB-551). The idle-lead heartbeat * uses this to stand aside: while the push loop is actively nudging the lead, a concurrent * heartbeat injection would start a second competing turn in the same pane — racing loops multiply - * turns and context burn. "Active" means a schedule exists in {@link #activeTargets}; the set is - * bounded by what has been triggered, not by any persistent state. + * turns and context burn. "Active" means a schedule exists in {@link #activeTargets} or + * {@link #activeLeads} (CB-588 ticket nudges are a second source of pane injections the heartbeat + * must equally stand aside for); the sets are bounded by what has been triggered, not by any + * persistent state. */ public boolean isActive() { - return !activeTargets.isEmpty(); + return !activeTargets.isEmpty() || !activeLeads.isEmpty(); } /** Shut down the scheduler. Outstanding reminders are cancelled. */ public void stop() { scheduler.shutdownNow(); activeTargets.clear(); + activeLeads.clear(); + pendingTickets.clear(); } /** @see #stop() */ 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 0f4fe23..b257e50 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -714,6 +714,147 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); } + // --- CB-588: async ticket terminal nudges --------------------------------------------------- + // + // MessageService.reply's rendezvous fast path is exactly what an async ticket always takes + // (sendAsync registers a rendezvous waiter — see asyncTasksByWaiter), so it never reached + // ReplyPushLoop.onReplyQueued. These prove the ticket reaches ReplyPushLoop through the new + // onTicketTerminal entry point instead, with no bridge_poll from the lead first. + + private static final String LEAD = "term_lead"; + + /** A MessageService wired to a real ReplyPushLoop pointed at a private herdr fake for the lead. */ + private record PushWiring(MessageService service, FakeHerdr leadHerdr, + java.util.concurrent.ScheduledExecutorService scheduler) implements AutoCloseable { + @Override + public void close() { + scheduler.shutdownNow(); + } + } + + private PushWiring wireWithPushLoop(int maxReminders, long backoffMs) { + PrimaryRegistry registry = new PrimaryRegistry(null); + registry.recordDelegation(T, LEAD); + FakeHerdr leadHerdr = new FakeHerdr(); + AgentControl leadAgents = new AgentControl(leadHerdr); + var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor(); + ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs); + MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop); + return new PushWiring(service, leadHerdr, scheduler); + } + + private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException { + long deadline = System.currentTimeMillis() + 3000; + while (!leadHerdr.called("agent.prompt") && System.currentTimeMillis() < deadline) { + Thread.sleep(10); + } + assertTrue(leadHerdr.called("agent.prompt"), "expected a nudge in the lead's pane"); + } + + @Test + void anAsyncTicketThatFinishesNudgesTheLeadWithNoPriorPollCall() throws Exception { + try (var wiring = wireWithPushLoop(1, 50)) { + String ticket = wiring.service().sendAsync(T, "long task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "async result")); + + awaitNudge(wiring.leadHerdr()); + String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); + 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); + } + } + + @Test + void aFailedAsyncTicketAlsoNudgesAndSaysSo() throws Exception { + try (var wiring = wireWithPushLoop(1, 50)) { + wiring.service().sendAsync(T, "long task"); + awaitWaiting(); + wiring.service().abandon(T, "session released"); // a terminal failure phase + + awaitNudge(wiring.leadHerdr()); + String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); + assertTrue(nudge.toUpperCase().contains("FAILED"), "a failed ticket's nudge must say so: " + nudge); + } + } + + @Test + void anAlreadyCollectedTicketProducesNoNudge() throws Exception { + try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: poll before the first tick fires + String ticket = wiring.service().sendAsync(T, "long task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "async result")); + + MessageService.TaskView view = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE); + assertEquals(MessageService.Phase.DONE, view.phase()); + + Thread.sleep(400); // let the scheduled tick run — it must find nothing pending + assertFalse(wiring.leadHerdr().called("agent.prompt"), + "a ticket the lead already polled must never be nudged"); + } + } + + @Test + void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception { + try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires + String first = wiring.service().sendAsync(T, "first task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "first done")); + // Settle without polling: poll() itself marks a ticket collected (that's the point of + // anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion + // would collect the ticket before the coalescing this test checks ever gets a chance. + Thread.sleep(100); + + String second = wiring.service().sendAsync(T, "second task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "second done")); + Thread.sleep(100); + + awaitNudge(wiring.leadHerdr()); + Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge + long nudgeCount = wiring.leadHerdr().calls.stream() + .filter(c -> c.method().equals("agent.prompt")).count(); + assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two"); + String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); + assertTrue(nudge.contains(first) && nudge.contains(second), + "the coalesced nudge should name both tickets: " + nudge); + } + } + + @Test + void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception { + // `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw, + // and no nudge mechanism exists to fire regardless of how the ticket resolves. + String ticket = messages.sendAsync(T, "long task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "async result")); + + MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.DONE); + assertEquals(MessageService.Phase.DONE, view.phase()); + assertEquals("async result", view.reply()); + } + + private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket, + MessageService.Phase phase) throws Exception { + long deadline = System.currentTimeMillis() + 3000; + MessageService.TaskView view; + do { + view = svc.poll(ticket); + if (view.phase() == phase) { + return view; + } + Thread.sleep(5); + } while (System.currentTimeMillis() < deadline); + assertEquals(phase, view.phase()); + return view; + } + private void assertFailedTicket(String ticket, String reason) throws Exception { MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.FAILED); assertEquals(reason, view.detail()); 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 5051450..ca7da42 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -200,6 +200,167 @@ class ReplyPushLoopTest { assertTrue(nudge.contains("bridge_poll(target=term_worker)")); } + // --- CB-588: async ticket terminal nudges — decideTickets() logic -------------------------- + + @Test + void decideTicketsWithNothingPendingIsStop() { + agents = agentWithStatus("idle"); + assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0)); + } + + @Test + void decideTicketsAtCapIsStop() { + agents = agentWithStatus("idle"); + var loop = loop(2, 100_000); + loop.onTicketTerminal("task-1", WORKER, false); + assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2)); + } + + @Test + void decideTicketsUnderCapWithInjectableLeadIsInject() { + agents = agentWithStatus("idle"); + var loop = loop(5, 100_000); + loop.onTicketTerminal("task-1", WORKER, false); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0)); + } + + @Test + void decideTicketsUnderCapWithBusyLeadIsWaitBusy() { + agents = agentWithStatus("working"); + var loop = loop(5, 100_000); + loop.onTicketTerminal("task-1", WORKER, false); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0)); + } + + @Test + void decideTicketsWithoutAKnownLeadIsStop() { + 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)); + } + + // --- CB-588: onTicketTerminal integration --------------------------------------------------- + + @Test + void onTicketTerminalCausesExactlyOneNudgeNamingTheTicketAndThePollCall() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + + loop(1, 50).onTicketTerminal("task-1", WORKER, false); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent"); + assertEquals(1, rec.sendCount()); + 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"); + } + + @Test + void aFailedTicketNudgeSaysItWasAFailure() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + + loop(1, 50).onTicketTerminal("task-1", WORKER, true); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS)); + String nudge = rec.sentParams().getFirst().getValue().toString(); + assertTrue(nudge.toUpperCase().contains("FAILED"), "a failed ticket's nudge must say so: " + nudge); + } + + @Test + void severalTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + var loop = loop(1, 300); // backoff wide enough that both onTicketTerminal calls land first + + loop.onTicketTerminal("task-1", WORKER, false); + loop.onTicketTerminal("task-2", WORKER, true); // arrives while the schedule is already active + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one coalesced nudge should have been sent"); + assertEquals(1, rec.sendCount(), "two tickets finishing together must produce ONE nudge, not two"); + String nudge = rec.sentParams().getFirst().getValue().toString(); + assertTrue(nudge.contains("task-1") && nudge.contains("task-2"), + "the coalesced nudge should name both tickets: " + nudge); + assertTrue(nudge.contains("2 tickets"), "the coalesced nudge should name the count: " + nudge); + } + + @Test + void onTicketTerminalIsIdempotentPerLeadWhileActive() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + rec.sendLatch = new CountDownLatch(1); + + var loop = loop(1, 100); + loop.onTicketTerminal("task-1", WORKER, false); + loop.onTicketTerminal("task-1", WORKER, false); // duplicate — should not start a second schedule + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS)); + Thread.sleep(200); + assertEquals(1, rec.sendCount(), "a duplicate onTicketTerminal for the same lead must not double-nudge"); + } + + @Test + void aTicketAlreadyCollectedProducesNoNudge() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + var loop = loop(1, 100); + + loop.onTicketTerminal("task-1", WORKER, false); + loop.ticketCollected("task-1"); // the lead polled before the first tick fired + + Thread.sleep(300); // let the scheduled tick run + assertEquals(0, rec.sendCount(), "an already-collected ticket must never be nudged"); + } + + @Test + void ticketNudgesSendUpToCapThenStop() throws Exception { + int cap = 2; + var rec = recordingClient(); + agents = new AgentControl(rec); + rec.sendLatch = new CountDownLatch(cap); + + loop(cap, 50).onTicketTerminal("task-1", WORKER, false); + + assertTrue(rec.sendLatch.await(5, TimeUnit.SECONDS), cap + " ticket nudges should have fired"); + Thread.sleep(300); + assertEquals(cap, rec.sendCount(), "exactly " + cap + " ticket nudges (cap=" + cap + ")"); + } + + @Test + void isActiveReflectsALiveTicketScheduleForTheHeartbeatStandDown() { + var rec = recordingClient(); + agents = new AgentControl(rec); + ReplyPushLoop loop = loop(1, 100_000); // long backoff so the tick cannot fire mid-test + + loop.onTicketTerminal("task-1", WORKER, false); + assertTrue(loop.isActive(), "an active ticket-reminder schedule must also stand off the heartbeat"); + + loop.stop(); + assertFalse(loop.isActive(), "stopping clears the active ticket schedule too"); + } + + @Test + void ticketNudgeFormatIsCorrect() { + String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1"); + assertTrue(single.contains("Ticket task-1")); + assertTrue(single.contains("bridge_poll(ticket=task-1)")); + + String multi = ReplyPushLoop.TICKETS_NUDGE_FORMAT.formatted(2, "", "task-1, task-2"); + assertTrue(multi.contains("2 tickets")); + assertTrue(multi.contains("bridge_poll(ticket=...)")); + } + + // --- CB-307 nudge path is unchanged (regression) -------------------------------------------- + + @Test + 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"); + } + // --- metrics (CB-512) ---------------------------------------------------------------------- @Test @@ -233,6 +394,33 @@ class ReplyPushLoopTest { assertEquals(0, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered")); } + @Test + void successfulTicketNudgeIncrementsDelivered() throws Exception { + var rec = recordingClient(); + agents = new AgentControl(rec); + Metrics metrics = new Metrics(); + + loop(1, 50, metrics).onTicketTerminal("task-1", WORKER, false); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent"); + Thread.sleep(200); + assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered"), + "a successfully sent ticket nudge must count as delivered, same metric as CB-307"); + } + + @Test + void ticketReminderCapIncrementsExhausted() { + agents = agentWithStatus("idle"); + Metrics metrics = new Metrics(); + var loop = loop(2, 100_000, metrics); + loop.onTicketTerminal("task-1", WORKER, false); + + assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2)); + + assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"), + "hitting the ticket reminder cap must count as exhausted"); + } + // --- helpers ------------------------------------------------------------------------------- private ReplyPushLoop loop() { From 6ebad2a91f81dae6fcd1d50fdc45e819bae27323 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 16:41:23 +0200 Subject: [PATCH 2/3] CB-588 follow-up: reclaim a pruned ticket's pendingTickets entry too MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit tasks is the sole authority on whether a ticket exists, but pruneTerminalTickets dropped entries from it without telling ReplyPushLoop. ticketCollected(ticket) was only ever called from MessageService.poll's terminal branch, which pruneTerminalTickets short-circuits past once a ticket is gone (poll returns null at the top). A ticket the lead never polled — or one the reminder cap already gave up on — was pruned from tasks but never collected in ReplyPushLoop.pendingTickets, so it rode along on every later nudge to the same lead forever, naming a ticket bridge_poll could no longer find, and the map itself never shrank. pruneTerminalTickets now calls pushLoop.ticketCollected for every ticket it actually removes (guarded on pushLoop != null), so a pending nudge entry lives exactly as long as its ticket is pollable. Reused tasks as the only removal trigger rather than adding a second live query back into MessageService — no new source of truth. Added an injectable clock (LongSupplier nowNanos, defaulting to System::nanoTime) to MessageService, mirroring the SessionManager/ SessionReaper nowNanos seam, so a test can cross the 10-minute TICKET_TTL_NANOS deterministically instead of sleeping for real. Confirmed the new regression test fails against the prior pruneTerminalTickets (a stale ticket rides along on a later coalesced nudge) before restoring the fix. --- .../dev/ltms/bridged/msg/MessageService.java | 49 ++++++++++++--- .../ltms/bridged/msg/MessageServiceTest.java | 63 ++++++++++++++++++- 2 files changed, 103 insertions(+), 9 deletions(-) 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 41d9170..eaa2bb3 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -20,6 +20,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.locks.ReentrantLock; +import java.util.function.LongSupplier; /** * The blocking delegation feature (CB-104): deliver {@code content} into a worker and block until @@ -52,8 +53,12 @@ public final class MessageService { */ private static final long ASYNC_TIMEOUT_MS = 30 * 60 * 1_000L; - /** How long a finished (terminal) ticket is retained for polling before it is pruned. */ - private static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L; + /** + * How long a finished (terminal) ticket is retained for polling before it is pruned. Package- + * private (not {@code private}) so a test can advance an injected clock past it deterministically + * instead of duplicating the magic number or sleeping for real. + */ + static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L; /** Outcome of a blocking send. */ public enum Outcome { @@ -161,12 +166,13 @@ public final class MessageService { private static final class Task { private final String target; private final CompletableFuture future = new CompletableFuture<>(); - private final long createdNanos = System.nanoTime(); + private final long createdNanos; private volatile Reply question; private volatile String turnId; - private Task(String target) { + private Task(String target, long createdNanos) { this.target = target; + this.createdNanos = createdNanos; } } @@ -176,6 +182,9 @@ public final class MessageService { private final ReplyInbox inbox; private final ReplyPushLoop pushLoop; private final Metrics metrics; // CB-502: nullable — no registry in unit tests + // CB-588: injectable so pruneTerminalTickets' 10-minute TICKET_TTL_NANOS can be exercised in a + // test without a real wait — same seam SessionManager already uses for its idle reaper (nowNanos). + private final LongSupplier nowNanos; private final ConcurrentHashMap sessionLocks = new ConcurrentHashMap<>(); private final ConcurrentHashMap tasks = new ConcurrentHashMap<>(); /** Async task that owns each exact forward rendezvous waiter. */ @@ -208,12 +217,19 @@ public final class MessageService { */ public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, ReplyPushLoop pushLoop, Metrics metrics) { + this(agents, injector, rendezvous, inbox, pushLoop, metrics, System::nanoTime); + } + + /** Test constructor with an injectable clock (CB-588: exercise the ticket-prune TTL without a real wait). */ + MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, + ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) { this.agents = agents; this.injector = injector; this.rendezvous = rendezvous; this.inbox = inbox; this.pushLoop = pushLoop; this.metrics = metrics; + this.nowNanos = nowNanos; } /** Create with an explicit {@link ReplyInbox} and no push loop. */ @@ -573,7 +589,7 @@ public final class MessageService { */ public String sendAsync(String target, String content, Runnable onAccepted) { String ticket = "task-" + ticketSeq.incrementAndGet(); - Task task = new Task(target); + Task task = new Task(target, nowNanos.getAsLong()); tasks.put(ticket, task); if (pushLoop != null) { // CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a @@ -659,10 +675,27 @@ public final class MessageService { } } - /** Drop finished tickets older than the TTL so the registry cannot grow without bound. */ + /** + * Drop finished tickets older than the TTL so {@link #tasks} cannot grow without bound. + * + *

{@code tasks} is the sole authority on whether a ticket still exists — {@link #poll} returns + * {@code null} the instant a ticket is gone from here, before it ever reaches the terminal branch + * that calls {@link ReplyPushLoop#ticketCollected}. Without telling the push loop about a prune + * too, its own {@code pendingTickets} entry would outlive the ticket it names: an unpolled ticket + * (or one the reminder cap already gave up on) is pruned here but never collected there, so it + * lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead + * — naming a ticket {@code bridge_poll} can no longer find (CB-588 follow-up). + */ private void pruneTerminalTickets() { - long cutoff = System.nanoTime() - TICKET_TTL_NANOS; - tasks.values().removeIf(t -> t.future.isDone() && t.createdNanos < cutoff); + long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS; + tasks.entrySet().removeIf(e -> { + Task t = e.getValue(); + boolean expired = t.future.isDone() && t.createdNanos < cutoff; + if (expired && pushLoop != null) { + pushLoop.ticketCollected(e.getKey()); + } + return expired; + }); } /** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */ 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 b257e50..67057f5 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -733,13 +733,18 @@ class MessageServiceTest { } private PushWiring wireWithPushLoop(int maxReminders, long backoffMs) { + return wireWithPushLoop(maxReminders, backoffMs, System::nanoTime); + } + + /** As above, with an injectable clock (CB-588 follow-up: exercise pruneTerminalTickets' TTL). */ + private PushWiring wireWithPushLoop(int maxReminders, long backoffMs, java.util.function.LongSupplier nowNanos) { PrimaryRegistry registry = new PrimaryRegistry(null); registry.recordDelegation(T, LEAD); FakeHerdr leadHerdr = new FakeHerdr(); AgentControl leadAgents = new AgentControl(leadHerdr); var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor(); ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs); - MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop); + MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos); return new PushWiring(service, leadHerdr, scheduler); } @@ -840,6 +845,62 @@ class MessageServiceTest { assertEquals("async result", view.reply()); } + /** + * CB-588 follow-up: {@code tasks} is the sole authority on whether a ticket exists, and + * {@code pruneTerminalTickets} drops entries from it once {@link MessageService#TICKET_TTL_NANOS} + * elapses. Before this test, that prune never told {@code ReplyPushLoop} — its own + * {@code pendingTickets} entry for a pruned, never-collected ticket had no remover at all, so it + * rode along on every later nudge to the same lead, naming a ticket {@code bridge_poll} could no + * longer find. Uses the injectable clock (mirroring {@code SessionManager}'s {@code nowNanos} seam + * for its idle reaper) to cross the 10-minute TTL without a real wait. + */ + @Test + void aPrunedTicketIsReclaimedFromThePushLoopNotLeakedForever() throws Exception { + java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L); + // maxReminders=1 + a short backoff: the stale ticket gets its one legitimate reminder, then + // decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix) + // pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of + // the bug report, reached deterministically rather than by timing it against a live tick. + try (var wiring = wireWithPushLoop(1, 50, clock::get)) { + String stale = wiring.service().sendAsync(T, "first task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "stale result")); + + // Let the reminder loop fire its one nudge and hit the cap (STOP removes it from + // activeLeads; pendingTickets is untouched either way — that asymmetry is the bug). + awaitNudge(wiring.leadHerdr()); + Thread.sleep(300); + assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale), + "sanity: the stale ticket's own reminder must have fired first"); + + // Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly. + clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1)); + + // A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes. + String fresh = wiring.service().sendAsync(T, "second task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "fresh result")); + + // The fresh ticket restarts the (now-dormant) reminder loop with its own nudge. + long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count(); + long deadline = System.currentTimeMillis() + 3000; + while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before + && System.currentTimeMillis() < deadline) { + Thread.sleep(10); + } + String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); + assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge); + assertFalse(latestNudge.contains(stale), + "a pruned ticket must never be named in a later nudge — it is gone and bridge_poll " + + "on it would return nothing: " + latestNudge); + + // And bridge_poll(ticket=stale) really does return nothing now — the nudge would have lied. + assertNull(wiring.service().poll(stale), "the pruned ticket must actually be gone, not just unmentioned"); + } + } + private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket, MessageService.Phase phase) throws Exception { long deadline = System.currentTimeMillis() + 3000; From ac044e75730a28a6b221af251129f7de7e8f0b2a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 17:02:05 +0200 Subject: [PATCH 3/3] CB-588 round 3: close the STOP-vs-onTicketTerminal race, pin nudge polarity, fix comment - ReplyPushLoop.stopOrRestartTicketLoop: after releasing a lead's active-schedule slot on STOP, restart only if a ticket landed that the pre-decision snapshot did not already account for. A naive "restart on any pending ticket" version was tried first and reverted: it defeated the reminder cap by restarting forever on a stale, never-collected ticket (broke successfulTicketNudgeIncrementsDelivered and ticketNudgesSendUpToCapThenStop). Diffing against a pendingBefore snapshot distinguishes a genuine race arrival from stale cap-exhausted backlog. - Two new ReplyPushLoopTest cases exercise stopOrRestartTicketLoop directly (now package-private) rather than forcing the underlying thread race: aTicketStillPendingWhenTheLoopStopsIsNotStranded (the race must restart) and aStaleUncollectedTicketAtCapDoesNotRestartTheLoop (the cap must still hold). Both were verified to fail against deliberately-reverted versions of the fix before being restored to green. - MessageServiceTest: pin the success-path nudge test's negative direction too (must not contain "FAILED"), not just the failure-path test. - MessageService: complete the whenComplete comment's list of completion paths (TIMED_OUT/BUSY/BACKEND_EXHAUSTED via finishAsyncTask, completeExceptionally on throw, and answer() -> finishAsyncTask(turnId, result)). --- .../dev/ltms/bridged/msg/MessageService.java | 12 ++-- .../dev/ltms/bridged/msg/ReplyPushLoop.java | 55 +++++++++++++++++-- .../ltms/bridged/msg/MessageServiceTest.java | 2 + .../ltms/bridged/msg/ReplyPushLoopTest.java | 44 +++++++++++++++ 4 files changed, 105 insertions(+), 8 deletions(-) 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");