Files
fleetd/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java
T
Dai Ha ac044e7573
CI / build (pull_request) Successful in 1m19s
CI / contract (pull_request) Successful in 1m31s
CB-588 round 3: close the STOP-vs-onTicketTerminal race, pin nudge polarity, fix comment
- ReplyPushLoop.stopOrRestartTicketLoop: after releasing a lead's active-schedule
  slot on STOP, restart only if a ticket landed that the pre-decision snapshot
  did not already account for. A naive "restart on any pending ticket" version
  was tried first and reverted: it defeated the reminder cap by restarting
  forever on a stale, never-collected ticket (broke
  successfulTicketNudgeIncrementsDelivered and ticketNudgesSendUpToCapThenStop).
  Diffing against a pendingBefore snapshot distinguishes a genuine race arrival
  from stale cap-exhausted backlog.
- Two new ReplyPushLoopTest cases exercise stopOrRestartTicketLoop directly
  (now package-private) rather than forcing the underlying thread race:
  aTicketStillPendingWhenTheLoopStopsIsNotStranded (the race must restart) and
  aStaleUncollectedTicketAtCapDoesNotRestartTheLoop (the cap must still hold).
  Both were verified to fail against deliberately-reverted versions of the fix
  before being restored to green.
- MessageServiceTest: pin the success-path nudge test's negative direction too
  (must not contain "FAILED"), not just the failure-path test.
- MessageService: complete the whenComplete comment's list of completion paths
  (TIMED_OUT/BUSY/BACKEND_EXHAUSTED via finishAsyncTask, completeExceptionally
  on throw, and answer() -> finishAsyncTask(turnId, result)).
2026-08-15 17:02:05 +02:00

551 lines
22 KiB
Java

package dev.ltms.bridged.msg;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
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;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.*;
/**
* Unit tests for {@link ReplyPushLoop}: decision logic, nudge injection, idempotency,
* bounded reminders, and stop conditions.
*
* <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.
*/
class ReplyPushLoopTest {
private static final String PRIMARY = "term_primary";
private static final String WORKER = "term_worker";
private static final ObjectMapper MAPPER = new ObjectMapper();
private PrimaryRegistry registry;
private AgentControl agents;
private InMemoryReplyInbox inbox;
private ScheduledExecutorService scheduler;
@BeforeEach
void setUp() {
registry = new PrimaryRegistry(PRIMARY);
inbox = new InMemoryReplyInbox();
inbox.own(WORKER); // CB-520: the inbox only peeks/acks targets it owns
scheduler = Executors.newSingleThreadScheduledExecutor();
}
@AfterEach
void tearDown() {
scheduler.shutdownNow();
}
// --- 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));
}
@Test
void decideWithEmptyInboxIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
}
@Test
void decideAtCapIsStop() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
}
@Test
void decideUnderCapWithInjectablePrimaryIsInject() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
}
@Test
void decideUnderCapWithBlockedPrimaryIsInject() {
agents = agentWithStatus("blocked");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
"BLOCKED is injectable");
}
@Test
void decideUnderCapWithDonePrimaryIsInject() {
agents = agentWithStatus("done");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
"DONE is injectable");
}
@Test
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
agents = agentWithStatus("working");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
}
@Test
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
agents = agentWithStatus("unknown");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
}
@Test
void decideStopsAfterInboxIsEmptied() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
inbox.ack(WORKER, "m1");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@Test
void injectablePrimaryCausesExactlyOneNudge() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
loop(1, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (1 agent.prompt call) should have been sent");
// Exactly one nudge = exactly 1 agent.prompt call (it submits itself)
assertEquals(1, rec.sendCount());
assertTrue(rec.sentParams().stream()
.anyMatch(e -> e.getValue().toString().contains("bridge_poll")),
"nudge text should contain bridge_poll");
}
@Test
void onReplyQueuedIsIdempotentPerTarget() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 100);
loop.onReplyQueued(WORKER);
loop.onReplyQueued(WORKER); // second call — should be a no-op
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"expected exactly one nudge (1 prompt)");
Thread.sleep(200);
assertEquals(1, rec.sendCount(),
"second onReplyQueued must not trigger another nudge");
}
@Test
void sendsUpToCapThenStops() throws Exception {
int cap = 2;
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
rec.sendLatch = new CountDownLatch(cap);
loop(cap, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(5, TimeUnit.SECONDS),
cap + " nudges (" + cap + " prompts) should have fired");
Thread.sleep(300);
assertEquals(cap, rec.sendCount(),
"exactly " + cap + " agent.prompt calls (cap=" + cap + ")");
}
// --- nudge format --------------------------------------------------------------------------
@Test
void isActiveReflectsALiveScheduleForTheHeartbeatStandDown() {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
ReplyPushLoop loop = loop(1, 100_000); // long backoff so the tick cannot fire mid-test
loop.onReplyQueued(WORKER);
assertTrue(loop.isActive(), "CB-551: the heartbeat must stand aside while a reminder is live");
loop.stop();
assertFalse(loop.isActive(), "stopping clears the active schedule");
}
@Test
void nudgeFormatIsCorrect() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("Worker term_worker"));
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
void successfulNudgeIncrementsDelivered() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
loop(1, 50, metrics).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (1 agent.prompt call) should have been sent");
// The delivered count is bumped on the scheduler thread right after the send that releases
// the latch — settle briefly so the counter is published before we read it.
Thread.sleep(200);
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered"),
"a successfully sent nudge must count as delivered");
}
@Test
void reminderCapIncrementsExhausted() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
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() {
return loop(5, 100);
}
private ReplyPushLoop loop(int maxReminders, long backoffMs) {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs);
}
private ReplyPushLoop loop(int maxReminders, long backoffMs, Metrics metrics) {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics);
}
private static AgentControl agentWithStatus(String status) {
return new AgentControl(new FakeHerdrClient(status));
}
/** Non-recording (single-threaded) fake — safe for decide() tests. */
private static final class FakeHerdrClient implements HerdrClient {
private final String agentStatus;
FakeHerdrClient(String agentStatus) {
this.agentStatus = agentStatus;
}
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", PRIMARY)
.put("agent_status", agentStatus));
}
return MAPPER.createObjectNode();
}
@Override
public void close() {
}
}
/**
* Thread-safe recording fake that counts agent.prompt calls (protocol 19: one nudge = one
* prompt). Uses synchronized access so the scheduler thread and test thread never race.
*/
private static final class RecordingHerdrClient implements HerdrClient {
private final List<Map.Entry<String, Object>> calls =
Collections.synchronizedList(new ArrayList<>());
volatile CountDownLatch sendLatch = new CountDownLatch(1);
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", PRIMARY)
.put("agent_status", "idle")); // recording double is always injectable
}
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() {
}
}
private static RecordingHerdrClient recordingClient() {
return new RecordingHerdrClient();
}
}