CB-588: nudge the lead when an async ticket reaches a terminal phase
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.
This commit is contained in:
@@ -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());
|
||||
|
||||
@@ -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() {
|
||||
|
||||
Reference in New Issue
Block a user