Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d4f93a7b13 | |||
| 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
|
||||
|
||||
@@ -255,7 +255,18 @@ public interface PeerLauncher {
|
||||
*
|
||||
* <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.
|
||||
* 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).
|
||||
*
|
||||
* @throws IllegalArgumentException if the decision names an unknown profile
|
||||
*/
|
||||
|
||||
@@ -17,15 +17,70 @@ 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). 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.
|
||||
*
|
||||
* <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,
|
||||
* this only avoids sending input into a shell still mid-startup.
|
||||
*/
|
||||
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 +91,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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user