stopHandler =
(exchange, req) -> {
@@ -585,6 +613,54 @@ public final class FleetMcp {
return null;
}
+ /**
+ * {@code fleet_send} carrying a {@code coordId} (CB-635): a message to a PEER LEAD, published to
+ * that lead's durable mailbox on the shared coordination broker. This is the only lead→lead path
+ * that crosses hosts — the existing pane-injection route can only reach a lead whose herdr socket
+ * this daemon shares.
+ *
+ * {@code coordId} is mutually exclusive with {@code sessionId} and {@code turnId}: those two
+ * address a worker session owned by this daemon, a coord-id addresses a lead owned by
+ * another one, and there is no sensible reading of a call that sets both. Rejected by name rather
+ * than resolved by precedence, so a caller that meant the other route learns it instead of having
+ * its message quietly go somewhere else.
+ *
+ *
The publish is synchronous and confirmed by the broker, so the result is a real delivery
+ * receipt rather than a hopeful one. Its failure — nobody owns {@code coordId}'s mailbox, the
+ * broker nacked, or the confirm timed out — arrives as {@link IllegalStateException} and is
+ * turned into a tool error naming the coord-id. It is never allowed to escape as a crash: an
+ * unreachable peer is an ordinary outcome of addressing a fleet you do not control.
+ *
+ * @param leadChannel this daemon's channel, or {@code null} when no coordinator is configured
+ */
+ static McpSchema.CallToolResult sendToLead(LeadChannel leadChannel, String coordId, String content,
+ String sessionId, String turnId) {
+ if (!isBlank(sessionId) || !isBlank(turnId)) {
+ String conflict = !isBlank(sessionId) ? "sessionId" : "turnId";
+ return error("coordId and " + conflict + " are mutually exclusive: coordId addresses a peer "
+ + "LEAD on another daemon over the coordination broker, while " + conflict
+ + " addresses a worker session on this one. Pass exactly one.");
+ }
+ if (isBlank(content)) {
+ return error("content is required");
+ }
+ if (leadChannel == null) {
+ return error("lead coordination is not configured (no coordinator: block) — cannot send to "
+ + "peer lead \"" + coordId + "\". Add a coordinator: block with a shared broker uri "
+ + "and this daemon's selfId, then restart bridged.");
+ }
+ LeadMessage msg = new LeadMessage(UUID.randomUUID().toString(), leadChannel.selfCoordId(),
+ coordId, content);
+ try {
+ leadChannel.publish(coordId, msg);
+ } catch (IllegalStateException e) {
+ return error("cannot deliver to peer lead \"" + coordId + "\": " + e.getMessage()
+ + ". Check that a daemon is running with coordinator.selfId=\"" + coordId
+ + "\" and is connected to the same coordination broker.");
+ }
+ return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")");
+ }
+
/** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
if (!isBlank(target)) {
@@ -870,6 +946,22 @@ public final class FleetMcp {
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map leads, String selfTerm) {
+ return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null);
+ }
+
+ /**
+ * As above, additionally reporting this daemon's own lead coordination id (CB-635) when one is
+ * configured and its channel opened. There is no peer-discovery surface yet — a lead addresses a
+ * peer by a coord-id it was told — so this row exists to answer the one question the operator
+ * cannot answer any other way: what is MY coord-id, the one a peer must use to reach me. It is
+ * omitted entirely when no coordinator is configured, so an ordinary fleet's output is unchanged.
+ *
+ * @param selfCoordId this daemon's coord-id, or {@code null} when lead coordination is off
+ */
+ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
+ CapacitySource capacity, HealthCoverageSource healthCoverage,
+ QuarantineSource quarantine, Map leads, String selfTerm,
+ String selfCoordId) {
try {
Map live = workers.list().stream()
.map(Agent.class::cast)
@@ -888,6 +980,9 @@ public final class FleetMcp {
Map result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
+ if (selfCoordId != null && !selfCoordId.isBlank()) {
+ result.put("coordinator", Map.of("selfId", selfCoordId, "configured", true));
+ }
if (capacity.available()) result.put("capacity", profiles.stream()
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
capacity.clock().getAsLong(), quarantine)).toList());
@@ -1015,7 +1110,9 @@ public final class FleetMcp {
"Delegate a task to a worker session. By default blocks until the worker replies and "
+ "returns its reply (or a 'still working / queued' note on timeout). Pass wait:false "
+ "for a long task to return a ticket immediately, then poll it with fleet_poll. To "
- + "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId.",
+ + "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId. "
+ + "To message a PEER LEAD on another daemon — possibly another host — pass its "
+ + "coordId instead; that is coordination, never a task.",
objectSchema(Map.of(
"sessionId", stringProp("The worker session id (herdr terminal_id) to delegate to"),
"content", stringProp("The task/message to send to the worker (or your answer, with turnId)"),
@@ -1023,7 +1120,11 @@ public final class FleetMcp {
"wait", Map.of("type", "boolean",
"description", "Block for the reply (default true); false returns a ticket to poll"),
"turnId", stringProp("When answering a worker's fleet_ask, its question turnId — "
- + "routes your answer back into the same turn (omit for a normal delegation)")),
+ + "routes your answer back into the same turn (omit for a normal delegation)"),
+ "coordId", stringProp("A peer LEAD's coordination id — delivers content to that "
+ + "lead's durable mailbox on the shared coordination broker, which works "
+ + "across hosts. Mutually exclusive with sessionId and turnId. Your own "
+ + "coordId is reported by fleet_list.")),
List.of("content")));
}
diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java
new file mode 100644
index 0000000..07664a9
--- /dev/null
+++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java
@@ -0,0 +1,40 @@
+package dev.ltms.fleet.msg;
+
+import java.util.List;
+
+/**
+ * The lead-to-lead message channel this daemon speaks, as its callers need it — one lead's own
+ * mailbox: publish to a peer's coord-id, look at what has arrived for me, and ack what I have
+ * delivered.
+ *
+ * Extracted from {@link LeadMailbox} purely as a seam. {@code LeadMailbox} is the one production
+ * implementation and owns a live AMQP connection, so a test that wanted to exercise the routing in
+ * {@code FleetMcp} or the delivery in {@link LeadCoordLoop} would have had to stand up a broker —
+ * which is exactly the kind of test that gets tagged {@code contract} and then does not run. With
+ * this interface both of those are hermetic: they inject a fake channel and assert on what was
+ * published, peeked and acked.
+ *
+ *
Note what is not here: {@code drain()} and {@code close()}. Draining is a convenience
+ * over peek+ack that no caller on this seam uses, and closing is the owner's job — {@code Fleetd}
+ * holds the concrete {@link LeadMailbox} for its shutdown hook and hands only this narrower view to
+ * everyone else.
+ */
+public interface LeadChannel {
+
+ /**
+ * Send {@code m} to {@code toCoordId}'s mailbox, blocking until the broker confirms it is
+ * durably queued. Throws {@link IllegalStateException} when it is not — unroutable (nobody owns
+ * that coord-id), nacked, or unconfirmed within the implementation's timeout. A caller must
+ * report that as a failed send, never as a delivered one.
+ */
+ void publish(String toCoordId, LeadMessage m);
+
+ /** Non-destructive FIFO snapshot of the messages held for this daemon's own coord-id. */
+ List peek();
+
+ /** Drop {@code msgId} from the held set and ack it on the broker. A no-op if it is not held. */
+ void ack(String msgId);
+
+ /** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
+ String selfCoordId();
+}
diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java
new file mode 100644
index 0000000..ad86a63
--- /dev/null
+++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java
@@ -0,0 +1,193 @@
+package dev.ltms.fleet.msg;
+
+import dev.ltms.fleet.herdr.AgentControl;
+import dev.ltms.fleet.herdr.AgentStatus;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.function.Supplier;
+
+/**
+ * The receive half of lead-to-lead messaging: a bounded background loop that takes what has arrived
+ * in this daemon's own {@link LeadChannel} mailbox and types it into the local lead's herdr pane.
+ *
+ * {@code fleet_send{coordId}} is the send half — it publishes to a peer daemon's mailbox and
+ * returns. Nothing on the receiving side reads that mailbox on its own, because the peer lead is an
+ * interactive agent, not a service that polls; this loop is what closes the gap.
+ *
+ *
Status-gated, exactly like {@link ReplyPushLoop}. A pane may only be injected
+ * into at a turn boundary ({@link AgentStatus#injectable()} — idle, blocked or done); pasting into
+ * a live turn corrupts it. So a tick that finds the lead busy simply does nothing and comes back
+ * later.
+ *
+ *
Ack only after delivery. A message is acked — removed from the broker — only
+ * once {@link AgentControl#send} has actually put it in the pane. Anything not delivered (no lead
+ * pane resolvable, lead mid-turn, herdr threw) stays unacked and is retried on the next tick, and
+ * survives a daemon restart because the broker still holds it. The cost of that choice is a
+ * possible duplicate — the send lands and the ack does not — which is the right way round: a peer
+ * lead seeing a message twice is a nuisance, a peer lead never seeing it at all is the failure this
+ * whole path exists to remove.
+ *
+ *
One message per tick. The loop delivers at most one held message per tick even
+ * when several are waiting. Injecting a second one immediately would mean acting on a status read
+ * taken before the first injection: that first paste starts a turn, and herdr does not
+ * report the pane as {@code working} the instant it does. Waiting for the next tick means every
+ * delivery is gated on a status read that already saw the previous one. A backlog therefore drains
+ * one message per {@code intervalMs}, in FIFO order.
+ */
+public final class LeadCoordLoop {
+
+ private static final Logger log = LoggerFactory.getLogger(LeadCoordLoop.class);
+
+ /** How an arriving peer message is rendered into the lead's pane — the sender's coord-id, then its text. */
+ static final String DELIVERY_FORMAT = "[lead %s] %s";
+
+ private final LeadChannel channel;
+ private final AgentControl agents;
+ private final Supplier