Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cf8da1d5fa | |||
| 3982ace544 | |||
| 48877315ca | |||
| b6db9c31f5 |
@@ -137,7 +137,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
|
||||
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
|
||||
| Answer a peer lead that messaged you | `fleet_reply{content}` — the one case a lead replies |
|
||||
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
|
||||
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
|
||||
| Tear down a member | `fleet_stop{paneId}` |
|
||||
|
||||
@@ -163,15 +163,21 @@ The traffic between leads is coordination and nothing else:
|
||||
3. **Verify a peer exactly as you verify yourself.** Peer status buys nothing: check the claim
|
||||
against the code, and re-run the build. A peer's correction gets the same treatment — right or
|
||||
wrong on the evidence, not on who said it. Neither of you merges the other's work unreviewed.
|
||||
**N observations are N data points only if they differ in the axis you are trusting.** This cuts
|
||||
both ways. N *failures* blamed on one cause are one data point when the cases share what you are
|
||||
not varying. N *agreeing measurements* are also one data point when they share an instrument —
|
||||
two hosts, two operators and the same formula is one formula, not two confirmations.
|
||||
4. **Ask a peer to read your project addendum.** Your addendum is instruction surface: every future
|
||||
session on your host obeys it, and a wrong one is obeyed just as faithfully as a right one. The
|
||||
author is the worst reader of their own qualifier placement — measured here, one addendum carried
|
||||
two defects and a non-author found both. If you have no peer, at least re-read it asking "which
|
||||
sentence goes false first, and would a reader reach the caveat before acting?"
|
||||
|
||||
Being messaged by a peer does not make you its worker: answer with `fleet_reply`, and push back on
|
||||
the substance if it is wrong. A peer that simply complies has thrown away the reason there are two of
|
||||
you.
|
||||
Being messaged by a peer does not make you its worker: answer the way you would open —
|
||||
`fleet_send{coordId}` for another daemon, `fleet_send{sessionId}` on this host — and push back on
|
||||
the substance if it is wrong. `fleet_reply` resolves a member's blocked `fleet_send`; a peer's
|
||||
coord-id message is durable and non-blocking, so there is nothing for it to resolve. A peer that
|
||||
simply complies has thrown away the reason there are two of you.
|
||||
|
||||
### Member (worker or architect) — the turn contract
|
||||
|
||||
|
||||
@@ -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<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> 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.
|
||||
*
|
||||
* <p>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
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -77,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)));
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
@@ -136,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.
|
||||
@@ -183,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);
|
||||
@@ -215,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");
|
||||
@@ -269,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(
|
||||
@@ -278,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)));
|
||||
}
|
||||
|
||||
@@ -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<InboxMessage> 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
|
||||
@@ -351,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"),
|
||||
@@ -364,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");
|
||||
@@ -415,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)));
|
||||
}
|
||||
@@ -1254,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");
|
||||
|
||||
|
||||
Reference in New Issue
Block a user