Stage-1 walking skeleton: bridged Java daemon (herdr client + REST + guard)

Maven/Java 25 module under bridged/. End-to-end verified against live
herdr 0.7.0 (protocol 14): GET /healthz and GET /sessions serve real
workspace data through the socket client.

- herdr client (CB-101): Unix-socket JSON-RPC via UnixDomainSocketAddress.
  Two contract facts pinned by tests against the real daemon:
  ids MUST be strings, and herdr is one-shot per connection
  (connection-per-call, which also makes the client lock-free).
- subscription guard: worker base_url must be on the off-subscription
  allowlist (gx00.gw, ollama.ltms.dev); primary env must carry no base_url.
- config (CB-106): Jackson YAML + Logback; example grounded in ltms-local.
- REST app (CB-104 start): injectable HerdrClient so acceptance tests run
  on an ephemeral port with a fake herdr, no daemon/Claude in the loop.
- tests: 17 unit/acceptance (mvn test) + 3 contract (mvn test -Pcontract).
This commit is contained in:
Dai Ha
2026-07-12 20:00:02 +02:00
parent b415003245
commit 93cfdb640f
18 changed files with 1035 additions and 0 deletions
@@ -0,0 +1,45 @@
package dev.ltms.bridged;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.rest.BridgedApp;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.file.Path;
/**
* {@code bridged} entry point. Wires the real herdr socket client to the REST app and
* starts listening. Before anything else it asserts its own environment is clean —
* {@code bridged} is not a Claude process and must never carry a base_url.
*/
public final class Bridged {
private static final Logger log = LoggerFactory.getLogger(Bridged.class);
public static void main(String[] args) {
Path configPath = Path.of(args.length > 0 ? args[0] : "bridged.yaml");
BridgedConfig cfg = BridgedConfig.load(configPath);
// The primary/host env that launched bridged must not be tainted.
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
guard.assertPrimaryClean(System.getenv());
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
? Path.of(cfg.herdrSocket())
: UnixSocketHerdrClient.defaultSocketPath();
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
Runtime.getRuntime().addShutdownHook(new Thread(herdr::close));
Javalin app = new BridgedApp(herdr).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
}
private Bridged() {
}
}
@@ -0,0 +1,83 @@
package dev.ltms.bridged.config;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Set;
/**
* {@code bridged} configuration, loaded from a YAML file (see
* {@code bridged.example.yaml}). Unknown keys are ignored so config can grow ahead
* of the code.
*
* @param bind REST/MCP listen host:port
* @param herdrSocket path to herdr's Unix socket ({@code null} → client default)
* @param worker worker-spawn settings
* @param guard subscription-boundary allowlist
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record BridgedConfig(
Bind bind,
String herdrSocket,
Worker worker,
Guard guard) {
@JsonIgnoreProperties(ignoreUnknown = true)
public record Bind(String host, int port) {
public Bind {
if (host == null || host.isBlank()) host = "127.0.0.1";
if (port <= 0) port = 8080;
}
}
/**
* @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
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Worker(String profile, String baseUrl, String model) {
}
/**
* Subscription boundary. Only these hosts may back a worker's
* {@code ANTHROPIC_BASE_URL}; the primary must carry none.
*
* @param offSubscriptionHosts hostnames allowed for worker base_urls
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Guard(List<String> offSubscriptionHosts) {
public Guard {
offSubscriptionHosts = offSubscriptionHosts == null ? List.of() : List.copyOf(offSubscriptionHosts);
}
public Set<String> hostSet() {
return Set.copyOf(offSubscriptionHosts);
}
}
private static final ObjectMapper YAML = new ObjectMapper(new YAMLFactory());
/** Load and validate config from {@code path}. */
public static BridgedConfig load(Path path) {
try {
BridgedConfig cfg = YAML.readValue(Files.readString(path), BridgedConfig.class);
return cfg.withDefaults();
} catch (IOException e) {
throw new UncheckedIOException("cannot read bridged config at " + path, e);
}
}
/** Fill in nested defaults so callers never see nulls for structural fields. */
public BridgedConfig withDefaults() {
Bind b = bind != null ? bind : new Bind(null, 0);
Guard g = guard != null ? guard : new Guard(List.of());
return new BridgedConfig(b, herdrSocket, worker, g);
}
}
@@ -0,0 +1,8 @@
package dev.ltms.bridged.guard;
/** Thrown when the subscription boundary would be violated. Never swallow this. */
public class GuardException extends RuntimeException {
public GuardException(String message) {
super(message);
}
}
@@ -0,0 +1,64 @@
package dev.ltms.bridged.guard;
import java.net.URI;
import java.net.URISyntaxException;
import java.util.Map;
import java.util.Set;
/**
* Enforces the non-negotiable subscription boundary:
*
* <ul>
* <li>A <b>worker</b> must egress to an off-subscription host — its
* {@code ANTHROPIC_BASE_URL} host must be on the allowlist
* (Stage-1: {@code gx00.gw}, {@code ollama.ltms.dev}).</li>
* <li>The <b>primary</b> must never carry {@code ANTHROPIC_BASE_URL}; a set value
* means its traffic would leave the subscription. That is a hard stop.</li>
* </ul>
*
* Both checks throw {@link GuardException} on violation. {@code bridged} calls
* {@link #assertWorker} before spawning a worker and {@link #assertPrimaryClean}
* against its own environment at startup.
*/
public final class SubscriptionGuard {
private final Set<String> offSubscriptionHosts;
public SubscriptionGuard(Set<String> offSubscriptionHosts) {
this.offSubscriptionHosts = Set.copyOf(offSubscriptionHosts);
}
/** The worker's base_url must resolve to an allowlisted off-subscription host. */
public void assertWorker(String baseUrl) {
if (baseUrl == null || baseUrl.isBlank()) {
throw new GuardException("worker has no ANTHROPIC_BASE_URL — refusing to spawn "
+ "a session that would bill the subscription");
}
String host = hostOf(baseUrl);
if (host == null) {
throw new GuardException("worker ANTHROPIC_BASE_URL is not a valid URL: " + baseUrl);
}
if (!offSubscriptionHosts.contains(host)) {
throw new GuardException("worker base_url host '" + host + "' is not on the "
+ "off-subscription allowlist " + offSubscriptionHosts
+ " — refusing to spawn");
}
}
/** The primary's environment must not contain {@code ANTHROPIC_BASE_URL}. */
public void assertPrimaryClean(Map<String, String> env) {
String v = env.get("ANTHROPIC_BASE_URL");
if (v != null && !v.isBlank()) {
throw new GuardException("primary environment is tainted: ANTHROPIC_BASE_URL="
+ v + " — the primary must run on the subscription, never a base_url");
}
}
private static String hostOf(String baseUrl) {
try {
return new URI(baseUrl).getHost();
} catch (URISyntaxException e) {
return null;
}
}
}
@@ -0,0 +1,38 @@
package dev.ltms.bridged.herdr;
import com.fasterxml.jackson.databind.JsonNode;
/**
* Client face onto the herdr daemon (protocol 14, herdr 0.7.0).
*
* <p>This is the ONLY thing in {@code bridged} that speaks to herdr. Every method
* maps to a herdr JSON-RPC call over its Unix domain socket. Requests are
* newline-delimited JSON with a <em>string</em> id; responses carry either a
* {@code result} object (whose {@code type} field discriminates the payload) or an
* {@code error} object.
*
* <p>Higher layers ({@code bridged}'s policy brain, REST endpoints, MCP adapters)
* depend on this interface, not on the socket. Tests substitute a fake; the
* {@code contract}-tagged suite exercises the real implementation against a running
* herdr to catch protocol drift.
*/
public interface HerdrClient extends AutoCloseable {
/**
* Invoke a herdr method and return its {@code result} node.
*
* @param method herdr method name, e.g. {@code "ping"}, {@code "workspace.list"}
* @param params params object (may be {@code null} → sent as {@code {}}); serialized by Jackson
* @return the {@code result} node of the response
* @throws HerdrException on transport failure or an {@code error} envelope
*/
JsonNode call(String method, Object params) throws HerdrException;
/** Convenience for parameterless calls. */
default JsonNode call(String method) throws HerdrException {
return call(method, null);
}
@Override
void close();
}
@@ -0,0 +1,70 @@
package dev.ltms.bridged.herdr;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import java.nio.charset.StandardCharsets;
/**
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 14).
*
* <p>Split out from the socket so the framing rules — the ones that actually bit us
* during the spike (id MUST be a string; response carries {@code result} or
* {@code error}, never {@code jsonrpc}) — are unit-testable without a live daemon.
*/
final class HerdrCodec {
private final ObjectMapper mapper;
HerdrCodec(ObjectMapper mapper) {
this.mapper = mapper;
}
/** Build one request frame: a single JSON object terminated by {@code '\n'}. */
byte[] encode(String id, String method, Object params) {
ObjectNode req = mapper.createObjectNode();
req.put("jsonrpc", "2.0");
req.put("id", id); // string id — herdr rejects integer ids with invalid_request
req.put("method", method);
JsonNode p = params == null ? mapper.createObjectNode() : mapper.valueToTree(params);
req.set("params", p);
try {
String line = mapper.writeValueAsString(req) + "\n";
return line.getBytes(StandardCharsets.UTF_8);
} catch (JsonProcessingException e) {
throw new HerdrException("failed to encode herdr request for method " + method, e);
}
}
/**
* Parse one response frame and return its {@code result} node.
*
* @throws HerdrException if the frame is an {@code error} envelope or is malformed
*/
JsonNode decodeResult(String line) {
JsonNode root;
try {
root = mapper.readTree(line);
} catch (JsonProcessingException e) {
throw new HerdrException("malformed herdr response: " + trim(line), e);
}
JsonNode error = root.get("error");
if (error != null && !error.isNull()) {
String code = error.path("code").asText(null);
String message = error.path("message").asText("unknown herdr error");
throw new HerdrException("herdr error [" + code + "]: " + message, code, null);
}
JsonNode result = root.get("result");
if (result == null || result.isNull()) {
throw new HerdrException("herdr response has neither result nor error: " + trim(line));
}
return result;
}
private static String trim(String s) {
String t = s.strip();
return t.length() > 200 ? t.substring(0, 200) + "…" : t;
}
}
@@ -0,0 +1,30 @@
package dev.ltms.bridged.herdr;
/**
* Raised when a herdr call fails: transport error, or an {@code error} envelope
* returned by the daemon. {@link #code()} carries herdr's error code
* (e.g. {@code invalid_request}) when the failure came back as a protocol error,
* or {@code null} for transport-level failures.
*/
public class HerdrException extends RuntimeException {
private final String code;
public HerdrException(String message) {
this(message, null, null);
}
public HerdrException(String message, Throwable cause) {
this(message, null, cause);
}
public HerdrException(String message, String code, Throwable cause) {
super(message, cause);
this.code = code;
}
/** herdr protocol error code, or {@code null} if this was a transport failure. */
public String code() {
return code;
}
}
@@ -0,0 +1,122 @@
package dev.ltms.bridged.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.net.StandardProtocolFamily;
import java.net.UnixDomainSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.SocketChannel;
import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.util.concurrent.atomic.AtomicLong;
/**
* {@link HerdrClient} over herdr's Unix domain socket using the JDK's
* {@link UnixDomainSocketAddress} + {@link SocketChannel} (no native/JNI dependency).
*
* <p><b>Connection-per-call.</b> The contract test against herdr 0.7.0 established
* that herdr serves <em>one</em> request/response per connection and then closes it —
* a second write on the same socket gets a broken pipe. So every {@link #call} opens a
* fresh connection, writes one frame, reads one line, and closes. A pleasant
* consequence: with no shared socket there is no shared read state, so the client is
* safe to call concurrently from many virtual threads with no locking.
*
* <p>Streaming methods ({@code events.subscribe}) keep their own long-lived connection
* and land in a later ticket; they do not reuse this request/response path.
*/
public final class UnixSocketHerdrClient implements HerdrClient {
private static final Logger log = LoggerFactory.getLogger(UnixSocketHerdrClient.class);
/** herdr's default socket, overridable by {@code HERDR_SOCKET_PATH}. */
public static Path defaultSocketPath() {
String override = System.getenv("HERDR_SOCKET_PATH");
if (override != null && !override.isBlank()) {
return Path.of(override);
}
return Path.of(System.getProperty("user.home"), ".config", "herdr", "herdr.sock");
}
private final Path socketPath;
private final HerdrCodec codec;
private final AtomicLong ids = new AtomicLong(1);
private UnixSocketHerdrClient(Path socketPath, ObjectMapper mapper) {
this.socketPath = socketPath;
this.codec = new HerdrCodec(mapper);
}
/** Client for the default socket with a fresh {@link ObjectMapper}. */
public static UnixSocketHerdrClient connect() {
return connect(defaultSocketPath(), new ObjectMapper());
}
/**
* Client for {@code socketPath}. Does not hold a connection open (herdr is
* one-shot per connection); connectivity surfaces on the first {@link #call}.
*/
public static UnixSocketHerdrClient connect(Path socketPath, ObjectMapper mapper) {
log.debug("herdr client bound to socket {}", socketPath);
return new UnixSocketHerdrClient(socketPath, mapper);
}
@Override
public JsonNode call(String method, Object params) {
String id = Long.toString(ids.getAndIncrement());
byte[] frame = codec.encode(id, method, params);
try (SocketChannel ch = SocketChannel.open(StandardProtocolFamily.UNIX)) {
ch.connect(UnixDomainSocketAddress.of(socketPath));
writeFully(ch, ByteBuffer.wrap(frame));
return codec.decodeResult(readLine(ch));
} catch (IOException e) {
throw new HerdrException("herdr call '" + method + "' failed at transport (socket "
+ socketPath + ", is herdr running?)", e);
}
}
private static void writeFully(SocketChannel ch, ByteBuffer buf) throws IOException {
while (buf.hasRemaining()) {
ch.write(buf);
}
}
/** Read up to and including the first {@code '\n'}, returning the line without it. */
private static String readLine(SocketChannel ch) throws IOException {
ByteBuffer readBuf = ByteBuffer.allocate(64 * 1024);
StringBuilder sb = new StringBuilder();
while (true) {
int nl = indexOfNewline(sb);
if (nl >= 0) {
return sb.substring(0, nl);
}
readBuf.clear();
int n = ch.read(readBuf);
if (n == -1) {
if (sb.length() > 0) {
return sb.toString(); // herdr closed after a complete, unterminated frame
}
throw new IOException("herdr closed the connection with no response");
}
readBuf.flip();
sb.append(StandardCharsets.UTF_8.decode(readBuf));
}
}
private static int indexOfNewline(CharSequence s) {
for (int i = 0; i < s.length(); i++) {
if (s.charAt(i) == '\n') {
return i;
}
}
return -1;
}
/** No persistent resources to release; present for the {@link AutoCloseable} contract. */
@Override
public void close() {
}
}
@@ -0,0 +1,73 @@
package dev.ltms.bridged.rest;
import com.fasterxml.jackson.databind.JsonNode;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import io.javalin.Javalin;
import io.javalin.http.Context;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
/**
* 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.
*/
public final class BridgedApp {
private final HerdrClient herdr;
public BridgedApp(HerdrClient herdr) {
this.herdr = herdr;
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
public Javalin build() {
Javalin app = Javalin.create(cfg -> cfg.showJavalinBanner = false);
app.get("/healthz", this::healthz);
app.get("/sessions", this::sessions);
return app;
}
/** Liveness + herdr reachability. 200 when herdr answers ping, 503 otherwise. */
private void healthz(Context ctx) {
try {
JsonNode pong = herdr.call("ping");
ctx.status(200).json(Map.of(
"status", "ok",
"herdr", Map.of(
"version", pong.path("version").asText(""),
"protocol", pong.path("protocol").asInt())));
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
"status", "degraded",
"herdr", "unreachable",
"detail", e.getMessage()));
}
}
/**
* 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.
*/
private void sessions(Context ctx) {
JsonNode result = herdr.call("workspace.list");
List<Map<String, Object>> out = new ArrayList<>();
for (JsonNode w : result.path("workspaces")) {
out.add(Map.of(
"id", w.path("workspace_id").asText(""),
"label", w.path("label").asText(""),
"focused", w.path("focused").asBoolean(false),
"paneCount", w.path("pane_count").asInt(),
"agentStatus", w.path("agent_status").asText("unknown")));
}
ctx.status(200).json(Map.of("sessions", out));
}
}