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