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 ->
|
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||||
Thread.ofVirtual().name("bridge-push-").unstarted(r));
|
Thread.ofVirtual().name("bridge-push-").unstarted(r));
|
||||||
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
|
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
|
||||||
System::nanoTime, pushScheduler, maxReminders, backoffMs);
|
pushScheduler, maxReminders, backoffMs);
|
||||||
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, pushLoop);
|
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, pushLoop);
|
||||||
|
|
||||||
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
|
// 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.ConcurrentHashMap;
|
||||||
import java.util.concurrent.ScheduledExecutorService;
|
import java.util.concurrent.ScheduledExecutorService;
|
||||||
import java.util.concurrent.TimeUnit;
|
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
|
* 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 PrimaryRegistry primaryRegistry;
|
||||||
private final AgentControl agents;
|
private final AgentControl agents;
|
||||||
private final ReplyInbox inbox;
|
private final ReplyInbox inbox;
|
||||||
private final LongSupplier clock;
|
|
||||||
private final ScheduledExecutorService scheduler;
|
private final ScheduledExecutorService scheduler;
|
||||||
private final int maxReminders;
|
private final int maxReminders;
|
||||||
private final long backoffMs;
|
private final long backoffMs;
|
||||||
@@ -40,12 +38,11 @@ public final class ReplyPushLoop {
|
|||||||
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
|
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||||
LongSupplier clock, ScheduledExecutorService scheduler,
|
ScheduledExecutorService scheduler,
|
||||||
int maxReminders, long backoffMs) {
|
int maxReminders, long backoffMs) {
|
||||||
this.primaryRegistry = primaryRegistry;
|
this.primaryRegistry = primaryRegistry;
|
||||||
this.agents = agents;
|
this.agents = agents;
|
||||||
this.inbox = inbox;
|
this.inbox = inbox;
|
||||||
this.clock = clock;
|
|
||||||
this.scheduler = scheduler;
|
this.scheduler = scheduler;
|
||||||
this.maxReminders = maxReminders;
|
this.maxReminders = maxReminders;
|
||||||
this.backoffMs = backoffMs;
|
this.backoffMs = backoffMs;
|
||||||
@@ -116,10 +113,8 @@ public final class ReplyPushLoop {
|
|||||||
injectNudge(target, reminderCount);
|
injectNudge(target, reminderCount);
|
||||||
scheduleNext(target, reminderCount + 1);
|
scheduleNext(target, reminderCount + 1);
|
||||||
}
|
}
|
||||||
case WAIT_BUSY -> {
|
// Re-check after the configured backoff; the primary may become injectable soon.
|
||||||
// Re-check after the configured backoff; the primary may become injectable soon.
|
case WAIT_BUSY -> scheduleNext(target, reminderCount);
|
||||||
scheduleNext(target, reminderCount);
|
|
||||||
}
|
|
||||||
case STOP -> {
|
case STOP -> {
|
||||||
activeTargets.remove(target);
|
activeTargets.remove(target);
|
||||||
log.debug("push: reminder loop ended for {}", 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.JsonNode;
|
||||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import dev.ltms.bridged.herdr.AgentControl;
|
import dev.ltms.bridged.herdr.AgentControl;
|
||||||
import dev.ltms.bridged.herdr.AgentStatus;
|
|
||||||
import dev.ltms.bridged.herdr.HerdrClient;
|
import dev.ltms.bridged.herdr.HerdrClient;
|
||||||
import dev.ltms.bridged.mcp.PrimaryRegistry;
|
import dev.ltms.bridged.mcp.PrimaryRegistry;
|
||||||
import org.junit.jupiter.api.AfterEach;
|
import org.junit.jupiter.api.AfterEach;
|
||||||
@@ -18,8 +17,6 @@ import java.util.concurrent.CountDownLatch;
|
|||||||
import java.util.concurrent.Executors;
|
import java.util.concurrent.Executors;
|
||||||
import java.util.concurrent.ScheduledExecutorService;
|
import java.util.concurrent.ScheduledExecutorService;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.atomic.AtomicReference;
|
|
||||||
import java.util.function.LongSupplier;
|
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.*;
|
import static org.junit.jupiter.api.Assertions.*;
|
||||||
|
|
||||||
@@ -40,14 +37,12 @@ class ReplyPushLoopTest {
|
|||||||
private PrimaryRegistry registry;
|
private PrimaryRegistry registry;
|
||||||
private AgentControl agents;
|
private AgentControl agents;
|
||||||
private InMemoryReplyInbox inbox;
|
private InMemoryReplyInbox inbox;
|
||||||
private LongSupplier clock;
|
|
||||||
private ScheduledExecutorService scheduler;
|
private ScheduledExecutorService scheduler;
|
||||||
|
|
||||||
@BeforeEach
|
@BeforeEach
|
||||||
void setUp() {
|
void setUp() {
|
||||||
registry = new PrimaryRegistry(PRIMARY);
|
registry = new PrimaryRegistry(PRIMARY);
|
||||||
inbox = new InMemoryReplyInbox();
|
inbox = new InMemoryReplyInbox();
|
||||||
clock = System::nanoTime;
|
|
||||||
scheduler = Executors.newSingleThreadScheduledExecutor();
|
scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -62,7 +57,7 @@ class ReplyPushLoopTest {
|
|||||||
void decideWithoutPrimaryIsStop() {
|
void decideWithoutPrimaryIsStop() {
|
||||||
agents = agentWithStatus("idle");
|
agents = agentWithStatus("idle");
|
||||||
var loop = new ReplyPushLoop(
|
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));
|
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -129,7 +124,7 @@ class ReplyPushLoopTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
void injectablePrimaryCausesExactlyOneNudge() throws Exception {
|
void injectablePrimaryCausesExactlyOneNudge() throws Exception {
|
||||||
var rec = recordingClient("idle");
|
var rec = recordingClient();
|
||||||
agents = new AgentControl(rec);
|
agents = new AgentControl(rec);
|
||||||
inbox.publish(WORKER, "m1", "hello");
|
inbox.publish(WORKER, "m1", "hello");
|
||||||
|
|
||||||
@@ -147,7 +142,7 @@ class ReplyPushLoopTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
void onReplyQueuedIsIdempotentPerTarget() throws Exception {
|
void onReplyQueuedIsIdempotentPerTarget() throws Exception {
|
||||||
var rec = recordingClient("idle");
|
var rec = recordingClient();
|
||||||
agents = new AgentControl(rec);
|
agents = new AgentControl(rec);
|
||||||
inbox.publish(WORKER, "m1", "hello");
|
inbox.publish(WORKER, "m1", "hello");
|
||||||
|
|
||||||
@@ -165,7 +160,7 @@ class ReplyPushLoopTest {
|
|||||||
@Test
|
@Test
|
||||||
void sendsUpToCapThenStops() throws Exception {
|
void sendsUpToCapThenStops() throws Exception {
|
||||||
int cap = 2;
|
int cap = 2;
|
||||||
var rec = recordingClient("idle");
|
var rec = recordingClient();
|
||||||
agents = new AgentControl(rec);
|
agents = new AgentControl(rec);
|
||||||
inbox.publish(WORKER, "m1", "hello");
|
inbox.publish(WORKER, "m1", "hello");
|
||||||
rec.sendLatch = new CountDownLatch(cap * 2);
|
rec.sendLatch = new CountDownLatch(cap * 2);
|
||||||
@@ -195,7 +190,7 @@ class ReplyPushLoopTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private ReplyPushLoop loop(int maxReminders, long backoffMs) {
|
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) {
|
private static AgentControl agentWithStatus(String status) {
|
||||||
@@ -231,22 +226,17 @@ class ReplyPushLoopTest {
|
|||||||
* so the scheduler thread and test thread never race.
|
* so the scheduler thread and test thread never race.
|
||||||
*/
|
*/
|
||||||
private static final class RecordingHerdrClient implements HerdrClient {
|
private static final class RecordingHerdrClient implements HerdrClient {
|
||||||
private final String agentStatus;
|
|
||||||
private final List<Map.Entry<String, Object>> calls =
|
private final List<Map.Entry<String, Object>> calls =
|
||||||
Collections.synchronizedList(new ArrayList<>());
|
Collections.synchronizedList(new ArrayList<>());
|
||||||
volatile CountDownLatch sendLatch = new CountDownLatch(2);
|
volatile CountDownLatch sendLatch = new CountDownLatch(2);
|
||||||
|
|
||||||
RecordingHerdrClient(String agentStatus) {
|
|
||||||
this.agentStatus = agentStatus;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public JsonNode call(String method, Object params) {
|
public JsonNode call(String method, Object params) {
|
||||||
if ("agent.get".equals(method)) {
|
if ("agent.get".equals(method)) {
|
||||||
return MAPPER.createObjectNode()
|
return MAPPER.createObjectNode()
|
||||||
.set("agent", MAPPER.createObjectNode()
|
.set("agent", MAPPER.createObjectNode()
|
||||||
.put("terminal_id", PRIMARY)
|
.put("terminal_id", PRIMARY)
|
||||||
.put("agent_status", agentStatus));
|
.put("agent_status", "idle")); // recording double is always injectable
|
||||||
}
|
}
|
||||||
if ("agent.send".equals(method)) {
|
if ("agent.send".equals(method)) {
|
||||||
calls.add(Map.entry(method, params));
|
calls.add(Map.entry(method, params));
|
||||||
@@ -268,7 +258,7 @@ class ReplyPushLoopTest {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private static RecordingHerdrClient recordingClient(String status) {
|
private static RecordingHerdrClient recordingClient() {
|
||||||
return new RecordingHerdrClient(status);
|
return new RecordingHerdrClient();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user