fleetd #361: close the lead-coordination visibility gap
Lead-to-lead AMQP coordination had a send half with tools and a receive
half without. This closes three blind spots:
- LeadChannel gains inspect(coordId) -> MailboxState(exists, pending,
consumers), implemented in LeadMailbox with a throwaway probe channel
(never the long-lived publish/consume channels) so a passive-declare
404 on a missing queue can never take down publish() on the same
instance.
- FleetConfig.Coordinator gains peers: List<String> (defaults to empty,
blank entries dropped) so a daemon can declare which peer coord-ids
it expects to reach.
- fleet_list reports coordination state via a new CoordinationSource
(own coord-id, own mailbox state, held messages as msgId/from/preview
only, and one row per configured peer with reachability/pending/
consumers), following the existing OutageSource/QuarantineSource
"Source record with none()" idiom instead of growing listFleet's
overload chain by another positional parameter. Every peer probe is
bounded by a 1.5s timeout on a virtual-thread pool and degrades to
absent rather than ever slowing or failing fleet_list.
- fleet_send{coordId}'s success text now says "durably confirmed by the
broker" instead of "delivered", and warns (while still reporting
success) when the target mailbox has zero consumers attached.
Tests: hermetic unit tests for exists/absent/zero-consumer/old-config-
no-peers-key/new fleet_list shape using FakeLeadChannel, plus a
@Tag("contract") LeadMailboxTest.inspectingAMissingMailboxNeverBreaks
PublishOnTheSameInstance proving the invariant against a real broker.
This commit is contained in:
@@ -821,10 +821,17 @@ guard:
|
||||
# across every daemon sharing this vhost.
|
||||
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
|
||||
# Default 32 when omitted.
|
||||
# peers → fleetd #361: the coord-ids of the OTHER daemons on this vhost, declared by the
|
||||
# operator (the daemon never guesses). fleet_list reports each one's live reachability
|
||||
# (a passive queue check, never a presence protocol) alongside this daemon's own
|
||||
# mailbox state. Omit, or leave empty, for a daemon with no known peers yet — an
|
||||
# undeclared peer can still reach you and be reached by fleet_send, it just will not
|
||||
# show up as a row in fleet_list.
|
||||
# coordinator:
|
||||
# uriEnv: LEAD_COORD_URI
|
||||
# selfId: mac-opus
|
||||
# prefetch: 32
|
||||
# peers: [fleet01-lead]
|
||||
|
||||
# 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
|
||||
|
||||
@@ -652,7 +652,12 @@ public final class Fleetd {
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
|
||||
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
|
||||
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
|
||||
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
|
||||
// leadMailbox above), so peers follows the same rule rather than half hot-reloading.
|
||||
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers());
|
||||
|
||||
// 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
|
||||
|
||||
@@ -888,12 +888,18 @@ public record FleetConfig(
|
||||
* {@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}.
|
||||
* @param peers fleetd #361: the coord-ids the operator declares as this daemon's peers — the
|
||||
* daemon never guesses who else exists. {@code fleet_list} reports each one's
|
||||
* live reachability. Blank entries are dropped; {@code null} ⇒ an empty list, so
|
||||
* a config written before this field existed still parses unchanged.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
|
||||
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch, List<String> peers) {
|
||||
|
||||
public Coordinator {
|
||||
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
|
||||
peers = peers == null ? List.of()
|
||||
: peers.stream().filter(p -> p != null && !p.isBlank()).toList();
|
||||
}
|
||||
|
||||
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
|
||||
|
||||
@@ -44,6 +44,7 @@ import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiFunction;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -101,6 +102,8 @@ public final class FleetMcp {
|
||||
private final LeadSeatSource leadSeats;
|
||||
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
|
||||
private final LeadChannel leadChannel;
|
||||
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
|
||||
private final List<String> peers;
|
||||
|
||||
/** 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,
|
||||
@@ -164,6 +167,32 @@ public final class FleetMcp {
|
||||
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
|
||||
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
|
||||
* coord-ids the operator has declared as peers ({@code coordinator.peers}). Bundled as its own
|
||||
* Source, the same idiom as {@link OutageSource}/{@link QuarantineSource}/{@link LeadSeatSource},
|
||||
* so {@code listFleet}'s already-long overload chain gains exactly one new required parameter
|
||||
* instead of a further bare positional argument.
|
||||
*
|
||||
* <p>Every mailbox look this triggers goes through {@link LeadChannel#inspect}, which is
|
||||
* specified to run on its own disposable channel — never the channel {@link LeadChannel#publish}
|
||||
* or the consume loop depends on — so a peer that happens to be down, or the coordination broker
|
||||
* itself being unreachable, can never take {@code fleet_send}/{@code LeadCoordLoop}'s own path
|
||||
* down with it. See {@code FleetMcp.probe} for the additional timeout bound on top of that.
|
||||
*
|
||||
* @param leadChannel this daemon's own channel, or {@code null} when no coordinator is configured
|
||||
* @param peers the coord-ids declared under {@code coordinator.peers}, or empty
|
||||
*/
|
||||
public record CoordinationSource(LeadChannel leadChannel, List<String> peers) {
|
||||
public CoordinationSource {
|
||||
peers = peers == null ? List.of() : List.copyOf(peers);
|
||||
}
|
||||
|
||||
/** Inert source — no coordinator row is ever reported. */
|
||||
public static CoordinationSource none() { return new CoordinationSource(null, List.of()); }
|
||||
}
|
||||
|
||||
/**
|
||||
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
|
||||
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
|
||||
@@ -211,19 +240,34 @@ public final class FleetMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}). This is what
|
||||
* {@code Fleetd.main} actually wires up.
|
||||
*
|
||||
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
|
||||
* the feature, never a defaulting overload (the same rule {@code quarantine} and
|
||||
* {@code outage} follow).
|
||||
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}).
|
||||
*/
|
||||
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, OutageSource outage,
|
||||
LeadSeatSource leadSeats) {
|
||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
|
||||
healthCoverage, quarantine, leadChannel, outage, leadSeats, List.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, with fleetd #361 {@code coordinator.peers} (see {@link CoordinationSource}). This is
|
||||
* what {@code Fleetd.main} actually wires up.
|
||||
*
|
||||
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
|
||||
* the feature, never a defaulting overload (the same rule {@code quarantine} and
|
||||
* {@code outage} follow).
|
||||
* @param peers the coord-ids declared under {@code coordinator.peers}; empty when unset or
|
||||
* when {@code leadChannel} is {@code null}.
|
||||
*/
|
||||
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, OutageSource outage,
|
||||
LeadSeatSource leadSeats, List<String> peers) {
|
||||
this.leadChannel = leadChannel;
|
||||
this.peers = peers == null ? List.of() : List.copyOf(peers);
|
||||
this.capacity = capacity;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
this.outage = Objects.requireNonNull(outage, "outage");
|
||||
@@ -361,7 +405,7 @@ public final class FleetMcp {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
leadChannel == null ? null : leadChannel.selfCoordId());
|
||||
new CoordinationSource(leadChannel, peers));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -697,6 +741,15 @@ public final class FleetMcp {
|
||||
* 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.
|
||||
*
|
||||
* <p><strong>fleetd #361: the success text is honest about what "durably confirmed" does and
|
||||
* does not mean.</strong> The broker's publisher confirm proves the message is durably queued —
|
||||
* it says nothing about whether the peer's pane has, or ever will, receive it. After a
|
||||
* successful publish this looks at the target mailbox's consumer count (via
|
||||
* {@link LeadChannel#inspect}, bounded and never allowed to fail the call — see {@link #probe})
|
||||
* and appends a warning when it is zero: that is the observable form of "nobody is reading this
|
||||
* right now". A zero-consumer publish is still reported as a SUCCESS, never an error — the
|
||||
* message is safely queued and will be read once a daemon owning that coord-id connects.
|
||||
*
|
||||
* @param leadChannel this daemon's channel, or {@code null} when no coordinator is configured
|
||||
*/
|
||||
static McpSchema.CallToolResult sendToLead(LeadChannel leadChannel, String coordId, String content,
|
||||
@@ -724,7 +777,16 @@ public final class FleetMcp {
|
||||
+ ". 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() + ")");
|
||||
String result = "published to peer lead \"" + coordId + "\"'s mailbox and durably confirmed "
|
||||
+ "by the broker (msgId " + msg.msgId() + ").";
|
||||
LeadChannel.MailboxState state = probe(leadChannel, coordId);
|
||||
if (state.exists() && state.consumers() == 0) {
|
||||
result += " Warning: that mailbox currently has NO consumers attached — nobody is reading "
|
||||
+ "it right now. The message is safely queued and will be delivered once a daemon "
|
||||
+ "with coordinator.selfId=\"" + coordId + "\" is running and connected; until then "
|
||||
+ "it will not reach that lead's pane.";
|
||||
}
|
||||
return text(result);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1110,7 +1172,8 @@ 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);
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
|
||||
CoordinationSource.none());
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1119,33 +1182,33 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, null);
|
||||
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* As above, additionally reporting this daemon's own lead coordination state (CB-637, fleetd
|
||||
* #361) when a coordinator is configured and its channel opened — see {@link #coordinatorView}
|
||||
* for the shape. 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
|
||||
* @param coordination this daemon's lead channel plus its declared peers, or
|
||||
* {@link CoordinationSource#none()} 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) {
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
|
||||
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm, String selfCoordId) {
|
||||
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
|
||||
@@ -1153,7 +1216,7 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
String selfCoordId) {
|
||||
CoordinationSource coordination) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
@@ -1175,8 +1238,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));
|
||||
Map<String, Object> coordinatorRow = coordinatorView(coordination);
|
||||
if (coordinatorRow != null) {
|
||||
result.put("coordinator", coordinatorRow);
|
||||
}
|
||||
if (capacity.available()) result.put("capacity", profiles.stream()
|
||||
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
|
||||
@@ -1187,6 +1251,102 @@ public final class FleetMcp {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: the {@code coordinator} row — this daemon's own coord-id and mailbox state, the
|
||||
* messages currently held for it, and the live reachability of every operator-declared peer.
|
||||
* {@code null} (the row is then omitted entirely) when lead coordination is off, so an ordinary
|
||||
* fleet's {@code fleet_list} output is byte-identical to before this feature existed.
|
||||
*
|
||||
* <p>Every peer/self mailbox look goes through {@link #probe}, which bounds each
|
||||
* {@link LeadChannel#inspect} call to {@link #PEER_PROBE_TIMEOUT_MS} and never lets it throw —
|
||||
* a coordination broker that is down or slow degrades this row toward "unreachable"/"unknown"
|
||||
* counts, it can never make {@code fleet_list} itself slow or fail. {@code held} comes from
|
||||
* {@link LeadChannel#peek}, a pure in-memory read with no broker round trip, so it is never
|
||||
* subject to that bound.
|
||||
*/
|
||||
private static Map<String, Object> coordinatorView(CoordinationSource coordination) {
|
||||
LeadChannel channel = coordination.leadChannel();
|
||||
if (channel == null) {
|
||||
return null;
|
||||
}
|
||||
String selfId = channel.selfCoordId();
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("selfId", selfId);
|
||||
row.put("configured", true);
|
||||
LeadChannel.MailboxState self = probe(channel, selfId);
|
||||
Map<String, Object> mailbox = new LinkedHashMap<>();
|
||||
mailbox.put("pending", self.pending());
|
||||
mailbox.put("consumers", self.consumers());
|
||||
row.put("mailbox", mailbox);
|
||||
row.put("held", channel.peek().stream().map(FleetMcp::heldView).toList());
|
||||
row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList());
|
||||
return row;
|
||||
}
|
||||
|
||||
/** One held-for-me message: enough to identify it and see roughly what it says, never the whole body. */
|
||||
private static Map<String, Object> heldView(LeadMessage m) {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("msgId", m.msgId());
|
||||
row.put("from", m.from());
|
||||
row.put("preview", preview(m.content()));
|
||||
return row;
|
||||
}
|
||||
|
||||
/** Cap a held message's content to a short preview — {@code fleet_list} must never dump a full body. */
|
||||
private static final int HELD_PREVIEW_MAX_CHARS = 80;
|
||||
|
||||
private static String preview(String content) {
|
||||
if (content == null) {
|
||||
return "";
|
||||
}
|
||||
return content.length() <= HELD_PREVIEW_MAX_CHARS
|
||||
? content
|
||||
: content.substring(0, HELD_PREVIEW_MAX_CHARS) + "…";
|
||||
}
|
||||
|
||||
/** One declared peer's row: its coord-id, whether its mailbox exists, and — only then — its counts. */
|
||||
private static Map<String, Object> peerView(LeadChannel channel, String coordId) {
|
||||
LeadChannel.MailboxState state = probe(channel, coordId);
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("coordId", coordId);
|
||||
row.put("reachable", state.exists());
|
||||
if (state.exists()) {
|
||||
row.put("pending", state.pending());
|
||||
row.put("consumers", state.consumers());
|
||||
}
|
||||
return row;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: how long {@code fleet_list} waits on any single {@link LeadChannel#inspect} call
|
||||
* before giving up on it — see {@link #probe}.
|
||||
*/
|
||||
private static final long PEER_PROBE_TIMEOUT_MS = 1_500L;
|
||||
|
||||
/**
|
||||
* Dedicated pool for {@link LeadChannel#inspect} calls so a slow one blocks only its own virtual
|
||||
* thread, never the MCP request thread calling {@code fleet_list}.
|
||||
*/
|
||||
private static final java.util.concurrent.ExecutorService PEER_PROBE_POOL =
|
||||
java.util.concurrent.Executors.newThreadPerTaskExecutor(Thread.ofVirtual().name("fleet-peer-probe-", 0).factory());
|
||||
|
||||
/**
|
||||
* fleetd #361: {@link LeadChannel#inspect}, bounded to {@link #PEER_PROBE_TIMEOUT_MS} and never
|
||||
* allowed to throw or hang the caller — a coordination broker that is unreachable or slow
|
||||
* degrades to {@link LeadChannel.MailboxState#absent} rather than making {@code fleet_list}
|
||||
* slow or failing it. {@code inspect} itself is already specified to never throw, but this is
|
||||
* the seam that also survives an implementation that does, or one that blocks indefinitely on a
|
||||
* dead connection.
|
||||
*/
|
||||
private static LeadChannel.MailboxState probe(LeadChannel channel, String coordId) {
|
||||
try {
|
||||
return PEER_PROBE_POOL.submit(() -> channel.inspect(coordId))
|
||||
.get(PEER_PROBE_TIMEOUT_MS, TimeUnit.MILLISECONDS);
|
||||
} catch (Exception e) {
|
||||
return LeadChannel.MailboxState.absent(coordId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that fleetd
|
||||
* may stop the member: the bridge has capacity facts but no work list, and choosing work needs
|
||||
@@ -1377,8 +1537,12 @@ public final class FleetMcp {
|
||||
+ "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.")),
|
||||
+ "across hosts. Mutually exclusive with sessionId and turnId. Success means "
|
||||
+ "the message is durably queued and confirmed by the broker, with a warning "
|
||||
+ "if that mailbox has no consumers attached right now (queued, but nobody is "
|
||||
+ "reading it yet) — it does not mean the peer's pane has seen it. Your own "
|
||||
+ "coordId, this daemon's mailbox state, and every coordinator.peers entry's "
|
||||
+ "live reachability are reported by fleet_list's coordinator row.")),
|
||||
List.of("content")));
|
||||
}
|
||||
|
||||
@@ -1486,7 +1650,13 @@ public final class FleetMcp {
|
||||
+ "for a quarantined profile's credential (see fleet_profiles), whatever its "
|
||||
+ "maxLoad/live — with credentialId and quarantinedForSeconds naming the "
|
||||
+ "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing "
|
||||
+ "for N seconds'.",
|
||||
+ "for N seconds'. When lead-to-lead coordination is configured, a 'coordinator' "
|
||||
+ "object reports this daemon's own coord-id ('selfId') and mailbox state "
|
||||
+ "('mailbox': pending/consumers), the messages currently held for it ('held': "
|
||||
+ "msgId/from/preview, never the full body), and one row per coordinator.peers "
|
||||
+ "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — "
|
||||
+ "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' "
|
||||
+ "array above. It is omitted entirely when no coordinator is configured.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
|
||||
@@ -4,8 +4,8 @@ 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.
|
||||
* mailbox: publish to a peer's coord-id, look at what has arrived for me, ack what I have
|
||||
* delivered, and (fleetd #361) inspect any coord-id's mailbox from the outside without owning it.
|
||||
*
|
||||
* <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
|
||||
@@ -37,4 +37,40 @@ public interface LeadChannel {
|
||||
|
||||
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
|
||||
String selfCoordId();
|
||||
|
||||
/**
|
||||
* A non-destructive look at {@code coordId}'s mailbox — does it exist, how many messages are
|
||||
* waiting on it, and how many consumers are attached — without owning, consuming, or otherwise
|
||||
* changing it. {@code consumers == 0} on an existing mailbox is the observable form of "nobody
|
||||
* is reading this right now": a publish to it will sit queued rather than reach a pane.
|
||||
*
|
||||
* <p>Returns {@link MailboxState#absent(String)}, never throws, when {@code coordId}'s mailbox
|
||||
* does not exist or the look otherwise fails (broker unreachable, timed out) — this is a
|
||||
* best-effort fact-finding call, not an operation a caller must handle failing.
|
||||
*
|
||||
* <p><strong>Must never share fate with {@link #publish} or {@link #peek}/{@link #ack}.</strong>
|
||||
* fleetd #361: in AMQP 0-9-1 a passive queue declare of a queue that does not exist closes the
|
||||
* channel it was declared on with a 404. An implementation backed by a real broker connection
|
||||
* must inspect on a channel it can afford to lose — never the channel {@link #publish} or the
|
||||
* consume loop depends on — so that looking at a peer that happens to be down can never break
|
||||
* this daemon's own send or receive path.
|
||||
*/
|
||||
MailboxState inspect(String coordId);
|
||||
|
||||
/**
|
||||
* The result of {@link #inspect}. {@code exists} is {@code false} for a mailbox nobody has ever
|
||||
* declared (or when the look could not be completed) — {@code pending}/{@code consumers} are
|
||||
* meaningless in that case and always read {@code 0}.
|
||||
*
|
||||
* @param coordId the coord-id inspected
|
||||
* @param exists whether the mailbox's queue is currently declared on the broker
|
||||
* @param pending messages ready for delivery but not yet in a consumer's hands (0 when absent)
|
||||
* @param consumers how many consumers are attached (0 when absent, or when nobody is reading it)
|
||||
*/
|
||||
record MailboxState(String coordId, boolean exists, int pending, int consumers) {
|
||||
/** The mailbox does not exist, or the look could not be completed. */
|
||||
public static MailboxState absent(String coordId) {
|
||||
return new MailboxState(coordId, false, 0, 0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -259,6 +259,41 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: look at {@code coordId}'s mailbox on a fresh, immediately-closed throwaway
|
||||
* channel — never {@link #channel} (consume/ack) or {@link #publishChannel} (publish). A
|
||||
* passive queue declare of a queue that does not exist closes the channel it was declared on
|
||||
* with a 404; using a disposable probe channel means that closure can never touch either
|
||||
* long-lived channel this instance depends on for {@link #publish} or the consume loop.
|
||||
*/
|
||||
@Override
|
||||
public MailboxState inspect(String coordId) {
|
||||
String queue = queueName(coordId);
|
||||
Channel probe;
|
||||
try {
|
||||
probe = connection.createChannel();
|
||||
} catch (IOException e) {
|
||||
log.debug("lead mailbox inspect: cannot open a probe channel for {}: {}", coordId, e.toString());
|
||||
return MailboxState.absent(coordId);
|
||||
}
|
||||
try {
|
||||
AMQP.Queue.DeclareOk declared = probe.queueDeclarePassive(queue);
|
||||
return new MailboxState(coordId, true, declared.getMessageCount(), declared.getConsumerCount());
|
||||
} catch (IOException e) {
|
||||
// Missing queue (404) or any other declare failure: the broker (or the client library)
|
||||
// has already closed `probe` for us — report "does not exist" rather than throw.
|
||||
return MailboxState.absent(coordId);
|
||||
} finally {
|
||||
try {
|
||||
if (probe.isOpen()) {
|
||||
probe.close();
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.debug("lead mailbox inspect: probe channel close for {}: {}", coordId, e.toString());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
|
||||
public List<LeadMessage> drain() {
|
||||
List<LeadMessage> snapshot = peek();
|
||||
|
||||
@@ -80,7 +80,7 @@ class FleetdLeadMailboxSelectionTest {
|
||||
@Test
|
||||
void opensTheMailboxWhenAUriAndSelfIdAreConfigured() {
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null);
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
|
||||
Fleetd.openLeadMailbox(coordinator, Map.of(), opener);
|
||||
|
||||
@@ -94,7 +94,7 @@ class FleetdLeadMailboxSelectionTest {
|
||||
void honoursUriEnvOverALiteralUri() {
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator("amqp://stale:stale@old:5672/x", "COORD_URI",
|
||||
"mac-opus", 8);
|
||||
"mac-opus", 8, null);
|
||||
|
||||
Fleetd.openLeadMailbox(coordinator, Map.of("COORD_URI", RESOLVED_URI), opener);
|
||||
|
||||
@@ -105,7 +105,7 @@ class FleetdLeadMailboxSelectionTest {
|
||||
@Test
|
||||
void turnsOffWhenUriEnvDoesNotResolve() {
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null);
|
||||
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
|
||||
|
||||
@@ -116,7 +116,7 @@ class FleetdLeadMailboxSelectionTest {
|
||||
void warnsAndStaysOffWhenSelfIdIsMissing() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null);
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
|
||||
|
||||
@@ -131,7 +131,7 @@ class FleetdLeadMailboxSelectionTest {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
opener.unreachable = true;
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null);
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
|
||||
"a down coordination broker turns the feature off; it must never take the daemon down");
|
||||
|
||||
+2
-2
@@ -104,7 +104,7 @@ class ConfigRefTopLevelReportingCoverageTest {
|
||||
v.put("configReload", new FleetConfig.ConfigReload(true, 10));
|
||||
v.put("quarantineCooldownSeconds", 1800);
|
||||
v.put("memberCredentials", null);
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1, null));
|
||||
v.put("worktreeGroup", "group-a");
|
||||
v.put("memberLoginShell", null);
|
||||
assertNamesMatchComponents(v);
|
||||
@@ -144,7 +144,7 @@ class ConfigRefTopLevelReportingCoverageTest {
|
||||
v.put("configReload", new FleetConfig.ConfigReload(false, 20));
|
||||
v.put("quarantineCooldownSeconds", 3600);
|
||||
v.put("memberCredentials", null);
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2, null));
|
||||
v.put("worktreeGroup", "group-b");
|
||||
v.put("memberLoginShell", null);
|
||||
assertNamesMatchComponents(v);
|
||||
|
||||
@@ -1189,7 +1189,7 @@ class FleetConfigTest {
|
||||
@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);
|
||||
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null, 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")),
|
||||
@@ -1200,12 +1200,61 @@ class FleetConfigTest {
|
||||
"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);
|
||||
"amqp://guest:guest@127.0.0.1:5672/coord", null, 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());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: the live block ships as just {@code uriEnv} + {@code selfId} (see
|
||||
* {@code Fleetd.example.yaml} / the operator's real {@code fleetd.yaml}, gitignored). That exact
|
||||
* shape, with no {@code peers:} key at all, must keep parsing unchanged after this field is added.
|
||||
*/
|
||||
@Test
|
||||
void coordinatorBlockWithNoPeersKeyStillParses(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("coordinator-no-peers.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
coordinator:
|
||||
uriEnv: COORD_AMQP_URI
|
||||
selfId: mac
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertNotNull(cfg.coordinator());
|
||||
assertEquals("mac", cfg.coordinator().selfId());
|
||||
assertEquals(List.of(), cfg.coordinator().peers(), "no peers: key means no configured peers, never null");
|
||||
}
|
||||
|
||||
@Test
|
||||
void coordinatorPeersParsesAndDropsBlankEntries(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("coordinator-peers.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
coordinator:
|
||||
selfId: mac
|
||||
uri: amqp://guest:guest@127.0.0.1:5672/coord
|
||||
peers:
|
||||
- fleet01
|
||||
- ""
|
||||
- fleet02
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertEquals(List.of("fleet01", "fleet02"), cfg.coordinator().peers(),
|
||||
"a blank peer entry must be dropped, never kept as an empty coord-id");
|
||||
}
|
||||
|
||||
@Test
|
||||
void coordinatorPeersDefaultsToEmptyWhenConstructedWithNull() {
|
||||
FleetConfig.Coordinator c = new FleetConfig.Coordinator(
|
||||
"amqp://guest:guest@127.0.0.1:5672/coord", null, "mac", null, null);
|
||||
assertEquals(List.of(), c.peers(), "a null peers list must default to empty, never NPE downstream");
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentWorktreeGroupLeavesItNull(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-worktree-group.yaml");
|
||||
|
||||
+1
-1
@@ -92,7 +92,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
|
||||
v.put("memberCredentials", new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
|
||||
List.of("git"), List.of("git", "ssh"), null));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3, null));
|
||||
v.put("worktreeGroup", "group-guard");
|
||||
v.put("memberLoginShell", "/bin/zsh");
|
||||
assertNamesMatchComponents(v);
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.msg.FakeLeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadMessage;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -93,4 +94,35 @@ class FleetMcpLeadCoordTest {
|
||||
assertTrue(FleetMcp.sendToLead(channel, PEER, " ", null, null).isError());
|
||||
assertEquals(0, channel.published().size());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: "delivered" overstated what publish actually proves — the broker's confirm means
|
||||
* durably queued, not read. A zero-consumer target is the observable form of "this will not
|
||||
* reach a pane right now", so the (still-successful) result must say so.
|
||||
*/
|
||||
@Test
|
||||
void warnsWhenThePeerMailboxHasNoConsumersButStillReportsSuccess() {
|
||||
var channel = new FakeLeadChannel(SELF)
|
||||
.withMailbox(PEER, new LeadChannel.MailboxState(PEER, true, 0, 0));
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
|
||||
|
||||
assertFalse(res.isError(), "a zero-consumer mailbox is still a successful, durably-queued publish");
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("durably confirmed"), out);
|
||||
assertTrue(out.toLowerCase().contains("no consumers"), () -> "must warn nobody is reading it: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void staysQuietAboutConsumersWhenThePeerMailboxHasOne() {
|
||||
var channel = new FakeLeadChannel(SELF)
|
||||
.withMailbox(PEER, new LeadChannel.MailboxState(PEER, true, 0, 1));
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
|
||||
|
||||
assertFalse(res.isError());
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("durably confirmed"), out);
|
||||
assertFalse(out.toLowerCase().contains("no consumers"), () -> "a consumer IS attached: " + out);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,6 +10,9 @@ import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.msg.FakeLeadChannel;
|
||||
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.session.FakeWorktrees;
|
||||
@@ -576,17 +579,23 @@ class FleetMcpTest {
|
||||
void listReportsThisDaemonsOwnCoordIdWhenLeadCoordinationIsOn() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", new LeadChannel.MailboxState("mac-opus", true, 0, 1));
|
||||
|
||||
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");
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
|
||||
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.
|
||||
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
|
||||
// own mailbox state (a self-diagnosis: is my own consumer actually attached?).
|
||||
assertTrue(out.contains("\"coordinator\""), out);
|
||||
assertTrue(out.contains("\"selfId\":\"mac-opus\""), out);
|
||||
assertTrue(out.contains("\"mailbox\":{\"pending\":0,\"consumers\":1}"), out);
|
||||
assertTrue(out.contains("\"held\":[]"), out);
|
||||
assertTrue(out.contains("\"peers\":[]"), out);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -601,6 +610,49 @@ class FleetMcpTest {
|
||||
"an ordinary fleet's output must be unchanged by this feature");
|
||||
}
|
||||
|
||||
@Test
|
||||
void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
String longContent = "x".repeat(200);
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", longContent));
|
||||
|
||||
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(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"msgId\":\"m1\""), out);
|
||||
assertTrue(out.contains("\"from\":\"fleet01-lead\""), out);
|
||||
assertFalse(out.contains(longContent), "fleet_list must never dump a held message's full body: " + out);
|
||||
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void listReportsEachDeclaredPeersLiveReachability() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("fleet01-lead", new LeadChannel.MailboxState("fleet01-lead", true, 2, 1));
|
||||
// "fleet02-lead" is declared as a peer but never configured on the fake — inspect() falls
|
||||
// back to MailboxState.absent, exactly as a real down peer would report.
|
||||
|
||||
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(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead")));
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"reachable\":true,\"pending\":2,\"consumers\":1"), out);
|
||||
assertTrue(out.contains("\"coordId\":\"fleet02-lead\",\"reachable\":false"), out);
|
||||
assertFalse(out.contains("\"coordId\":\"fleet02-lead\",\"reachable\":false,\"pending\""),
|
||||
"pending/consumers must be omitted, not faked as zero, for an unreachable peer: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void capacityUsesThePlacementLiveCount() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
@@ -911,7 +963,8 @@ class FleetMcpTest {
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 2, profile -> 3,
|
||||
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
|
||||
FleetMcp.CoordinationSource.none()));
|
||||
|
||||
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
|
||||
assertTrue(out.contains("\"live\":2"), out);
|
||||
@@ -934,7 +987,8 @@ class FleetMcpTest {
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 3,
|
||||
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
|
||||
FleetMcp.CoordinationSource.none()));
|
||||
|
||||
assertTrue(out.contains("\"live\":0"), out);
|
||||
assertTrue(out.contains("\"free\":3"), "the real gate never subtracts the lead's seat: " + out);
|
||||
@@ -988,7 +1042,7 @@ class FleetMcpTest {
|
||||
String out = textOf(FleetMcp.listFleet(composite, sm, null,
|
||||
new FleetMcp.CapacitySource(liveCount, p -> profiles.get(p).maxLoad(), profiles::keySet, () -> 0),
|
||||
new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(),
|
||||
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
|
||||
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", FleetMcp.CoordinationSource.none()));
|
||||
int reportedFree = extractInt(out, "free");
|
||||
|
||||
for (int i = 0; i < reportedFree; i++) {
|
||||
@@ -1017,7 +1071,7 @@ class FleetMcpTest {
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
|
||||
() -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", null));
|
||||
Map.of(), "", FleetMcp.CoordinationSource.none()));
|
||||
|
||||
assertTrue(out.contains("\"free\":2"), out);
|
||||
assertFalse(out.contains("leadSeats"), "no lead shares this profile's credential: " + out);
|
||||
|
||||
@@ -146,7 +146,7 @@ class MemberEnvAllowListTest {
|
||||
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
|
||||
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
|
||||
null, null, null, null,
|
||||
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults();
|
||||
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null, null)).withDefaults();
|
||||
}
|
||||
|
||||
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
|
||||
|
||||
@@ -3,6 +3,8 @@ package dev.ltms.fleet.msg;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
/**
|
||||
* Hermetic stand-in for {@link LeadChannel}: an in-memory mailbox that records what was published
|
||||
@@ -24,11 +26,19 @@ public final class FakeLeadChannel implements LeadChannel {
|
||||
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;
|
||||
/** Canned {@link #inspect} results by coord-id — absent for any coord-id not configured here. */
|
||||
private final Map<String, MailboxState> mailboxes = new ConcurrentHashMap<>();
|
||||
|
||||
public FakeLeadChannel(String selfCoordId) {
|
||||
this.selfCoordId = selfCoordId;
|
||||
}
|
||||
|
||||
/** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */
|
||||
public FakeLeadChannel withMailbox(String coordId, MailboxState state) {
|
||||
mailboxes.put(coordId, state);
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make every publish fail as an unreachable peer would. */
|
||||
public FakeLeadChannel failPublishWith(String message) {
|
||||
this.publishFailure = new IllegalStateException(message);
|
||||
@@ -65,6 +75,11 @@ public final class FakeLeadChannel implements LeadChannel {
|
||||
return selfCoordId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public MailboxState inspect(String coordId) {
|
||||
return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId));
|
||||
}
|
||||
|
||||
public List<LeadMessage> published() {
|
||||
return List.copyOf(published);
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ 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.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
@@ -159,6 +160,82 @@ class LeadMailboxTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumer() throws Exception {
|
||||
// A LeadMailbox declares AND consumes its own queue the moment open() returns (see own()),
|
||||
// so inspecting a coord-id this same process owns must always find exactly one consumer.
|
||||
String self = coordId("lead-inspect-self");
|
||||
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
|
||||
LeadChannel.MailboxState state = mailbox.inspect(self);
|
||||
assertEquals(self, state.coordId());
|
||||
assertTrue(state.exists(), "this daemon owns and has declared this exact queue");
|
||||
assertEquals(0, state.pending(), "nothing has been published to it yet");
|
||||
assertEquals(1, state.consumers(), "the mailbox's own constructor already attached a consumer");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception {
|
||||
String nobody = coordId("lead-inspect-nobody");
|
||||
try (LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-caller"))) {
|
||||
LeadChannel.MailboxState state = mailbox.inspect(nobody);
|
||||
assertEquals(LeadChannel.MailboxState.absent(nobody), state,
|
||||
"a queue nobody has ever declared must report absent, never throw");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void inspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymore() throws Exception {
|
||||
// Publish into a mailbox this test owns, then never consume from it, to prove `pending` and
|
||||
// `consumers` really come off the broker rather than off this process's own in-memory state.
|
||||
String to = coordId("lead-inspect-pending");
|
||||
String observerId = coordId("lead-inspect-observer");
|
||||
try (LeadMailbox owner = LeadMailbox.open(uri(), to);
|
||||
LeadMailbox observer = LeadMailbox.open(uri(), observerId)) {
|
||||
owner.publish(to, new LeadMessage("m1", "lead-from", to, "sitting in the queue"));
|
||||
awaitPeek(owner); // make sure the broker has actually enqueued it before inspecting
|
||||
} // `owner` closes here: its consumer disconnects, but the durable, unacked message stays queued.
|
||||
|
||||
try (LeadMailbox observer = LeadMailbox.open(uri(), coordId("lead-inspect-observer-2"))) {
|
||||
// The broker requeues `owner`'s unacked delivery asynchronously once its connection drops,
|
||||
// so poll rather than assume the very first passive declare already sees the settled state.
|
||||
LeadChannel.MailboxState state = awaitInspect(observer, to, s -> s.consumers() == 0);
|
||||
assertTrue(state.exists());
|
||||
assertEquals(0, state.consumers(), "the only owner just closed — nobody is reading this anymore");
|
||||
assertEquals(1, state.pending(), "the unacked message must be requeued, never dropped");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361's central invariant, proved rather than assumed: a passive queue declare of a
|
||||
* missing queue closes ITS channel with a 404 in AMQP 0-9-1. {@link LeadMailbox#inspect} is
|
||||
* specified to run on its own disposable channel for exactly this reason — this test is the one
|
||||
* that actually exercises the failure mode and shows {@link LeadMailbox#publish} on the SAME
|
||||
* instance is unaffected by it.
|
||||
*/
|
||||
@Test
|
||||
void inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance() throws Exception {
|
||||
String self = coordId("lead-invariant-self");
|
||||
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
|
||||
// Miss on a queue that has never existed — this is exactly the 404-closes-the-channel case.
|
||||
LeadChannel.MailboxState missed = mailbox.inspect(coordId("lead-invariant-nobody-home"));
|
||||
assertFalse(missed.exists());
|
||||
|
||||
// publish() must still work on THIS SAME instance: if inspect() had reused `publishChannel`
|
||||
// (or `channel`), the broker's 404 would have closed it out from underneath publish().
|
||||
LeadMessage sent = new LeadMessage("after-miss", "lead-from", self, "still alive");
|
||||
mailbox.publish(self, sent);
|
||||
List<LeadMessage> got = awaitPeek(mailbox);
|
||||
assertEquals(1, got.size(), "publish must still reach this mailbox's own queue after a missed inspect");
|
||||
assertEquals("after-miss", got.getFirst().msgId());
|
||||
|
||||
// And a second inspect() — of a mailbox that DOES exist this time — must also still work,
|
||||
// proving the miss did not wedge inspect() itself either.
|
||||
LeadChannel.MailboxState self2 = mailbox.inspect(self);
|
||||
assertTrue(self2.exists());
|
||||
}
|
||||
}
|
||||
|
||||
/** 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 {
|
||||
@@ -170,4 +247,18 @@ class LeadMailboxTest {
|
||||
}
|
||||
return msgs;
|
||||
}
|
||||
|
||||
/** Poll inspect(coordId) until it satisfies {@code done}, or ~10s elapse (broker state settles async). */
|
||||
@SuppressWarnings("BusyWait")
|
||||
private static LeadChannel.MailboxState awaitInspect(
|
||||
LeadMailbox observer, String coordId, java.util.function.Predicate<LeadChannel.MailboxState> done)
|
||||
throws InterruptedException {
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
|
||||
LeadChannel.MailboxState state = observer.inspect(coordId);
|
||||
while (!done.test(state) && System.nanoTime() < deadline) {
|
||||
Thread.sleep(50);
|
||||
state = observer.inspect(coordId);
|
||||
}
|
||||
return state;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user