CB-102: native agent.* worker spawn (env-injected, guard-checked)

Spike decided the worker south side in favour of herdr's native agent.*
namespace over pane+send_text. Proven against live herdr 0.7.0:
agent.start takes a first-class env map that reaches the process
environment (ANTHROPIC_BASE_URL confirmed via agent.read), and herdr
tracks each worker's Claude session UUID itself.

- AgentControl: start/send/read/get/status/list + pane.close over agent.*.
- Agent/AgentStatus: projection of herdr agent records (session UUID,
  injectable status gate).
- WorkerService: build worker env (base_url/token/model/config dir),
  assertWorker BEFORE any herdr call, then agent.start.
- REST: GET /agents (discovery by session UUID), POST /workers (201, or
  403 subscription_boundary), DELETE /workers/{paneId}.
- Full stack smoke-tested live: POST->guard->agent.start->new pane,
  GET /agents lists it, DELETE closes it.
- Tests: 21 unit/acceptance + 4 contract (incl. a live end-to-end probe
  spawn that proves env injection and cleans up its pane).

No agent.stop in herdr (use pane.close); herdr id must be a string;
one request per connection — all pinned by contract tests.
This commit is contained in:
Dai Ha
2026-07-12 20:18:18 +02:00
parent 1d558a3692
commit 2a815ee64e
10 changed files with 512 additions and 57 deletions
@@ -2,8 +2,10 @@ package dev.ltms.bridged;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -34,7 +36,10 @@ public final class Bridged {
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
Runtime.getRuntime().addShutdownHook(new Thread(herdr::close));
Javalin app = new BridgedApp(herdr).build();
AgentControl agents = new AgentControl(herdr);
WorkerService workers = new WorkerService(agents, guard, cfg.worker(), System::getenv);
Javalin app = new BridgedApp(herdr, workers).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
@@ -37,12 +37,22 @@ public record BridgedConfig(
}
/**
* @param profile ccs profile a worker is spawned under (Stage-1: {@code ltms-local})
* @param baseUrl the off-subscription endpoint the worker's launch line sets
* @param model model alias to request from that endpoint
* @param profile ccs profile a worker is spawned under (Stage-1: {@code ltms-local})
* @param baseUrl the off-subscription endpoint injected as {@code ANTHROPIC_BASE_URL}
* @param model model alias, injected as {@code ANTHROPIC_MODEL} (may be {@code null})
* @param configDir {@code CLAUDE_CONFIG_DIR} so the worker inherits the profile's
* skills/MCP/hooks (may be {@code null})
* @param tokenEnv name of the host env var holding the worker's auth token; its value
* is injected as {@code ANTHROPIC_AUTH_TOKEN} (never stored in config)
* @param argv launch command; defaults to {@code ["claude"]}
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Worker(String profile, String baseUrl, String model) {
public record Worker(String profile, String baseUrl, String model,
String configDir, String tokenEnv, List<String> argv) {
public Worker {
argv = (argv == null || argv.isEmpty()) ? List.of("claude") : List.copyOf(argv);
tokenEnv = (tokenEnv == null || tokenEnv.isBlank()) ? "BRIDGED_WORKER_TOKEN" : tokenEnv;
}
}
/**
@@ -0,0 +1,42 @@
package dev.ltms.bridged.herdr;
import com.fasterxml.jackson.databind.JsonNode;
/**
* A herdr-tracked agent (a Claude session, or any spawned command). Projected from the
* {@code agent} node returned by {@code agent.start}/{@code agent.get}/{@code agent.list}.
*
* @param terminalId herdr's stable handle — the {@code target} for send/read/get
* @param paneId pane handle — the argument to {@code pane.close}
* @param workspaceId owning workspace
* @param sessionId the agent's own session id (Claude's session UUID), or {@code null}
* before it has registered one (e.g. immediately after start)
* @param agentType agent kind, e.g. {@code "claude"} (the launch label for spawned probes)
* @param status current lifecycle state
*/
public record Agent(
String terminalId,
String paneId,
String workspaceId,
String sessionId,
String agentType,
AgentStatus status) {
/** Project a herdr {@code agent} node. Tolerates the start-time shape (no session yet). */
public static Agent from(JsonNode a) {
JsonNode session = a.get("agent_session");
String sessionId = session != null && session.hasNonNull("value")
? session.get("value").asText()
: null;
// start returns "name" (the launch label); list/get return "agent" (the kind).
String type = a.hasNonNull("agent") ? a.get("agent").asText()
: a.path("name").asText(null);
return new Agent(
a.path("terminal_id").asText(null),
a.path("pane_id").asText(null),
a.path("workspace_id").asText(null),
sessionId,
type,
AgentStatus.fromWire(a.path("agent_status").asText(null)));
}
}
@@ -0,0 +1,81 @@
package dev.ltms.bridged.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
/**
* Domain layer over herdr's native {@code agent.*} namespace — the worker south side.
* Chosen in the CB-102 spike over the pane + {@code send_text} fallback because
* {@code agent.start} takes a first-class {@code env} map (clean, guard-checked
* subscription injection) and herdr tracks each worker's Claude session UUID itself.
*
* <p>Every method is one herdr call through the injected {@link HerdrClient}, so this
* layer is unit-testable with a fake and contract-tested against a live daemon.
*/
public final class AgentControl {
private final HerdrClient herdr;
public AgentControl(HerdrClient herdr) {
this.herdr = herdr;
}
/**
* Spawn an agent. {@code env} is applied to the process environment verbatim — this
* is where a worker's {@code ANTHROPIC_BASE_URL} lives, and the ONLY place it should.
*
* @param name label/kind for herdr status detection (e.g. {@code "claude"})
* @param argv launch command, e.g. {@code ["claude"]}
* @param env process environment additions ({@code ANTHROPIC_BASE_URL}, token, …)
*/
public Agent start(String name, List<String> argv, Map<String, String> env) {
JsonNode result = herdr.call("agent.start", Map.of(
"name", name,
"argv", argv,
"env", env));
return Agent.from(result.get("agent"));
}
/** Deliver {@code text} to an agent (its next prompt input). */
public void send(String target, String text) {
herdr.call("agent.send", Map.of("target", target, "text", text));
}
/**
* Read an agent's terminal.
*
* @param source one of {@code visible|recent|recent_unwrapped|detection}
*/
public String read(String target, String source) {
JsonNode result = herdr.call("agent.read", Map.of("target", target, "source", source));
return result.path("read").path("text").asText("");
}
/** Current agent record (status, session UUID, pane). */
public Agent get(String target) {
return Agent.from(herdr.call("agent.get", Map.of("target", target)).get("agent"));
}
/** Just the lifecycle status — what the status-gated injector checks before send. */
public AgentStatus status(String target) {
return get(target).status();
}
/** All agents herdr tracks — the discovery surface ("what workers exist"). */
public List<Agent> list() {
JsonNode result = herdr.call("agent.list");
List<Agent> out = new ArrayList<>();
for (JsonNode a : result.path("agents")) {
out.add(Agent.from(a));
}
return out;
}
/** Tear a worker down (there is no agent.stop — close its pane). */
public void close(String paneId) {
herdr.call("pane.close", Map.of("pane_id", paneId));
}
}
@@ -0,0 +1,29 @@
package dev.ltms.bridged.herdr;
/**
* A herdr agent's lifecycle state, as reported by {@code agent_status}. Drives the
* status-gated injector: a worker is safe to inject into only when {@link #IDLE} or
* {@link #BLOCKED}, never mid-turn ({@link #WORKING}).
*/
public enum AgentStatus {
IDLE,
WORKING,
BLOCKED,
UNKNOWN;
/** Map herdr's wire string ({@code idle|working|blocked|unknown}) to the enum. */
public static AgentStatus fromWire(String s) {
if (s == null) return UNKNOWN;
return switch (s.toLowerCase()) {
case "idle" -> IDLE;
case "working" -> WORKING;
case "blocked" -> BLOCKED;
default -> UNKNOWN;
};
}
/** Whether {@code bridged} may inject a message now without stepping on a live turn. */
public boolean injectable() {
return this == IDLE || this == BLOCKED;
}
}
@@ -1,30 +1,36 @@
package dev.ltms.bridged.rest;
import com.fasterxml.jackson.databind.JsonNode;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
import io.javalin.http.Context;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
/**
* The REST surface — {@code bridged}'s contract, and its testability seam. Every
* The REST surface — {@code bridged}'s contract and its testability seam. Every
* feature is reachable here without Claude or MCP in the loop, so each is an
* acceptance test against plain HTTP. MCP tools (later) are thin adapters over these
* same endpoints and are validated by parity, not by re-implementing behaviour.
*
* <p>Built from an injected {@link HerdrClient} so tests can supply a fake and run on
* an ephemeral port; {@code main} supplies the real Unix-socket client.
* <p>Built from injected collaborators so tests supply fakes and run on an ephemeral
* port; {@code main} supplies the real Unix-socket client and worker service.
*/
public final class BridgedApp {
private final HerdrClient herdr;
private final WorkerService workers;
public BridgedApp(HerdrClient herdr) {
public BridgedApp(HerdrClient herdr, WorkerService workers) {
this.herdr = herdr;
this.workers = workers;
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -32,6 +38,9 @@ public final class BridgedApp {
Javalin app = Javalin.create(cfg -> cfg.showJavalinBanner = false);
app.get("/healthz", this::healthz);
app.get("/sessions", this::sessions);
app.get("/agents", this::agents);
app.post("/workers", this::spawnWorker);
app.delete("/workers/{paneId}", this::stopWorker);
return app;
}
@@ -52,11 +61,7 @@ public final class BridgedApp {
}
}
/**
* Sessions view, derived from herdr {@code workspace.list}. Stage-1 maps one
* workspace → one session summary; later tickets enrich this with the primary/
* worker role and the subscription-guard verdict per pane.
*/
/** Sessions view derived from herdr {@code workspace.list} (one workspace → one row). */
private void sessions(Context ctx) {
JsonNode result = herdr.call("workspace.list");
List<Map<String, Object>> out = new ArrayList<>();
@@ -70,4 +75,37 @@ public final class BridgedApp {
}
ctx.status(200).json(Map.of("sessions", out));
}
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
private void agents(Context ctx) {
ctx.status(200).json(Map.of("agents", workers.list().stream().map(BridgedApp::view).toList()));
}
/** Spawn a guard-checked worker. 403 if the base_url would breach the subscription boundary. */
private void spawnWorker(Context ctx) {
try {
Agent worker = workers.spawn();
ctx.status(201).json(view(worker));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
}
}
/** Tear a worker down by pane id. */
private void stopWorker(Context ctx) {
workers.stop(ctx.pathParam("paneId"));
ctx.status(204);
}
/** Stable JSON projection of an agent (null-safe for the start-time shape). */
private static Map<String, Object> view(Agent a) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("terminalId", a.terminalId());
m.put("paneId", a.paneId());
m.put("workspaceId", a.workspaceId());
m.put("sessionId", a.sessionId());
m.put("agentType", a.agentType());
m.put("status", a.status().name().toLowerCase());
return m;
}
}
@@ -0,0 +1,77 @@
package dev.ltms.bridged.worker;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
/**
* Spawns and lists worker sessions — the safe path from a delegation request to a
* running off-subscription Claude.
*
* <p>The spawn sequence encodes the subscription boundary: build the worker env with
* {@code ANTHROPIC_BASE_URL}, assert that host is on the allowlist <em>before</em>
* touching herdr, and only then {@code agent.start}. A worker's base_url lives in the
* env map handed to herdr and nowhere else; {@code bridged}'s own environment is never
* mutated.
*/
public final class WorkerService {
private static final Logger log = LoggerFactory.getLogger(WorkerService.class);
private final AgentControl agents;
private final SubscriptionGuard guard;
private final BridgedConfig.Worker cfg;
private final Function<String, String> env; // host env lookup (injectable for tests)
public WorkerService(AgentControl agents, SubscriptionGuard guard,
BridgedConfig.Worker cfg, Function<String, String> env) {
this.agents = agents;
this.guard = guard;
this.cfg = cfg;
this.env = env;
}
/** Spawn a worker for the configured profile. Guard runs before any herdr call. */
public Agent spawn() {
String baseUrl = cfg.baseUrl();
guard.assertWorker(baseUrl); // hard stop before we spawn anything
Map<String, String> workerEnv = new LinkedHashMap<>();
workerEnv.put("ANTHROPIC_BASE_URL", baseUrl);
putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model());
putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir());
String token = env.apply(cfg.tokenEnv());
putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", token);
String name = "claude";
List<String> argv = cfg.argv();
log.info("spawning worker profile={} base_url={} argv={}", cfg.profile(), baseUrl, argv);
Agent worker = agents.start(name, argv, workerEnv);
log.info("worker started pane={} terminal={}", worker.paneId(), worker.terminalId());
return worker;
}
/** All herdr-tracked agents — discovery for "what workers exist". */
public List<Agent> list() {
return agents.list();
}
/** Tear a worker down by pane id. */
public void stop(String paneId) {
agents.close(paneId);
}
private static void putIfPresent(Map<String, String> m, String k, String v) {
if (v != null && !v.isBlank()) {
m.put(k, v);
}
}
}