Files
fleetd/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java
T
Dai Ha e5038c6d13
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 56s
CB-551: idle-lead heartbeat — nudge an idle lead back to work on a timer
2026-08-13 21:10:21 +02:00

200 lines
8.7 KiB
Java

package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
* 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 ScheduledExecutorService scheduler;
private final int maxReminders;
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/** Track targets that have an active schedule. */
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
int maxReminders, long backoffMs) {
this(primaryRegistry, agents, inbox, scheduler, maxReminders, backoffMs, null);
}
/** As above, with a metric registry (CB-512) so push outcomes are counted. */
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
int maxReminders, long backoffMs, Metrics metrics) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.inbox = inbox;
this.scheduler = scheduler;
this.maxReminders = maxReminders;
this.backoffMs = backoffMs;
this.metrics = metrics;
}
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(BridgedMetrics.PUSH_NUDGES, "outcome", outcome);
}
}
// --- 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) {
// CB-532: the destination is per-delegation — the lead that sent this worker its work, not
// "the primary". With two leads orchestrating one fleet the singular question has no right
// answer, and answering it anyway interrupted whichever lead happened to call bridge_send
// first with results it never asked for.
var nudgeTarget = primaryRegistry.nudgeTargetFor(target);
if (nudgeTarget.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, stopping reminder", 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);
countNudge("exhausted");
return Action.STOP;
}
String leadTerminal = nudgeTarget.get();
AgentStatus status;
try {
status = agents.status(leadTerminal);
} catch (RuntimeException e) {
log.debug("push: status check failed for lead {}, will retry", leadTerminal, e);
return Action.WAIT_BUSY;
}
if (status.injectable()) {
return Action.INJECT;
}
log.debug("push: lead {} is {} (not injectable), waiting", leadTerminal, 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);
}
// 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);
}
}
}
/** Send the nudge and log the event. */
private void injectNudge(String target, int reminderCount) {
// Re-read rather than threading it down from decide(): the delegating lead can change
// between the decision and the injection, and the nudge should follow the current one.
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: lead for {} disappeared before the nudge could be sent", target);
return;
}
String leadTerminal = lead.get();
String nudge = NUDGE_FORMAT.formatted(target, target);
try {
agents.send(leadTerminal, nudge);
log.debug("push: nudge {}/{} sent to lead {} for target {}",
reminderCount + 1, maxReminders, leadTerminal, target);
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} for target {} (reminder {}/{}): {}",
leadTerminal, 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 -----------------------------------------------------------------------------
/**
* Whether any reminder loop is currently active for some target (CB-551). The idle-lead heartbeat
* uses this to stand aside: while the push loop is actively nudging the lead, a concurrent
* heartbeat injection would start a second competing turn in the same pane — racing loops multiply
* turns and context burn. "Active" means a schedule exists in {@link #activeTargets}; the set is
* bounded by what has been triggered, not by any persistent state.
*/
public boolean isActive() {
return !activeTargets.isEmpty();
}
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
}
/** @see #stop() */
public void close() {
stop();
}
}