Compare commits

..

5 Commits

Author SHA1 Message Date
Dai Ha edabccd885 CB-637: relabel lead-comms (was colliding CB-635) + document coordId route in the CLAUDE.md tool table
The lead-to-lead wiring shipped with CB-635 in its comments, but CB-635 is
already the broker.uriEnv / unreachable-broker work. Relabel the mailbox +
fleet_send{coordId} + receive loop to CB-637 so a ticket number names one
feature. Add the cross-host peer-lead row to the primary intent->tool table
(kept byte-identical with the wiki template).
2026-08-24 17:48:51 +02:00
Dai Ha 1dbe3a03fc Merge lead-comms-wiring: wire LeadMailbox into the daemon (fleet_send{coordId} + receive loop) 2026-08-24 17:42:11 +02:00
Dai Ha 6058b8472b lead comms: wire LeadMailbox into the daemon so lead-to-lead messages flow
CI / contract (pull_request) Successful in 41s
CI / build (pull_request) Failing after 1m22s
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 <from>] <text>"
  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.
2026-08-24 17:39:29 +02:00
Dai Ha f29968c334 Merge autocompact-window: per-profile autoCompactWindow for claude-code + opencode members 2026-08-24 17:32:53 +02:00
Dai Ha 7b98cca967 LeadMailbox: durable leader-to-leader mailbox over a shared coordination vhost
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Failing after 1m35s
Adds the broker-side mechanism for lead-to-lead messages across daemons/hosts
(unit 1 of 2): a LeadMessage envelope carrying from/to coord-ids, an
AMQP-backed LeadMailbox modeled closely on AmqpReplyInbox (consume-and-hold,
deferred manual ack, confirm-mode publish, recovery handling), and a new
optional coordinator: config block (separate vhost from broker:, leader
traffic only). Config parsing + accessors only — FleetMcp/Fleetd/Injector/
MessageService and the send path are untouched; wiring is a separate ticket.
2026-08-24 17:23:42 +02:00
16 changed files with 1747 additions and 8 deletions
+2 -1
View File
@@ -112,7 +112,8 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
| Delegate (blocking) | `fleet_send{sessionId, content}` |
| Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` |
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
| Message a **peer lead** | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| 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 |
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
| Tear down a member | `fleet_stop{paneId}` |
+17
View File
@@ -667,6 +667,23 @@ guard:
# uriEnv: LAVINMQ_URI
# prefetch: 32
# Shared cross-host LEADER coordination broker. OMIT this block to leave lead-to-lead messaging
# off entirely (config-only in this ticket — nothing here wires it into a live LeadMailbox yet).
# This is a SEPARATE AMQP vhost from `broker:` above: member/worker inboxes always stay on the
# per-fleet `broker:` vhost, and this vhost carries only leader-to-leader traffic, so two fleets
# whose members must never see each other can still share one coordination vhost for their leads.
# uriEnv → name of a host env var holding the coordination AMQP URI, same convention as
# broker.uriEnv (keeps the credential out of fleetd.yaml). Wins over `uri` when set.
# selfId → this daemon's own lead coord-id — the name its mailbox is owned under
# (lead.<selfId>.inbox), e.g. "mac-opus" or "fleet01-lead". Must be globally unique
# across every daemon sharing this vhost.
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
# Default 32 when omitted.
# coordinator:
# uriEnv: LEAD_COORD_URI
# selfId: mac-opus
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
# pane — status-gated (only when injectable, never mid-turn) and bounded. Ack = drain: the loop
@@ -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-637: 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-637: 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-637: 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-637: 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-637: 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-637: 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.
*
* <p>Every "off" path returns {@code null}, and each says why at the level it deserves:
*
* <ul>
* <li>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.</li>
* <li>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.</li>
* <li>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.</li>
* <li>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.</li>
* </ul>
*/
static LeadMailbox openLeadMailbox(FleetConfig.Coordinator coordinator, Map<String, String> 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.
@@ -6,6 +6,7 @@ import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import dev.ltms.fleet.msg.AmqpReplyInbox;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.placement.PlacementPolicies;
import org.slf4j.Logger;
@@ -69,6 +70,11 @@ import java.util.Set;
* @param memberCredentials deny-by-default policy (CB-596) for which of the operator's own host
* credentials a spawned member's pane inherits. {@code null} (the block
* omitted) blocks nothing — see {@link MemberCredentials}.
* @param coordinator shared cross-host leader coordination broker: a SEPARATE AMQP vhost from
* {@link #broker} used only for lead-to-lead traffic (member/worker inboxes
* stay on {@code broker}'s vhost). {@code null} → no lead mailbox is opened.
* Config parsing + accessors only — nothing here wires it into a live
* {@code LeadMailbox}; that is a separate ticket. See {@link Coordinator}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -89,7 +95,20 @@ public record FleetConfig(
Auth auth,
ConfigReload configReload,
Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
MemberCredentials memberCredentials,
Coordinator coordinator) {
/** Back-compat form before the {@code coordinator:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, null);
}
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
@@ -690,6 +709,82 @@ public record FleetConfig(
}
}
/**
* Shared cross-host leader coordination broker: a {@link LeadMailbox} lets two leads on
* different daemons — possibly different hosts — exchange durable messages, which a herdr pane
* injection (how {@code fleet_send} reaches a lead today) cannot do at all. Its mere presence is
* config only in this ticket: nothing here opens a live {@code LeadMailbox} yet, that wiring is
* a separate ticket.
*
* <p>Deliberately a SEPARATE vhost from {@link Broker}, not a reuse of it. {@link Broker} is
* per-fleet — its queues are named by worker session id, and two fleets sharing one broker stay
* isolated by vhost (see {@code Two fleets share one LavinMQ}). Leader coordination is meant to
* cross exactly that boundary: two independently-owned fleets' leads talking to each other. Using
* the same vhost would either leak member traffic across the fleet boundary this is meant to
* cross, or force every fleet's members onto one shared vhost to get leader coordination — a
* second, dedicated vhost keeps "member inboxes stay fleet-local" true while still letting leads
* reach across fleets.
*
* @param uri AMQP connection URI for the coordination vhost, e.g.
* {@code amqp://guest:guest@127.0.0.1:5672/coord}. Blank/{@code null} ⇒ the
* coordinator block is treated as unconfigured. Ignored when {@code uriEnv} is set.
* @param uriEnv name of a host env var holding the AMQP URI, same convention as
* {@link Broker#uriEnv()} — keeps the credential out of the config file. Wins
* over {@code uri} whenever set. Blank/{@code null} ⇒ ignored.
* @param selfId this daemon's own lead coord-id — the name its {@code LeadMailbox} is owned
* under ({@code lead.<selfId>.inbox}), e.g. {@code "mac-opus"}. Blank/
* {@code null} ⇒ kept as {@code null} (no self id configured).
* @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒
* {@link LeadMailbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
public Coordinator {
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
}
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
public boolean hasUriEnv() {
return uriEnv != null && !uriEnv.isBlank();
}
/**
* True when a usable coordination broker URI is configured (an empty block does not enable
* it). Honors {@code uriEnv} first: if it names a variable that is unset or blank, the
* coordinator is <em>not</em> configured — a bare {@code uri} is only consulted when no
* {@code uriEnv} is set.
*/
public boolean isConfigured() {
return effectiveUri() != null;
}
/**
* The effective AMQP URI to connect with. {@code uriEnv} wins when set (both over
* {@code uri} and alone). When {@code uriEnv} names a variable that is unset or blank,
* returns {@code null} rather than falling back to {@code uri} — an operator who moved to
* the secret store must not silently drop back onto a stale clear-text URI. Returns the
* literal {@code uri} when no {@code uriEnv} is configured.
*/
public String effectiveUri() {
return effectiveUri(System.getenv());
}
/** As {@link #effectiveUri()}, reading the variable value from {@code env} (the injection seam). */
public String effectiveUri(Map<String, String> env) {
if (hasUriEnv()) {
String value = env.get(uriEnv);
return (value != null && !value.isBlank()) ? value : null;
}
return (uri != null && !uri.isBlank()) ? uri : null;
}
/** The prefetch to use, defaulting to {@link LeadMailbox#DEFAULT_PREFETCH} when unset. */
public int prefetchOrDefault() {
return (prefetch != null && prefetch > 0) ? prefetch : LeadMailbox.DEFAULT_PREFETCH;
}
}
/**
* Optional pinned primary terminal config (CB-307). When present with a non-blank
* {@code terminal}, the bridge uses this as the primary's herdr identity instead of
@@ -1224,7 +1319,7 @@ public record FleetConfig(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials");
"memberCredentials", "coordinator");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -1840,9 +1935,11 @@ public record FleetConfig(
// and an upgrade must not change what a running deployment's members inherit.
MemberCredentials mc = memberCredentials != null ? memberCredentials
: new MemberCredentials(null, List.of(), List.of());
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
return new FleetConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc);
quarantineCooldown, mc, coordinator);
}
/**
@@ -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-637: 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<String, Integer> liveCount, Function<String, Integer> 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-637). {@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-637: 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<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -585,6 +613,54 @@ public final class FleetMcp {
return null;
}
/**
* {@code fleet_send} carrying a {@code coordId} (CB-637): 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.
*
* <p>{@code coordId} is mutually exclusive with {@code sessionId} and {@code turnId}: those two
* address a worker session owned by <em>this</em> 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.
*
* <p>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<String, String> 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-637) 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<String, String> leads, String selfTerm,
String selfCoordId) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -888,6 +980,9 @@ public final class FleetMcp {
Map<String, Object> 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")));
}
@@ -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.
*
* <p>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.
*
* <p>Note what is <em>not</em> 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<LeadMessage> 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();
}
@@ -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.
*
* <p>{@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.
*
* <p><strong>Status-gated, exactly like {@link ReplyPushLoop}.</strong> 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.
*
* <p><strong>Ack only after delivery.</strong> 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.
*
* <p><strong>One message per tick.</strong> 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 <em>before</em> 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<Map<String, String>> 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<Map<String, String>> 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<LeadMessage> 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:
*
* <ol>
* <li>a lead whose configured name equals {@code selfCoordId} — the explicit, unambiguous case;</li>
* <li>otherwise the sole lead, when this daemon recognises exactly one;</li>
* <li>otherwise nothing, and the message waits.</li>
* </ol>
*
* <p>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<String, String> 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;
}
}
@@ -0,0 +1,431 @@
package dev.ltms.fleet.msg;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.NavigableMap;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/**
* Durable, AMQP-backed mailbox for lead-to-lead messages across daemons — including daemons on
* different hosts, where a herdr pane injection (how {@code fleet_send} reaches a lead today)
* cannot reach at all. The broker is the only medium two independently-owned daemons share, which
* is exactly why {@link AmqpReplyInbox}'s javadoc already calls out "one gateway may publish to an
* agent owned by another gateway" (CB-308 federation) as the reason {@code publish} and
* {@code own}/consume are separate operations there — this class leans on the same split.
*
* <p><strong>Single-target, unlike {@link AmqpReplyInbox}.</strong> {@code AmqpReplyInbox}
* multiplexes many workers' reply queues under one gateway connection. A {@code LeadMailbox}
* instance is simpler: it owns exactly <em>one</em> queue — this daemon's own
* {@code lead.<selfCoordId>.inbox} — declared and consumed the moment it is constructed. There is
* no {@code own}/{@code release} pair to call separately; a daemon either runs a {@code LeadMailbox}
* for its own coord-id, or it does not run one at all.
*
* <p><strong>Consume-and-hold with deferred manual ack</strong> — same mapping as
* {@code AmqpReplyInbox}. The constructor declares the durable queue and starts a manual-ack
* consumer that pulls persistent messages into an in-memory {@code held} map (keyed by
* {@link LeadMessage#msgId()}) but does not ack them. {@link #peek} returns a non-destructive
* snapshot; {@link #ack} acks the broker delivery-tag and drops the entry. A message that is never
* acked (a crash, a bounce) survives — the broker redelivers it to the next connection that owns
* the queue.
*
* <p><strong>Publishing does not imply owning.</strong> {@link #publish} sends to
* {@code lead.<toCoordId>.inbox} over a dedicated confirm-mode channel; it never declares that
* queue as owned and never attaches a consumer to it. A sender that has never opened its own
* {@code LeadMailbox} for {@code toCoordId} can still publish to it, exactly as CB-308 federation
* requires. Publish blocks for the broker's publisher confirm (persistent delivery, {@code
* mandatory=true}) and throws {@link IllegalStateException} on an unroutable return, a nack, or a
* timeout — the caller must not report success for a black-holed message.
*
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
* 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 LeadChannel, AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class);
private static final String QUEUE_PREFIX = "lead.";
private static final String QUEUE_SUFFIX = ".inbox";
/** The prefetch used when a caller does not pass an explicit value to {@link #open(String, String, int)}. */
public static final int DEFAULT_PREFETCH = 32;
/** How long {@link #publish} waits for its publisher confirm before failing the call. */
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
private static final ObjectMapper MAPPER = new ObjectMapper();
private final Connection connection;
private final String selfCoordId;
private final Channel channel;
/** All consume-channel operations (declare/consume/ack) serialize on this — a Channel is not thread-safe. */
private final Object channelLock = new Object();
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
/**
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
* so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
*/
private final Channel publishChannel;
private final Object publishChannelLock = new Object();
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the deserialized envelope. */
private record Held(long deliveryTag, LeadMessage message) {}
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
private static final class Pending {
final String msgId;
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
volatile boolean returned;
Pending(String msgId) {
this.msgId = msgId;
}
}
/**
* Connect to {@code uri} (the shared cross-host coordination vhost, e.g.
* {@code amqp://guest:guest@127.0.0.1:5672/coord}) and own {@code selfCoordId}'s mailbox, with
* {@link #DEFAULT_PREFETCH}.
*/
public static LeadMailbox open(String uri, String selfCoordId) {
return open(uri, selfCoordId, DEFAULT_PREFETCH);
}
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new LeadMailbox(factory.newConnection("bridged-lead-mailbox"), selfCoordId, prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
}
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId) {
this(connection, selfCoordId, DEFAULT_PREFETCH);
}
/** As above, with an explicit prefetch (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId, int prefetch) {
this.connection = connection;
this.selfCoordId = selfCoordId;
try {
this.channel = connection.createChannel();
// Bound the held backlog — must be set before basicConsume.
this.channel.basicQos(prefetch);
this.publishChannel = connection.createChannel();
this.publishChannel.confirmSelect();
this.publishChannel.addReturnListener(this::onReturn);
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
own();
} catch (IOException e) {
throw new IllegalStateException("cannot open AMQP channel", e);
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// confirm still in flight when the connection dropped is equally stale — fail it now rather
// than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
synchronized (held) {
held.clear();
}
failPendingPublishesOnRecovery();
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
}
@Override
public void handleRecoveryStarted(Recoverable recoverable) {
// no-op: we act once recovery completes
}
});
}
}
/** Declare + consume this daemon's own {@code lead.<selfCoordId>.inbox}. Called once, at construction. */
private void own() throws IOException {
String queue = queueName(selfCoordId);
synchronized (channelLock) {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack
}
log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId);
}
/**
* Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher
* confirm for it lands. Does <em>not</em> imply owning or consuming {@code toCoordId}'s queue.
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
* 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 {
body = MAPPER.writeValueAsBytes(msg);
} catch (JsonProcessingException e) {
throw new IllegalStateException("cannot serialize lead message " + msg.msgId(), e);
}
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.messageId(msg.msgId())
.deliveryMode(2) // persistent — survives a broker restart
.contentType("application/json")
.build();
Pending pending = new Pending(msg.msgId());
long seq;
synchronized (publishChannelLock) {
seq = publishChannel.getNextPublishSeqNo();
pendingBySeq.put(seq, pending);
pendingByMsgId.put(msg.msgId(), pending);
try {
publishChannel.basicPublish("", queueName(toCoordId), true, props, body);
} catch (IOException e) {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msg.msgId(), pending);
throw new IllegalStateException("cannot publish lead message to " + queueName(toCoordId), e);
}
}
try {
pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
} catch (TimeoutException e) {
throw new IllegalStateException("publish confirm for lead message " + msg.msgId() + " to "
+ queueName(toCoordId) + " timed out after " + CONFIRM_TIMEOUT_MS
+ "ms — broker may be unreachable or overloaded", e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting publish confirm for " + msg.msgId(), e);
} finally {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msg.msgId(), pending);
}
}
/** 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<LeadMessage> peek() {
synchronized (held) {
return held.values().stream().map(Held::message).toList();
}
}
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
public List<LeadMessage> drain() {
List<LeadMessage> snapshot = peek();
snapshot.forEach(m -> ack(m.msgId()));
return snapshot;
}
/** 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) {
h = held.remove(msgId);
}
if (h == null) {
return; // never held (or already acked) — no-op
}
try {
synchronized (channelLock) {
channel.basicAck(h.deliveryTag(), false);
}
} catch (IOException e) {
// Ack didn't reach the broker: restore the entry so a later ack (or a redelivery after
// reconnect) can retry. Keeps the at-least-once contract — a message is never silently lost.
synchronized (held) {
held.putIfAbsent(msgId, h);
}
throw new IllegalStateException("cannot ack lead message " + msgId, e);
}
}
private DeliverCallback deliverCallback() {
return (_, delivery) -> {
long tag = delivery.getEnvelope().getDeliveryTag();
LeadMessage msg;
try {
msg = MAPPER.readValue(delivery.getBody(), LeadMessage.class);
} catch (IOException e) {
// A malformed body can never be dedup-keyed or handed to a caller; ack it so the
// broker does not redeliver it forever, and log loudly since this should never happen
// for a producer that only ever calls publish(String, LeadMessage).
log.warn("dropping malformed lead-mailbox delivery (tag {}): {}", tag, e.toString());
synchronized (channelLock) {
channel.basicAck(tag, false);
}
return;
}
boolean duplicate;
synchronized (held) {
if (held.containsKey(msg.msgId())) {
duplicate = true;
} else {
held.put(msg.msgId(), new Held(tag, msg));
duplicate = false;
}
}
if (duplicate) {
// Redelivered duplicate: ack the new tag and drop it so the broker stops resending.
synchronized (channelLock) {
channel.basicAck(tag, false);
}
}
};
}
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm. */
private void onReturn(Return r) {
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
if (pending != null) {
pending.returned = true;
} else {
log.warn("AMQP return for lead message {} (routingKey={}, {} {}) with no matching in-flight publish"
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
r.getReplyText());
}
}
private void onAck(long seq, boolean multiple) {
resolveConfirm(seq, multiple, true);
}
private void onNack(long seq, boolean multiple) {
resolveConfirm(seq, multiple, false);
}
/**
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
* a return means "confirmed but never routed", not "durably queued".
*/
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
NavigableMap<Long, Pending> covered = multiple
? pendingBySeq.headMap(seq, true)
: pendingBySeq.subMap(seq, true, seq, true);
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
if (ack && !pending.returned) {
pending.confirmed.complete(null);
} else if (ack) {
pending.confirmed.completeExceptionally(new IllegalStateException(
"lead message " + pending.msgId + " was returned as unroutable (mailbox not owned)"));
} else {
pending.confirmed.completeExceptionally(new IllegalStateException(
"broker nacked publish of lead message " + pending.msgId));
}
}
}
/**
* Fail every publish still awaiting its confirm — their sequence numbers are stale after
* recovery. Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while
* it takes its sequence number and registers its {@link Pending} — see
* {@code AmqpReplyInbox.failPendingPublishesOnRecovery}'s javadoc for the full race analysis this
* mirrors. Package-private only so a unit test can drive it directly without a live broker
* reconnect.
*/
void failPendingPublishesOnRecovery() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP connection recovered mid-publish; confirm status of lead message "
+ pending.msgId + " is unknown"));
}
}
}
/**
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving
* it to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it.
*/
private void failPendingPublishesOnClose() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"lead mailbox closed while publish of lead message " + pending.msgId
+ " was still awaiting its confirm"));
}
}
}
/** The durable queue name a coord-id's mailbox lives on: {@code lead.<coordId>.inbox}. */
public static String queueName(String coordId) {
return QUEUE_PREFIX + coordId + QUEUE_SUFFIX;
}
@Override
public void close() {
failPendingPublishesOnClose();
try {
channel.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox channel close: {}", e.toString());
}
try {
publishChannel.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox publish channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox connection close: {}", e.toString());
}
}
}
@@ -0,0 +1,22 @@
package dev.ltms.fleet.msg;
/**
* Wire envelope for a lead-to-lead message carried over {@link LeadMailbox}.
*
* <p>Unlike {@link ReplyInbox.InboxMessage} (a worker→primary reply, addressed only by the single
* gateway that owns the worker), a lead message crosses independently-owned daemons — possibly on
* different hosts — so it carries an explicit sender ({@code from}) as well as the recipient
* ({@code to}): the recipient needs the sender's coord-id to reply back.
*
* <p>{@code from} and {@code to} are globally-unique lead coordination ids (e.g. {@code "mac-opus"},
* {@code "fleet01-lead"}) — NOT herdr terminal ids. A herdr terminal id is meaningful only on the
* host that owns it, so it cannot address a lead running on another daemon; a coord-id is chosen
* by configuration ({@code coordinator.selfId}) precisely so it means the same thing everywhere.
*
* @param msgId idempotency id; a redelivered duplicate (at-least-once delivery) is deduped on this
* @param from the sending lead's coord-id
* @param to the receiving lead's coord-id — identifies the mailbox this message is held on
* @param content the message text
*/
public record LeadMessage(String msgId, String from, String to, String content) {
}
@@ -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-637: 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.
*
* <p>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<ILoggingEvent> captureFleetdLogs() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static String joined(ListAppender<ILoggingEvent> 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);
}
}
@@ -1,6 +1,7 @@
package dev.ltms.fleet.config;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.peer.MemberRole;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -975,6 +976,66 @@ class FleetConfigTest {
assertFalse(cfg.broker().isConfigured(), "an empty uri must not enable AMQP");
}
@Test
void absentCoordinatorBlockLeavesCoordinatorNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-coordinator.yaml");
Files.writeString(f, "bind:\n port: 8080\n");
FleetConfig cfg = FleetConfig.load(f);
assertNull(cfg.coordinator(), "no coordinator: block → null → no lead mailbox is opened");
}
@Test
void coordinatorBlockParses(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator.yaml");
Files.writeString(f, """
bind:
port: 8080
coordinator:
uri: amqp://guest:guest@127.0.0.1:5672/coord
selfId: mac-opus
prefetch: 16
""");
FleetConfig cfg = FleetConfig.load(f);
assertNotNull(cfg.coordinator());
assertTrue(cfg.coordinator().isConfigured(), "a non-blank uri enables the coordinator");
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", cfg.coordinator().uri());
assertEquals("mac-opus", cfg.coordinator().selfId());
assertEquals(16, cfg.coordinator().prefetchOrDefault());
}
@Test
void coordinatorBlockWithBlankUriStaysUnconfigured(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator-blank.yaml");
Files.writeString(f, "bind:\n port: 8080\ncoordinator:\n uri: \"\"\n");
FleetConfig cfg = FleetConfig.load(f);
assertNotNull(cfg.coordinator());
assertFalse(cfg.coordinator().isConfigured(), "an empty uri must not enable the coordinator");
assertNull(cfg.coordinator().selfId(), "a blank/absent selfId stays null, never coerced to empty");
}
@Test
void coordinatorEffectiveUriHonorsUriEnv() {
FleetConfig.Coordinator withEnv = new FleetConfig.Coordinator(
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null);
assertEquals("amqp://from-env@127.0.0.1:5672/coord",
withEnv.effectiveUri(Map.of("LEAD_COORD_URI", "amqp://from-env@127.0.0.1:5672/coord")),
"uriEnv wins over a literal uri when its variable resolves");
assertNull(withEnv.effectiveUri(Map.of()),
"an unset uriEnv variable must not fall back to the literal uri");
assertNull(withEnv.effectiveUri(Map.of("LEAD_COORD_URI", " ")),
"a blank uriEnv variable must not fall back to the literal uri");
FleetConfig.Coordinator noEnv = new FleetConfig.Coordinator(
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null);
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", noEnv.effectiveUri(Map.of()),
"the literal uri is used when no uriEnv is configured");
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
}
@Test
void absentPrimaryBlockLeavesPrimaryNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-primary.yaml");
@@ -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-637: 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.
*
* <p>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());
}
}
@@ -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();
@@ -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.
*
* <p>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.
*
* <p>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<LeadMessage> held = Collections.synchronizedList(new ArrayList<>());
private final List<LeadMessage> published = Collections.synchronizedList(new ArrayList<>());
private final List<String> 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<LeadMessage> 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<LeadMessage> published() {
return List.copyOf(published);
}
public List<String> acked() {
return List.copyOf(acked);
}
}
@@ -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}.
*
* <p>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<String, String> leads) {
return new LeadCoordLoop(channel, new AgentControl(herdr), () -> leads, NO_SCHEDULER, 3_000L);
}
private static List<FakeHerdr.Call> 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<String, Object> params = (Map<String, Object>) 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<String, String>();
leads.put("term_other", "some-other-lead");
leads.put(LEAD_TERM, SELF);
loop(channel, herdr, leads).tick();
@SuppressWarnings("unchecked")
Map<String, Object> params = (Map<String, Object>) 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<String, String>();
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");
}
}
@@ -0,0 +1,173 @@
package dev.ltms.fleet.msg;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.testcontainers.containers.RabbitMQContainer;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Contract test for {@link LeadMailbox} against a REAL broker — same approach as
* {@code AmqpReplyInboxContractTest}, which this mirrors: a Testcontainers RabbitMQ locally, or an
* externally-provisioned broker in CI via {@code AMQP_URI}. Tagged {@code contract} so it is
* excluded from {@code mvn test}/{@code mvn clean install} (which stay hermetic and need no
* Docker); run it with Docker present via {@code mvn test -Pcontract}.
*
* <p>Proves the mechanism this ticket adds: a {@link LeadMessage} published to
* {@code lead.<to>.inbox} is received with {@code from}/{@code to}/{@code content} intact, and
* {@link LeadMailbox#ack} removes it — the same publish→peek→ack roundtrip
* {@code AmqpReplyInboxContractTest} proves for {@link AmqpReplyInbox}, adapted to this class's
* single-owned-mailbox shape (no {@code own}/{@code release} — the mailbox for {@code selfCoordId}
* is owned the moment {@link LeadMailbox#open} returns).
*/
@Tag("contract")
// disabledWithoutDocker=false: on the CI path (AMQP_URI set) no container is started and the class
// must still run against the external broker even though the runner has no Docker.
@Testcontainers(disabledWithoutDocker = false)
class LeadMailboxTest {
private static final String EXTERNAL_URI = System.getenv("AMQP_URI");
private static final RabbitMQContainer BROKER =
new RabbitMQContainer(DockerImageName.parse("rabbitmq:3.13-management"));
private static final AtomicLong SEQ = new AtomicLong();
// No @Container: the JUnit 5 extension would force-start it even when AMQP_URI is set. Start it
// manually only on the local (no-external-broker) path; Ryuk reaps it on JVM exit.
@BeforeAll
static void startBrokerUnlessExternal() {
if (EXTERNAL_URI == null) {
BROKER.start();
}
}
private static String uri() {
if (EXTERNAL_URI != null) {
return EXTERNAL_URI;
}
// No trailing slash: an empty path is vhost "", which does not exist — omitting it selects
// the default vhost "/".
return "amqp://guest:guest@" + BROKER.getHost() + ":" + BROKER.getAmqpPort();
}
/** A fresh coord-id per test run so parallel/repeat runs never collide on the same queue. */
private static String coordId(String prefix) {
return prefix + "-" + System.nanoTime() + "-" + SEQ.incrementAndGet();
}
@Test
void publishThenPeekThenAckRoundTrip() throws Exception {
String to = coordId("lead-to");
String from = "lead-from";
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
LeadMessage sent = new LeadMessage("m1", from, to, "hello peer lead");
inbox.publish(to, sent);
List<LeadMessage> got = awaitPeek(inbox);
assertEquals(1, got.size(), "the published message should be held for drain");
assertEquals("m1", got.getFirst().msgId());
assertEquals(from, got.getFirst().from());
assertEquals(to, got.getFirst().to());
assertEquals("hello peer lead", got.getFirst().content());
inbox.ack("m1");
assertTrue(inbox.peek().isEmpty(), "an acked message is dropped");
}
}
@Test
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
String to = coordId("lead-dedup");
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
inbox.publish(to, new LeadMessage("dup", "lead-from", to, "first"));
awaitPeek(inbox);
inbox.publish(to, new LeadMessage("dup", "lead-from", to, "second")); // same msgId — no-op
Thread.sleep(500); // give any erroneous second delivery time to land
List<LeadMessage> got = inbox.peek();
assertEquals(1, got.size(), "a repeated msgId must not double-queue");
assertEquals("first", got.getFirst().content(), "the first payload wins");
}
}
@Test
void unackedMessageSurvivesRestartAndIsRedelivered() throws Exception {
String to = coordId("lead-durable");
// First "process life": publish, see it held, but crash before acking.
try (LeadMailbox first = LeadMailbox.open(uri(), to)) {
first.publish(to, new LeadMessage("persist-1", "lead-from", to, "survive me"));
assertEquals(1, awaitPeek(first).size());
// no ack — simulate a java -jar bounce with the message still pending
}
// Second "process life": a fresh connection owning the same mailbox must be redelivered it.
try (LeadMailbox second = LeadMailbox.open(uri(), to)) {
List<LeadMessage> got = awaitPeek(second);
assertEquals(1, got.size(), "an unacked persistent message is redelivered after restart");
assertEquals("persist-1", got.getFirst().msgId());
assertEquals("survive me", got.getFirst().content());
second.ack("persist-1");
}
// Third life: once acked, it is gone for good.
try (LeadMailbox third = LeadMailbox.open(uri(), to)) {
Thread.sleep(500);
assertTrue(third.peek().isEmpty(), "an acked message does not come back on the next restart");
}
}
@Test
void publishDoesNotRequireTheSenderToOwnTheTargetMailbox() throws Exception {
// CB-308 federation: a sender that never opened its own LeadMailbox for `to` can still
// publish to it — publish must not imply ownership. Only the owner ever consumes here.
String to = coordId("lead-federated");
String senderId = coordId("lead-sender");
try (LeadMailbox owner = LeadMailbox.open(uri(), to);
LeadMailbox sender = LeadMailbox.open(uri(), senderId)) {
sender.publish(to, new LeadMessage("m1", senderId, to, "from a federated peer"));
List<LeadMessage> got = awaitPeek(owner);
assertEquals(1, got.size(), "only the owner's mailbox should receive the message");
assertEquals(senderId, got.getFirst().from());
assertTrue(sender.peek().isEmpty(), "the sender must not also hold a copy — it never owns `to`");
}
}
@Test
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
// Publish to a coord-id whose mailbox was never opened by anyone: the queue is never
// declared, so the default-exchange route to lead.<to>.inbox does not exist and the broker
// must return the publish.
String to = coordId("lead-nobody-home");
try (LeadMailbox sender = LeadMailbox.open(uri(), coordId("lead-sender"))) {
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> sender.publish(to, new LeadMessage("m1", "lead-from", to, "nobody home")));
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
"expected an unroutable-publish failure, got: " + ex.getMessage());
}
}
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
@SuppressWarnings("BusyWait")
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
List<LeadMessage> msgs = inbox.peek();
while (msgs.isEmpty() && System.nanoTime() < deadline) {
Thread.sleep(50);
msgs = inbox.peek();
}
return msgs;
}
}