fleetd #801: drop the two mayDrainPoll-defaulting send/answer overloads
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m57s

They had no production caller — only test call sites — and the default
pointed the unsafe way: true means naming fleet_poll{target}, which is
the exact receipt this ticket is fixing. Keep one signature each and
pass mayDrainPoll explicitly at every test call site instead.
This commit is contained in:
Dai Ha
2026-10-07 06:12:57 +02:00
parent 24782fe329
commit c8e31e1195
2 changed files with 23 additions and 43 deletions
@@ -1031,17 +1031,6 @@ public final class FleetMcp {
? "[fleet_send from observer " + caller.terminal() + "]\n" + content : content;
}
/**
* As {@link #send(MessageService, String, String, Long, Runnable, Set, String, boolean)},
* for a caller that may always drain the target's inbox — the production handler instead
* passes the caller's real {@code Authz.Action.DRAIN} grant.
*/
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
String callerOwner) {
return send(messages, sessionId, content, timeoutMs, onAccepted, profiles, callerOwner, true);
}
/**
* {@code fleet_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.
@@ -1073,15 +1062,6 @@ public final class FleetMcp {
}
}
/**
* As {@link #answer(MessageService, String, String, Long, String, boolean)}, for a caller
* that may always drain the target's inbox.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
String callerOwner) {
return answer(messages, turnId, content, timeoutMs, callerOwner, true);
}
/**
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
@@ -78,7 +78,7 @@ class FleetMcpTest {
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null));
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -104,7 +104,7 @@ class FleetMcpTest {
void sendThenReplyRoundTrips() throws Exception {
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null));
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null, true));
// 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).
@@ -182,7 +182,7 @@ class FleetMcpTest {
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null, true));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
@@ -280,7 +280,7 @@ class FleetMcpTest {
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null, true));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -298,7 +298,7 @@ class FleetMcpTest {
@Test
void aDifferentCallersMcpAnswerIsRefusedForABlockingSendButTheRealOwnerSucceeds() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner"));
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner", true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -313,12 +313,12 @@ class FleetMcpTest {
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker", true);
assertTrue(hijacked.isError(), "a caller that did not open this turn must get an error, not an answer");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner", true));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -357,14 +357,14 @@ class FleetMcpTest {
String turnId = asking.turnId();
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L,
"worker:term_attacker");
"worker:term_attacker", true);
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
"a refused answer must not advance the async ticket's phase");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey()));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey(), true));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -383,7 +383,7 @@ class FleetMcpTest {
@Test
void sendTimesOutWithAWorkingNote() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, true);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
@@ -395,7 +395,7 @@ class FleetMcpTest {
*/
@Test
void sendTimesOutQueuedNamesFleetPollTargetAndNotTheBareWordRetry() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, true);
String text = textOf(res);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("fleet_poll{target=\"term_a\"}"), "got: " + text);
@@ -406,7 +406,7 @@ class FleetMcpTest {
@Test
void sendTimesOutWorkingNamesFleetPollTargetAndNotTheBareWordRetry() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null));
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -427,7 +427,7 @@ class FleetMcpTest {
@Test
void sendTimesOutBusyNamesFleetPollTargetAndNotTheBareWordRetry() throws Exception {
CompletableFuture<McpSchema.CallToolResult> first = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "first", 5000L, null, Set.of(), null));
() -> FleetMcp.send(messages, T, "first", 5000L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -435,7 +435,7 @@ class FleetMcpTest {
}
assertTrue(rendezvous.isWaiting(T), "the first send should hold the session lock");
McpSchema.CallToolResult busy = FleetMcp.send(messages, T, "second", 100L, null, Set.of(), null);
McpSchema.CallToolResult busy = FleetMcp.send(messages, T, "second", 100L, null, Set.of(), null, true);
String text = textOf(busy);
assertNotEquals(Boolean.TRUE, busy.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("fleet_poll{target=\"" + T + "\"}"), "got: " + text);
@@ -459,7 +459,7 @@ class FleetMcpTest {
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null));
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -496,7 +496,7 @@ class FleetMcpTest {
@Test
void answerTimesOutWithoutDrainGrantNamesNoRoute() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null));
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -523,13 +523,13 @@ class FleetMcpTest {
@Test
void sendRejectsMissingArgs() {
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null, true).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null, true).isError());
}
@Test
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null);
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null, true);
McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
assertTrue(blocking.isError());
@@ -548,7 +548,7 @@ class FleetMcpTest {
assertSendRoundTrips("term_live_member", profiles);
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null);
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null, true);
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
}
@@ -640,7 +640,7 @@ class FleetMcpTest {
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(
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null));
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -662,7 +662,7 @@ class FleetMcpTest {
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null, true));
// The worker's ask returns the answer — it resumes the same turn.
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
@@ -688,7 +688,7 @@ class FleetMcpTest {
@Test
void answerToAStaleTurnIsAnError() {
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null);
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null, true);
assertTrue(res.isError());
assertTrue(textOf(res).contains("no longer open"), textOf(res));
}