diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index aad56e7..b552fdf 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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). diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 714d803..9845dad 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -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; + } } /** diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 046dadd..7a15666 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -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 sessionLocks = new ConcurrentHashMap<>(); private final ConcurrentHashMap 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 } diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java new file mode 100644 index 0000000..8168cd2 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -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. + * + *

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. + * + *

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 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(); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java new file mode 100644 index 0000000..1b122bc --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -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. + * + *

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> 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> sentParams() { + return List.copyOf(calls); + } + + @Override + public void close() { + } + } + + private static RecordingHerdrClient recordingClient(String status) { + return new RecordingHerdrClient(status); + } +}