From 081893fb32f3d13b7bc9614479e563787df9f1d0 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 13 Jul 2026 15:52:53 +0200 Subject: [PATCH] =?UTF-8?q?CB-103:=20status-gated=20injector=20=E2=80=94?= =?UTF-8?q?=20per-worker=20FIFO,=20safe-window=20delivery,=20single=20writ?= =?UTF-8?q?er?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../main/java/dev/ltms/bridged/Bridged.java | 14 +- .../dev/ltms/bridged/inject/Injector.java | 182 +++++++++++++++++ .../dev/ltms/bridged/inject/StatusPoller.java | 84 ++++++++ .../dev/ltms/bridged/herdr/FakeHerdr.java | 30 ++- .../dev/ltms/bridged/inject/InjectorTest.java | 185 ++++++++++++++++++ 5 files changed, 490 insertions(+), 5 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/inject/Injector.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java 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()); + } +}