fleetd #391: require reply caller role
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 1m45s

This commit is contained in:
Dai Ha
2026-09-10 09:13:16 +07:00
parent 3982ace544
commit cf8da1d5fa
2 changed files with 11 additions and 16 deletions
@@ -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)) {
@@ -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<McpSchema.CallToolResult> 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");