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..fe1a0b9 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. */ @@ -191,7 +200,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) { @@ -206,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. */ @@ -571,8 +589,24 @@ 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 + // worker paused in bridge_ask leaves it running, per finishAsyncTask's own contract — so + // 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); + }); + } asyncExecutor.submit(() -> { try { Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task); @@ -610,6 +644,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); @@ -640,10 +679,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/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index 890a0c0..77aefb5 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,12 @@ import dev.ltms.bridged.metrics.Metrics; 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; +import java.util.stream.Collectors; /** * Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own @@ -23,11 +26,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 +58,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 +196,211 @@ 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) { + Set pendingBefore = pendingIdsFor(lead); + var action = decideTickets(lead, reminderCount); + switch (action) { + case INJECT -> { + injectTicketNudge(lead, reminderCount); + scheduleTicketTick(lead, reminderCount + 1); + } + case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount); + case STOP -> stopOrRestartTicketLoop(lead, pendingBefore); + } + } + + /** 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 + // 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..07a0a33 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,210 @@ 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) { + 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, null, nowNanos); + 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); + assertFalse(nudge.toUpperCase().contains("FAILED"), + "a successfully-replied ticket's nudge must not say it failed: " + 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()); + } + + /** + * 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; + 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..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; @@ -200,6 +201,210 @@ 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 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"); + 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 +438,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() {