Files
fleetd/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java
T
Dai Ha 541df87272
CI / build (pull_request) Failing after 59s
CI / contract (pull_request) Successful in 1m10s
CB-578 stage A: classify a usage-limit refusal instead of a completed reply
A backend that refuses on a subscription usage limit leaves the pane healthy but
the turn ends with no bridge_reply; the completion fallback used to scrape and
hand that refusal back as if it were a real answer. CompletionResolver now
matches the scrape against a per-profile exhaustedPattern (config, never a
vendor string) and resolves the send as Rendezvous.Kind/Outcome.BACKEND_EXHAUSTED
with a reason carrying the matched line, kept distinct from GONE/WORKER_FAILED.
A profile with no pattern configured is unaffected. Coverage is logged at
startup via CompletionResolver.coverage(...), naming which profiles have a
pattern and which don't, following FleetHealthMonitor.coverage's pattern.
2026-08-15 09:55:06 +02:00

578 lines
26 KiB
Java

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;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.peer.PeerLauncher;
import io.javalin.Javalin;
import io.javalin.http.Context;
import jakarta.servlet.http.HttpServlet;
import org.eclipse.jetty.servlet.ServletHolder;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
* 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 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 {
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
/** Blocking window for a worker's bridge_ask (CB-205); the worker's MCP client caps its own call. */
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 MemberPresence 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, MemberPresence 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, MemberPresence 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. */
public Javalin build() {
Javalin app = Javalin.create(cfg -> {
cfg.showJavalinBanner = false;
if (mcpServlet != null) {
// The MCP server shares the daemon's port; Jetty routes /mcp to its servlet.
cfg.jetty.modifyServletContextHandler(h ->
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("/members", this::listMembers); // CB-304: registry roster + live herdr status
app.get("/profiles", this::profiles); // configured backend profiles
app.post("/members", this::spawnMember); // optional ?role=&profile= or {"role":…,"profile":…}
app.delete("/members/{paneId}", this::stopMember);
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking, wait:false, or answer via turnId)
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker)
app.get("/sessions/{id}/replies", this::drainReplies); // drain reply inbox (CB-307)
app.post("/sessions/{id}/ask", this::askMessage); // bridge_ask (worker → primary, CB-205)
app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status
app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send
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 {
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} (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")) {
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));
}
/** 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 listMembers(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> out = sessions.roster().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
ctx.status(200).json(Map.of("workers", out));
}
/** 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()));
}
/**
* Spawn a guard-checked worker. An optional {@code profile} (query param or {@code {"profile":…}}
* body) picks which configured profile; omitted → the default. 403 if the base_url would breach
* the subscription boundary, 400 for an unknown profile.
*/
private void spawnMember(Context ctx) {
if (!allow(ctx, Authz.Action.SPAWN, null)) {
return;
}
String role = ctx.queryParam("role");
String profile = ctx.queryParam("profile");
String cwd = ctx.queryParam("cwd");
String worktree = ctx.queryParam("worktree");
String ticket = ctx.queryParam("ticket");
if (profile == null || profile.isBlank() || cwd == null || cwd.isBlank()
|| worktree == null || worktree.isBlank()) {
try {
String body = ctx.body();
if (!body.isBlank()) {
JsonNode b = mapper.readTree(body);
if (role == null || role.isBlank()) role = b.path("role").asText(null);
if (profile == null || profile.isBlank()) profile = b.path("profile").asText(null);
if (cwd == null || cwd.isBlank()) cwd = b.path("cwd").asText(null);
if (worktree == null || worktree.isBlank()) worktree = b.path("worktree").asText(null);
if (ticket == null || ticket.isBlank()) ticket = b.path("ticket").asText(null);
}
} catch (Exception ignored) {
// A malformed/empty body just means "no overrides" → fall through to defaults.
}
}
WorktreeRequest wt = worktreeRequest(worktree, ticket);
MemberRole memberRole;
try {
memberRole = (role == null || role.isBlank()) ? MemberRole.DEV : MemberRole.parse(role);
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_role", "detail", e.getMessage()));
return;
}
try {
// No MCP caller over REST, so callerCwd and ownerTerminal are null.
MemberSession member = sessions.acquire(blankToNull(profile), memberRole,
blankToNull(cwd), null, null, wt);
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
} catch (PeerUnreachableException e) {
ctx.status(502).json(Map.of("error", "spawn_timeout", "detail", e.getMessage()));
}
}
private static WorktreeRequest worktreeRequest(String worktree, String ticket) {
if (worktree == null || worktree.isBlank() || "false".equalsIgnoreCase(worktree)) {
return null;
}
if ("true".equalsIgnoreCase(worktree)) {
if (ticket == null || ticket.isBlank()) {
throw new IllegalArgumentException("worktree=true requires a ticket slug");
}
return new WorktreeRequest(ticket, null);
}
return new WorktreeRequest(worktree, null);
}
private static String blankToNull(String s) {
return (s == null || s.isBlank()) ? null : s;
}
/** Tear a worker down by pane id. */
private void stopMember(Context ctx) {
String paneId = ctx.pathParam("paneId");
if (!allow(ctx, Authz.Action.STOP, paneId)) {
return;
}
sessions.release(paneId);
ctx.status(204);
}
/**
* The blocking delegation call (CB-104): inject {@code content} into the worker via the
* status-gated injector and block until the worker returns a structured {@code bridge_reply}.
* Times out with a typed 202 (working / queued / busy) rather than an error — the message may
* still land.
*/
private void sendMessage(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, Authz.Action.SEND, id)) {
return;
}
String content;
String turnId;
long timeout;
boolean wait;
try {
JsonNode body = mapper.readTree(ctx.body());
content = body.path("content").asText("");
turnId = body.path("turnId").asText(null);
timeout = body.path("timeoutMs").asLong(DEFAULT_MESSAGE_TIMEOUT_MS);
wait = body.path("wait").asBoolean(true); // default: block for the reply (CB-104)
} catch (Exception e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
if (content.isBlank()) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required"));
return;
}
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
// Answering a worker's bridge_ask (CB-205): always blocks, and derives the worker from turnId.
if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
return;
}
if (!wait) {
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
String ticket = messages.sendAsync(id, content);
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
try {
writeReply(ctx, id, messages.send(id, content, timeout), timeout);
} catch (HerdrException e) {
herdrError(ctx, e);
}
}
/**
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a
* bridge_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
* (with its {@code turnId}); a stale answer a 409; every other non-terminal outcome a typed 202.
*/
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout) {
switch (reply.outcome()) {
case QUESTION -> ctx.status(202).json(Map.of(
"sessionId", id, "status", "question",
"question", reply.text(), "turnId", reply.turnId()));
case STALE_TURN -> ctx.status(409).json(Map.of(
"sessionId", id, "error", "stale_turn",
"detail", "that question is no longer open (timed out or already answered)"));
case REPLIED, COMPLETED_UNREPLIED -> {
// replySource distinguishes a structured bridge_reply from the CB-106 completion
// fallback (a scrape of the worker's transcript when it finished without replying).
String source = reply.outcome() == MessageService.Outcome.REPLIED ? "reply" : "transcript";
ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text(), "replySource", source));
}
default -> ctx.status(202).json(Map.of(
"sessionId", id,
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
default -> "done"; // unreachable (terminal outcomes handled above)
},
"detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
&& reply.text() != null
? reply.text()
: "no reply within " + timeout + "ms; poll status or retry"));
}
}
/**
* A worker's mid-turn question ({@code bridge_ask}, CB-205) — surfaces to the primary's open
* blocking send and blocks until it answers. 200 with the answer, 409 if no delegation is open,
* 202 if the primary stayed silent.
*/
private void askMessage(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, Authz.Action.ASK, id)) {
return;
}
String question;
long timeout;
try {
JsonNode body = mapper.readTree(ctx.body());
question = body.path("question").asText("");
timeout = body.path("timeoutMs").asLong(DEFAULT_ASK_TIMEOUT_MS);
} catch (Exception e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
if (question.isBlank()) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "question is required"));
return;
}
timeout = Math.clamp(timeout, 1, MAX_ASK_TIMEOUT_MS);
MessageService.AskResult r = messages.ask(id, question, timeout);
switch (r.outcome()) {
case ANSWERED -> ctx.status(200).json(Map.of("sessionId", id, "answered", true, "answer", r.answer()));
case NO_WAITER -> ctx.status(409).json(Map.of(
"sessionId", id, "error", "no_pending_send",
"detail", "no primary is awaiting this turn to answer a question"));
case TIMED_OUT -> ctx.status(202).json(Map.of(
"sessionId", id, "status", "no_answer",
"detail", "the primary did not answer within " + timeout + "ms"));
}
}
/**
* The worker's structured reply ({@code bridge_reply}) — resolves the blocking send awaiting
* on this session, or queues the reply in the inbox when no send is open (CB-307).
*/
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("");
} catch (Exception e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
messages.reply(id, content);
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
}
/**
* Drain the reply inbox for a worker session — peek + ack any replies that arrived when no send
* was open. At-least-once: draining removes them from the inbox so a subsequent read returns
* nothing; an in-flight failure between the drain and the caller's processing re-surfaces them.
*/
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(
"msgId", m.msgId(),
"content", m.content())).toList()));
}
/**
* Live lifecycle status of a worker (MCP `bridge_status` wraps this in CB-105), plus its
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
* is also true during boot.
*/
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,
"status", messages.status(id).name().toLowerCase(),
"ready", presence.isPresent(id)));
} catch (HerdrException e) {
herdrError(ctx, e);
}
}
/** 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)"));
return;
}
Map<String, Object> body = new LinkedHashMap<>();
body.put("ticket", v.ticket());
body.put("phase", v.phase().name().toLowerCase());
if (v.reply() != null) {
body.put("reply", v.reply());
body.put("replySource", v.replySource());
}
if (v.detail() != null) {
body.put("detail", v.detail());
}
ctx.status(200).json(body);
}
/** Map a herdr failure: unknown target → 404, anything else → 502 (herdr is upstream). */
private static void herdrError(Context ctx, HerdrException e) {
if (e.code() != null && e.code().endsWith("_not_found")) {
ctx.status(404).json(Map.of("error", "session_not_found", "detail", e.getMessage()));
} else {
ctx.status(502).json(Map.of("error", "herdr_error", "detail", e.getMessage()));
}
}
/** 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("tabId", a.tabId());
m.put("sessionId", a.sessionId());
m.put("agentType", a.agentType());
m.put("status", a.status().name().toLowerCase());
return m;
}
/** CB-301 projection of an authoritative bridge-owned session. */
private static Map<String, Object> view(MemberSession s) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("terminalId", s.terminalId());
m.put("paneId", s.paneId());
m.put("profile", s.profile());
m.put("cwd", s.cwd());
m.put("ownerTerminal", s.ownerTerminal());
m.put("state", s.state().name().toLowerCase());
if (s.worktree() != null) {
m.put("worktree", s.worktree());
}
if (s.branch() != null) {
m.put("branch", s.branch());
}
return m;
}
}