Files
fleetd/bridged/src/main/java/dev/ltms/bridged/herdr/UnixSocketHerdrClient.java
T
Dai Ha 93cfdb640f 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).
2026-07-12 20:00:02 +02:00

123 lines
4.7 KiB
Java

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() {
}
}