fleetd #801: drop the two mayDrainPoll-defaulting send/answer overloads
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:
@@ -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));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user