From 3982ace544e06eb060e2b84aad76bf2336984d5a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 08:43:54 +0700 Subject: [PATCH 1/2] fleetd #391: refuse lead fleet replies --- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 20 ++++++++-- .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 40 ++++++++++++++++++- 2 files changed, 54 insertions(+), 6 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 84a28ca..c2e4f22 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -346,7 +346,7 @@ public final class FleetMcp { String self = callerTerminal(exchange); McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_reply", req.arguments()), self); if (denied != null) return denied; - return reply(messages, self, str(req.arguments(), "content")); + return reply(messages, self, principal(exchange).role(), str(req.arguments(), "content")); }; // fleet_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION. BiFunction askHandler = @@ -868,19 +868,26 @@ public final class FleetMcp { /** * {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send * or — when no send is open — queueing the reply in the inbox for later drain (CB-307). - * {@code callerTerminal} is resolved from the connection (never an argument); a {@code null} - * means the caller is not a known worker (e.g. the primary called it by mistake). + * {@code callerTerminal} and {@code callerRole} are resolved from the connection (never an + * argument). A {@code null} terminal means the caller is not a known worker. A PRIMARY with a + * terminal is a lead and must use {@code fleet_send}, because reply has no peer-lead route. * *

fleetd #365: the result text names which of those actually happened * ({@link MessageService.ReplyOutcome#description()}) instead of the single word "delivered" * for both — a queued reply is a real success, but it is not the same fact as one that resolved * a live waiter, and the caller could not previously tell them apart. */ - static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) { + static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, Role callerRole, String content) { if (callerTerminal == null) { return error("fleet_reply is for workers only — could not identify the calling worker " + "from the connection"); } + if (callerRole == Role.PRIMARY) { + return error("fleet_reply has no route to a peer lead. Use fleet_send{coordId: ...} for a peer on another " + + "daemon or fleet_send{sessionId: ...} for a peer on this host. fleet_reply resolves a member's " + + "blocked fleet_send, and a peer's coord-id message is durable and non-blocking, so there is " + + "nothing for it to resolve."); + } // fleetd #302: isBlank, not == null, to match fleet_send's own guard above. MessageService // .reply now REJECTS blank content, and this handler is a bare BiFunction with no try/catch // around it — so a whitespace-only fleet_reply would leave here as an uncaught @@ -893,6 +900,11 @@ public final class FleetMcp { return text(outcome.description()); } + /** Test helper for worker replies that bypass the MCP transport context. */ + static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) { + return reply(messages, callerTerminal, Role.WORKER, content); + } + /** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */ static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) { if (isBlank(target) || isBlank(msgId)) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 743fb5e..db85d58 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -3,6 +3,7 @@ package dev.ltms.fleet.mcp; import dev.ltms.fleet.auth.CallerResolver; import dev.ltms.fleet.auth.MemberRegistry; import dev.ltms.fleet.auth.Principal; +import dev.ltms.fleet.auth.Role; import dev.ltms.fleet.config.FleetConfig; import dev.ltms.fleet.guard.SubscriptionGuard; import dev.ltms.fleet.herdr.AgentControl; @@ -28,6 +29,7 @@ import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementPolicies; import io.modelcontextprotocol.spec.McpSchema; import dev.ltms.fleet.msg.InMemoryReplyInbox; +import dev.ltms.fleet.msg.ReplyInbox; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -110,7 +112,7 @@ class FleetMcpTest { // fleetd #365: a resolved live send must read distinctly from a merely-queued reply — // see replyWithNoPendingSendIsQueuedNotError below for the other case. - McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM"); + McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "LGTM"); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply)); McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS); @@ -331,7 +333,7 @@ class FleetMcpTest { void replyWithNoPendingSendIsQueuedNotError() { // CB-307: a reply with no open send is now queued in the inbox, not an error. // fleetd #365: it must also no longer claim "delivered" — nothing was waiting for it. - McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan"); + McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan"); assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error"); assertEquals(MessageService.ReplyOutcome.QUEUED.description(), textOf(res)); @@ -341,6 +343,40 @@ class FleetMcpTest { assertEquals("orphan", drained.getFirst().content()); } + @Test + void replyFromLeadIsRefusedBeforeItCanPublishToTheWorkerInbox() { + ReplyInbox inboxThatRejectsPublishes = new ReplyInbox() { + @Override public void own(String target) { } + @Override public void release(String target) { } + @Override public void publish(String target, String msgId, String content) { + fail("a lead fleet_reply must not publish to the worker inbox"); + } + @Override public List peek(String target) { return List.of(); } + @Override public void ack(String target, String msgId) { } + }; + MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(), + inboxThatRejectsPublishes); + + McpSchema.CallToolResult res = assertDoesNotThrow( + () -> FleetMcp.reply(leadMessages, "term_lead", Role.PRIMARY, "peer reply")); + + assertTrue(res.isError()); + assertEquals("fleet_reply has no route to a peer lead. Use fleet_send{coordId: ...} for a peer on another " + + "daemon or fleet_send{sessionId: ...} for a peer on this host. fleet_reply resolves a member's " + + "blocked fleet_send, and a peer's coord-id message is durable and non-blocking, so there is " + + "nothing for it to resolve.", + textOf(res)); + } + + @Test + void replyFromUnidentifiedCallerKeepsItsOwnError() { + McpSchema.CallToolResult res = FleetMcp.reply(messages, null, Role.PRIMARY, "reply"); + + assertTrue(res.isError()); + assertEquals("fleet_reply is for workers only — could not identify the calling worker from the connection", + textOf(res)); + } + @Test void replyWithBlankContentIsACleanToolErrorNotAnUncaughtException() { // fleetd #302: MessageService.reply now REJECTS blank content by throwing. fleet_reply's -- 2.52.0 From cf8da1d5fa68982782d3bf229a79c7d7b184f8f2 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 09:13:16 +0700 Subject: [PATCH 2/2] fleetd #391: require reply caller role --- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 5 ----- .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 22 +++++++++---------- 2 files changed, 11 insertions(+), 16 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index c2e4f22..56a7dae 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -900,11 +900,6 @@ public final class FleetMcp { return text(outcome.description()); } - /** Test helper for worker replies that bypass the MCP transport context. */ - static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) { - return reply(messages, callerTerminal, Role.WORKER, content); - } - /** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */ static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) { if (isBlank(target) || isBlank(msgId)) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index db85d58..da985f9 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -79,7 +79,7 @@ class FleetMcpTest { Thread.sleep(5); } assertTrue(rendezvous.isWaiting(target), "send should be accepted for " + target); - FleetMcp.reply(messages, target, "received"); + FleetMcp.reply(messages, target, Role.WORKER, "received"); assertEquals("received", textOf(send.get(6, TimeUnit.SECONDS))); } @@ -138,7 +138,7 @@ class FleetMcpTest { } assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter"); - McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM"); + McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "async LGTM"); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply)); // Poll until the async send completes and reports the reply. @@ -185,7 +185,7 @@ class FleetMcpTest { Thread.sleep(5); } assertTrue(rendezvous.isWaiting("term_a")); - FleetMcp.reply(messages, "term_a", "done"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "done"); assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS))); McpSchema.CallToolResult done = FleetMcp.poll(messages, ticket, null); @@ -217,7 +217,7 @@ class FleetMcpTest { // fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session // released before it replied" reason. This used to land in the inbox instead (see the old // assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug. - FleetMcp.reply(messages, "term_a", "finished after timeout"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "finished after timeout"); assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null))); assertTrue(messages.drainReplies("term_a").isEmpty(), "the reply completed its own ticket directly and never touched the inbox"); @@ -271,7 +271,7 @@ class FleetMcpTest { } assertEquals(MessageService.Phase.FAILED, second.phase()); - FleetMcp.reply(messages, "term_a", "late reply"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "late reply"); assertEquals("late reply", messages.drainReplies("term_a").getFirst().content()); CompletableFuture answer = CompletableFuture.supplyAsync( @@ -280,7 +280,7 @@ class FleetMcpTest { while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { Thread.sleep(5); } - FleetMcp.reply(messages, "term_a", "done"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "done"); assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS))); } @@ -387,7 +387,7 @@ class FleetMcpTest { // has always used isBlank for exactly this reason. for (String blank : new String[] {null, "", " ", "\n\t"}) { McpSchema.CallToolResult res = assertDoesNotThrow( - () -> FleetMcp.reply(messages, "term_a", blank), + () -> FleetMcp.reply(messages, "term_a", Role.WORKER, blank), "blank content must be refused as a tool error, never thrown out of the handler"); assertEquals(Boolean.TRUE, res.isError(), "blank content is an error result"); assertTrue(textOf(res).contains("content is required"), @@ -400,7 +400,7 @@ class FleetMcpTest { @Test void bridgePollWithTargetDrainsReplies() { // A reply with no open send queues it in the inbox. - FleetMcp.reply(messages, "term_a", "queued-msg"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "queued-msg"); // fleet_poll with target drains the inbox. McpSchema.CallToolResult res = FleetMcp.poll(messages, null, "term_a"); @@ -451,7 +451,7 @@ class FleetMcpTest { Thread.sleep(5); } assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter"); - McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "done"); + McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "done"); assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply)); assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS))); } @@ -1290,13 +1290,13 @@ class FleetMcpTest { @Test void bridgeAckRemovesSpecificReply() { // Queue a reply and capture its msgId. - FleetMcp.reply(messages, "term_a", "orphan"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan"); var before = messages.drainReplies("term_a"); assertEquals(1, before.size(), "one reply in the inbox"); String msgId = before.getFirst().msgId(); // Publish the same reply again and ack it via fleet_ack surface. - FleetMcp.reply(messages, "term_a", "orphan-again"); + FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan-again"); var peeked = messages.drainReplies("term_a"); assertEquals(1, peeked.size(), "one fresh reply in the inbox"); -- 2.52.0