CB-307 Increment 2: ReplyPushLoop — status-gated push loop (mechanism b)

- ReplyPushLoop: dedicated scheduled loop for nudging the primary
- Package-private decide() method for pure decision logic (unit-testable)
- Status-gated injection via AgentControl.status().injectable()
- Bounded reminders (cap + backoff), idempotent per target
- Wire into MessageService.reply after inbox.publish on no-waiter branch
- Wire into Bridged.main (constructor + shutdown hook)
- Primary config record updated with pushReminders/pushBackoffMs knobs
- Tests: decide() matrix, nudge injection, idempotency, cap enforcement
This commit is contained in:
Dai Ha
2026-07-19 10:16:24 +02:00
parent a1aecbf4fc
commit f756933879
5 changed files with 487 additions and 8 deletions
@@ -21,6 +21,7 @@ import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.ReplyInbox;
import dev.ltms.bridged.msg.ReplyPushLoop;
import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.SessionManager;
@@ -32,6 +33,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.file.Path;
import java.util.concurrent.Executors;
/**
* {@code bridged} entry point. Wires the real herdr socket client to the REST app and
@@ -131,15 +133,24 @@ public final class Bridged {
replyInbox = new InMemoryReplyInbox();
log.info("reply inbox: in-memory (soft-state)");
}
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox);
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
PrimaryRegistry primaryRegistry = new PrimaryRegistry(
cfg.primary() != null ? cfg.primary().terminal() : null);
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
// open bridge_send. Uses its own lightweight scheduled executor, separate from the injector.
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-push-").unstarted(r));
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
System::nanoTime, 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.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
PrimaryRegistry primaryRegistry = new PrimaryRegistry(
cfg.primary() != null ? cfg.primary().terminal() : null);
// Cast to ClaudeCodeLauncher: BridgeMcp is not yet migrated to PeerLauncher (Stage A scope).
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, (ClaudeCodeLauncher) workers, sessions, identity, presence,
primaryRegistry);
@@ -151,6 +162,7 @@ public final class Bridged {
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
poller.stop();
messages.close();
pushLoop.close();
mcp.close();
if (reaper != null) reaper.stop();
// Release the broker connection last among message resources (no-op for the in-memory inbox).
@@ -197,10 +197,26 @@ public record BridgedConfig(
* deriving it from the MCP connection. Useful when the primary runs off-host or in a
* non-herdr terminal where connection-derived identity is unavailable.
*
* @param terminal the primary's herdr {@code terminal_id} ({@code null}/blank → derive)
* @param terminal the primary's herdr {@code terminal_id} ({@code null}/blank → derive)
* @param pushReminders max reminder nudges before giving up (default 5)
* @param pushBackoffMs delay between reminders in ms (default 15000)
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Primary(String terminal) {
public record Primary(String terminal, Integer pushReminders, Integer pushBackoffMs) {
/** Legacy constructor with just a terminal — both push knobs default. */
public Primary(String terminal) {
this(terminal, null, null);
}
/** @return configured reminder cap, or 5 */
public int remindersOrDefault() {
return pushReminders != null ? pushReminders : 5;
}
/** @return configured backoff in ms, or 15000 */
public long backoffMsOrDefault() {
return pushBackoffMs != null ? pushBackoffMs.longValue() : 15_000L;
}
}
/**
@@ -150,18 +150,31 @@ public final class MessageService {
private final Injector injector;
private final Rendezvous rendezvous;
private final ReplyInbox inbox;
private final ReplyPushLoop pushLoop;
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
/** Create with an explicit {@link ReplyInbox} (Stage 1: {@link InMemoryReplyInbox}). */
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) {
/**
* Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}.
*
* @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply
* branch ({@link #reply}) so it can nudge the primary to drain the inbox
*/
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
ReplyInbox inbox, ReplyPushLoop pushLoop) {
this.agents = agents;
this.injector = injector;
this.rendezvous = rendezvous;
this.inbox = inbox;
this.pushLoop = pushLoop;
}
/** Create with an explicit {@link ReplyInbox} and no push loop. */
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) {
this(agents, injector, rendezvous, inbox, null);
}
/** Backward-compatible constructor that uses a default {@link InMemoryReplyInbox}. */
@@ -190,6 +203,9 @@ public final class MessageService {
return true; // a live send took it — unchanged fast path
}
inbox.publish(session, UUID.randomUUID().toString(), content);
if (pushLoop != null) {
pushLoop.onReplyQueued(session);
}
return true; // held, not lost
}
@@ -0,0 +1,161 @@
package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import org.slf4j.Logger;
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
* herdr pane when a worker reply lands with no live {@code bridge_send} to resolve it.
*
* <p>The loop is triggered by {@link #onReplyQueued(String)} (called from
* {@link MessageService#reply} after the durable inbox publish). It checks four conditions
* at each tick via {@link #decide(String, int)}, then either injects a drain nudge,
* waits for the primary to become injectable, or stops reminding.
*
* <p>Bounded: at most {@link #maxReminders} nudges per target, with a configurable backoff
* between them. The reply is never lost — the durable inbox is the backstop.
*/
public final class ReplyPushLoop {
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
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;
/** Track targets that have an active schedule. */
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
LongSupplier clock, 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;
}
// --- decision logic (package-private for unit-testing) -------------------------------------
/** The action the loop should take for a target at the given reminder count. */
enum Action { INJECT, WAIT_BUSY, STOP }
/**
* Pure decision function: examine the current state and return what the loop should do.
*
* @param target the worker session (target terminal id)
* @param reminderCount how many nudges have been sent so far for this target
* @return the action the caller should take
*/
Action decide(String target, int reminderCount) {
if (!primaryRegistry.isKnown()) {
log.debug("push: primary unknown, stopping reminder for {}", target);
return Action.STOP;
}
if (inbox.peek(target).isEmpty()) {
log.debug("push: inbox empty for {}, stopping reminder", target);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: reminder cap ({}) reached for {}, stopping", maxReminders, target);
return Action.STOP;
}
var primaryTerminal = primaryRegistry.primaryTerminal().orElseThrow();
AgentStatus status;
try {
status = agents.status(primaryTerminal);
} catch (RuntimeException e) {
log.debug("push: status check failed for primary {}, will retry", primaryTerminal, e);
return Action.WAIT_BUSY;
}
if (status.injectable()) {
return Action.INJECT;
}
log.debug("push: primary {} is {} (not injectable), waiting", primaryTerminal, status);
return Action.WAIT_BUSY;
}
// --- public entrypoint ---------------------------------------------------------------------
/**
* Called when a reply is queued for {@code target}. Idempotent per target: a second call while
* a schedule is active is a no-op. The schedule nudges the primary, then schedules a follow-up
* check (reminder on backoff, or re-check on WAIT_BUSY), until the inbox is empty or the cap
* is reached.
*/
public void onReplyQueued(String target) {
if (activeTargets.putIfAbsent(target, Boolean.TRUE) != null) {
log.debug("push: already active for {}, ignoring duplicate trigger", target);
return; // already scheduled
}
log.debug("push: starting reminder loop for {}", target);
scheduleNext(target, 0);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String target, int reminderCount) {
var action = decide(target, reminderCount);
switch (action) {
case INJECT -> {
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);
}
case STOP -> {
activeTargets.remove(target);
log.debug("push: reminder loop ended for {}", target);
}
}
}
/** Send the nudge and log the event. */
private void injectNudge(String target, int reminderCount) {
var primaryTerminal = primaryRegistry.primaryTerminal().orElseThrow();
String nudge = NUDGE_FORMAT.formatted(target, target);
try {
agents.send(primaryTerminal, nudge);
log.debug("push: nudge {}/{} sent to primary {} for target {}",
reminderCount + 1, maxReminders, primaryTerminal, target);
} catch (RuntimeException e) {
log.warn("push: failed to nudge primary {} for target {} (reminder {}/{}): {}",
primaryTerminal, target, reminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String target, int nextReminderCount) {
scheduler.schedule(() -> tick(target, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
}
// --- lifecycle -----------------------------------------------------------------------------
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
}
/** @see #stop() */
public void close() {
stop();
}
}
@@ -0,0 +1,274 @@
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;
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.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.*;
/**
* 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 LongSupplier clock;
private ScheduledExecutorService scheduler;
@BeforeEach
void setUp() {
registry = new PrimaryRegistry(PRIMARY);
inbox = new InMemoryReplyInbox();
clock = System::nanoTime;
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, clock, 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("idle");
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
loop(1, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (2 agent.send calls) should have been sent");
// Exactly one nudge = exactly 2 agent.send calls (text + submit)
assertEquals(2, 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("idle");
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 (2 sends)");
Thread.sleep(200);
assertEquals(2, rec.sendCount(),
"second onReplyQueued must not trigger another nudge");
}
@Test
void sendsUpToCapThenStops() throws Exception {
int cap = 2;
var rec = recordingClient("idle");
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
rec.sendLatch = new CountDownLatch(cap * 2);
loop(cap, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(5, TimeUnit.SECONDS),
cap + " nudges (" + (cap * 2) + " sends) should have fired");
Thread.sleep(300);
assertEquals(cap * 2, rec.sendCount(),
"exactly " + (cap * 2) + " agent.send calls (cap=" + cap + ")");
}
// --- nudge format --------------------------------------------------------------------------
@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)"));
}
// --- helpers -------------------------------------------------------------------------------
private ReplyPushLoop loop() {
return loop(5, 100);
}
private ReplyPushLoop loop(int maxReminders, long backoffMs) {
return new ReplyPushLoop(registry, agents, inbox, clock, scheduler, maxReminders, backoffMs);
}
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.send calls. Uses synchronized access
* 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));
}
if ("agent.send".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(String status) {
return new RecordingHerdrClient(status);
}
}