diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index b552fdf..21bff47 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -144,7 +144,7 @@ public final class Bridged { var pushScheduler = Executors.newSingleThreadScheduledExecutor(r -> Thread.ofVirtual().name("bridge-push-").unstarted(r)); var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox, - System::nanoTime, pushScheduler, maxReminders, backoffMs); + pushScheduler, maxReminders, backoffMs); MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, pushLoop); // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. 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 8168cd2..ac61679 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -9,7 +9,6 @@ import org.slf4j.LoggerFactory; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.function.LongSupplier; /** * Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own @@ -31,7 +30,6 @@ public final class ReplyPushLoop { private final PrimaryRegistry primaryRegistry; private final AgentControl agents; private final ReplyInbox inbox; - private final LongSupplier clock; private final ScheduledExecutorService scheduler; private final int maxReminders; private final long backoffMs; @@ -40,12 +38,11 @@ public final class ReplyPushLoop { private final ConcurrentHashMap activeTargets = new ConcurrentHashMap<>(); public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, - LongSupplier clock, ScheduledExecutorService scheduler, + ScheduledExecutorService scheduler, int maxReminders, long backoffMs) { this.primaryRegistry = primaryRegistry; this.agents = agents; this.inbox = inbox; - this.clock = clock; this.scheduler = scheduler; this.maxReminders = maxReminders; this.backoffMs = backoffMs; @@ -116,10 +113,8 @@ public final class ReplyPushLoop { injectNudge(target, reminderCount); scheduleNext(target, reminderCount + 1); } - case WAIT_BUSY -> { - // Re-check after the configured backoff; the primary may become injectable soon. - scheduleNext(target, reminderCount); - } + // Re-check after the configured backoff; the primary may become injectable soon. + case WAIT_BUSY -> scheduleNext(target, reminderCount); case STOP -> { activeTargets.remove(target); log.debug("push: reminder loop ended for {}", target); 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 1b122bc..e2dd5ec 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -3,7 +3,6 @@ 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.AgentStatus; import dev.ltms.bridged.herdr.HerdrClient; import dev.ltms.bridged.mcp.PrimaryRegistry; import org.junit.jupiter.api.AfterEach; @@ -18,8 +17,6 @@ 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.AtomicReference; -import java.util.function.LongSupplier; import static org.junit.jupiter.api.Assertions.*; @@ -40,14 +37,12 @@ class ReplyPushLoopTest { private PrimaryRegistry registry; private AgentControl agents; private InMemoryReplyInbox inbox; - private LongSupplier clock; private ScheduledExecutorService scheduler; @BeforeEach void setUp() { registry = new PrimaryRegistry(PRIMARY); inbox = new InMemoryReplyInbox(); - clock = System::nanoTime; scheduler = Executors.newSingleThreadScheduledExecutor(); } @@ -62,7 +57,7 @@ class ReplyPushLoopTest { void decideWithoutPrimaryIsStop() { agents = agentWithStatus("idle"); var loop = new ReplyPushLoop( - new PrimaryRegistry(null), agents, inbox, clock, scheduler, 5, 100); + new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100); assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0)); } @@ -129,7 +124,7 @@ class ReplyPushLoopTest { @Test void injectablePrimaryCausesExactlyOneNudge() throws Exception { - var rec = recordingClient("idle"); + var rec = recordingClient(); agents = new AgentControl(rec); inbox.publish(WORKER, "m1", "hello"); @@ -147,7 +142,7 @@ class ReplyPushLoopTest { @Test void onReplyQueuedIsIdempotentPerTarget() throws Exception { - var rec = recordingClient("idle"); + var rec = recordingClient(); agents = new AgentControl(rec); inbox.publish(WORKER, "m1", "hello"); @@ -165,7 +160,7 @@ class ReplyPushLoopTest { @Test void sendsUpToCapThenStops() throws Exception { int cap = 2; - var rec = recordingClient("idle"); + var rec = recordingClient(); agents = new AgentControl(rec); inbox.publish(WORKER, "m1", "hello"); rec.sendLatch = new CountDownLatch(cap * 2); @@ -195,7 +190,7 @@ class ReplyPushLoopTest { } private ReplyPushLoop loop(int maxReminders, long backoffMs) { - return new ReplyPushLoop(registry, agents, inbox, clock, scheduler, maxReminders, backoffMs); + return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs); } private static AgentControl agentWithStatus(String status) { @@ -231,22 +226,17 @@ class ReplyPushLoopTest { * so the scheduler thread and test thread never race. */ private static final class RecordingHerdrClient implements HerdrClient { - private final String agentStatus; private final List> calls = Collections.synchronizedList(new ArrayList<>()); volatile CountDownLatch sendLatch = new CountDownLatch(2); - RecordingHerdrClient(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)); + .put("agent_status", "idle")); // recording double is always injectable } if ("agent.send".equals(method)) { calls.add(Map.entry(method, params)); @@ -268,7 +258,7 @@ class ReplyPushLoopTest { } } - private static RecordingHerdrClient recordingClient(String status) { - return new RecordingHerdrClient(status); + private static RecordingHerdrClient recordingClient() { + return new RecordingHerdrClient(); } }