diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java
index c87e43e..a59d7d3 100644
--- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java
+++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java
@@ -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 {}",
diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java
new file mode 100644
index 0000000..5cdb47d
--- /dev/null
+++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java
@@ -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.
+ *
+ *
Per target it holds a FIFO queue and delivers at most one message per turn:
+ * 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.
+ *
+ *
Polling caveat. This is driven by {@link #onStatus} sampling (a
+ * {@code StatusPoller}), not by reliable status edges. 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 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 delivered) {
+ }
+
+ /** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
+ private static final class Target {
+ final Deque 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.
+ *
+ * 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 enqueue(String target, String text) {
+ CompletableFuture 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 after 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 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;
+ synchronized (t) {
+ pending = new ArrayList<>(t.queue);
+ t.queue.clear();
+ }
+ for (Pending p : pending) {
+ p.delivered().completeExceptionally(cause);
+ }
+ }
+}
diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java b/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java
new file mode 100644
index 0000000..43dfa64
--- /dev/null
+++ b/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java
@@ -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.
+ *
+ * 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 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();
+ }
+}
diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java
index 1e1b548..b5f8e5e 100644
--- a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java
+++ b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java
@@ -23,6 +23,8 @@ public final class FakeHerdr implements HerdrClient {
private int agentNameTakenFor = 0;
private int workerTabPaneCount = 1;
private String paneCloseErrorCode = null;
+ private String agentSendErrorCode = null;
+ private volatile String agentStatus = "idle"; // what agent.get reports
public FakeHerdr healthy(boolean h) {
this.healthy = h;
@@ -47,8 +49,16 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
- private long callCount(String method) {
- return calls.stream().filter(c -> c.method().equals(method)).count();
+ /** Set the {@code agent_status} that {@code agent.get} reports (drives the injector). */
+ public FakeHerdr agentStatus(String status) {
+ this.agentStatus = status;
+ return this;
+ }
+
+ /** Make {@code agent.send} fail with this herdr error code. */
+ public FakeHerdr agentSendFailsWith(String code) {
+ this.agentSendErrorCode = code;
+ return this;
}
/** Seed an additional workspace into {@code workspace.list} (e.g. a pre-existing worker space). */
@@ -65,7 +75,7 @@ public final class FakeHerdr implements HerdrClient {
public Call lastCall(String method) {
return calls.stream().filter(c -> c.method().equals(method))
- .reduce((a, b) -> b).orElseThrow();
+ .reduce((_, b) -> b).orElseThrow();
}
@Override
@@ -86,8 +96,20 @@ public final class FakeHerdr implements HerdrClient {
{"terminal_id":"term_a","agent":"claude","agent_status":"idle",
"agent_session":{"kind":"id","value":"sess-1111"},
"workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"}]}""");
+ case "agent.send" -> {
+ if (agentSendErrorCode != null) {
+ throw new HerdrException("herdr error [" + agentSendErrorCode + "]: agent.send failed",
+ agentSendErrorCode, null);
+ }
+ yield mapper.readTree("{\"type\":\"ok\"}");
+ }
+ case "agent.get" -> mapper.readTree(("""
+ {"type":"agent_info","agent":{"terminal_id":"term_a","agent":"claude",
+ "agent_status":"%s","workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"}}""")
+ .formatted(agentStatus));
case "agent.start" -> {
- if (callCount("agent.start") <= agentNameTakenFor) {
+ long starts = calls.stream().filter(c -> c.method().equals("agent.start")).count();
+ if (starts <= agentNameTakenFor) {
throw new HerdrException(
"herdr error [agent_name_taken]: agent name already used",
"agent_name_taken", null);
diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java
new file mode 100644
index 0000000..7d0797d
--- /dev/null
+++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java
@@ -0,0 +1,185 @@
+package dev.ltms.bridged.inject;
+
+import dev.ltms.bridged.herdr.AgentControl;
+import dev.ltms.bridged.herdr.AgentStatus;
+import dev.ltms.bridged.herdr.FakeHerdr;
+import dev.ltms.bridged.herdr.HerdrException;
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.*;
+
+/**
+ * Golden-transcript tests for the status-gated injector (CB-103). Delivery rules are driven
+ * deterministically by feeding {@code onStatus}, so no timing or real polling is involved.
+ */
+class InjectorTest {
+
+ private static final String T = "term_a";
+
+ private final FakeHerdr herdr = new FakeHerdr();
+ private final Injector injector = new Injector(new AgentControl(herdr));
+
+ /** Text of every agent.send, in order. */
+ @SuppressWarnings("unchecked")
+ private List sent() {
+ return herdr.calls.stream()
+ .filter(c -> c.method().equals("agent.send"))
+ .map(c -> ((Map) c.params()).get("text").toString())
+ .toList();
+ }
+
+ @Test
+ void deliversWhenIdle() {
+ CompletableFuture f = injector.enqueue(T, "hello");
+ assertFalse(f.isDone(), "not delivered until an injectable status arrives");
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertTrue(f.isDone());
+ assertEquals(List.of("hello"), sent());
+ }
+
+ @Test
+ void holdsWhileWorkingThenDeliversOnIdle() {
+ injector.enqueue(T, "later");
+ injector.onStatus(T, AgentStatus.WORKING);
+ assertEquals(List.of(), sent(), "must not inject mid-turn");
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("later"), sent());
+ }
+
+ @Test
+ void blockedIsInjectableButUnknownIsNot() {
+ injector.enqueue(T, "answer");
+ injector.onStatus(T, AgentStatus.UNKNOWN);
+ assertEquals(List.of(), sent(), "unknown status is not safe to inject");
+ injector.onStatus(T, AgentStatus.BLOCKED);
+ assertEquals(List.of("answer"), sent(), "blocked worker can be answered");
+ }
+
+ @Test
+ void twoRapidDeliveriesNeverInterleave() {
+ injector.enqueue(T, "m1");
+ injector.enqueue(T, "m2");
+
+ // First idle window delivers only m1, even if idle is observed twice before pickup.
+ injector.onStatus(T, AgentStatus.IDLE);
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("m1"), sent(), "second message must wait for the turn to complete");
+
+ // Worker picks up m1 (works), then returns idle → m2 delivers.
+ injector.onStatus(T, AgentStatus.WORKING);
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("m1", "m2"), sent());
+ }
+
+ @Test
+ void transientUnknownDoesNotReleaseThePickupLatch() {
+ injector.enqueue(T, "m1");
+ injector.enqueue(T, "m2");
+ injector.onStatus(T, AgentStatus.IDLE); // m1 sent, awaiting pickup
+ assertEquals(List.of("m1"), sent());
+
+ injector.onStatus(T, AgentStatus.UNKNOWN); // a detection glitch is NOT a pickup
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("m1"), sent(), "unknown must not let the next message interleave the turn");
+
+ injector.onStatus(T, AgentStatus.WORKING); // real pickup
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("m1", "m2"), sent());
+ }
+
+ @Test
+ void missedPickupEdgeIsReleasedByGraceSoTheQueueNeverWedges() {
+ injector.enqueue(T, "m1");
+ injector.enqueue(T, "m2");
+ injector.onStatus(T, AgentStatus.IDLE); // m1 sent
+ assertEquals(List.of("m1"), sent());
+
+ // WORKING is never sampled (turn faster than the poll). The latch must release after
+ // the grace window so m2 is delivered rather than wedged forever.
+ for (int i = 0; i < 20; i++) injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("m1", "m2"), sent(), "safety valve must eventually deliver m2");
+ }
+
+ @Test
+ void fifoOrderAcrossManyTurns() {
+ injector.enqueue(T, "a");
+ injector.enqueue(T, "b");
+ injector.enqueue(T, "c");
+ for (int i = 0; i < 3; i++) {
+ injector.onStatus(T, AgentStatus.IDLE); // deliver one
+ injector.onStatus(T, AgentStatus.WORKING); // pickup
+ }
+ injector.onStatus(T, AgentStatus.IDLE);
+ assertEquals(List.of("a", "b", "c"), sent());
+ }
+
+ @Test
+ void activeWhileQueuedOrInFlightThenQuietAfterPickup() {
+ assertTrue(injector.activeTargets().isEmpty());
+ injector.enqueue(T, "x");
+ assertEquals(java.util.Set.of(T), injector.activeTargets(), "active while a message is queued");
+
+ injector.onStatus(T, AgentStatus.IDLE); // delivers; still in-flight (awaiting pickup)
+ assertEquals(java.util.Set.of(T), injector.activeTargets(),
+ "stays active so the poller can observe the worker pick the message up");
+
+ injector.onStatus(T, AgentStatus.WORKING); // pickup observed → in-flight cleared
+ assertTrue(injector.activeTargets().isEmpty(), "quiet once queue is empty and pickup is seen");
+ }
+
+ @Test
+ void sendFailureDropsMessageAndFailsItsFuture() {
+ FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
+ Injector inj = new Injector(new AgentControl(failing));
+ CompletableFuture f = inj.enqueue(T, "boom");
+
+ inj.onStatus(T, AgentStatus.IDLE);
+ assertTrue(f.isCompletedExceptionally());
+ assertTrue(inj.activeTargets().isEmpty(), "poisoned message is dropped, not left blocking the queue");
+ }
+
+ @Test
+ void dropFailsPendingWaiters() {
+ CompletableFuture f = injector.enqueue(T, "orphan");
+ injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
+ assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
+ }
+
+ @Test
+ void pollerDeliversToAnIdleWorker() throws Exception {
+ // End-to-end through the poller: idle worker → message delivered without manual onStatus.
+ FakeHerdr idle = new FakeHerdr().agentStatus("idle");
+ Injector inj = new Injector(new AgentControl(idle));
+ StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10);
+ poller.start();
+ try {
+ CompletableFuture delivered = inj.enqueue(T, "via-poller");
+ delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send
+ } finally {
+ poller.stop();
+ }
+ assertEquals(List.of("via-poller"), idle.calls.stream()
+ .filter(c -> c.method().equals("agent.send"))
+ .map(c -> {
+ @SuppressWarnings("unchecked")
+ Map p = (Map) c.params();
+ return p.get("text").toString();
+ }).toList());
+ }
+
+ @Test
+ void deliveredFutureCarriesSendFailure() {
+ FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
+ Injector inj = new Injector(new AgentControl(failing));
+ CompletableFuture f = inj.enqueue(T, "boom");
+ inj.onStatus(T, AgentStatus.IDLE);
+ ExecutionException ex = assertThrows(ExecutionException.class, f::get);
+ assertInstanceOf(HerdrException.class, ex.getCause());
+ }
+}