From cf8da1d5fa68982782d3bf229a79c7d7b184f8f2 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 09:13:16 +0700 Subject: [PATCH] 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");