Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha d4f93a7b13 fleetd #449: fix stale herdr protocol 14 javadocs/assertion, diagnose and fix the timing-raced AgentControlContractTest, select contract tests by tag in CI
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 2m9s
- HerdrClient.java, HerdrCodec.java, HerdrContractTest.java: the herdr port to
  protocol 19 (CB-521) left the client javadoc and the contract test's own
  assertion still saying protocol 14 / herdr 0.7.0. Updated to 19 / 0.8.0 and
  renamed pingReturnsProtocol14 -> pingReturnsProtocol19. Verified the
  assertion is real by temporarily changing the expected value to 20 (fails),
  then restoring 19 (passes).

- AgentControlContractTest.java: tabCreateInjectsEnvIntoTheSeedShell was
  failing, not skipping, on a host with a live herdr socket. Diagnosed with a
  temporary instrumented run (not committed) that polled the pane every
  200ms before and after sending input: the seed shell reliably takes ~2.5s
  to reach its prompt (measured 3x), while the test's fixed 1000ms sleep
  raced that startup. Input typed too early was swallowed by the shell's own
  startup, leaving the typed line followed by the "Restored session" banner
  and no command output — indistinguishable at a glance from the env map
  never reaching the shell. Once the shell was actually ready, the injected
  env value showed up in ~200ms, ruling out an env-seam defect. Replaced both
  fixed sleeps with bounded polling on the actual conditions (pane text
  settling, then the expected output appearing). Ran the fixed test 3x
  standalone, all green.

- .gitea/workflows/ci.yml: the "Contract tests" step ran exactly one class by
  name (-Dtest=AmqpReplyInboxContractTest), silently excluding every other
  @Tag("contract") test from CI including the herdr ones above -- which is
  how the stale protocol 14 assertion went unnoticed. Changed to
  -Dgroups=contract, which selects the whole tagged group and picks up
  future contract tests automatically.
2026-09-10 17:14:09 +07:00
9 changed files with 91 additions and 78 deletions
+10 -4
View File
@@ -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
@@ -17,7 +17,6 @@ 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;
@@ -595,23 +594,6 @@ 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();
@@ -253,27 +253,26 @@ 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>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.
* <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).
*
* @throws IllegalArgumentException if the decision names an unknown profile
*/
PeerHandle spawn(SpawnRequest req, PlacementDecision decision);
default PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return spawn(req.withProfile(decision.profile()));
}
/**
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
@@ -20,7 +20,6 @@ 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;
@@ -124,11 +123,6 @@ 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,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());
}
}
@@ -23,7 +23,6 @@ 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;
@@ -1037,11 +1036,6 @@ 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();
@@ -1603,11 +1597,6 @@ 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");
@@ -1672,11 +1661,6 @@ 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();
@@ -1883,11 +1867,6 @@ 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");