Merge CB-572: reject a configured profile name as a send target
A lead sent to sessionId "sol" — a profile name, not a terminal id. The bridge accepted it, handed out a ticket, failed 60s later inside the injector, and still reported the ticket as pending 20 minutes on. The sender never learned anything and a whole delegation was lost. bridge_send now rejects a target that exactly matches a configured profile name, on both the blocking and the wait:false path, before any ticket is issued. The error names the value and points at bridge_list. The check is deliberately narrow. A target absent from the member roster may still be a peer lead's terminal or a herdr-owned pane, so only a value the bridge can prove is a profile is refused. profiles is a required parameter on both send methods — the earlier revision kept overloads that defaulted it to an empty set, which is the same silent disable shape as CB-561.
This commit is contained in:
@@ -33,6 +33,7 @@ import jakarta.servlet.http.HttpServlet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@@ -151,8 +152,8 @@ public final class BridgeMcp {
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
|
||||
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
|
||||
return Boolean.FALSE.equals(a.get("wait"))
|
||||
? sendAsync(messages, target, content, onAccepted)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted);
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles())
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
|
||||
})
|
||||
// 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.
|
||||
@@ -379,21 +380,22 @@ public final class BridgeMcp {
|
||||
|
||||
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
||||
|
||||
/** {@code bridge_send}: delegate {@code content} to a worker session and block for its reply. */
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, Long timeoutMs) {
|
||||
return send(messages, sessionId, content, timeoutMs, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #send(MessageService, String, String, Long)}, wiring an accepted-delivery hook
|
||||
* {@code bridge_send}: delegate {@code content} to a worker session and block for its reply.
|
||||
* The configured profiles are required so a profile name can never bypass target validation.
|
||||
*
|
||||
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
|
||||
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
|
||||
*/
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
|
||||
Long timeoutMs, Runnable onAccepted) {
|
||||
Long timeoutMs, Runnable onAccepted, Set<String> profiles) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
McpSchema.CallToolResult targetError = profileTargetError(sessionId, profiles);
|
||||
if (targetError != null) {
|
||||
return targetError;
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
try {
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
|
||||
@@ -463,24 +465,33 @@ public final class BridgeMcp {
|
||||
/**
|
||||
* {@code bridge_send} with {@code wait:false}: delegate {@code content} and return a ticket
|
||||
* immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content) {
|
||||
return sendAsync(messages, sessionId, content, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #sendAsync(MessageService, String, String)}, wiring the accepted-delivery hook
|
||||
* The configured profiles are required so a profile name can never bypass target validation.
|
||||
*
|
||||
* This wires the accepted-delivery hook
|
||||
* (CB-548) so an async flooding send records delegator ownership exactly once it is accepted.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted) {
|
||||
Runnable onAccepted, Set<String> profiles) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
McpSchema.CallToolResult targetError = profileTargetError(sessionId, profiles);
|
||||
if (targetError != null) {
|
||||
return targetError;
|
||||
}
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted);
|
||||
return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket);
|
||||
}
|
||||
|
||||
/** A configured profile is never a send target; other unknown values may be herdr-owned panes. */
|
||||
private static McpSchema.CallToolResult profileTargetError(String sessionId, Set<String> profiles) {
|
||||
if (profiles.contains(sessionId)) {
|
||||
return error("unknown send target \"" + sessionId + "\": it is a configured profile name, not a "
|
||||
+ "session id. Call bridge_list to find a member or lead sessionId.");
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/** {@code bridge_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
|
||||
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
|
||||
if (!isBlank(target)) {
|
||||
|
||||
@@ -53,6 +53,18 @@ class BridgeMcpTest {
|
||||
return ((McpSchema.TextContent) r.content().getFirst()).text();
|
||||
}
|
||||
|
||||
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> BridgeMcp.send(messages, target, "hi", 4000L, null, profiles));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(target), "send should be accepted for " + target);
|
||||
BridgeMcp.reply(messages, target, "received");
|
||||
assertEquals("received", textOf(send.get(6, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
private static ClaudeCodeLauncher workerService(FakeHerdr h, String baseUrl, Set<String> allow) {
|
||||
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||
"ltms-local", baseUrl, "coder", null, "BRIDGED_WORKER_TOKEN", null,
|
||||
@@ -69,7 +81,7 @@ class BridgeMcpTest {
|
||||
void sendThenReplyRoundTrips() throws Exception {
|
||||
// bridge_send blocks; bridge_reply resolves it with the worker's structured answer.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L));
|
||||
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
|
||||
|
||||
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
|
||||
// queues in the inbox if no waiter is open, which would break the round-trip).
|
||||
@@ -91,7 +103,7 @@ class BridgeMcpTest {
|
||||
@Test
|
||||
void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception {
|
||||
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by bridge_poll.
|
||||
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it");
|
||||
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
|
||||
assertNotEquals(Boolean.TRUE, accepted.isError());
|
||||
String out = textOf(accepted);
|
||||
assertTrue(out.contains("ticket="), out);
|
||||
@@ -129,15 +141,40 @@ class BridgeMcpTest {
|
||||
|
||||
@Test
|
||||
void sendTimesOutWithAWorkingNote() {
|
||||
McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L);
|
||||
McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L, null, Set.of());
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
|
||||
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsMissingArgs() {
|
||||
assertTrue(BridgeMcp.send(messages, null, "hi", null).isError());
|
||||
assertTrue(BridgeMcp.send(messages, "term_a", " ", null).isError());
|
||||
assertTrue(BridgeMcp.send(messages, null, "hi", null, null, Set.of()).isError());
|
||||
assertTrue(BridgeMcp.send(messages, "term_a", " ", null, null, Set.of()).isError());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
|
||||
McpSchema.CallToolResult blocking = BridgeMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"));
|
||||
McpSchema.CallToolResult async = BridgeMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
|
||||
|
||||
assertTrue(blocking.isError());
|
||||
assertTrue(async.isError());
|
||||
assertTrue(textOf(blocking).contains("sol"));
|
||||
assertTrue(textOf(blocking).contains("configured profile name"));
|
||||
assertTrue(textOf(blocking).contains("bridge_list"));
|
||||
assertFalse(textOf(async).contains("ticket="));
|
||||
}
|
||||
|
||||
@Test
|
||||
void sendAllowsPeerLeadMemberAndUnclassifiedTargets() throws Exception {
|
||||
Set<String> profiles = Set.of("sol");
|
||||
|
||||
assertSendRoundTrips("term_peer_lead", profiles);
|
||||
assertSendRoundTrips("term_live_member", profiles);
|
||||
|
||||
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
|
||||
McpSchema.CallToolResult result = BridgeMcp.send(messages, "external-pane", "hi", 10L, null, profiles);
|
||||
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -173,7 +210,7 @@ class BridgeMcpTest {
|
||||
void askThenAnswerRoundTrips() throws Exception {
|
||||
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L));
|
||||
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L, null, Set.of()));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
|
||||
Reference in New Issue
Block a user