CB-307: primary-gate cleanup of push-loop (drop dead clock, IDE 0/0)
Primary verification pass over the delegated push-loop delivery: - Remove the unused LongSupplier clock threaded into ReplyPushLoop (timing is the scheduler's; the field was never read) from the component, Bridged wiring, and both test call sites. - Collapse the single-statement WAIT_BUSY switch arm (redundant block). - Drop now-dead test scaffolding: the always-"idle" recordingClient param and unused AgentStatus/AtomicReference imports. IDE diagnostics 0/0 on all changed files; mvn clean install green (242 tests, 0 failures).
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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<String, Boolean> 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);
|
||||
|
||||
@@ -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<Map.Entry<String, Object>> 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();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user