From 3d10ed385c9c17a4ca2c5a86b0aeff528ad9e59d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 16:10:22 +0200 Subject: [PATCH] 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() {