Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c1ca6273fc | |||
| e54e3d87ea | |||
| 3c5873dfe2 | |||
| e29227d5f4 | |||
| 2af13ab1ff | |||
| 0788d84be8 | |||
| 20c1094cbf | |||
| 9011c59b9f | |||
| c11ad71ed0 | |||
| bdcf285265 | |||
| cfebc575ea | |||
| 5289eb509f | |||
| 703a05db41 |
@@ -426,7 +426,8 @@ public final class FleetMcp {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
new CoordinationSource(leadChannel, peers));
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
coordinatorVisibleTo(principal(exchange)));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -579,6 +580,20 @@ public final class FleetMcp {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439: only the primary may read {@code fleet_list}'s {@code coordinator} row —
|
||||
* lead-to-lead coordination state (coord-ids, mailbox facts, held-message previews), never the
|
||||
* roster. Split out of the {@code fleet_list} handler, same reason as {@link #denyFor} and
|
||||
* {@link #recordPrimarySingleton}: the decision must be unit-testable without fabricating an
|
||||
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather than
|
||||
* inlining the check, so a future edit cannot silently pass a literal instead of asking who
|
||||
* called ({@code FleetMcpAuthzTest.theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}
|
||||
* reads the source and asserts the handler calls this method by name, not a literal).
|
||||
*/
|
||||
static boolean coordinatorVisibleTo(Principal caller) {
|
||||
return caller.isPrimary();
|
||||
}
|
||||
|
||||
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
|
||||
private static String callerTerminal(McpSyncServerExchange exchange) {
|
||||
Object v = exchange.transportContext().get(CALLER_TERMINAL);
|
||||
@@ -994,7 +1009,11 @@ public final class FleetMcp {
|
||||
if (isBlank(target) || isBlank(msgId)) {
|
||||
return error("target and msgId are required");
|
||||
}
|
||||
messages.ackReply(target, msgId);
|
||||
if (!messages.ackReply(target, msgId)) {
|
||||
return error(msgId + " is not in " + target + "'s reply inbox (wrong id, wrong target, "
|
||||
+ "or already acked). Held lead-to-lead (peer) mail cannot be acked this way — "
|
||||
+ "read it with fleet_poll{coordId}.");
|
||||
}
|
||||
return text("acknowledged " + msgId);
|
||||
}
|
||||
|
||||
@@ -1340,12 +1359,44 @@ public final class FleetMcp {
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
|
||||
/**
|
||||
* As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}).
|
||||
*
|
||||
* <p>Assumes the caller is the primary — every wrapper overload above delegates here without
|
||||
* carrying a caller identity, which is exactly right for them: they exist for call sites (and
|
||||
* unit tests) that have no {@link Principal} to hand over, and this preserves their pre-#439
|
||||
* behavior unchanged. The one call site that has a real caller ({@code fleet_list}'s MCP
|
||||
* handler) uses {@link #listFleet(PeerLauncher, SessionManager, MessageService, CapacitySource,
|
||||
* HealthCoverageSource, QuarantineSource, OutageSource, LeadSeatSource, Map, String,
|
||||
* CoordinationSource, boolean)} instead, so it can pass the true answer.
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, leads, selfTerm, coordination, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, gated by the caller's role (fleetd #439). The {@code coordinator} row is
|
||||
* lead-to-lead coordination state — coordination between orchestrators, not roster
|
||||
* observation — so it is assembled and included only when {@code callerIsPrimary} is
|
||||
* {@code true}. A worker or an architect gets a result with the {@code coordinator} key
|
||||
* <strong>absent</strong>, never an empty or redacted one, and never pays the cost of
|
||||
* {@link #coordinatorView} probing peer mailboxes for a row it will not receive.
|
||||
*
|
||||
* @param callerIsPrimary whether the {@code fleet_list} caller is the primary; only the MCP
|
||||
* handler computes this from the real connection (see
|
||||
* {@code Principal#isPrimary()}) — every other overload passes
|
||||
* {@code true}
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
@@ -1367,9 +1418,14 @@ public final class FleetMcp {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("leads", leadRows); result.put("members", out);
|
||||
result.put("healthCoverage", healthCoverage.value().get());
|
||||
Map<String, Object> coordinatorRow = coordinatorView(coordination);
|
||||
if (coordinatorRow != null) {
|
||||
result.put("coordinator", coordinatorRow);
|
||||
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
|
||||
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
|
||||
// key is absent rather than present-and-empty.
|
||||
if (callerIsPrimary) {
|
||||
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,
|
||||
@@ -1766,7 +1822,10 @@ public final class FleetMcp {
|
||||
+ "has processed a reply and wants to confirm it, leaving other pending replies "
|
||||
+ "in the inbox for later drain.",
|
||||
objectSchema(Map.of(
|
||||
"target", stringProp("Worker session id whose inbox to ack from"),
|
||||
"target", stringProp("Worker session id whose inbox to ack from. Must name a "
|
||||
+ "reply actually queued for it — an id in no inbox, or a coord-id "
|
||||
+ "(peer held mail, read with fleet_poll{coordId} instead), errors "
|
||||
+ "rather than reporting a false success"),
|
||||
"msgId", stringProp("The message id to acknowledge")),
|
||||
List.of("target", "msgId")));
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.PeerHandle;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -594,6 +595,23 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
req.sessionName(), spawned.agentSessionId(), spawned.receipt());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
*
|
||||
* <p>fleetd #450: re-enters {@link #spawn(SpawnRequest)} with {@code decision}'s profile named
|
||||
* explicitly. This is the re-entering form the interface javadoc describes for a launcher with
|
||||
* no placement concept of its own — an instance of this class spawns a single adapter's own
|
||||
* profile set by explicit name only ({@link #place}/{@link #defaultProfileFor} are unoverridden
|
||||
* here and just wrap {@link #defaultProfile()}); it does no quarantine/cool-off/maxLoad/model-off
|
||||
* filtering of its own to re-apply. That filtering lives one layer up, in {@code
|
||||
* CompositePeerLauncher}, which is the launcher that routes across more than one profile and
|
||||
* therefore overrides this method with the routing form instead.
|
||||
*/
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
|
||||
/** The herdr daemon that owns this launcher's pane coordinates. */
|
||||
public HerdrClient herdr() {
|
||||
return agents.herdr();
|
||||
|
||||
@@ -394,17 +394,17 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
public boolean ack(String target, String msgId) {
|
||||
var perTarget = held.get(target);
|
||||
if (perTarget == null || perTarget == RELEASED) {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
Held h;
|
||||
synchronized (perTarget) {
|
||||
h = perTarget.remove(msgId);
|
||||
}
|
||||
if (h == null) {
|
||||
return; // never held (or already acked) — no-op
|
||||
return false; // never held (or already acked) — no-op
|
||||
}
|
||||
try {
|
||||
synchronized (channelLock) {
|
||||
@@ -418,6 +418,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
throw new IllegalStateException("cannot ack reply " + msgId + " on " + queueName(target), e);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
private DeliverCallback deliverCallback(String target) {
|
||||
|
||||
@@ -59,16 +59,17 @@ public final class InMemoryReplyInbox implements ReplyInbox {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
public boolean ack(String target, String msgId) {
|
||||
if (!owned.contains(target)) {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
var perTarget = store.get(target);
|
||||
if (perTarget != null) {
|
||||
//noinspection SynchronizationOnLocalVariableOrMethodParameter
|
||||
synchronized (perTarget) {
|
||||
perTarget.remove(msgId);
|
||||
}
|
||||
if (perTarget == null) {
|
||||
return false;
|
||||
}
|
||||
//noinspection SynchronizationOnLocalVariableOrMethodParameter
|
||||
synchronized (perTarget) {
|
||||
return perTarget.remove(msgId) != null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -826,9 +826,13 @@ public final class MessageService {
|
||||
/**
|
||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
*
|
||||
* @return {@code true} if an entry was actually removed, {@code false} if {@code msgId} was not
|
||||
* in {@code target}'s inbox (wrong id, wrong target, or already acked). The caller —
|
||||
* {@link dev.ltms.fleet.mcp.FleetMcp#ack} — must not report success on {@code false}.
|
||||
*/
|
||||
public void ackReply(String target, String msgId) {
|
||||
inbox.ack(target, msgId);
|
||||
public boolean ackReply(String target, String msgId) {
|
||||
return inbox.ack(target, msgId);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -46,6 +46,13 @@ public interface ReplyInbox {
|
||||
/** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */
|
||||
List<InboxMessage> peek(String target);
|
||||
|
||||
/** Remove the reply {@code msgId} for {@code target} once the primary has taken it. No-op if absent. */
|
||||
void ack(String target, String msgId);
|
||||
/**
|
||||
* Remove the reply {@code msgId} for {@code target} once the primary has taken it.
|
||||
*
|
||||
* @return {@code true} if an entry was actually removed, {@code false} if there was nothing to
|
||||
* remove (unknown {@code target}, unowned {@code target}, or a {@code msgId} not held for
|
||||
* it). A {@code false} is not an error — acking a {@code target} this daemon does not own is
|
||||
* part of the normal contract, not a failure.
|
||||
*/
|
||||
boolean ack(String target, String msgId);
|
||||
}
|
||||
|
||||
@@ -168,6 +168,23 @@ public interface PeerLauncher {
|
||||
* answer for a launcher with no role-pool concept of its own (e.g. a single {@code
|
||||
* HerdrPeerLauncher} adapter, which is never reached this way in production: {@code
|
||||
* CompositePeerLauncher} always fronts it and resolves roles itself).
|
||||
*
|
||||
* <p>fleetd #453: this default is deliberately <em>not</em> abstract — unlike {@link
|
||||
* #spawn(SpawnRequest, PlacementDecision)} (fleetd #450), there is no live defect in inheriting
|
||||
* it today, and the only current single-adapter implementer ({@code HerdrPeerLauncher}) is
|
||||
* correct to do so. But it stays correct only as long as that holds: <strong>if a launcher ever
|
||||
* routes more than one profile per role, it MUST override this method</strong>, or every role
|
||||
* silently resolves to {@link #defaultProfile()} with no error and no log line. {@code
|
||||
* HerdrPeerLauncher.spawn(SpawnRequest, PlacementDecision)} — the override in {@code
|
||||
* dev.ltms.fleet.member}, not the declaration below — names this method and {@link #place}
|
||||
* explicitly as "unoverridden here" for exactly this reason. Read it before adding role-pool
|
||||
* routing to any {@code HerdrPeerLauncher} subclass.
|
||||
*
|
||||
* <p>Who is forced to read which paragraph, because it is not symmetric. A new class that
|
||||
* implements this interface directly must write a body for {@link #spawn(SpawnRequest,
|
||||
* PlacementDecision)}, which is abstract here, so it lands on this javadoc. A subclass of
|
||||
* {@code HerdrPeerLauncher} does not: that class already implements the method, and the
|
||||
* subclass inherits the body. So for a subclass this paragraph is advice, not a gate.
|
||||
*/
|
||||
default String defaultProfileFor(MemberRole role) {
|
||||
return defaultProfile();
|
||||
@@ -228,6 +245,13 @@ public interface PeerLauncher {
|
||||
* placement condition — the right answer for a launcher with no pool or placement-policy
|
||||
* concept of its own, matching {@link #defaultProfileFor}'s own default.
|
||||
*
|
||||
* <p>fleetd #453: same reasoning as {@link #defaultProfileFor}'s own #453 note — this default
|
||||
* is deliberately not abstract (no live defect today, correct for the sole single-adapter
|
||||
* implementer), but <strong>a launcher that ever routes more than one profile per role MUST
|
||||
* override this method too</strong>, or placement silently ignores {@code role} for it. See
|
||||
* {@code HerdrPeerLauncher.spawn(SpawnRequest, PlacementDecision)}'s javadoc, which names this
|
||||
* method as "unoverridden here" and why that is correct only for a single-profile adapter.
|
||||
*
|
||||
* @throws RuntimeException (implementation-specific, typically a placement exception) if no
|
||||
* candidate in {@code role}'s pool is currently placeable
|
||||
*/
|
||||
@@ -253,26 +277,27 @@ public interface PeerLauncher {
|
||||
* matching profile; passing a request that names a <em>different</em>, explicit profile than
|
||||
* the decision it is paired with is a caller bug this method does not attempt to detect.
|
||||
*
|
||||
* <p>Default implementation for a launcher with no placement concept of its own: delegates to
|
||||
* {@link #spawn(SpawnRequest)} with the decision's profile named explicitly — its only spawn
|
||||
* contract, since there is no separate routing path to honor. This default is correct ONLY for
|
||||
* a launcher that spawns a single profile of its own (e.g. {@code HerdrPeerLauncher}), where
|
||||
* the explicit-profile branch it re-enters and the routing branch {@link #place} would have
|
||||
* used are the same thing. <strong>A launcher that routes across more than one profile — the
|
||||
* way {@code CompositePeerLauncher} routes across every configured adapter — MUST override
|
||||
* this method instead of inheriting this default.</strong> Re-entering {@link
|
||||
* #spawn(SpawnRequest)} re-applies that single-argument method's explicit-profile checks
|
||||
* ({@code enforceNotQuarantined}, {@code enforceNotCoolingOff}, {@code enforceMaxLoad}, {@code
|
||||
* enforceModelEnabled} in {@code CompositePeerLauncher}), which can refuse the very profile
|
||||
* {@link #place} just chose, if the underlying placement state moved in the window between the
|
||||
* {@link #place} call and this one — the exact window this method and {@link PlacementDecision}
|
||||
* exist to close (fleetd #444).
|
||||
* <p>No default implementation (fleetd #450): the two correct bodies disagree on purpose, so an
|
||||
* implementer must choose one rather than silently inherit whichever this interface happened to
|
||||
* provide. An implementer with no placement concept of its own — spawns a single profile, e.g.
|
||||
* {@code HerdrPeerLauncher} — should delegate to {@link #spawn(SpawnRequest)} with the decision's
|
||||
* profile named explicitly, since there is no separate routing path to honor there: the
|
||||
* explicit-profile branch it re-enters and the routing branch {@link #place} would have used are
|
||||
* the same thing. <strong>A launcher that routes across more than one profile — the way {@code
|
||||
* CompositePeerLauncher} routes across every configured adapter — MUST NOT re-enter {@link
|
||||
* #spawn(SpawnRequest)}.</strong> Doing so re-applies that single-argument method's
|
||||
* explicit-profile checks ({@code enforceNotQuarantined}, {@code enforceNotCoolingOff}, {@code
|
||||
* enforceMaxLoad}, {@code enforceModelEnabled} in {@code CompositePeerLauncher}), which can
|
||||
* refuse the very profile {@link #place} just chose, if the underlying placement state moved in
|
||||
* the window between the {@link #place} call and this one — the exact window this method and
|
||||
* {@link PlacementDecision} exist to close (fleetd #444). Before #450 this was a {@code default}
|
||||
* method that only {@code CompositePeerLauncher} overrode; a future placement-doing launcher
|
||||
* could have inherited the re-entering body silently and never known. Making it abstract turns
|
||||
* that silent inheritance into a compile error.
|
||||
*
|
||||
* @throws IllegalArgumentException if the decision names an unknown profile
|
||||
*/
|
||||
default PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
PeerHandle spawn(SpawnRequest req, PlacementDecision decision);
|
||||
|
||||
/**
|
||||
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
|
||||
|
||||
@@ -20,6 +20,7 @@ import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
@@ -123,6 +124,11 @@ class FleetdBackendErrorSinkTest {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of();
|
||||
|
||||
@@ -18,12 +18,34 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
* the throwaway space down.
|
||||
*
|
||||
* <p>The seed shell's own startup (restoring its session, printing its banner) is asynchronous
|
||||
* and its length is not a fleetd contract — measured here at ~2.5s on one host (fleetd #449). A
|
||||
* fixed sleep before typing raced that startup: input typed before the shell reached its prompt
|
||||
* was swallowed by the shell's own startup, and the pane showed the typed line followed by the
|
||||
* startup banner with no command output at all — indistinguishable, at a glance, from the env
|
||||
* map never reaching the shell. So this polls for a real signal (the pane's visible text
|
||||
* settling, then the expected output appearing) instead of guessing a sleep length.
|
||||
* and its length is not a fleetd contract — measured here at ~2.5s on one host (fleetd #449).
|
||||
* The old version used two fixed sleeps: 1000ms before typing, then 800ms before reading. It
|
||||
* failed, and the pane showed the typed line followed by the startup banner with no command
|
||||
* output at all — which looks, at a glance, exactly like the env map never reaching the shell.
|
||||
* So this polls for a real signal instead of guessing a sleep length.
|
||||
*
|
||||
* <p><strong>What was measured, and what was not.</strong> Polling fixes it: 5 standalone runs
|
||||
* green. The load-bearing half is {@link #waitForText}, and one cell proves it. Keep the old
|
||||
* 1000ms write sleep and change only the read — the 800ms fixed sleep becomes a 5s poll — and
|
||||
* the test goes from 0 of 3 passing to 3 of 3. Removing the write wait instead
|
||||
* ({@link #SHELL_READY_TIMEOUT_MS} set to 0, so input is typed at once) also passes 3 of 3. So
|
||||
* the cause is the 800ms READ deadline, not the 1000ms write delay. The old version fails every
|
||||
* time, not sometimes, so "race" is the wrong word for it. The earlier explanation — that input
|
||||
* typed before the prompt is swallowed by the shell's startup — is not supported by anything
|
||||
* measured here. Please do not repeat it: if it were true, typing at 0ms would be worse than
|
||||
* typing at 1000ms, and it is not.
|
||||
*
|
||||
* <p><strong>Where this was measured.</strong> A 12-core macOS host, load average 2.6 to 5.9,
|
||||
* on commit 20c1094. The same four cells were also run under load and gave the same answer, but
|
||||
* that run is not clean evidence: the load average climbed from 7 to 50 while the cells ran, and
|
||||
* the old version ran last, at the top of that climb. Above about load 20 everything here is
|
||||
* slow for reasons that have nothing to do with this seam. So read the claim as "measured near
|
||||
* idle on a 12-core host", and nothing stronger. If this test fails on a smaller or busier
|
||||
* machine, raise {@link #OUTPUT_TIMEOUT_MS} before you suspect the seam.
|
||||
*
|
||||
* <p>{@link #waitUntilSettled} is kept as cheap insurance against that swallow case, not because
|
||||
* anyone showed it was needed. If you want to delete it, the honest test is whether you can make
|
||||
* this test fail by typing early. Nobody has managed that yet.
|
||||
*
|
||||
* <p>Tagged {@code contract}; run with {@code mvn test -Pcontract}.
|
||||
*/
|
||||
@@ -48,8 +70,9 @@ class AgentControlContractTest {
|
||||
/**
|
||||
* Poll {@code pane.read} until two consecutive reads come back identical — the shell's own
|
||||
* startup output (restore banner, prompt) has stopped changing — or {@code timeoutMs} elapses.
|
||||
* Never asserts by itself; the caller's own assertion is what actually verifies the outcome,
|
||||
* this only avoids sending input into a shell still mid-startup.
|
||||
* Never asserts by itself; the caller's own assertion is what actually verifies the outcome.
|
||||
* Setting {@code timeoutMs} to 0 skips the wait entirely and the test still passes here, so
|
||||
* treat this as insurance rather than as the fix — see the class javadoc.
|
||||
*/
|
||||
private static String waitUntilSettled(UnixSocketHerdrClient herdr, String paneId, long timeoutMs)
|
||||
throws InterruptedException {
|
||||
|
||||
@@ -187,6 +187,71 @@ class FleetMcpAuthzTest {
|
||||
"no CallerResolver supplied ⇒ authorization not enforced (legacy behaviour)");
|
||||
}
|
||||
|
||||
// --- fleetd #439: who may see fleet_list's coordinator row ----------------------------------
|
||||
|
||||
/**
|
||||
* fleetd #439: {@link FleetMcp#coordinatorVisibleTo} is the whole policy decision for
|
||||
* {@code fleet_list}'s {@code coordinator} row — lead-to-lead coordination state, not roster
|
||||
* observation. Only the primary may see it; a worker, an architect, and (the case the previous
|
||||
* pass of this ticket did not cover) an anonymous caller must all be refused.
|
||||
*/
|
||||
@Test
|
||||
void onlyThePrimaryMaySeeTheCoordinatorRow() {
|
||||
assertTrue(FleetMcp.coordinatorVisibleTo(PRIMARY), "the primary must see its own coordination state");
|
||||
assertFalse(FleetMcp.coordinatorVisibleTo(WORKER_A), "a worker must not see lead-to-lead coordination state");
|
||||
assertFalse(FleetMcp.coordinatorVisibleTo(ARCH_DESIGN),
|
||||
"an architect holds READ today, but that must not extend to coordinator");
|
||||
assertFalse(FleetMcp.coordinatorVisibleTo(ANON), "authenticated as nothing must not see it either");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439 / PR #462 review finding M2: the predicate above can be perfectly correct while
|
||||
* the one production call site (the {@code fleet_list} MCP handler) never actually asks it —
|
||||
* a literal {@code true} compiles, and the whole suite stayed green under that mutation because
|
||||
* every existing test drives {@link FleetMcp#listFleet} directly and supplies the boolean
|
||||
* itself. This test reads {@code FleetMcp.java}'s own source (same idiom as {@link
|
||||
* #toolsTheServerRegisters()} / {@link #everyRegisteredToolHasItsHandlerActionPinned()}) and
|
||||
* asserts the handler's call passes {@code coordinatorVisibleTo(principal(exchange))} — not a
|
||||
* literal {@code true} or {@code false} — as {@code listFleet}'s trailing argument.
|
||||
*
|
||||
* <p>Anchored on argument position, not a bare substring search: {@code true} appears many
|
||||
* times elsewhere in this file for unrelated reasons, so a plain {@code contains("true")}
|
||||
* check would prove nothing. The pattern requires the literal text immediately before the
|
||||
* closing {@code );} of the {@code listFleet(} call inside the handler block to be exactly
|
||||
* {@code coordinatorVisibleTo(principal(exchange))}.
|
||||
*/
|
||||
@Test
|
||||
void theFleetListHandlerActuallyConsultsCoordinatorVisibleTo() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
// Isolate the fleet_list handler block: from its declaration up to the next handler's
|
||||
// declaration. A change to variable naming would break this scrape loudly (see the control
|
||||
// assertion just below), rather than silently reporting "no violation found".
|
||||
int start = source.indexOf("listHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("stopHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
|
||||
// fails, the anchors above moved and the assertion below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("listFleet("),
|
||||
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
|
||||
+ "the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
Pattern trailingArg = Pattern.compile(
|
||||
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
|
||||
Pattern.DOTALL);
|
||||
Matcher m = trailingArg.matcher(handlerBlock);
|
||||
assertTrue(m.find(), "could not locate listFleet(...)'s trailing boolean argument in the "
|
||||
+ "listHandler block -- the call shape changed, update this test's anchor: " + handlerBlock);
|
||||
String trailing = m.group(1);
|
||||
assertEquals("coordinatorVisibleTo(principal(exchange))", trailing,
|
||||
"the fleet_list handler must ask coordinatorVisibleTo(principal(exchange)) who is "
|
||||
+ "calling, not pass a literal boolean -- found: " + trailing);
|
||||
}
|
||||
|
||||
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -352,7 +352,7 @@ class FleetMcpTest {
|
||||
fail("a lead fleet_reply must not publish to the worker inbox");
|
||||
}
|
||||
@Override public List<InboxMessage> peek(String target) { return List.of(); }
|
||||
@Override public void ack(String target, String msgId) { }
|
||||
@Override public boolean ack(String target, String msgId) { return false; }
|
||||
};
|
||||
MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
inboxThatRejectsPublishes);
|
||||
@@ -676,6 +676,95 @@ class FleetMcpTest {
|
||||
"an ordinary fleet's output must be unchanged by this feature");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439: a worker calling {@code fleet_list} must get a result with the {@code
|
||||
* coordinator} key <strong>absent</strong> -- not an empty object, not a redacted one -- even
|
||||
* though lead coordination is fully configured and would otherwise report a row. This drives
|
||||
* the same {@code callerIsPrimary} value the MCP handler computes ({@code
|
||||
* Principal.worker(...).isPrimary()}), so it pins the real production boolean, not a literal.
|
||||
*/
|
||||
@Test
|
||||
void listOmitsTheCoordinatorKeyEntirelyForAWorkerEvenWhenLeadCoordinationIsOn() {
|
||||
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", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
|
||||
|
||||
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(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""), "a worker must never see the coordinator key at all: " + out);
|
||||
assertFalse(out.contains("mac-opus"), "no fragment of the coordinator row may leak either: " + out);
|
||||
assertTrue(out.contains("\"leads\""), "the rest of the result must still be present: " + out);
|
||||
assertTrue(out.contains("\"members\""), out);
|
||||
assertTrue(out.contains("\"healthCoverage\""), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439 acceptance criterion 2: an architect gets exactly the same treatment as a worker.
|
||||
* This is a real, executed test (not just reasoning by analogy) -- it drives the actual
|
||||
* {@code Principal.architect(...).isPrimary()} value the production handler would compute for
|
||||
* an architect caller, through the same gate a worker's call goes through.
|
||||
*/
|
||||
@Test
|
||||
void listOmitsTheCoordinatorKeyEntirelyForAnArchitectToo() {
|
||||
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", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
|
||||
|
||||
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(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""), "an architect must never see the coordinator key either: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439 acceptance criterion 3: the primary's {@code fleet_list} is byte-for-byte
|
||||
* unchanged by this fix. Proven by comparing the new gated overload (with {@code
|
||||
* callerIsPrimary=true}, exactly what the MCP handler passes for the primary) against the
|
||||
* pre-#439 overload that always assembled the row -- if the gate changed anything for a
|
||||
* primary caller, these two strings would differ.
|
||||
*/
|
||||
@Test
|
||||
void listIsByteForByteUnchangedForThePrimaryCaller() {
|
||||
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", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
|
||||
String preExisting = textOf(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 gatedAsPrimary = textOf(FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
|
||||
|
||||
assertEquals(preExisting, gatedAsPrimary,
|
||||
"a primary caller must see byte-for-byte the same result as before this fix");
|
||||
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"mailbox\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"heldCount\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"heldDurable\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"held\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"peers\""), gatedAsPrimary);
|
||||
}
|
||||
|
||||
@Test
|
||||
void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
@@ -1386,6 +1475,10 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void bridgeAckReturnsConfirmationForValidArgs() {
|
||||
// fleet_ack only reports success for a msgId actually queued in the target's inbox
|
||||
// (fleetd #437) — publish one via the inbox directly rather than asserting on a
|
||||
// fabricated id nothing ever queued.
|
||||
inbox.publish("term_a", "msg-1", "queued reply");
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "msg-1");
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
assertTrue(textOf(res).contains("msg-1"), "response should mention the msgId");
|
||||
@@ -1398,6 +1491,26 @@ class FleetMcpTest {
|
||||
assertTrue(FleetMcp.ack(messages, " ", "msg-1").isError());
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckOfAnIdInNoInboxIsAnError() {
|
||||
// fleetd #437: fleet_ack used to say "acknowledged <msgId>" for a message it never
|
||||
// touched, because nothing in the chain reported hit vs. miss. "never-queued" is in no
|
||||
// inbox at all, so this must error rather than claim success.
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "never-queued");
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("never-queued"), textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckOfACoordIdTargetIsAnErrorNamingFleetPoll() {
|
||||
// A coord-id names a peer lead's held mailbox (LeadChannel/LeadMailbox), never a
|
||||
// worker's ReplyInbox — fleet_ack has no route to it and must say so, pointing at
|
||||
// fleet_poll{coordId} instead of reporting a false "acknowledged".
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "coord-some-peer", "msg-1");
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("fleet_poll{coordId}"), textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckRemovesSpecificReply() {
|
||||
// Queue a reply and capture its msgId.
|
||||
@@ -1411,8 +1524,18 @@ class FleetMcpTest {
|
||||
var peeked = messages.drainReplies("term_a");
|
||||
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
|
||||
|
||||
// ackReply works (no-op since published with a different UUID, but callable).
|
||||
assertDoesNotThrow(() -> messages.ackReply("term_a", msgId));
|
||||
// fleetd #437: msgId was already drained above (a fresh UUID each publish), so it is no
|
||||
// longer in the inbox — ackReply must now report that miss instead of pretending to ack.
|
||||
assertFalse(messages.ackReply("term_a", msgId));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckRemovingARealQueuedReplyReportsSuccessAndRemovesIt() {
|
||||
// The worker path must not change behaviour: acking a reply that IS still in the inbox
|
||||
// still succeeds and still removes it (fleetd #437).
|
||||
inbox.publish("term_a", "real-1", "still queued");
|
||||
assertTrue(messages.ackReply("term_a", "real-1"), "ack of a real queued reply must report true");
|
||||
assertTrue(inbox.peek("term_a").isEmpty(), "the acked reply must be gone from the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -16,6 +16,7 @@ import java.util.List;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
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;
|
||||
|
||||
@@ -221,6 +222,31 @@ class AmqpReplyInboxContractTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackReportsHitVsMissAgainstARealBroker() throws Exception {
|
||||
// fleetd #437: fleet_ack said "acknowledged <msgId>" for a message it never touched,
|
||||
// because ReplyInbox.ack() (void) could not tell a hit from a miss. Pin the fixed
|
||||
// boolean contract against a real broker — the adapter fleetd actually runs live.
|
||||
String target = "worker-ack-contract-" + System.nanoTime();
|
||||
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
|
||||
inbox.own(target);
|
||||
|
||||
// Never held for this target at all: must report false, not throw.
|
||||
assertFalse(inbox.ack(target, "never-held"),
|
||||
"acking a msgId never held for an owned target must report false");
|
||||
|
||||
// A real message: first ack removes it and reports true...
|
||||
inbox.publish(target, "m1", "ack me");
|
||||
assertEquals(1, awaitPeek(inbox, target).size(), "the published reply should be held");
|
||||
assertTrue(inbox.ack(target, "m1"), "acking a held reply must report true");
|
||||
assertTrue(inbox.peek(target).isEmpty(), "an acked reply is dropped");
|
||||
|
||||
// ...and the second ack of the SAME msgId has nothing left to remove: false.
|
||||
assertFalse(inbox.ack(target, "m1"),
|
||||
"acking the same msgId twice must report false the second time");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void confirmedPublishDeliversNormally() throws Exception {
|
||||
String target = "worker-confirm-" + System.nanoTime();
|
||||
|
||||
@@ -41,20 +41,20 @@ class InMemoryReplyInboxTest {
|
||||
@Test
|
||||
void ackRemovesTheMessage() {
|
||||
inbox.publish("term_a", "m1", "hello");
|
||||
inbox.ack("term_a", "m1");
|
||||
assertTrue(inbox.ack("term_a", "m1"), "fleetd #437: ack of a real entry must report true");
|
||||
assertTrue(inbox.peek("term_a").isEmpty(), "after ack, the message is gone");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackForUnknownMsgIdIsNoOp() {
|
||||
inbox.publish("term_a", "m1", "hello");
|
||||
inbox.ack("term_a", "no-such-id"); // no-op
|
||||
assertFalse(inbox.ack("term_a", "no-such-id"), "fleetd #437: a miss must report false"); // no-op
|
||||
assertEquals(1, inbox.peek("term_a").size(), "the published message is still there");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackForUnknownTargetIsNoOp() {
|
||||
inbox.ack("no-such-target", "m1"); // no-op, should not throw
|
||||
assertFalse(inbox.ack("no-such-target", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -167,7 +167,7 @@ class InMemoryReplyInboxTest {
|
||||
@Test
|
||||
void peekAndAckAreNoOpsForUnownedTarget() {
|
||||
assertTrue(inbox.peek("term_not_owned").isEmpty());
|
||||
inbox.ack("term_not_owned", "m1"); // no-op, should not throw
|
||||
assertFalse(inbox.ack("term_not_owned", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -23,6 +23,7 @@ import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
@@ -1036,6 +1037,11 @@ class SessionManagerTest {
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return delegate.spawn(req, decision);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
@@ -1597,6 +1603,11 @@ class SessionManagerTest {
|
||||
throw new UnsupportedOperationException("not reachable — the capability check refuses first");
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
throw new UnsupportedOperationException("not reachable — the capability check refuses first");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of("stub-profile");
|
||||
@@ -1661,6 +1672,11 @@ class SessionManagerTest {
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return delegate.spawn(req, decision);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
@@ -1867,6 +1883,11 @@ class SessionManagerTest {
|
||||
return handle;
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of("lazy");
|
||||
|
||||
Reference in New Issue
Block a user