CB-103: status-gated injector — per-worker FIFO, safe-window delivery, single writer

This commit is contained in:
Dai Ha
2026-07-13 15:52:53 +02:00
parent 7fccdbeeca
commit 081893fb32
5 changed files with 490 additions and 5 deletions
@@ -5,6 +5,8 @@ import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
@@ -22,7 +24,10 @@ public final class Bridged {
private static final Logger log = LoggerFactory.getLogger(Bridged.class);
public static void main(String[] args) {
/** How often the injector samples a busy worker's status while it has queued work. */
private static final long INJECT_POLL_MILLIS = 250;
static void main(String[] args) {
Path configPath = Path.of(args.length > 0 ? args[0] : "bridged.yaml");
BridgedConfig cfg = BridgedConfig.load(configPath);
@@ -41,6 +46,13 @@ public final class Bridged {
WorkspaceControl spaces = new WorkspaceControl(herdr);
WorkerService workers = new WorkerService(agents, spaces, guard, cfg.worker(), System::getenv);
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then.
Injector injector = new Injector(agents);
StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS);
poller.start();
Runtime.getRuntime().addShutdownHook(new Thread(poller::stop));
Javalin app = new BridgedApp(herdr, workers).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
@@ -0,0 +1,182 @@
package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
/**
* The status-gated injector (CB-103): the single writer that delivers a message into a
* worker only when it is safe — {@code idle} or {@code blocked}, never mid-turn.
*
* <p>Per target it holds a FIFO queue and delivers <strong>at most one message per turn</strong>:
* after a send it waits for the worker to pick the message up (go {@code working}) before
* delivering the next, so two rapid deliveries never interleave into one turn. Each
* status-check-then-send for a target is serialized on the target's monitor, closing the
* TOCTOU window between "is it idle?" and "send" — one writer per worker.
*
* <p><strong>Polling caveat.</strong> This is driven by {@link #onStatus} sampling (a
* {@code StatusPoller}), not by reliable status <em>edges</em>. A turn can begin and end
* entirely between two polls, so the {@code working} pickup may never be sampled. To avoid
* wedging a queue forever, an awaited pickup is released after {@link #PICKUP_GRACE_POLLS}
* consecutive injectable samples (the worker has plainly moved on). Only {@code working} — not
* a transient {@code unknown} — counts as a real pickup, so a detection glitch can't prematurely
* release the latch. Perfectly reliable turn boundaries require a herdr {@code events.subscribe}
* stream; that is the intended upgrade and would replace only the sampling, not this queue.
*/
public final class Injector {
private static final Logger log = LoggerFactory.getLogger(Injector.class);
/**
* How many consecutive injectable samples (with no {@code working} in between) after a send
* before we assume the turn completed unobserved and release the pickup latch. At the
* default 250ms poll interval this is a ~2s grace — far longer than a worker takes to start
* a turn, so it only fires on a genuinely missed pickup edge.
*/
private static final int PICKUP_GRACE_POLLS = 8;
private final AgentControl agents;
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
public Injector(AgentControl agents) {
this.agents = agents;
}
/** A pending message and the future that completes when it has been delivered. */
private record Pending(String text, CompletableFuture<Void> delivered) {
}
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
private static final class Target {
final Deque<Pending> queue = new ArrayDeque<>();
boolean awaitingPickup; // sent a message, waiting for the worker to pick it up
int injectableSincePickup; // consecutive injectable samples while awaitingPickup
synchronized void add(Pending p) {
queue.add(p);
}
}
/**
* Queue {@code text} for delivery to {@code target} (a herdr {@code terminal_id}). Returns
* immediately with a future that completes when the message is actually sent — the worker
* may be mid-turn, in which case delivery waits for the next injectable window.
*
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the
* target" and "queue the message" and orphan it in a target it just removed.
*/
public CompletableFuture<Void> enqueue(String target, String text) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(text, delivered);
targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
return t;
});
return delivered;
}
/**
* Feed a fresh status observation for {@code target}. Delivers the head of the queue iff the
* worker is injectable and no earlier message is still awaiting pickup. Serialized per target
* so the check and the send cannot race another delivery to the same worker; the delivered
* future is completed <em>after</em> the monitor is released so a caller's continuation never
* runs on the poller thread while it holds the lock.
*/
public void onStatus(String target, AgentStatus status) {
Target t = targets.get(target);
if (t == null) return;
Pending sent = null;
RuntimeException sendError = null;
synchronized (t) {
if (status == AgentStatus.WORKING) {
// Definitive pickup: the worker is busy on our last message.
t.awaitingPickup = false;
t.injectableSincePickup = 0;
} else if (status.injectable()) { // IDLE or BLOCKED
if (t.awaitingPickup && ++t.injectableSincePickup >= PICKUP_GRACE_POLLS) {
// Pickup edge was never sampled (turn faster than the poll, or status lag).
// The worker has plainly moved on — release the latch rather than wedge.
t.awaitingPickup = false;
t.injectableSincePickup = 0;
}
if (!t.awaitingPickup) {
Pending p = t.queue.peek();
if (p != null) {
try {
agents.send(target, p.text());
t.queue.poll();
t.awaitingPickup = true;
t.injectableSincePickup = 0;
sent = p;
} catch (RuntimeException e) {
// Delivery failed at herdr; drop the poisoned message and surface it
// rather than blocking the queue behind it.
t.queue.poll();
sent = p;
sendError = e;
}
}
}
}
// UNKNOWN (and any other non-injectable, non-working): do nothing — neither a safe
// window nor a reliable pickup signal, so we must not deliver or release the latch.
// Reclaim the entry once the worker is fully quiescent, so the map cannot grow without
// bound across many short-lived workers.
if (t.queue.isEmpty() && !t.awaitingPickup) {
targets.remove(target, t);
}
}
if (sent != null) {
if (sendError != null) {
log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage());
sent.delivered().completeExceptionally(sendError);
} else {
sent.delivered().complete(null);
}
}
}
/** Targets the poller must keep sampling: those with a queued message or an awaited pickup. */
public Set<String> activeTargets() {
return targets.entrySet().stream()
.filter(e -> {
synchronized (e.getValue()) {
return !e.getValue().queue.isEmpty() || e.getValue().awaitingPickup;
}
})
.map(java.util.Map.Entry::getKey)
.collect(Collectors.toSet());
}
/**
* Forget a target whose worker is gone, failing every still-queued message so awaiting
* callers unblock instead of hanging forever. Futures are completed after the monitor is
* released.
*/
public void drop(String target, Throwable cause) {
Target t = targets.remove(target);
if (t == null) return;
List<Pending> pending;
synchronized (t) {
pending = new ArrayList<>(t.queue);
t.queue.clear();
}
for (Pending p : pending) {
p.delivered().completeExceptionally(cause);
}
}
}
@@ -0,0 +1,84 @@
package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.HerdrException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Set;
/**
* Drives the {@link Injector} by sampling each active worker's {@code agent_status} and
* feeding it in. A single virtual-thread loop polls only targets that have work outstanding,
* so an idle bridge does no herdr traffic.
*
* <p>This is the Stage-1 gate signal. It can later be replaced (or fronted) by a herdr
* {@code events.subscribe} stream without touching the {@link Injector} — the injector only
* consumes {@code onStatus} calls, however they are produced.
*/
public final class StatusPoller {
private static final Logger log = LoggerFactory.getLogger(StatusPoller.class);
private final AgentControl agents;
private final Injector injector;
private final long intervalMillis;
private volatile boolean running;
private Thread thread;
public StatusPoller(AgentControl agents, Injector injector, long intervalMillis) {
this.agents = agents;
this.injector = injector;
this.intervalMillis = intervalMillis;
}
/** Start the polling loop on a virtual thread. Idempotent. */
public synchronized void start() {
if (running) return;
running = true;
thread = Thread.ofVirtual().name("status-poller").start(this::loop);
log.info("status poller started (interval {}ms)", intervalMillis);
}
private void loop() {
while (running) {
Set<String> active = injector.activeTargets();
for (String target : active) {
if (!running) return;
try {
AgentStatus status = agents.status(target);
injector.onStatus(target, status);
} catch (HerdrException e) {
// The worker's agent is gone — stop trying and unblock its waiters.
if (e.code() != null && e.code().endsWith("_not_found")) {
log.debug("target {} gone; dropping its queue", target);
injector.drop(target, e);
} else {
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
}
} catch (RuntimeException e) {
// Never let one target's unexpected error (e.g. an odd agent.get shape) kill
// the single poller thread and stall injection for every worker.
log.warn("unexpected error polling {}; skipping this round", target, e);
}
}
sleep();
}
}
private void sleep() {
try {
Thread.sleep(intervalMillis);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
running = false;
}
}
/** Stop the polling loop. Idempotent. */
public synchronized void stop() {
running = false;
if (thread != null) thread.interrupt();
}
}