From 6058b8472b7013620b0a63d709be32dc42638b8b Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 24 Aug 2026 17:39:29 +0200 Subject: [PATCH] lead comms: wire LeadMailbox into the daemon so lead-to-lead messages flow MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The sibling ticket landed the mechanism (LeadMailbox, LeadMessage, the coordinator: config block) but nothing opened it, nothing sent through it, and nothing read it. This is the wiring. - LeadChannel: a small interface LeadMailbox now implements (publish/peek/ack plus a selfCoordId() accessor). It exists so FleetMcp and the receive loop can be tested with a fake instead of a live broker. LeadMailbox's AMQP logic is untouched — the diff is the implements clause, four @Override marks and the accessor. - Fleetd.openLeadMailbox: opens this daemon's mailbox after the reply inbox is selected, with the same env-injected seam selectReplyInbox uses. Every "off" path returns null and the daemon still starts: no coordinator block (silent), a uriEnv that does not resolve (INFO), a configured broker with no selfId (WARN — a mailbox is named after the coord-id that owns it), or a broker that refuses at boot (WARN, credentials stripped). Closed in the ordered shutdown hook, after the loop that reads it has stopped. - fleet_send{coordId}: publishes a LeadMessage(from=selfCoordId, to=coordId) to the peer's mailbox and returns the broker-confirmed receipt. coordId is mutually exclusive with sessionId/turnId and is rejected by name rather than resolved by precedence. An unroutable/nacked/timed-out publish comes back as a tool error naming the coordId, never a crash. The worker send/reply path is not touched. - LeadCoordLoop: the receive half. Each tick peeks the mailbox, resolves the local lead pane, and — only at a turn boundary — injects "[lead ] " and acks. Anything not delivered stays unacked and is retried, so a message is never dropped; one message per tick, so every delivery is gated on a status read that already saw the previous one. - fleet_list reports {selfId, configured} when coordination is on, so an operator can find the coord-id a peer must use to reach them. Omitted entirely when it is off. Tests: 20 new hermetic tests (no broker) across routing, delivery and startup selection. mvn clean install: Tests run: 944, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 111 +++++++++- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 107 +++++++++- .../java/dev/ltms/fleet/msg/LeadChannel.java | 40 ++++ .../dev/ltms/fleet/msg/LeadCoordLoop.java | 193 ++++++++++++++++++ .../java/dev/ltms/fleet/msg/LeadMailbox.java | 11 +- .../fleet/FleetdLeadMailboxSelectionTest.java | 144 +++++++++++++ .../ltms/fleet/mcp/FleetMcpLeadCoordTest.java | 96 +++++++++ .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 29 +++ .../dev/ltms/fleet/msg/FakeLeadChannel.java | 75 +++++++ .../dev/ltms/fleet/msg/LeadCoordLoopTest.java | 150 ++++++++++++++ 10 files changed, 951 insertions(+), 5 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java create mode 100644 bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java create mode 100644 bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java create mode 100644 bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java create mode 100644 bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java create mode 100644 bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java diff --git a/bridged/src/main/java/dev/ltms/fleet/Fleetd.java b/bridged/src/main/java/dev/ltms/fleet/Fleetd.java index 65dd744..076a946 100644 --- a/bridged/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/bridged/src/main/java/dev/ltms/fleet/Fleetd.java @@ -31,6 +31,9 @@ import dev.ltms.fleet.mcp.LsofPeerPidLookup; import dev.ltms.fleet.mcp.LsofProcessCwdLookup; import dev.ltms.fleet.msg.AmqpReplyInbox; import dev.ltms.fleet.msg.InMemoryReplyInbox; +import dev.ltms.fleet.msg.LeadChannel; +import dev.ltms.fleet.msg.LeadCoordLoop; +import dev.ltms.fleet.msg.LeadMailbox; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.msg.ReplyInbox; @@ -60,6 +63,7 @@ import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; @@ -79,6 +83,14 @@ public final class Fleetd { /** CB-504: how long to wait at startup for herdr's socket before serving degraded. */ private static final long HERDR_WAIT_SECONDS = 30; + /** + * CB-635: how often the lead coordination loop looks for peer messages. A few seconds — slow + * enough that an idle fleet is not polling a broker in a tight loop, fast enough that a peer + * lead's message is not left sitting once the local lead reaches a turn boundary. The mailbox + * pushes into the loop's held set on its own consumer thread, so this interval bounds only the + * pane delivery, never the receive. + */ + private static final long LEAD_COORD_INTERVAL_MS = 3_000L; private static final long HERDR_WAIT_POLL_MILLIS = 500; /** @@ -375,6 +387,12 @@ public final class Fleetd { // unusable), bridged stays soft-state on the in-memory inbox. The AMQP inbox owns a broker // connection, so keep the reference to close it in the ordered shutdown hook. final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open); + // CB-635: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate + // broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator: + // block this is null and every lead path below is simply not wired, which is exactly the + // behaviour before this ticket. It owns a broker connection, so keep the reference for the + // ordered shutdown hook. + final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open); // CB-307: learn the primary's terminal from orchestration tool calls (or pin from config). // The pin also feeds CallerResolver below: a primary running inside a herdr pane would // otherwise resolve as a worker and be refused every orchestration tool. @@ -503,7 +521,27 @@ public final class Fleetd { new FleetMcp.QuarantineSource(profile -> { var configured = config.get().profiles().get(profile); return configured == null ? null : configured.effectiveCredentialId(); - }, quarantine)); + }, quarantine), + leadMailbox); + + // CB-635: the receive half. Only constructed when a lead mailbox actually opened — with no + // coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is + // created and no thread runs. It reads the SAME live lead supplier the injector's + // deliverability gate does, so a lead found by the tab scan after startup is reachable + // without a restart. + final LeadCoordLoop leadCoordLoop; + final ScheduledExecutorService leadCoordSchedulerRef; + if (leadMailbox != null) { + var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r -> + Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r)); + leadCoordLoop = new LeadCoordLoop(leadMailbox, agents, leads, leadCoordScheduler, + LEAD_COORD_INTERVAL_MS); + leadCoordLoop.start(); + leadCoordSchedulerRef = leadCoordScheduler; + } else { + leadCoordLoop = null; + leadCoordSchedulerRef = null; + } // CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an // upgraded daemon behaves exactly as before — the file is read once at boot and never again. @@ -524,6 +562,8 @@ public final class Fleetd { messages.close(); pushLoop.close(); if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler + if (leadCoordLoop != null) leadCoordLoop.close(); // CB-635: stop delivering peer-lead messages + if (leadCoordSchedulerRef != null) leadCoordSchedulerRef.shutdownNow(); if (healthMonitor != null) healthMonitor.stop(); if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file mcp.close(); @@ -536,6 +576,15 @@ public final class Fleetd { log.debug("reply inbox close: {}", e.toString()); } } + // CB-635: the coordination connection goes with it — after the loop that reads it has + // stopped, so no tick can be mid-ack against a closed channel. + if (leadMailbox != null) { + try { + leadMailbox.close(); + } catch (Exception e) { + log.debug("lead mailbox close: {}", e.toString()); + } + } herdr.close(); })); @@ -575,6 +624,66 @@ public final class Fleetd { ReplyInbox open(String uri, int prefetch); } + /** Injection seam for {@link #openLeadMailbox}: production binds {@link LeadMailbox#open}. */ + @FunctionalInterface + interface LeadMailboxOpener { + LeadMailbox open(String uri, String selfCoordId, int prefetch); + } + + /** + * CB-635: open this daemon's lead-to-lead mailbox, or return {@code null} to leave the feature + * off. Package-private and env-injected for the same reason as {@link #selectReplyInbox}: the + * selection is then testable without a broker or a mutable process environment. + * + *

Every "off" path returns {@code null}, and each says why at the level it deserves: + * + *

    + *
  • no {@code coordinator:} block — silent. Lead coordination is opt-in; an operator who + * never configured it does not need to be told it is off on every boot.
  • + *
  • a block whose {@code uriEnv} does not resolve — INFO, the same "you moved to the secret + * store and the variable is not there" case {@code selectReplyInbox} warns about.
  • + *
  • a configured broker but no {@code selfId} — WARN. This one is a half-finished config: a + * mailbox is named after the coord-id that owns it, so with no id there is no queue to own + * and no {@code from} to send as. Loud, because the operator plainly intended the feature.
  • + *
  • the broker refuses at boot — WARN, and carry on. Mirrors {@code openAmqpOrFallback}: a + * coordination broker that is down must never take a whole fleet's daemon with it, and the + * fleet still works exactly as it did before this feature existed.
  • + *
+ */ + static LeadMailbox openLeadMailbox(FleetConfig.Coordinator coordinator, Map env, + LeadMailboxOpener opener) { + if (coordinator == null) { + return null; // opt-in: nothing configured, nothing to say + } + String uri = coordinator.effectiveUri(env); + if (uri == null) { + log.info("lead coordination: OFF — coordinator{} has no usable broker uri", + coordinator.uriEnv() == null ? "" : ".uriEnv=" + coordinator.uriEnv()); + return null; + } + if (coordinator.selfId() == null || coordinator.selfId().isBlank()) { + log.warn("coordinator.selfId is unset — lead coordination is OFF. A lead mailbox is the " + + "queue named after the coord-id that owns it, so with no id there is nothing to " + + "own and no sender identity to publish as. Set coordinator.selfId to a name that " + + "is unique across every daemon sharing {} and restart bridged.", + stripCredentials(uri)); + return null; + } + try { + LeadMailbox mailbox = opener.open(uri, coordinator.selfId(), coordinator.prefetchOrDefault()); + log.info("lead coordination: ON as coord-id {} (prefetch={})", + coordinator.selfId(), coordinator.prefetchOrDefault()); + return mailbox; + } catch (IllegalStateException e) { + log.warn("cannot reach the AMQP coordination broker ({}) — lead-to-lead messaging is OFF " + + "for this process lifetime. fleet_send{{coordId}} will report it as " + + "not configured, and peer messages already queued stay on the broker until " + + "a restart picks them up. Reason: {}", + stripCredentials(uri), reasonOf(e)); + return null; + } + } + /** * CB-151/152: pick the reply inbox. A usable broker — a literal {@code uri}, or a {@code * uriEnv} whose variable resolves (both read from {@code env}) — selects the durable AMQP inbox. diff --git a/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index fe67ef2..dd1bcae 100644 --- a/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -11,6 +11,8 @@ import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; import dev.ltms.fleet.inject.MemberPresence; import dev.ltms.fleet.herdr.HerdrException; +import dev.ltms.fleet.msg.LeadChannel; +import dev.ltms.fleet.msg.LeadMessage; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.peer.PeerUnreachableException; @@ -37,6 +39,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Set; +import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.function.BiFunction; import java.util.function.Function; @@ -97,6 +100,8 @@ public final class FleetMcp { private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; private final QuarantineSource quarantine; + /** CB-635: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */ + private final LeadChannel leadChannel; /** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */ public record CapacitySource(Function liveCount, Function maxLoad, @@ -131,6 +136,21 @@ public final class FleetMcp { ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine) { + this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity, + healthCoverage, quarantine, null); + } + + /** + * As above, with this daemon's lead-to-lead channel (CB-635). {@code leadChannel} is + * {@code null} whenever no {@code coordinator:} block is configured or its broker could not be + * reached at boot — cross-daemon lead messaging is simply off, and {@code fleet_send{coordId}} + * says so rather than failing obscurely. + */ + public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, LeadChannel leadChannel) { + this.leadChannel = leadChannel; this.capacity = capacity; this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); this.healthCoverage = healthCoverage; @@ -177,6 +197,13 @@ public final class FleetMcp { String target = str(a, "sessionId"); String content = str(a, "content"); String turnId = str(a, "turnId"); + String coordId = str(a, "coordId"); + if (coordId != null && !coordId.isBlank()) { + // CB-635: a peer LEAD on another daemon, addressed by coord-id over the shared + // coordination broker. Checked before the turnId branch so a call that sets both + // is rejected as the conflict it is, rather than silently taking one route. + return sendToLead(leadChannel, coordId, content, target, turnId); + } if (turnId != null && !turnId.isBlank()) { // Answering a worker's fleet_ask (CB-205): resolve its blocked question and // block for the worker's reply as it resumes the same turn. This is the same @@ -257,7 +284,8 @@ public final class FleetMcp { if (denied != null) return denied; return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, callers == null ? Map.of() : callers.leads(), - callerTerminal(exchange)); + callerTerminal(exchange), + leadChannel == null ? null : leadChannel.selfCoordId()); }; BiFunction 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> leads; + private final ScheduledExecutorService scheduler; + private final long intervalMs; + + private volatile boolean running; + + /** + * @param channel this daemon's own lead mailbox + * @param agents herdr control, for the status gate and the pane injection + * @param leads live {@code terminal_id → name} view of the leads this daemon recognises — + * read through the supplier on every tick, never snapshotted, so a lead found by + * the tab scan after startup becomes reachable without a restart + * @param scheduler the loop's own scheduler; the caller owns its shutdown + * @param intervalMs how long between ticks + */ + public LeadCoordLoop(LeadChannel channel, AgentControl agents, Supplier> leads, + ScheduledExecutorService scheduler, long intervalMs) { + this.channel = channel; + this.agents = agents; + this.leads = leads; + this.scheduler = scheduler; + this.intervalMs = intervalMs; + } + + /** Begin ticking. Idempotent-ish: calling it twice would schedule two chains, so call it once. */ + public void start() { + running = true; + log.info("lead coordination: delivering peer messages for coord-id {} every {}ms", + channel.selfCoordId(), intervalMs); + scheduleNext(); + } + + /** Stop ticking. In-flight work finishes; nothing further is scheduled. */ + public void close() { + running = false; + } + + private void scheduleNext() { + if (!running) { + return; + } + scheduler.schedule(this::tickAndReschedule, intervalMs, TimeUnit.MILLISECONDS); + } + + private void tickAndReschedule() { + try { + tick(); + } catch (RuntimeException e) { + // Never let one bad tick end the chain — the next one re-reads everything from scratch. + log.warn("lead coordination tick failed: {}", e.toString()); + } + scheduleNext(); + } + + /** + * One tick: deliver at most one held peer message into the local lead's pane and ack it. + * Package-private so a test drives it directly rather than waiting on the scheduler. + */ + void tick() { + List held; + try { + held = channel.peek(); + } catch (RuntimeException e) { + log.debug("lead coordination: cannot read the mailbox this tick: {}", e.toString()); + return; + } + if (held.isEmpty()) { + return; + } + String lead = resolveLocalLead(); + if (lead == null) { + // Left unacked on purpose: the broker keeps holding it until a lead pane exists. + log.debug("lead coordination: {} message(s) waiting but no local lead pane to deliver to", + held.size()); + return; + } + AgentStatus status; + try { + status = agents.status(lead); + } catch (RuntimeException e) { + log.debug("lead coordination: status check failed for lead {}, retrying next tick", lead, e); + return; + } + if (!status.injectable()) { + log.debug("lead coordination: lead {} is {} (not injectable), holding {} message(s)", + lead, status, held.size()); + return; + } + LeadMessage msg = held.getFirst(); + try { + agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content())); + } catch (RuntimeException e) { + // Not delivered, so not acked — the broker still has it for the next tick. + log.warn("lead coordination: failed to deliver message {} from {} to lead {}: {}", + msg.msgId(), msg.from(), lead, e.toString()); + return; + } + try { + channel.ack(msg.msgId()); + } catch (RuntimeException e) { + // Delivered but not acked: it will be redelivered, which the javadoc calls out as the + // deliberate direction of this trade. + log.warn("lead coordination: delivered message {} but could not ack it: {}", + msg.msgId(), e.toString()); + return; + } + log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead); + } + + /** + * Which local pane a peer's message is for. The mailbox's {@code selfCoordId} is this daemon's + * one lead identity, so there is exactly one right answer — this only has to find it: + * + *

    + *
  1. a lead whose configured name equals {@code selfCoordId} — the explicit, unambiguous case;
  2. + *
  3. otherwise the sole lead, when this daemon recognises exactly one;
  4. + *
  5. otherwise nothing, and the message waits.
  6. + *
+ * + *

Step 3 is deliberate rather than a guess-the-lead fallback. Picking one of several leads + * arbitrarily would type a peer's message into a pane it was not addressed to, and the message + * would then be acked and gone. Leaving it held costs a delay and nothing else. + */ + private String resolveLocalLead() { + Map known = leads.get(); + if (known.isEmpty()) { + return null; + } + String self = channel.selfCoordId(); + for (var entry : known.entrySet()) { + if (entry.getValue() != null && entry.getValue().equals(self)) { + return entry.getKey(); + } + } + if (known.size() == 1) { + return known.keySet().iterator().next(); + } + log.warn("lead coordination: {} leads are known and none is named \"{}\" — cannot tell which " + + "pane a peer message is for; name one lead after coordinator.selfId to fix this", + known.size(), self); + return null; + } +} diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java index 4ea4d6a..f3c8e68 100644 --- a/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -61,7 +61,7 @@ import java.util.concurrent.TimeoutException; * redelivery) and any publish still awaiting its confirm is failed rather than left to idle out * the confirm timeout against a sequence number that means nothing on the new channel. */ -public final class LeadMailbox implements AutoCloseable { +public final class LeadMailbox implements LeadChannel, AutoCloseable { private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class); @@ -195,6 +195,7 @@ public final class LeadMailbox implements AutoCloseable { * confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish, * not a lost-and-forgotten one. */ + @Override public void publish(String toCoordId, LeadMessage msg) { byte[] body; try { @@ -239,7 +240,14 @@ public final class LeadMailbox implements AutoCloseable { } } + /** The coord-id whose mailbox this instance owns — the {@code from} of everything it publishes. */ + @Override + public String selfCoordId() { + return selfCoordId; + } + /** Non-destructive FIFO snapshot of this mailbox's currently-held messages. */ + @Override public List peek() { synchronized (held) { return held.values().stream().map(Held::message).toList(); @@ -254,6 +262,7 @@ public final class LeadMailbox implements AutoCloseable { } /** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */ + @Override public void ack(String msgId) { Held h; synchronized (held) { diff --git a/bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java b/bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java new file mode 100644 index 0000000..9cac1eb --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java @@ -0,0 +1,144 @@ +package dev.ltms.fleet; + +import ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.msg.LeadMailbox; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-635: the daemon decides whether lead-to-lead messaging is on in + * {@link Fleetd#openLeadMailbox}, not in the config record — so testing + * {@code Coordinator.isConfigured()} alone would pass even if {@code Fleetd} never honoured it. + * These drive the real selection with an injected env map and an injected opener, so no broker is + * involved and no process environment is mutated. + * + *

The invariant every case shares: the feature turns itself OFF, never takes the daemon down. + * A fleet whose coordination broker is missing, half-configured or unreachable must still start and + * still work exactly as it did before this feature existed. + */ +class FleetdLeadMailboxSelectionTest { + + private static final String SECRET = "c00rdPw"; + private static final String RESOLVED_URI = "amqp://user:" + SECRET + "@coord.example:5672/coord"; + + /** Fake opener: records what it was offered, or fails as an unreachable broker would. */ + private static final class RecordingOpener implements Fleetd.LeadMailboxOpener { + String offeredUri; + String offeredSelfId; + int offeredPrefetch = -1; + boolean unreachable; + + @Override + public LeadMailbox open(String uri, String selfCoordId, int prefetch) { + this.offeredUri = uri; + this.offeredSelfId = selfCoordId; + this.offeredPrefetch = prefetch; + if (unreachable) { + throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, + new java.net.ConnectException("Connection refused")); + } + // A real LeadMailbox needs a live connection; nothing here dereferences the result + // beyond a null check, so the "reachable" cases assert on what was OFFERED instead. + return null; + } + } + + private static ListAppender captureFleetdLogs() { + Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + return appender; + } + + private static String joined(ListAppender appender, Level level) { + return appender.list.stream().filter(e -> e.getLevel() == level) + .map(ILoggingEvent::getFormattedMessage).reduce("", (a, b) -> a + "\n" + b); + } + + @Test + void noCoordinatorBlockLeavesTheFeatureOffSilently() { + var appender = captureFleetdLogs(); + var opener = new RecordingOpener(); + + assertNull(Fleetd.openLeadMailbox(null, Map.of(), opener)); + + assertNull(opener.offeredUri, "nothing configured means nothing is opened"); + assertEquals("", joined(appender, Level.WARN), + "an opt-in feature nobody asked for must not warn on every boot"); + } + + @Test + void opensTheMailboxWhenAUriAndSelfIdAreConfigured() { + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null); + + Fleetd.openLeadMailbox(coordinator, Map.of(), opener); + + assertEquals(RESOLVED_URI, opener.offeredUri); + assertEquals("mac-opus", opener.offeredSelfId); + assertEquals(LeadMailbox.DEFAULT_PREFETCH, opener.offeredPrefetch, + "an unset prefetch takes the mailbox's own default, not zero"); + } + + @Test + void honoursUriEnvOverALiteralUri() { + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator("amqp://stale:stale@old:5672/x", "COORD_URI", + "mac-opus", 8); + + Fleetd.openLeadMailbox(coordinator, Map.of("COORD_URI", RESOLVED_URI), opener); + + assertEquals(RESOLVED_URI, opener.offeredUri, "the secret store wins over clear text"); + assertEquals(8, opener.offeredPrefetch); + } + + @Test + void turnsOffWhenUriEnvDoesNotResolve() { + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null); + + assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener)); + + assertNull(opener.offeredUri); + } + + @Test + void warnsAndStaysOffWhenSelfIdIsMissing() { + var appender = captureFleetdLogs(); + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null); + + assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener)); + + assertNull(opener.offeredUri, "a mailbox with no owning coord-id has no queue to declare"); + String warns = joined(appender, Level.WARN); + assertTrue(warns.contains("coordinator.selfId"), () -> "say which key is missing: " + warns); + assertFalse(warns.contains(SECRET), () -> "the URI's password must never be logged: " + warns); + } + + @Test + void warnsAndStaysOffWhenTheBrokerIsUnreachableAtBoot() { + var appender = captureFleetdLogs(); + var opener = new RecordingOpener(); + opener.unreachable = true; + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null); + + assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener), + "a down coordination broker turns the feature off; it must never take the daemon down"); + + String warns = joined(appender, Level.WARN); + assertTrue(warns.contains("coord.example"), () -> "name the host that failed: " + warns); + assertFalse(warns.contains(SECRET), () -> "with credentials stripped: " + warns); + assertTrue(warns.contains("Connection refused"), () -> "and the real reason: " + warns); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java new file mode 100644 index 0000000..10fe555 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java @@ -0,0 +1,96 @@ +package dev.ltms.fleet.mcp; + +import dev.ltms.fleet.msg.FakeLeadChannel; +import dev.ltms.fleet.msg.LeadMessage; +import io.modelcontextprotocol.spec.McpSchema; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-635: the {@code fleet_send{coordId}} route — a message to a PEER LEAD on another daemon, + * published to its durable mailbox instead of typed into a pane this daemon can reach. + * + *

Hermetic: {@link FakeLeadChannel} replaces the AMQP-backed {@code LeadMailbox}, so these run + * with no broker. What is checked here is routing and refusal — which envelope goes out, and which + * calls are rejected before anything is sent. + */ +class FleetMcpLeadCoordTest { + + private static final String SELF = "mac-opus"; + private static final String PEER = "fleet01-lead"; + + private static String textOf(McpSchema.CallToolResult r) { + return ((McpSchema.TextContent) r.content().getFirst()).text(); + } + + @Test + void publishesAnEnvelopeAddressedFromThisDaemonToThePeer() { + var channel = new FakeLeadChannel(SELF); + + McpSchema.CallToolResult res = + FleetMcp.sendToLead(channel, PEER, "you own the auth layer, I own the config one", null, null); + + assertFalse(res.isError(), () -> "expected a success result, got: " + textOf(res)); + assertEquals(1, channel.published().size()); + LeadMessage sent = channel.published().getFirst(); + assertEquals(SELF, sent.from(), "the sender is this daemon's own coord-id, never an argument"); + assertEquals(PEER, sent.to()); + assertEquals("you own the auth layer, I own the config one", sent.content()); + assertNotNull(sent.msgId()); + assertFalse(sent.msgId().isBlank(), "the envelope carries an idempotency id for dedup on redelivery"); + assertTrue(textOf(res).contains(PEER), "the receipt names the peer it reached"); + } + + @Test + void rejectsCoordIdTogetherWithSessionId() { + var channel = new FakeLeadChannel(SELF); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", "term_worker", null); + + assertTrue(res.isError()); + assertTrue(textOf(res).contains("coordId"), () -> textOf(res)); + assertTrue(textOf(res).contains("sessionId"), () -> "the error must name the conflict: " + textOf(res)); + assertEquals(0, channel.published().size(), "an ambiguous call must send nothing at all"); + } + + @Test + void rejectsCoordIdTogetherWithTurnId() { + var channel = new FakeLeadChannel(SELF); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, "turn_7"); + + assertTrue(res.isError()); + assertTrue(textOf(res).contains("turnId"), () -> textOf(res)); + assertEquals(0, channel.published().size()); + } + + @Test + void saysSoPlainlyWhenLeadCoordinationIsNotConfigured() { + McpSchema.CallToolResult res = FleetMcp.sendToLead(null, PEER, "hi", null, null); + + assertTrue(res.isError()); + assertTrue(textOf(res).contains("lead coordination is not configured"), () -> textOf(res)); + assertTrue(textOf(res).contains("coordinator"), () -> "point at the config block to add: " + textOf(res)); + } + + @Test + void reportsAnUnreachablePeerAsAToolErrorNamingIt() { + var channel = new FakeLeadChannel(SELF) + .failPublishWith("lead message m1 was returned as unroutable (mailbox not owned)"); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); + + assertTrue(res.isError(), "a black-holed message must never be reported as delivered"); + assertTrue(textOf(res).contains(PEER), () -> "the error must name the coordId: " + textOf(res)); + assertTrue(textOf(res).contains("unroutable"), () -> "and carry the broker's reason: " + textOf(res)); + } + + @Test + void requiresContent() { + var channel = new FakeLeadChannel(SELF); + + assertTrue(FleetMcp.sendToLead(channel, PEER, " ", null, null).isError()); + assertEquals(0, channel.published().size()); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index f21fc6c..f99872c 100644 --- a/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -527,6 +527,35 @@ class FleetMcpTest { assertTrue(out.contains("\"liveStatus\":\"unknown\""), out); } + @Test + void listReportsThisDaemonsOwnCoordIdWhenLeadCoordinationIsOn() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, + FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), Map.of(), "", "mac-opus"); + + String out = textOf(res); + // There is no peer-discovery surface yet, so this row answers the one question an operator + // cannot answer any other way: which coord-id a peer must use to reach ME. + assertTrue(out.contains("\"coordinator\""), out); + assertTrue(out.contains("\"selfId\":\"mac-opus\""), out); + } + + @Test + void listOmitsTheCoordinatorRowWhenLeadCoordinationIsOff() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), ""); + + assertFalse(textOf(res).contains("coordinator"), + "an ordinary fleet's output must be unchanged by this feature"); + } + @Test void capacityUsesThePlacementLiveCount() { FakeHerdr h = new FakeHerdr(); diff --git a/bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java b/bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java new file mode 100644 index 0000000..38ee2e9 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java @@ -0,0 +1,75 @@ +package dev.ltms.fleet.msg; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * Hermetic stand-in for {@link LeadChannel}: an in-memory mailbox that records what was published + * and what was acked, with no broker anywhere. + * + *

Its whole point is that {@link LeadMailbox} — the production implementation — owns a live AMQP + * connection, so every test of the code AROUND it would otherwise need a broker and end up tagged + * {@code contract}. The broker round trip is covered once, by {@code LeadMailboxTest}; the routing + * ({@code FleetMcp}) and the delivery ({@link LeadCoordLoop}) are covered here. + * + *

Thread-safe: {@link LeadCoordLoop} calls it from its own scheduler thread while a test reads + * the recorded lists. + */ +public final class FakeLeadChannel implements LeadChannel { + + private final String selfCoordId; + private final List held = Collections.synchronizedList(new ArrayList<>()); + private final List published = Collections.synchronizedList(new ArrayList<>()); + private final List acked = Collections.synchronizedList(new ArrayList<>()); + /** When set, every {@link #publish} throws it — the unroutable/nacked/timed-out peer. */ + private volatile IllegalStateException publishFailure; + + public FakeLeadChannel(String selfCoordId) { + this.selfCoordId = selfCoordId; + } + + /** Make every publish fail as an unreachable peer would. */ + public FakeLeadChannel failPublishWith(String message) { + this.publishFailure = new IllegalStateException(message); + return this; + } + + /** Put a message in this mailbox as if a peer had sent it. */ + public FakeLeadChannel hold(LeadMessage m) { + held.add(m); + return this; + } + + @Override + public void publish(String toCoordId, LeadMessage m) { + if (publishFailure != null) { + throw publishFailure; + } + published.add(m); + } + + @Override + public List peek() { + return List.copyOf(held); + } + + @Override + public void ack(String msgId) { + held.removeIf(m -> m.msgId().equals(msgId)); + acked.add(msgId); + } + + @Override + public String selfCoordId() { + return selfCoordId; + } + + public List published() { + return List.copyOf(published); + } + + public List acked() { + return List.copyOf(acked); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java b/bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java new file mode 100644 index 0000000..e1150b9 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java @@ -0,0 +1,150 @@ +package dev.ltms.fleet.msg; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.ScheduledExecutorService; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Unit tests for {@link LeadCoordLoop} — the receive half of lead-to-lead messaging. Hermetic + * throughout: a {@link FakeLeadChannel} stands in for the mailbox and {@link FakeHerdr} for the + * pane, so no broker and no herdr daemon is involved. The broker round trip has its own + * {@code contract}-tagged test on {@link LeadMailbox}. + * + *

Every test drives {@link LeadCoordLoop#tick()} directly rather than waiting on the scheduler: + * the schedule itself is one {@code scheduler.schedule} call, while the decisions worth pinning — + * inject or hold, ack or leave unacked — all live in the tick. + */ +class LeadCoordLoopTest { + + private static final String SELF = "mac-opus"; + private static final String PEER = "fleet01-lead"; + private static final String LEAD_TERM = "term_lead"; + + /** No scheduler is needed: nothing here calls start(). */ + private static final ScheduledExecutorService NO_SCHEDULER = null; + + private static LeadCoordLoop loop(LeadChannel channel, FakeHerdr herdr, Map leads) { + return new LeadCoordLoop(channel, new AgentControl(herdr), () -> leads, NO_SCHEDULER, 3_000L); + } + + private static List prompts(FakeHerdr herdr) { + return herdr.calls.stream().filter(c -> c.method().equals("agent.prompt")).toList(); + } + + @Test + void deliversAHeldMessageToTheLeadPaneAndAcksIt() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked")); + var herdr = new FakeHerdr().agentStatus("idle"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(1, prompts(herdr).size(), "an injectable lead must receive the peer's message"); + @SuppressWarnings("unchecked") + Map params = (Map) prompts(herdr).getFirst().params(); + assertEquals(LEAD_TERM, params.get("target"), "delivered to the resolved local lead pane"); + assertEquals("[lead " + PEER + "] the merge is blocked", params.get("text"), + "the sender's coord-id is carried into the pane — the lead must know who to answer"); + assertEquals(List.of("m1"), channel.acked(), "a delivered message is acked off the broker"); + assertTrue(channel.peek().isEmpty(), "and is no longer held"); + } + + @Test + void leavesTheMessageUnackedWhenTheLeadIsMidTurn() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("working"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(0, prompts(herdr).size(), "never paste into a live turn"); + assertEquals(List.of(), channel.acked(), "an undelivered message must NOT be acked"); + assertEquals(1, channel.peek().size(), "it stays held for the next tick"); + } + + @Test + void leavesTheMessageUnackedWhenNoLeadPaneIsKnown() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle"); + + loop(channel, herdr, Map.of()).tick(); + + assertEquals(0, prompts(herdr).size()); + assertEquals(List.of(), channel.acked(), "nowhere to deliver is not a reason to drop it"); + assertEquals(1, channel.peek().size()); + } + + @Test + void leavesTheMessageUnackedWhenHerdrRefusesTheInjection() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("agent_not_found"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(List.of(), channel.acked(), "a send that threw delivered nothing, so nothing is acked"); + assertEquals(1, channel.peek().size()); + } + + @Test + void resolvesTheLeadByNameWhenSeveralAreKnown() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle"); + // Two leads on this daemon; only one carries the coord-id the mailbox is owned as. + var leads = new java.util.LinkedHashMap(); + leads.put("term_other", "some-other-lead"); + leads.put(LEAD_TERM, SELF); + + loop(channel, herdr, leads).tick(); + + @SuppressWarnings("unchecked") + Map params = (Map) prompts(herdr).getFirst().params(); + assertEquals(LEAD_TERM, params.get("target"), "the lead named after coordinator.selfId wins"); + } + + @Test + void holdsWhenSeveralLeadsAreKnownAndNoneCarriesTheCoordId() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle"); + var leads = new java.util.LinkedHashMap(); + leads.put("term_one", "lead-one"); + leads.put("term_two", "lead-two"); + + loop(channel, herdr, leads).tick(); + + assertEquals(0, prompts(herdr).size(), + "guessing a pane would type a peer's message into the wrong lead and then ack it"); + assertEquals(List.of(), channel.acked()); + } + + @Test + void deliversOneMessagePerTickSoEachIsGatedOnItsOwnStatusRead() { + var channel = new FakeLeadChannel(SELF) + .hold(new LeadMessage("m1", PEER, SELF, "first")) + .hold(new LeadMessage("m2", PEER, SELF, "second")); + var herdr = new FakeHerdr().agentStatus("idle"); + var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF)); + + loop.tick(); + assertEquals(List.of("m1"), channel.acked(), "FIFO: the oldest goes first"); + assertEquals(1, prompts(herdr).size(), "the second waits for a fresh status read"); + + loop.tick(); + assertEquals(List.of("m1", "m2"), channel.acked()); + assertEquals(2, prompts(herdr).size()); + assertTrue(channel.peek().isEmpty()); + } + + @Test + void anEmptyMailboxNeverTouchesHerdr() { + var channel = new FakeLeadChannel(SELF); + var herdr = new FakeHerdr().agentStatus("idle"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(0, herdr.calls.size(), "an idle fleet must not poll a pane's status every tick"); + } +}