Compare commits
15 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c1ca6273fc | |||
| e54e3d87ea | |||
| 3c5873dfe2 | |||
| e29227d5f4 | |||
| 2af13ab1ff | |||
| 0788d84be8 | |||
| 20c1094cbf | |||
| 9011c59b9f | |||
| c11ad71ed0 | |||
| bdcf285265 | |||
| d4f93a7b13 | |||
| cfebc575ea | |||
| 822327eed5 | |||
| e4c703a51a | |||
| 3f036b2a62 |
+10
-4
@@ -87,12 +87,18 @@ jobs:
|
||||
apt-get update && apt-get install -y --no-install-recommends maven
|
||||
mvn -version
|
||||
|
||||
# The `contract` profile clears the default-excludes group, so the @Tag("contract") AMQP test
|
||||
# runs against the RabbitMQ service container (AMQP_URI). Pinned to the one contract test to
|
||||
# avoid re-running the unit suite already covered by the `build` job.
|
||||
# The `contract` profile clears the default-excludes group, so `-Dgroups=contract` runs every
|
||||
# @Tag("contract") test and nothing from the unit suite the `build` job already covered — a
|
||||
# tag selects the whole group, so a test added to it later runs here automatically. A prior
|
||||
# version of this step pinned `-Dtest=AmqpReplyInboxContractTest` by class name instead: that
|
||||
# silently excluded every other contract test (including the herdr ones) from CI, and nobody
|
||||
# noticed until the herdr protocol drifted out from under a test that never ran here
|
||||
# (fleetd #449). If this runner has no herdr socket, the herdr-backed tests in the group
|
||||
# skip on their own `assumeTrue` and only the broker-backed ones actually run — check the
|
||||
# step output rather than assuming which.
|
||||
- name: Contract tests
|
||||
working-directory: fleetd
|
||||
run: mvn -B -Pcontract test -Dtest=AmqpReplyInboxContractTest
|
||||
run: mvn -B -Pcontract test -Dgroups=contract
|
||||
|
||||
- name: Failing test output
|
||||
if: failure()
|
||||
|
||||
@@ -3,7 +3,7 @@ package dev.ltms.fleet.herdr;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
|
||||
/**
|
||||
* Client face onto the herdr daemon (protocol 14, herdr 0.7.0).
|
||||
* Client face onto the herdr daemon (protocol 19, herdr 0.8.0).
|
||||
*
|
||||
* <p>This is the ONLY thing in {@code fleetd} that speaks to herdr. Every method
|
||||
* maps to a herdr JSON-RPC call over its Unix domain socket. Requests are
|
||||
|
||||
@@ -8,7 +8,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
|
||||
/**
|
||||
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 14).
|
||||
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 19).
|
||||
*
|
||||
* <p>Split out from the socket so the framing rules — the ones that actually bit us
|
||||
* during the spike (id MUST be a string; response carries {@code result} or
|
||||
|
||||
@@ -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);
|
||||
@@ -1344,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)
|
||||
@@ -1371,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,
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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,15 +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.
|
||||
* <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();
|
||||
|
||||
@@ -17,15 +17,93 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
* SHELL directly (never {@code claude}, so no subscription/token involvement) and always tears
|
||||
* 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).
|
||||
* 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}.
|
||||
*/
|
||||
@Tag("contract")
|
||||
class AgentControlContractTest {
|
||||
|
||||
private static final long POLL_INTERVAL_MS = 150;
|
||||
/** Bound for the seed shell to settle: observed ~2.5s three times running; this leaves headroom. */
|
||||
private static final long SHELL_READY_TIMEOUT_MS = 8_000;
|
||||
/** Bound for the typed command's output to appear once the shell is ready: observed ~0.2s. */
|
||||
private static final long OUTPUT_TIMEOUT_MS = 5_000;
|
||||
|
||||
private boolean noSocket() {
|
||||
return !Files.exists(UnixSocketHerdrClient.defaultSocketPath());
|
||||
}
|
||||
|
||||
private static String readPane(UnixSocketHerdrClient herdr, String paneId) {
|
||||
return herdr.call("pane.read", Map.of("pane_id", paneId, "source", "visible"))
|
||||
.path("read").path("text").asText("");
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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 {
|
||||
long deadline = System.currentTimeMillis() + timeoutMs;
|
||||
String previous = null;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(POLL_INTERVAL_MS);
|
||||
String current = readPane(herdr, paneId);
|
||||
if (current.equals(previous) && !current.isBlank()) {
|
||||
return current;
|
||||
}
|
||||
previous = current;
|
||||
}
|
||||
return previous == null ? "" : previous;
|
||||
}
|
||||
|
||||
/** Poll {@code pane.read} until {@code needle} appears or {@code timeoutMs} elapses. */
|
||||
private static String waitForText(UnixSocketHerdrClient herdr, String paneId, String needle, long timeoutMs)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + timeoutMs;
|
||||
String last = "";
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
last = readPane(herdr, paneId);
|
||||
if (last.contains(needle)) {
|
||||
return last;
|
||||
}
|
||||
Thread.sleep(POLL_INTERVAL_MS);
|
||||
}
|
||||
return last;
|
||||
}
|
||||
|
||||
@Test
|
||||
void tabCreateInjectsEnvIntoTheSeedShell() throws Exception {
|
||||
assumeTrue(!noSocket(), "no herdr socket — skipping");
|
||||
@@ -36,15 +114,13 @@ class AgentControlContractTest {
|
||||
Map.of("ANTHROPIC_BASE_URL", "http://gx00.gw:8000"));
|
||||
try {
|
||||
assertNotNull(tab.rootPaneId(), "tab.create must return the seed pane");
|
||||
Thread.sleep(1000); // let the seed shell reach its prompt
|
||||
waitUntilSettled(herdr, tab.rootPaneId(), SHELL_READY_TIMEOUT_MS);
|
||||
herdr.call("pane.send_input", Map.of(
|
||||
"pane_id", tab.rootPaneId(),
|
||||
"text", "printf 'PROBE_BASE=[%s]\\n' \"$ANTHROPIC_BASE_URL\"",
|
||||
"keys", List.of("enter")));
|
||||
Thread.sleep(800);
|
||||
String visible = herdr.call("pane.read",
|
||||
Map.of("pane_id", tab.rootPaneId(), "source", "visible"))
|
||||
.path("read").path("text").asText("");
|
||||
String visible = waitForText(herdr, tab.rootPaneId(),
|
||||
"PROBE_BASE=[http://gx00.gw:8000]", OUTPUT_TIMEOUT_MS);
|
||||
assertTrue(visible.contains("PROBE_BASE=[http://gx00.gw:8000]"),
|
||||
"env map must reach the seed shell; saw: " + visible);
|
||||
} finally {
|
||||
|
||||
@@ -14,7 +14,7 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
* Contract test against a REAL running herdr. Tagged {@code contract} so it is
|
||||
* excluded from {@code mvn test}; run it with {@code mvn test -Pcontract}. It fails
|
||||
* loudly if herdr drifts from the protocol {@code fleetd} was built against
|
||||
* (0.7.0, protocol 14) — catching breakage that unit tests with canned frames cannot.
|
||||
* (0.8.0, protocol 19) — catching breakage that unit tests with canned frames cannot.
|
||||
*/
|
||||
@Tag("contract")
|
||||
class HerdrContractTest {
|
||||
@@ -24,13 +24,13 @@ class HerdrContractTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void pingReturnsProtocol14() {
|
||||
void pingReturnsProtocol19() {
|
||||
assumeTrue(Files.exists(socket()), "no herdr socket at " + socket() + " — skipping");
|
||||
try (UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect()) {
|
||||
JsonNode pong = herdr.call("ping");
|
||||
assertEquals("pong", pong.get("type").asText());
|
||||
assertEquals(14, pong.get("protocol").asInt(),
|
||||
"fleetd is built against herdr protocol 14");
|
||||
assertEquals(19, pong.get("protocol").asInt(),
|
||||
"fleetd is built against herdr protocol 19");
|
||||
assertFalse(pong.get("version").asText().isBlank());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -20,6 +20,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
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.PlacementException;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -1123,6 +1124,68 @@ class CompositePeerLauncherTest {
|
||||
assertEquals(0, adapter.spawnCount("b"), "routedProfileFor never spawns anything");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #444: {@link PlacementDecision} exists to close the window between {@link
|
||||
* CompositePeerLauncher#place} and {@link CompositePeerLauncher#spawn(SpawnRequest,
|
||||
* PlacementDecision)} — the placement state must be free to move in that window without the
|
||||
* held decision being re-checked against the new state. Every quarantine test above resolves
|
||||
* and spawns in one call, so none of them ever open that window; this test is the one that
|
||||
* does: "sol" is placed FIRST, while nothing is quarantined yet, and only THEN is its
|
||||
* credential quarantined, before the held decision is spawned.
|
||||
*
|
||||
* <p>This is the test that tells the real override apart from the alternative body the ticket
|
||||
* measured: routing {@code decision.profile()} straight to its adapter (the real override)
|
||||
* never re-runs {@code enforceNotQuarantined}, so the spawn against the held decision still
|
||||
* succeeds on sol. Re-entering {@code spawn(req.withProfile(decision.profile()))} instead
|
||||
* lands in the explicit-profile branch, which refuses a now-quarantined sol outright — before
|
||||
* this test existed, replacing the real override's body with that re-entering call left the
|
||||
* whole suite green.
|
||||
*/
|
||||
@Test
|
||||
void spawnHonorsAPlacementDecisionEvenAfterItsProfileIsQuarantinedInTheWindowAfterPlace() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"sol", stubWorker("sol", "shared-openai"),
|
||||
"b", stubWorker("b"));
|
||||
// The adapter's OWN fallback default is "b", deliberately different from the profile place()
|
||||
// decides ("sol") — see the note below on why this must not be "sol" too.
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "b", Set.of());
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
|
||||
|
||||
// 1. Resolve BEFORE anything is quarantined — sol (definition order first, fixed policy) wins.
|
||||
// composite's own defaultProfile ("sol", the constructor arg above) never enters this: the
|
||||
// pool poolFor(DEV) resolves to is never empty here, so place() only ever reads that field as
|
||||
// a fallback for an empty pool, which this test does not exercise.
|
||||
PlacementDecision decision = composite.place(MemberRole.DEV);
|
||||
assertEquals("sol", decision.profile(), "sanity: nothing is quarantined yet, so sol is placed");
|
||||
|
||||
// 2. Move the placement state IN THE WINDOW between place() and spawn() — sol's credential
|
||||
// is now quarantined. A fresh place()/spawn(req) pair would fall through to b instead; the
|
||||
// held decision must not be re-evaluated against this new state at all.
|
||||
quarantine.quarantine("shared-openai");
|
||||
|
||||
// 3. Spawn against the HELD decision, not a fresh resolve.
|
||||
SpawnRequest req = new SpawnRequest(null, null, null, null, null, MemberRole.DEV);
|
||||
PeerHandle handle = composite.spawn(req, decision);
|
||||
|
||||
assertEquals("sol", handle.profile(),
|
||||
"the decision from place() is honored even though sol is now quarantined");
|
||||
// A fixture whose adapter falls back to "sol" too would let an UNSTAMPED request (one
|
||||
// routed but never given req.withProfile("sol")) land on spawnCount("sol") == 1 by
|
||||
// COINCIDENCE, since StubLauncher.spawn falls back to its own defaultProfile whenever
|
||||
// req.profileName() is blank. Giving the adapter "b" as its fallback instead means only an
|
||||
// actually-stamped request can produce this count — an unstamped one would count against
|
||||
// "b" and this assertion would fail.
|
||||
assertEquals(1, adapter.spawnCount("sol"),
|
||||
"the request that reached the delegate actually carried sol as its profile "
|
||||
+ "(the adapter's own fallback default is 'b', so this can't happen by accident)");
|
||||
assertEquals(0, adapter.spawnCount("b"),
|
||||
"b must never be touched — neither as the decision's profile nor as an unstamped "
|
||||
+ "request's accidental fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aQuarantineLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -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