CB-5xx: Stage 5 hardening — auth, authz+audit, metrics, CI, supervision

Closes out single-host before the cross-host work. Sequenced BEFORE CB-308
deliberately: federation's own gating concern is the trust model, and it
inherits whatever identity shape lands here.

The finding this stage is built around: bridged had exactly ONE security
control, the loopback bind. ConnectionIdentity resolves a worker from its
connection (unforgeable), but every caller that was not a recognised worker
pane fell through to being treated as the PRIMARY -- the most privileged role
on the bus. Latent today; load-bearing the moment a bind widens.

CB-501 auth:
- Role/Principal/CallerResolver: connection identity first, bearer token
  second, ANONYMOUS third. Inverts the old default so absence of identity
  means nothing, not everything.
- Worker identity is never token-gated, so enabling auth cannot lock the
  fleet out of bridge_reply.
- Constant-time token compare (MessageDigest.isEqual).
- validateAuthExposure(): a non-loopback bind under loopback-trust now
  REFUSES TO START. Makes the dangerous config unrepresentable rather than
  merely documented.
- TLS terminates at a reverse proxy by design (D3), not in the JVM.

CB-505 authz + audit, enforced on BOTH entry paths:
- The docs describe MCP as "a thin adapter over the REST core"; at code level
  it is not. BridgeMcp calls MessageService directly, and /mcp is a raw
  servlet on Jetty's context handler that never traverses Javalin's before
  filter. Enforcing only at REST would have left /mcp open.
- Load-bearing rule is own-session-only: a worker may reply/ask only as
  itself. Structurally true over MCP already; over REST the session id in the
  URL path had simply been trusted.
- Audit: JSON lines to a dedicated appender, additivity=false. Never records
  message content -- this bus carries source and prompts.

CB-502 metrics: zero new dependencies. A ~150-line Prometheus text renderer
instead of the specced Micrometer, because this pom already hand-pins
jackson-annotations to reconcile Jackson 2/3, imports a Jetty BOM against
skew, and carries four accepted-CVE advisories -- and CLAUDE.md's mandated
dependency CVE gate could not be run (no JetBrains MCP server connected).
Instrumented at MessageService, the single funnel both surfaces share.

CB-503 CI: .gitea/workflows/ci.yml against the already-running Gitea runner.
Needs no contract-exclusion flag -- the pom's default-excludes profile
already sets excludedGroups=contract, so plain `mvn clean install` IS the
mock-socket surface. Provisions JDK 25 explicitly (runner default-jdk is older).

CB-504 supervision: launchd agent (the real target -- this host is macOS,
there is no systemd) plus a systemd unit for the Linux gateways CB-308 adds.
Ordering directives are advisory, so the actual fix is that startup now waits
up to 30s for the herdr socket and then serves degraded, instead of crashing
into a restart loop on a boot-order race.

Also fixes drift found while surveying:
- bridged.example.yaml documented spawn_ready_timeout_ms in snake_case; config
  binds via plain Jackson with ignoreUnknown, so uncommenting it would have
  been silently dropped and the default kept. Now camelCase, with a test that
  loads the shipped example and one that pins every documented knob's
  spelling -- no test had ever loaded that file.
- Added the 6 shipped-but-undocumented knobs (worktreeRoot, parityOverlay,
  gitTokenEnv, gitHostEnv, configDir, primary:).
- README "Next" listed bridge_ask and session lifecycle as upcoming; both
  shipped long ago.
- docs/CB-301-ext and docs/CB-402 status headers said "design"/"pre-
  implementation" for work already merged.

307 unit/acceptance tests green (was 266), mvn clean install BUILD SUCCESS.
Note: CLAUDE.md's per-file ide_diagnostics gate and the pom Mend.io CVE check
could not be run -- no JetBrains/intellij-index MCP server is connected this
session. mvn clean install is the only gate that ran.
This commit is contained in:
2026-07-29 22:29:26 +07:00
parent c9f0ca9359
commit 9daf1ec5ba
28 changed files with 2298 additions and 37 deletions
@@ -3,6 +3,8 @@ 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.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.herdr.PaneLocator;
import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.herdr.WorkspaceControl;
@@ -11,8 +13,11 @@ import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.TurnListener;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.mcp.BridgeMcp;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.mcp.LsofPeerPidLookup;
import dev.ltms.bridged.mcp.LsofProcessCwdLookup;
@@ -54,6 +59,10 @@ public final class Bridged {
/** How often the injector samples a busy worker's status while it has queued work. */
private static final long INJECT_POLL_MILLIS = 250;
/** CB-504: how long to wait at startup for herdr's socket before serving degraded. */
private static final long HERDR_WAIT_SECONDS = 30;
private static final long HERDR_WAIT_POLL_MILLIS = 500;
static void main(String[] args) {
Path configPath = Path.of(args.length > 0 ? args[0] : "bridged.yaml");
BridgedConfig cfg = BridgedConfig.load(configPath);
@@ -62,6 +71,12 @@ public final class Bridged {
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
guard.assertPrimaryClean(System.getenv());
// CB-501: refuse to start if the bind is wider than the auth mode can defend. Under
// loopback-trust, "not a known worker" means "the primary" — sound only because the OS
// refuses remote connections to a loopback socket. This throws rather than warns so the
// dangerous configuration cannot be reached by ignoring a log line.
cfg.validateAuthExposure();
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
? Path.of(cfg.herdrSocket())
: UnixSocketHerdrClient.defaultSocketPath();
@@ -97,9 +112,20 @@ public final class Bridged {
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs()));
}
PeerLauncher workers = new CompositePeerLauncher(adapters, cfg.defaultProfile());
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died with
// the previous process — reap those leaked orphans now, before we start serving.
workers.reapOrphanWorkers();
// CB-504: under supervision (launchd/systemd) bridged can start before herdr's socket
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
// the first thing that actually talks to herdr, so without this wait a boot-order race
// would crash the daemon into a restart loop. Wait, then degrade rather than die: serving
// with /healthz reporting "degraded" is strictly more useful than exiting.
if (awaitHerdr(herdr)) {
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
// with the previous process — reap those leaked orphans now, before we start serving.
workers.reapOrphanWorkers();
} else {
log.warn("herdr did not answer within {}s — starting anyway; /healthz will report "
+ "degraded until it comes up. Orphaned worker panes (if any) were NOT reaped.",
HERDR_WAIT_SECONDS);
}
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
@@ -176,14 +202,36 @@ public final class Bridged {
Thread.ofVirtual().name("bridge-push-").unstarted(r));
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
pushScheduler, maxReminders, backoffMs);
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, pushLoop);
// CB-502: the registry is built before the service so send/reply outcomes are counted at
// their single funnel rather than at each of the two caller-facing surfaces.
Metrics metrics = BridgedMetrics.create(sessions, replyInbox);
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
final CallerResolver callers;
if (cfg.auth().tokenMode()) {
String token = System.getenv(cfg.auth().tokenEnv());
if (token == null || token.isBlank()) {
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
+ " is unset or empty — export it before starting bridged");
}
callers = new CallerResolver(identity, true, token);
log.info("auth: token mode (bearer required for non-worker callers, env {})",
cfg.auth().tokenEnv());
} else {
callers = new CallerResolver(identity);
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
BridgeMcp mcp = new BridgeMcp(messages, workers, sessions, identity, presence,
primaryRegistry);
primaryRegistry, callers, metrics);
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
// open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr
@@ -206,12 +254,46 @@ public final class Bridged {
herdr.close();
}));
Javalin app = new BridgedApp(herdr, workers, sessions, messages, presence, mcp.servlet()).build();
Javalin app = new BridgedApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
}
/**
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
*
* @return true if herdr answered, false if it never did
*/
private static boolean awaitHerdr(HerdrClient herdr) {
long deadline = System.nanoTime() + HERDR_WAIT_SECONDS * 1_000_000_000L;
boolean waited = false;
while (true) {
try {
herdr.call("ping");
if (waited) {
log.info("herdr is up");
}
return true;
} catch (HerdrException e) {
if (System.nanoTime() >= deadline) {
return false;
}
if (!waited) {
log.info("waiting up to {}s for the herdr socket…", HERDR_WAIT_SECONDS);
waited = true;
}
try {
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
return false;
}
}
}
}
private Bridged() {
}
}
@@ -0,0 +1,63 @@
package dev.ltms.bridged.auth;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* Append-only record of privileged actions (CB-505).
*
* <p>Writes JSON lines to a dedicated {@code audit} logger — its own appender, separate from the
* chatty app log — so the trail stays greppable and can later be shipped without dragging debug
* noise along.
*
* <p><strong>Message content is never recorded.</strong> This bridge carries the user's source
* code, diffs, and prompts; an audit trail that quietly accumulated them would be a transcript
* archive wearing a security control's clothing. Records carry <em>who / what / against what /
* outcome</em> and correlation ids only.
*/
public final class AuditLog {
private static final Logger AUDIT = LoggerFactory.getLogger("audit");
private AuditLog() {
}
/** Record an allowed action. */
public static void allowed(Principal caller, Authz.Action action, String target) {
write(caller, action, target, "allowed", null);
}
/** Record a refused action and why. */
public static void denied(Principal caller, Authz.Action action, String target, String reason) {
write(caller, action, target, "denied", reason);
}
/** Record an action that was authorized but then failed downstream (guard, timeout, herdr). */
public static void failed(Principal caller, Authz.Action action, String target, String reason) {
write(caller, action, target, "failed", reason);
}
private static void write(Principal caller, Authz.Action action, String target,
String outcome, String reason) {
Principal c = caller != null ? caller : Principal.anonymous();
StringBuilder sb = new StringBuilder(160);
sb.append("{\"role\":\"").append(c.role()).append('"')
.append(",\"actor\":\"").append(esc(c.describe())).append('"')
.append(",\"pid\":").append(c.pid())
.append(",\"action\":\"").append(action).append('"')
.append(",\"target\":").append(target == null ? "null" : '"' + esc(target) + '"')
.append(",\"outcome\":\"").append(outcome).append('"');
if (reason != null) {
sb.append(",\"reason\":\"").append(esc(reason)).append('"');
}
sb.append('}');
// The appender supplies the timestamp, so it cannot disagree with the app log's clock.
AUDIT.info(sb.toString());
}
/** Minimal JSON string escaping — these values are ids and short reasons, never free text. */
private static String esc(String s) {
return s.replace("\\", "\\\\").replace("\"", "\\\"")
.replace("\n", "\\n").replace("\r", "\\r").replace("\t", "\\t");
}
}
@@ -0,0 +1,72 @@
package dev.ltms.bridged.auth;
/**
* The authorization table (CB-505), stated once and enforced on both entry paths.
*
* <p>Most of these rules are already true de facto — {@code BridgeMcp} derives a worker's identity
* from the connection rather than reading it from an argument, so a worker has never been able to
* reply <em>as</em> another worker over MCP. What was missing is that the REST surface trusted the
* session id in the URL path, and neither surface checked role at all. This class makes the
* invariant explicit and testable rather than emergent.
*/
public final class Authz {
private Authz() {
}
/** A privileged operation, named for the audit trail. */
public enum Action {
/** Spawn a worker peer. */
SPAWN,
/** Tear a worker peer down. */
STOP,
/** Deliver a turn to a session (or answer a worker's question). */
SEND,
/** A worker's terminal reply for its own turn. */
REPLY,
/** A worker's mid-turn question to the primary. */
ASK,
/** Collect held replies from a session's inbox. */
DRAIN,
/** Read-only observation: status, roster, profiles, task polling. */
READ,
/** Scrape the metrics endpoint. */
METRICS
}
/**
* Whether {@code caller} may perform {@code action} against {@code targetSession}.
*
* @param targetSession the session id in the request path; only consulted for the worker-scoped
* actions ({@code REPLY}, {@code ASK}), ignored otherwise, may be
* {@code null}
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
if (caller == null || caller.isAnonymous()) {
return false; // authenticated as nothing ⇒ authorized for nothing
}
return switch (action) {
// Orchestration is the primary's alone. A worker driving spawn/stop/send would be a
// worker escalating into the orchestrator role.
case SPAWN, STOP, SEND, DRAIN -> caller.isPrimary();
// The load-bearing rule: a worker acts only as itself. The primary is deliberately
// excluded — a reply/ask is a worker's own turn output, and letting the primary forge
// one would corrupt the rendezvous correlation it is itself waiting on.
case REPLY, ASK -> caller.ownsSession(targetSession);
// Observation is open to both authenticated roles: a worker legitimately polls its own
// status, and the roster carries no secrets.
case READ, METRICS -> caller.isPrimary() || caller.isWorker();
};
}
/**
* Why a request was refused, for the error body. Distinguishes "you are nobody" from "you are
* somebody, but not the right somebody" — the first is a credential problem (401), the second
* an authorization one (403).
*/
public static boolean isUnauthenticated(Principal caller) {
return caller == null || caller.isAnonymous();
}
}
@@ -0,0 +1,119 @@
package dev.ltms.bridged.auth;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
/**
* Resolves every caller to a {@link Principal}, for both entry paths into the core (CB-501).
*
* <p>There are two of them and they are not layered the way the docs suggest: {@code BridgeMcp}
* calls the service layer directly and is mounted as a raw servlet (so it never passes through a
* Javalin filter), while the REST routes historically resolved no identity at all. Both now
* delegate here, so the authorization rules are stated once instead of drifting apart.
*
* <p><strong>Resolution order</strong> — connection identity first, token second, nothing third:
* <ol>
* <li>A loopback peer PID that maps to a herdr worker pane ⇒ {@link Role#WORKER}. This is
* unforgeable (the OS reports the PID, herdr owns the PID→pane map) and is honoured
* regardless of auth mode, so enabling auth never breaks the fleet.</li>
* <li>Otherwise, under {@code token} mode, a valid bearer token ⇒ {@link Role#PRIMARY}.</li>
* <li>Otherwise, under {@code loopback-trust}, a loopback caller ⇒ {@link Role#PRIMARY}
* (the historical behaviour, now an explicit configured choice).</li>
* <li>Otherwise {@link Role#ANONYMOUS}.</li>
* </ol>
*/
public final class CallerResolver {
private final ConnectionIdentity identity;
private final boolean tokenMode;
private final byte[] expectedToken; // null unless tokenMode
/** Loopback-trust resolver: no token required, historical behaviour. */
public CallerResolver(ConnectionIdentity identity) {
this(identity, false, null);
}
/**
* @param identity connection-based worker identification
* @param tokenMode when true, a non-worker caller must present a valid bearer token
* @param token the expected bearer token; required (non-blank) when {@code tokenMode}
*/
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token) {
if (tokenMode && (token == null || token.isBlank())) {
throw new IllegalArgumentException(
"auth.mode=token requires a non-empty token; check that the env var named by "
+ "auth.tokenEnv is exported to the daemon's environment");
}
this.identity = identity;
this.tokenMode = tokenMode;
this.expectedToken = tokenMode ? token.getBytes(StandardCharsets.UTF_8) : null;
}
/**
* Resolve the caller of a request.
*
* @param remoteAddr the connection's remote address
* @param remotePort the connection's remote port (used for the peer-PID lookup)
* @param authorizationHeader the raw {@code Authorization} header, or {@code null}
*/
public Principal resolve(String remoteAddr, int remotePort, String authorizationHeader) {
ConnectionIdentity.Caller c = identity.resolve(remoteAddr, remotePort);
if (c.terminal() != null) {
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
}
if (tokenMode) {
return presentedTokenMatches(authorizationHeader)
? Principal.primary(c.pid())
: Principal.anonymous();
}
// loopback-trust: same-host callers that are not workers are the primary. A non-loopback
// caller is anonymous even here — and startup refuses that combination anyway
// (BridgedConfig.validateAuthExposure), so this is defence in depth, not the control.
return isLoopback(remoteAddr) ? Principal.primary(c.pid()) : Principal.anonymous();
}
/** The working directory of the calling process (CB-112 spawn cwd inheritance), or {@code null}. */
public String cwdForPid(long pid) {
return identity.cwdForPid(pid);
}
/** True when auth requires a bearer token of non-worker callers. */
public boolean tokenMode() {
return tokenMode;
}
private boolean presentedTokenMatches(String authorizationHeader) {
String presented = bearerValue(authorizationHeader);
if (presented == null) {
return false;
}
// Constant-time: MessageDigest.isEqual does not short-circuit on the first differing byte,
// so a token cannot be recovered a byte at a time by timing the response.
return MessageDigest.isEqual(presented.getBytes(StandardCharsets.UTF_8), expectedToken);
}
/** Extract the credential from {@code Authorization: Bearer <token>}, or {@code null}. */
private static String bearerValue(String header) {
if (header == null) {
return null;
}
String h = header.trim();
if (h.length() < 7 || !h.regionMatches(true, 0, "Bearer ", 0, 7)) {
return null;
}
String token = h.substring(7).trim();
return token.isEmpty() ? null : token;
}
private static boolean isLoopback(String remoteAddr) {
if (remoteAddr == null) {
return false;
}
return remoteAddr.equals("127.0.0.1") || remoteAddr.equals("::1")
|| remoteAddr.equals("0:0:0:0:0:0:0:1") || remoteAddr.startsWith("127.");
}
}
@@ -0,0 +1,57 @@
package dev.ltms.bridged.auth;
/**
* A resolved caller: its {@link Role}, and — for a worker — the herdr {@code terminal_id} that
* identifies which worker it is (CB-501).
*
* @param role what this caller is authorized to act as
* @param terminal the worker's herdr terminal id; {@code null} for {@code PRIMARY}/{@code ANONYMOUS}
* @param pid the connecting process id, or {@code -1} when not resolvable (audit context)
*/
public record Principal(Role role, String terminal, long pid) {
/** A caller authenticated as nothing — the default when no check establishes anything else. */
public static Principal anonymous() {
return new Principal(Role.ANONYMOUS, null, -1);
}
/** The orchestrating session. */
public static Principal primary(long pid) {
return new Principal(Role.PRIMARY, null, pid);
}
/** A worker peer, identified by its herdr pane. */
public static Principal worker(String terminal, long pid) {
return new Principal(Role.WORKER, terminal, pid);
}
public boolean isPrimary() {
return role == Role.PRIMARY;
}
public boolean isWorker() {
return role == Role.WORKER;
}
public boolean isAnonymous() {
return role == Role.ANONYMOUS;
}
/**
* Whether this caller may act <em>as</em> {@code sessionId} — the "own session only" rule that
* keeps one worker from replying or asking on another's behalf. Only a worker can own a
* session, and only its own.
*/
public boolean ownsSession(String sessionId) {
return isWorker() && terminal != null && terminal.equals(sessionId);
}
/** Short, non-sensitive description for audit lines and error details. */
public String describe() {
return switch (role) {
case WORKER -> "worker:" + terminal;
case PRIMARY -> "primary";
case ANONYMOUS -> "anonymous";
};
}
}
@@ -0,0 +1,29 @@
package dev.ltms.bridged.auth;
/**
* What a caller is allowed to be on the bus (CB-501).
*
* <p>The ordering matters conceptually: {@link #PRIMARY} is the <em>most</em> privileged role
* (it spawns, stops, sends to any session, and drains any inbox), not the least. Before CB-501
* the daemon reached {@code PRIMARY} by <em>failing</em> every other check — any caller that did
* not resolve to a known worker pane was treated as the primary. That is inverted here:
* {@link #ANONYMOUS} is the fallback, and {@code PRIMARY} must be established.
*/
public enum Role {
/**
* The orchestrating session. Established either by being a loopback caller that is not a
* worker pane (under {@code loopback-trust}) or by presenting a valid bearer token (under
* {@code token} mode).
*/
PRIMARY,
/**
* A worker peer, identified by its herdr pane. Unforgeable: derived from the connection's
* loopback peer PID via herdr's PID→pane map, never from a request argument.
*/
WORKER,
/** Authenticated as nothing. Authorized for nothing but {@code /healthz}. */
ANONYMOUS
}
@@ -36,6 +36,8 @@ import java.util.Set;
* @param primary optional pinned primary terminal config ({@code null} → derived from connection);
* a non-blank {@code terminal} seeds {@code PrimaryRegistry} and prevents
* connection-derived overrides, CB-307
* @param auth API authentication mode ({@code null} → {@code loopback-trust}, the
* historical behaviour), CB-501
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record BridgedConfig(
@@ -50,7 +52,8 @@ public record BridgedConfig(
Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs,
Broker broker,
Primary primary) {
Primary primary,
Auth auth) {
@JsonIgnoreProperties(ignoreUnknown = true)
public record Bind(String host, int port) {
@@ -256,6 +259,43 @@ public record BridgedConfig(
}
}
/**
* API authentication (CB-501). Governs how a caller that is <em>not</em> an on-host worker
* pane proves it is the primary.
*
* <p>Worker identity never depends on this block: a loopback peer PID that maps to a herdr
* pane is unforgeable and is always honoured (see
* {@link dev.ltms.bridged.mcp.ConnectionIdentity}). This only decides what happens for
* <em>everyone else</em>.
*
* @param mode {@code "loopback-trust"} (default) — any loopback caller that is not a known
* worker is the primary, no credential needed; this is the historical
* behaviour, now chosen explicitly rather than implied. {@code "token"} — such
* a caller must present {@code Authorization: Bearer <token>} or it is
* {@code ANONYMOUS} and authorized for nothing.
* @param tokenEnv name of the host env var holding the bearer token; the literal value is
* never stored in config. Defaults to {@code BRIDGED_API_TOKEN}. Only read
* when {@code mode} is {@code token}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Auth(String mode, String tokenEnv) {
/** Historical behaviour: loopback non-worker ⇒ primary, no credential. */
public static final String MODE_LOOPBACK_TRUST = "loopback-trust";
/** A non-worker caller must present a valid bearer token to be the primary. */
public static final String MODE_TOKEN = "token";
public Auth {
mode = (mode == null || mode.isBlank()) ? MODE_LOOPBACK_TRUST : mode.toLowerCase();
tokenEnv = (tokenEnv == null || tokenEnv.isBlank()) ? "BRIDGED_API_TOKEN" : tokenEnv;
}
/** True when a bearer token is required of every non-worker caller. */
public boolean tokenMode() {
return MODE_TOKEN.equals(mode);
}
}
/**
* Subscription boundary. Only these hosts may back a worker's
* {@code ANTHROPIC_BASE_URL}; the primary must carry none.
@@ -326,8 +366,43 @@ public record BridgedConfig(
Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null);
Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000;
Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300;
Auth a = auth != null ? auth : new Auth(null, null);
// broker is left as-is: null (or an empty/blank uri) keeps the in-memory soft-state inbox.
// primary is left as-is: null defaults to connection-derived identity.
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary);
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, a);
}
/**
* Reject a configuration whose network exposure outruns its authentication (CB-501).
*
* <p>{@code loopback-trust} means "any caller that is not a known worker pane is the primary" —
* safe only because the OS refuses non-local connections to a loopback bind. Widen
* {@code bind.host} without switching to {@code token} mode and that sentence becomes "any
* client that can reach this port is the primary", which is the most privileged role on the
* bus. Rather than document the hazard, make it unrepresentable: fail fast at startup.
*
* @throws IllegalStateException when a non-loopback bind is paired with {@code loopback-trust}
*/
public void validateAuthExposure() {
String host = bind().host();
if (isLoopbackBind(host) || auth().tokenMode()) {
return;
}
throw new IllegalStateException(
"refusing to start: bind.host=" + host + " is not loopback, but auth.mode="
+ auth().mode() + ". A non-loopback bind treats every unauthenticated "
+ "caller as the primary (spawn/stop/send/drain on any session). Set "
+ "auth.mode: token (with auth.tokenEnv) before exposing this port, or "
+ "bind to 127.0.0.1 and put a reverse proxy in front.");
}
/** True for the loopback addresses and the unspecified-but-local forms we treat as same-host. */
private static boolean isLoopbackBind(String host) {
if (host == null || host.isBlank()) {
return true; // Bind's own default is 127.0.0.1
}
String h = host.trim().toLowerCase();
return h.equals("127.0.0.1") || h.equals("::1") || h.equals("localhost")
|| h.startsWith("127.");
}
}
@@ -1,7 +1,14 @@
package dev.ltms.bridged.mcp;
import dev.ltms.bridged.auth.AuditLog;
import dev.ltms.bridged.auth.Authz;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.auth.Principal;
import dev.ltms.bridged.auth.Role;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.msg.MessageService;
@@ -56,13 +63,34 @@ public final class BridgeMcp {
static final String CALLER_TERMINAL = "callerTerminal";
/** Transport-context key under which the extractor stashes the caller's PID (for cwd inherit). */
static final String CALLER_PID = "callerPid";
/** Transport-context key under which the extractor stashes the resolved {@link Role} (CB-501). */
static final String CALLER_ROLE = "callerRole";
private final HttpServletStreamableServerTransportProvider transport;
private final McpSyncServer server;
private final CallerResolver authz; // CB-501: null → authorization not enforced (legacy)
private final Metrics metrics; // CB-502: null → auth failures not counted
/**
* Legacy constructor — no authorization. Retained so existing tests exercise tool behaviour
* without an auth fixture.
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence,
PrimaryRegistry primaryRegistry) {
this(messages, workers, sessions, identity, presence, primaryRegistry, null, null);
}
/**
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
* Jetty's context handler and never passes through Javalin's {@code before}
* filter, so the REST guard does not cover it.
* @param metrics registry for auth-failure counting; may be {@code null}
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence,
PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics) {
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
.jsonMapper(json)
@@ -72,17 +100,26 @@ public final class BridgeMcp {
// inherit the primary's cwd (CB-112). Any contact from a worker marks it available
// (CB-113) — its MCP initialize is the reliable "the agent is up" signal.
.contextExtractor(req -> {
ConnectionIdentity.Caller c = identity.resolve(req.getRemoteAddr(), req.getRemotePort());
presence.markPresent(c.terminal()); // no-op for the primary (null terminal)
// One resolution per call, shared with the REST surface via CallerResolver so
// the two paths cannot drift on who a caller is.
Principal p = callers != null
? callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
req.getHeader("Authorization"))
: legacyPrincipal(identity, req.getRemoteAddr(), req.getRemotePort());
presence.markPresent(p.terminal()); // no-op for the primary (null terminal)
return McpTransportContext.create(Map.of(
CALLER_TERMINAL, orEmpty(c.terminal()),
CALLER_PID, Long.toString(c.pid())));
CALLER_TERMINAL, orEmpty(p.terminal()),
CALLER_PID, Long.toString(p.pid()),
CALLER_ROLE, p.role().name()));
})
.build();
this.server = McpServer.sync(transport)
.serverInfo("bridge", "0.1.0")
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
.toolCall(sendTool(), (exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
str(req.arguments(), "sessionId"));
if (denied != null) return denied;
String caller = callerTerminal(exchange);
if (caller != null) primaryRegistry.record(caller);
Map<String, Object> a = req.arguments();
@@ -97,25 +134,44 @@ public final class BridgeMcp {
? sendAsync(messages, str(a, "sessionId"), str(a, "content"))
: send(messages, str(a, "sessionId"), str(a, "content"), timeoutMs(a));
})
// bridge_reply's identity is the CONNECTION, never an argument.
.toolCall(replyTool(), (exchange, req) ->
reply(messages, callerTerminal(exchange), str(req.arguments(), "content")))
// bridge_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
.toolCall(replyTool(), (exchange, req) -> {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self);
if (denied != null) return denied;
return reply(messages, self, str(req.arguments(), "content"));
})
// bridge_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
.toolCall(askTool(), (exchange, req) ->
ask(messages, callerTerminal(exchange), str(req.arguments(), "question"), timeoutMs(req.arguments())))
.toolCall(statusTool(), (_, req) ->
status(messages, str(req.arguments(), "sessionId")))
.toolCall(pollTool(), (_, req) -> {
.toolCall(askTool(), (exchange, req) -> {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self);
if (denied != null) return denied;
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
})
.toolCall(statusTool(), (exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return status(messages, str(req.arguments(), "sessionId"));
})
.toolCall(pollTool(), (exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
Map<String, Object> a = req.arguments();
return poll(messages, str(a, "ticket"), str(a, "target"));
})
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
.toolCall(ackTool(), (_, req) -> {
// Acking removes a reply from the inbox, so it is a drain, not a read.
.toolCall(ackTool(), (exchange, req) -> {
Map<String, Object> a = req.arguments();
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target"));
if (denied != null) return denied;
return ack(messages, str(a, "target"), str(a, "msgId"));
})
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
.toolCall(spawnTool(), (exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
if (denied != null) return denied;
String caller = callerTerminal(exchange);
if (caller != null) primaryRegistry.record(caller);
Map<String, Object> a = req.arguments();
@@ -126,10 +182,72 @@ public final class BridgeMcp {
return spawn(sessions, str(a, "profile"), str(a, "cwd"), callerCwd,
callerTerminal(exchange), worktreeRequest(a));
})
.toolCall(listTool(), (_, _) -> listWorkers(workers, sessions))
.toolCall(stopTool(), (_, req) -> stop(sessions, str(req.arguments(), "paneId")))
.toolCall(profilesTool(), (_, _) -> profiles(workers))
.toolCall(listTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return listWorkers(workers, sessions);
})
.toolCall(stopTool(), (exchange, req) -> {
String paneId = str(req.arguments(), "paneId");
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId);
if (denied != null) return denied;
return stop(sessions, paneId);
})
.toolCall(profilesTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return profiles(workers);
})
.build();
this.authz = callers;
this.metrics = metrics;
}
/**
* Pre-CB-501 identity: worker if the connection maps to a pane, otherwise the primary. Used
* only by the legacy constructor, where authorization is not enforced anyway.
*/
private static Principal legacyPrincipal(ConnectionIdentity identity, String addr, int port) {
ConnectionIdentity.Caller c = identity.resolve(addr, port);
return c.terminal() != null
? Principal.worker(c.terminal(), c.pid())
: Principal.primary(c.pid());
}
/** The caller reconstructed from the transport context. */
private static Principal principal(McpSyncServerExchange exchange) {
Object r = exchange.transportContext().get(CALLER_ROLE);
String terminal = callerTerminal(exchange);
long pid = callerPid(exchange);
if (r == null) {
// No role stashed (legacy path): fall back to the historical interpretation.
return terminal != null ? Principal.worker(terminal, pid) : Principal.primary(pid);
}
return new Principal(Role.valueOf(r.toString()), terminal, pid);
}
/**
* Gate a tool call on the CB-505 table. Returns {@code null} when the call may proceed, or the
* error result to return when it may not.
*/
private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action,
String target) {
if (authz == null) {
return null; // legacy: authorization not enforced
}
Principal caller = principal(exchange);
if (Authz.permits(caller, action, target)) {
if (action != Authz.Action.READ) {
AuditLog.allowed(caller, action, target);
}
return null;
}
String reason = Authz.isUnauthenticated(caller) ? "unauthenticated" : "forbidden";
AuditLog.denied(caller, action, target, reason);
if (metrics != null) {
metrics.inc(BridgedMetrics.AUTH_FAILURES, "reason", reason);
}
return error(reason + ": " + caller.describe() + " may not " + action);
}
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
@@ -0,0 +1,96 @@
package dev.ltms.bridged.metrics;
import dev.ltms.bridged.msg.ReplyInbox;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.WorkerSession;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* The daemon's metric definitions (CB-502) — one place where every series is named, described, and
* (for gauges) bound to live state.
*
* <p>The set is deliberately small: each series maps to a failure mode this project has actually
* hit, not to whatever was easy to count. The two worth watching in practice are
* {@code bridged_sends_total{outcome="completion_fallback"}} — a rising share means turn detection
* is degrading, the CB-115/116/118 failure family — and
* {@code bridged_push_nudges_total{outcome="exhausted"}}, which means the primary stopped draining
* its inbox and CB-307's active push gave up.
*/
public final class BridgedMetrics {
/** Counter: delegated sends by terminal outcome. */
public static final String SENDS = "bridged_sends_total";
/** Counter: worker replies by the path that carried them (rendezvous vs stranded-to-inbox). */
public static final String REPLIES = "bridged_replies_total";
/** Counter: push-loop nudges to the primary, by outcome. */
public static final String PUSH_NUDGES = "bridged_push_nudges_total";
/** Counter: spawn attempts by peer kind and outcome. */
public static final String SPAWNS = "bridged_spawns_total";
/** Counter: herdr socket calls by method and outcome. */
public static final String HERDR_CALLS = "bridged_herdr_calls_total";
/** Counter: rejected requests by reason (CB-501). */
public static final String AUTH_FAILURES = "bridged_auth_failures_total";
/** Gauge: session census by lifecycle state. */
public static final String SESSIONS = "bridged_sessions";
/** Gauge: undrained replies held per target. */
public static final String INBOX_DEPTH = "bridged_inbox_depth";
private BridgedMetrics() {
}
/**
* Build the registry with its help text and live gauges bound.
*
* @param sessions the authoritative session registry (census gauge)
* @param inbox the reply inbox; only used for a depth gauge when it can be inspected
*/
public static Metrics create(SessionManager sessions, ReplyInbox inbox) {
Metrics m = new Metrics();
m.describe(SENDS, "counter",
"Delegated sends by terminal outcome (replied|completion_fallback|timeout|failed).");
m.describe(REPLIES, "counter",
"Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held).");
m.describe(PUSH_NUDGES, "counter",
"CB-307 push-loop nudges to the primary (delivered|exhausted).");
m.describe(SPAWNS, "counter",
"Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected).");
m.describe(HERDR_CALLS, "counter",
"herdr socket calls by method and outcome — the dependency everything else rests on.");
m.describe(AUTH_FAILURES, "counter",
"Requests refused by CB-501/505 (unauthenticated|forbidden).");
m.describe(SESSIONS, "gauge",
"Registered worker sessions by lifecycle state.");
m.describe(INBOX_DEPTH, "gauge",
"Replies held for a target that the primary has not drained. Steady state is 0; "
+ "a target stuck above 0 means CB-307 delivery is not completing.");
// One gauge per state so a scrape shows the whole census even when a state is empty —
// an absent series and a zero series read very differently on a dashboard.
for (WorkerSession.State state : WorkerSession.State.values()) {
String label = state.name().toLowerCase();
m.gauge(SESSIONS, () -> countIn(sessions, state), "state", label);
}
// Depth is per live session, so the label set is only known at scrape time. peek() is the
// port's non-destructive read — scraping metrics must never ack a reply out of the inbox.
m.collector(INBOX_DEPTH, "target", () -> {
Map<String, Number> depths = new LinkedHashMap<>();
for (WorkerSession s : sessions.roster()) {
String target = s.terminalId();
if (target == null) {
continue;
}
depths.put(target, inbox.peek(target).size());
}
return depths;
});
return m;
}
private static long countIn(SessionManager sessions, WorkerSession.State state) {
return sessions.roster().stream().filter(s -> s.state() == state).count();
}
}
@@ -0,0 +1,181 @@
package dev.ltms.bridged.metrics;
import java.util.Map;
import java.util.NavigableMap;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.atomic.LongAdder;
import java.util.function.Supplier;
/**
* The daemon's metric registry and Prometheus text renderer (CB-502).
*
* <p>Deliberately dependency-free. The roadmap's tech-stack table specified Micrometer, but this
* pom already carries an unusual reconciliation burden (a hand-pinned {@code jackson-annotations}
* to make the MCP SDK's Jackson 3 coexist with our Jackson 2, a Jetty BOM import to stop version
* skew, and four documented accepted-CVE advisories), and the dependency CVE gate this project
* mandates could not be run when this landed. The metric set is small and fully known, and
* Prometheus text exposition is a stable, well-specified format — so the registry is ~100 lines
* here instead of a new transitive tree. {@code GET /metrics} is the swap seam if Micrometer's
* ecosystem is ever wanted.
*
* <p>Thread-safe: counters are {@link LongAdder} (built for contended increment), gauges are
* supplier-backed so they read live state at scrape time rather than needing to be pushed.
*/
public final class Metrics {
/** Counter series, keyed by the fully-rendered {@code name{labels}} sample id. */
private final NavigableMap<String, LongAdder> counters = new ConcurrentSkipListMap<>();
/** Gauge series, evaluated at scrape time. */
private final NavigableMap<String, Supplier<Number>> gauges = new ConcurrentSkipListMap<>();
/** Gauge families whose label set is only known at scrape time, keyed by metric name. */
private final NavigableMap<String, Collector> collectors = new ConcurrentSkipListMap<>();
/** HELP/TYPE metadata, keyed by bare metric name. */
private final Map<String, String[]> meta = new ConcurrentHashMap<>();
/** A gauge family whose series are discovered per scrape (one label, many values). */
private record Collector(String labelName, Supplier<Map<String, Number>> samples) {
}
/** Declare a metric's help text and type once, so the exposition carries HELP/TYPE lines. */
public Metrics describe(String name, String type, String help) {
meta.put(name, new String[]{type, help});
return this;
}
/** Increment a counter by one. */
public void inc(String name, String... labelPairs) {
add(name, 1, labelPairs);
}
/** Increment a counter by {@code delta}. */
public void add(String name, long delta, String... labelPairs) {
counters.computeIfAbsent(sample(name, labelPairs), _ -> new LongAdder()).add(delta);
}
/**
* Register a live gauge. The supplier is called at scrape time, so it reflects current state
* (session census, inbox depth) without anything having to remember to update it.
*/
public void gauge(String name, Supplier<Number> value, String... labelPairs) {
gauges.put(sample(name, labelPairs), value);
}
/**
* Register a gauge family whose label values are not known up front — inbox depth per target,
* for instance, where the set of targets changes as workers come and go. The supplier returns
* {@code labelValue → value} and is evaluated once per scrape.
*/
public void collector(String name, String labelName, Supplier<Map<String, Number>> samples) {
collectors.put(name, new Collector(labelName, samples));
}
/** Current value of a counter series — for assertions in tests. */
public long count(String name, String... labelPairs) {
LongAdder a = counters.get(sample(name, labelPairs));
return a == null ? 0 : a.sum();
}
/**
* Render the Prometheus text exposition format (version 0.0.4): optional {@code # HELP} and
* {@code # TYPE} lines per metric family, then one line per sample.
*/
public String render() {
StringBuilder out = new StringBuilder(1024);
String lastFamily = null;
for (Map.Entry<String, LongAdder> e : counters.entrySet()) {
lastFamily = emitHeader(out, e.getKey(), lastFamily);
out.append(e.getKey()).append(' ').append(e.getValue().sum()).append('\n');
}
for (Map.Entry<String, Supplier<Number>> e : gauges.entrySet()) {
lastFamily = emitHeader(out, e.getKey(), lastFamily);
Number v;
try {
v = e.getValue().get();
} catch (RuntimeException ex) {
continue; // a broken gauge must never break the whole scrape
}
if (v == null) {
continue;
}
out.append(e.getKey()).append(' ').append(format(v)).append('\n');
}
for (Map.Entry<String, Collector> e : collectors.entrySet()) {
Map<String, Number> samples;
try {
samples = e.getValue().samples().get();
} catch (RuntimeException ex) {
continue; // a broken collector must never break the whole scrape
}
if (samples == null || samples.isEmpty()) {
continue;
}
lastFamily = emitHeader(out, e.getKey(), lastFamily);
// Sort so repeated scrapes are byte-stable and diffable.
new java.util.TreeMap<>(samples).forEach((label, v) -> {
if (v != null) {
out.append(sample(e.getKey(), e.getValue().labelName(), label))
.append(' ').append(format(v)).append('\n');
}
});
}
return out.toString();
}
/** Emit HELP/TYPE when the sample starts a new metric family; returns the current family. */
private String emitHeader(StringBuilder out, String sampleId, String lastFamily) {
String family = familyOf(sampleId);
if (family.equals(lastFamily)) {
return lastFamily;
}
String[] m = meta.get(family);
if (m != null) {
out.append("# HELP ").append(family).append(' ').append(m[1]).append('\n');
out.append("# TYPE ").append(family).append(' ').append(m[0]).append('\n');
}
return family;
}
private static String familyOf(String sampleId) {
int brace = sampleId.indexOf('{');
return brace < 0 ? sampleId : sampleId.substring(0, brace);
}
/** Whole numbers render without a decimal point; everything else as-is. */
private static String format(Number v) {
double d = v.doubleValue();
return (d == Math.rint(d) && !Double.isInfinite(d))
? Long.toString((long) d)
: Double.toString(d);
}
/** Build the {@code name{k="v",k2="v2"}} sample id; labels are sorted for stable output. */
private static String sample(String name, String... labelPairs) {
if (labelPairs == null || labelPairs.length == 0) {
return name;
}
if (labelPairs.length % 2 != 0) {
throw new IllegalArgumentException("labels must be key/value pairs, got " + labelPairs.length);
}
NavigableMap<String, String> sorted = new java.util.TreeMap<>();
for (int i = 0; i < labelPairs.length; i += 2) {
sorted.put(labelPairs[i], labelPairs[i + 1] == null ? "" : labelPairs[i + 1]);
}
StringBuilder sb = new StringBuilder(name.length() + 16 * sorted.size());
sb.append(name).append('{');
boolean first = true;
for (Map.Entry<String, String> e : sorted.entrySet()) {
if (!first) {
sb.append(',');
}
first = false;
sb.append(e.getKey()).append("=\"").append(escapeLabel(e.getValue())).append('"');
}
return sb.append('}').toString();
}
/** Label values are escaped per the exposition format: backslash, quote, newline. */
private static String escapeLabel(String v) {
return v.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n");
}
}
@@ -3,6 +3,8 @@ package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -151,6 +153,7 @@ public final class MessageService {
private final Rendezvous rendezvous;
private final ReplyInbox inbox;
private final ReplyPushLoop pushLoop;
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
@@ -165,11 +168,23 @@ public final class MessageService {
*/
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
ReplyInbox inbox, ReplyPushLoop pushLoop) {
this(agents, injector, rendezvous, inbox, pushLoop, null);
}
/**
* As above, with a metric registry (CB-502). Instrumenting here rather than at the REST and MCP
* edges means both surfaces are counted by one piece of code and cannot drift.
*
* @param metrics nullable — when null, nothing is recorded
*/
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
ReplyInbox inbox, ReplyPushLoop pushLoop, Metrics metrics) {
this.agents = agents;
this.injector = injector;
this.rendezvous = rendezvous;
this.inbox = inbox;
this.pushLoop = pushLoop;
this.metrics = metrics;
}
/** Create with an explicit {@link ReplyInbox} and no push loop. */
@@ -200,15 +215,46 @@ public final class MessageService {
*/
public boolean reply(String session, String content) {
if (rendezvous.resolve(session, content)) {
count(BridgedMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
inbox.publish(session, UUID.randomUUID().toString(), content);
// A rising inbox share is the signal CB-307 exists to make visible: the worker finished but
// nobody was waiting, so delivery now depends on the push loop and a drain.
count(BridgedMetrics.REPLIES, "path", "inbox");
if (pushLoop != null) {
pushLoop.onReplyQueued(session);
}
return true; // held, not lost
}
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
private void count(String name, String... labels) {
if (metrics != null) {
metrics.inc(name, labels);
}
}
/** Count a send's terminal outcome and pass the reply through unchanged. */
private Reply recorded(Reply r) {
String label = sendOutcomeLabel(r.outcome());
if (label != null) {
count(BridgedMetrics.SENDS, "outcome", label);
}
return r;
}
/** Map a terminal send outcome to its metric label, or {@code null} for non-terminal ones. */
private static String sendOutcomeLabel(Outcome o) {
return switch (o) {
case REPLIED -> "replied";
case COMPLETED_UNREPLIED -> "completion_fallback";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
case WORKER_FAILED -> "failed";
case STALE_TURN, QUESTION -> null; // not a completed delegation
};
}
/**
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
* so that a subsequent drain or peek no longer returns it.
@@ -248,11 +294,12 @@ public final class MessageService {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
} catch (TimeoutException e) {
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
return new Reply(wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null);
return recorded(new Reply(
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
@@ -2,7 +2,12 @@ package dev.ltms.bridged.rest;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.bridged.auth.AuditLog;
import dev.ltms.bridged.auth.Authz;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.auth.Principal;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
@@ -43,23 +48,47 @@ public final class BridgedApp {
private static final long DEFAULT_ASK_TIMEOUT_MS = 55_000;
private static final long MAX_ASK_TIMEOUT_MS = 115_000;
/** Context attribute under which the resolved caller is stashed by the auth filter. */
private static final String CALLER = "bridged.caller";
private final HerdrClient herdr;
private final PeerLauncher workers;
private final SessionManager sessions; // CB-301: authoritative session registry
private final MessageService messages;
private final WorkerPresence presence; // CB-113: which workers are MCP-connected (available)
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
private final Metrics metrics; // CB-502: null → /metrics not exposed
private final ObjectMapper mapper = new ObjectMapper();
/**
* Legacy constructor — no identity resolution and no authorization, exactly as the REST surface
* behaved before CB-501. Retained so existing acceptance tests keep exercising handler
* behaviour without each needing an auth fixture.
*/
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, WorkerPresence presence,
HttpServlet mcpServlet) {
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
}
/**
* @param auth resolves each request's {@link Principal}; {@code null} disables authorization
* entirely (legacy). {@code main} always supplies one.
* @param metrics registry to instrument and expose at {@code GET /metrics}; {@code null} omits
* the endpoint
*/
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, WorkerPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
this.herdr = herdr;
this.workers = workers;
this.sessions = sessions;
this.messages = messages;
this.presence = presence;
this.mcpServlet = mcpServlet;
this.auth = auth;
this.metrics = metrics;
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -72,7 +101,18 @@ public final class BridgedApp {
h.addServlet(new ServletHolder(mcpServlet), "/mcp"));
}
});
// CB-501: resolve identity once per request, before any handler. /mcp does NOT pass through
// here — it is a raw servlet on Jetty's context handler — so BridgeMcp enforces separately
// against the same CallerResolver. Any check that lives in only one place is not a control.
if (auth != null) {
app.before(ctx -> ctx.attribute(CALLER,
auth.resolve(ctx.req().getRemoteAddr(), ctx.req().getRemotePort(),
ctx.header("Authorization"))));
}
app.get("/healthz", this::healthz);
if (metrics != null) {
app.get("/metrics", this::metrics);
}
app.get("/sessions", this::sessions);
app.get("/agents", this::agents);
app.get("/workers", this::listWorkers); // CB-304: registry roster + live herdr status
@@ -88,6 +128,54 @@ public final class BridgedApp {
return app;
}
/**
* Gate a handler on the CB-505 authorization table. Returns {@code true} when the request may
* proceed; otherwise writes the error response and returns {@code false}.
*
* <p>401 vs 403 is a real distinction here: 401 means "you presented no usable identity" (a
* credential problem the caller can fix), 403 means "you are authenticated, but this is not
* yours" (a worker reaching for another worker's session, or for orchestration).
*/
private boolean allow(Context ctx, Authz.Action action, String target) {
if (auth == null) {
return true; // legacy: authorization not enforced
}
Principal caller = ctx.attribute(CALLER);
if (Authz.permits(caller, action, target)) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
return true;
}
if (Authz.isUnauthenticated(caller)) {
AuditLog.denied(caller, action, target, "unauthenticated");
countAuthFailure("unauthenticated");
ctx.status(401).json(Map.of("error", "unauthenticated",
"detail", "present Authorization: Bearer <token>"));
} else {
AuditLog.denied(caller, action, target, "forbidden");
countAuthFailure("forbidden");
ctx.status(403).json(Map.of("error", "forbidden",
"detail", caller.describe() + " may not " + action + " on "
+ (target == null ? "this resource" : target)));
}
return false;
}
private void countAuthFailure(String reason) {
if (metrics != null) {
metrics.inc("bridged_auth_failures_total", "reason", reason);
}
}
/** Prometheus scrape endpoint (CB-502). */
private void metrics(Context ctx) {
if (!allow(ctx, Authz.Action.METRICS, null)) {
return;
}
ctx.status(200).contentType("text/plain; version=0.0.4; charset=utf-8").result(metrics.render());
}
/** Liveness + herdr reachability. 200 when herdr answers ping, 503 otherwise. */
private void healthz(Context ctx) {
try {
@@ -107,6 +195,9 @@ public final class BridgedApp {
/** Sessions view derived from herdr {@code workspace.list} (one workspace → one row). */
private void sessions(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
JsonNode result = herdr.call("workspace.list");
List<Map<String, Object>> out = new ArrayList<>();
for (JsonNode w : result.path("workspaces")) {
@@ -122,12 +213,18 @@ public final class BridgedApp {
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
private void agents(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(BridgedApp::view).toList()));
}
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
private void listWorkers(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.paneId() != null)
@@ -140,6 +237,9 @@ public final class BridgedApp {
/** The configured worker profiles and which one a no-argument spawn uses. */
private void profiles(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
ctx.status(200).json(Map.of(
"profiles", workers.profiles(),
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile()));
@@ -151,6 +251,9 @@ public final class BridgedApp {
* the subscription boundary, 400 for an unknown profile.
*/
private void spawnWorker(Context ctx) {
if (!allow(ctx, Authz.Action.SPAWN, null)) {
return;
}
String profile = ctx.queryParam("profile");
String cwd = ctx.queryParam("cwd");
String worktree = ctx.queryParam("worktree");
@@ -203,7 +306,11 @@ public final class BridgedApp {
/** Tear a worker down by pane id. */
private void stopWorker(Context ctx) {
sessions.release(ctx.pathParam("paneId"));
String paneId = ctx.pathParam("paneId");
if (!allow(ctx, Authz.Action.STOP, paneId)) {
return;
}
sessions.release(paneId);
ctx.status(204);
}
@@ -215,6 +322,9 @@ public final class BridgedApp {
*/
private void sendMessage(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, Authz.Action.SEND, id)) {
return;
}
String content;
String turnId;
long timeout;
@@ -296,6 +406,9 @@ public final class BridgedApp {
*/
private void askMessage(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, Authz.Action.ASK, id)) {
return;
}
String question;
long timeout;
try {
@@ -329,6 +442,12 @@ public final class BridgedApp {
*/
private void replyMessage(Context ctx) {
String id = ctx.pathParam("id");
// The rule that matters: a worker may reply only as itself. Over MCP this was already true
// structurally (identity comes from the connection, never an argument); over REST the path
// id was simply trusted, so this is where the invariant actually gets enforced.
if (!allow(ctx, Authz.Action.REPLY, id)) {
return;
}
String content;
try {
content = mapper.readTree(ctx.body()).path("content").asText("");
@@ -347,6 +466,9 @@ public final class BridgedApp {
*/
private void drainReplies(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, Authz.Action.DRAIN, id)) {
return;
}
var replies = messages.drainReplies(id);
ctx.status(200).json(Map.of("sessionId", id, "replies",
replies.stream().map(m -> Map.of(
@@ -362,6 +484,9 @@ public final class BridgedApp {
*/
private void sessionStatus(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, Authz.Action.READ, id)) {
return;
}
try {
ctx.status(200).json(Map.of(
"sessionId", id,
@@ -374,6 +499,9 @@ public final class BridgedApp {
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
private void taskStatus(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));