Merge lead-comms-wiring: wire LeadMailbox into the daemon (fleet_send{coordId} + receive loop)
This commit is contained in:
@@ -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.
|
||||
*
|
||||
* <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.
|
||||
|
||||
@@ -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<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-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<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-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.
|
||||
*
|
||||
* <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-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<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;
|
||||
}
|
||||
}
|
||||
@@ -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<LeadMessage> 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) {
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <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);
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
*
|
||||
* <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");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user