CB-590: collapse the CB-307 and CB-588 nudge schedules into one per lead
Both reply-queued and ticket-terminal nudges could independently decide to inject into the same lead pane in the same window, since they ran as two separate schedules keyed differently (worker target vs. lead) that never checked each other. Replace both with a single per-lead schedule (activeLeads) that drains pending reply targets and pending tickets together, sends at most one combined nudge per tick, and shares one reminder cap across both sources — so two injections into the same pane can no longer overlap, and work queued while the lead is busy is never lost, only deferred.
This commit is contained in:
@@ -20,6 +20,7 @@ import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -27,6 +28,12 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
* Unit tests for {@link ReplyPushLoop}: decision logic, nudge injection, idempotency,
|
||||
* bounded reminders, and stop conditions.
|
||||
*
|
||||
* <p>CB-590 collapsed the CB-307 reply-nudge schedule and the CB-588 ticket-nudge schedule into
|
||||
* one schedule per lead ({@link ReplyPushLoop#decide}), so most tests below register pending work
|
||||
* through the public entry points ({@code onReplyQueued} / {@code onTicketTerminal}) before
|
||||
* exercising {@code decide} directly, mirroring how the two entry points now share one decision
|
||||
* function keyed by the lead terminal rather than by worker target.
|
||||
*
|
||||
* <p>Uses a {@link RecordingHerdrClient} that synchronizes access to its call list so the
|
||||
* scheduler thread and test thread never have memory ordering issues. The {@code decide()}
|
||||
* tests use a simple client with no concurrency concern.
|
||||
@@ -58,38 +65,50 @@ class ReplyPushLoopTest {
|
||||
// --- decide() logic ------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void decideWithoutPrimaryIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = new ReplyPushLoop(
|
||||
new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0));
|
||||
void onReplyQueuedWithNoKnownLeadNeverStartsASchedule() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
|
||||
|
||||
loop.onReplyQueued(WORKER); // no lead known -> never registered, never scheduled
|
||||
|
||||
assertFalse(loop.isActive(), "no lead known means nothing to nudge yet");
|
||||
Thread.sleep(150);
|
||||
assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting");
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideWithEmptyInboxIsStop() {
|
||||
void decideWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideAtCapIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
|
||||
var loop = loop(2, 100);
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideUnderCapWithInjectablePrimaryIsInject() {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideUnderCapWithBlockedPrimaryIsInject() {
|
||||
agents = agentWithStatus("blocked");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
|
||||
"BLOCKED is injectable");
|
||||
}
|
||||
|
||||
@@ -97,7 +116,9 @@ class ReplyPushLoopTest {
|
||||
void decideUnderCapWithDonePrimaryIsInject() {
|
||||
agents = agentWithStatus("done");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
|
||||
"DONE is injectable");
|
||||
}
|
||||
|
||||
@@ -105,23 +126,29 @@ class ReplyPushLoopTest {
|
||||
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
|
||||
agents = agentWithStatus("working");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
|
||||
agents = agentWithStatus("unknown");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideStopsAfterInboxIsEmptied() {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
inbox.ack(WORKER, "m1");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
// --- onReplyQueued integration -------------------------------------------------------------
|
||||
@@ -201,12 +228,19 @@ class ReplyPushLoopTest {
|
||||
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
|
||||
@Test
|
||||
void repliesNudgeFormatIsCorrect() {
|
||||
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
|
||||
assertTrue(multi.contains("2 workers"));
|
||||
assertTrue(multi.contains("bridge_poll(target=...)"));
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
|
||||
|
||||
@Test
|
||||
void decideTicketsWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -214,7 +248,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(2, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -222,7 +256,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -230,7 +264,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("working");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -238,7 +272,7 @@ class ReplyPushLoopTest {
|
||||
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));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
|
||||
@@ -255,7 +289,7 @@ class ReplyPushLoopTest {
|
||||
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");
|
||||
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -346,11 +380,11 @@ class ReplyPushLoopTest {
|
||||
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
|
||||
// decide's 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 —
|
||||
// reliable, so this drives stopOrRestart — 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 was NOT part of the pre-decision snapshot (ticketsBefore=empty), standing in for one
|
||||
// that races in during the decision-to-release window.
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
@@ -358,9 +392,9 @@ class ReplyPushLoopTest {
|
||||
|
||||
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
|
||||
// Stand in for the scheduler thread reaching decide==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());
|
||||
loop.stopOrRestart(PRIMARY, Set.of(), 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");
|
||||
@@ -368,18 +402,19 @@ class ReplyPushLoopTest {
|
||||
|
||||
@Test
|
||||
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
|
||||
// The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) would be
|
||||
// The other direction of the same fix: restarting on ANY non-empty pending set 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.
|
||||
// #5: no spin / nudges stay bounded). task-1 here was already accounted for at decide time
|
||||
// (it is in ticketsBefore), 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"));
|
||||
loop.stopOrRestart(PRIMARY, Set.of(), 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");
|
||||
@@ -402,7 +437,52 @@ class ReplyPushLoopTest {
|
||||
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");
|
||||
"CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
|
||||
}
|
||||
|
||||
// --- CB-590: one schedule per lead — no overlap, no lost nudges -----------------------------
|
||||
|
||||
@Test
|
||||
void replyAndTicketForTheSameLeadCoalesceIntoOneSendNeverOverlapping() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
|
||||
|
||||
loop.onReplyQueued(WORKER);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
|
||||
Thread.sleep(300);
|
||||
assertEquals(1, rec.sendCount(),
|
||||
"a reply and a ticket for the same lead must coalesce onto ONE schedule — "
|
||||
+ "two nudge injections into the same lead pane must never overlap");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
|
||||
"the combined nudge must still mention the reply: " + nudge);
|
||||
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void aReplyQueuedWhileTheLeadIsBusyIsNotLostWhenATicketArrivesToo() throws Exception {
|
||||
// The lead is busy for its first two status checks, then becomes injectable. A reply is
|
||||
// queued while busy; a ticket for the same lead arrives before the lead frees up. Neither
|
||||
// may be dropped — deferred is fine, lost is not (acceptance criterion #2).
|
||||
var rec = new BusyThenIdleHerdrClient(2);
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(1, 50); // cap=1: WAIT_BUSY doesn't count against it, so exactly one send once injectable
|
||||
|
||||
loop.onReplyQueued(WORKER); // schedule starts, first tick(s) WAIT_BUSY
|
||||
loop.onTicketTerminal("task-1", WORKER, false); // coalesces onto the same waiting schedule
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
|
||||
"once the lead becomes injectable, the deferred work must still be nudged");
|
||||
Thread.sleep(200);
|
||||
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
|
||||
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
|
||||
}
|
||||
|
||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||
@@ -430,8 +510,10 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
Metrics metrics = new Metrics();
|
||||
var loop = loop(2, 100, metrics);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
|
||||
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the reminder cap must count as exhausted");
|
||||
@@ -459,7 +541,7 @@ class ReplyPushLoopTest {
|
||||
var loop = loop(2, 100_000, metrics);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
|
||||
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the ticket reminder cap must count as exhausted");
|
||||
@@ -547,4 +629,49 @@ class ReplyPushLoopTest {
|
||||
private static RecordingHerdrClient recordingClient() {
|
||||
return new RecordingHerdrClient();
|
||||
}
|
||||
|
||||
/**
|
||||
* Thread-safe fake that reports {@code working} (not injectable) for its first
|
||||
* {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
|
||||
* while the lead is busy is deferred, not dropped, once it becomes injectable.
|
||||
*/
|
||||
private static final class BusyThenIdleHerdrClient implements HerdrClient {
|
||||
private final List<Map.Entry<String, Object>> calls =
|
||||
Collections.synchronizedList(new ArrayList<>());
|
||||
private final AtomicInteger statusChecks = new AtomicInteger();
|
||||
private final int busyChecks;
|
||||
volatile CountDownLatch sendLatch = new CountDownLatch(1);
|
||||
|
||||
BusyThenIdleHerdrClient(int busyChecks) {
|
||||
this.busyChecks = busyChecks;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
String status = statusChecks.getAndIncrement() < busyChecks ? "working" : "idle";
|
||||
return MAPPER.createObjectNode()
|
||||
.set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", PRIMARY)
|
||||
.put("agent_status", status));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
calls.add(Map.entry(method, params));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
long sendCount() {
|
||||
return calls.size();
|
||||
}
|
||||
|
||||
List<Map.Entry<String, Object>> sentParams() {
|
||||
return List.copyOf(calls);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user