diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/ConnectionIdentity.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/ConnectionIdentity.java index 01a88a6..7ebb602 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/ConnectionIdentity.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/ConnectionIdentity.java @@ -31,6 +31,20 @@ public final class ConnectionIdentity { this.cwds = cwds; } + /** + * The {@link PaneLocator} this identity resolves callers against — fleetd #612 CB-185: lets a + * test drive the exact {@link PaneLocator} a real assembly wired up (e.g. {@code + * FleetdAssembly}'s {@code new ConnectionIdentity(new PaneLocator(herdr, memberHerdr), ...)}) + * directly with a chosen pid, bypassing the OS-dependent {@link PeerPidLookup} that {@link + * #resolve} otherwise goes through. A full HTTP round trip cannot exercise this: {@code + * LsofPeerPidLookup} excludes its own pid, and an in-process test client and server share one + * JVM pid, so {@code pidForLocalPort} always returns {@code -1} and {@link PaneLocator} never + * gets called at all. + */ + public PaneLocator panes() { + return panes; + } + /** * The caller resolved from the connection: its worker {@code terminal} (or {@code null} for the * primary / an off-host client), its {@code pid} (or {@code -1} if not resolvable), and whether diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 22cb8e2..82e3367 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -107,6 +107,13 @@ public final class FleetMcp { * untested identity heuristic. */ private final boolean authorizationEnforced; + /** + * fleetd #612 CB-185: kept as a field (rather than only captured by the {@code + * contextExtractor} closure built in the constructor) so a test can reach the exact {@link + * ConnectionIdentity} — and, through {@link ConnectionIdentity#panes()}, the exact {@link + * dev.ltms.fleet.herdr.PaneLocator} — that a real assembly wired up. See {@link #identity()}. + */ + private final ConnectionIdentity identity; private final Metrics metrics; // CB-502: null → auth failures not counted private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; @@ -396,6 +403,7 @@ public final class FleetMcp { Objects.requireNonNull(callers, "callers"); this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode") == AuthorizationMode.ENFORCED; + this.identity = identity; this.leadChannel = leadChannel; this.peers = peers == null ? List.of() : List.copyOf(peers); this.capacity = capacity; @@ -755,6 +763,15 @@ public final class FleetMcp { return transport; } + /** + * The {@link ConnectionIdentity} this server resolves every caller against — fleetd #612 + * CB-185: lets a test reach the exact {@link dev.ltms.fleet.herdr.PaneLocator} a real assembly + * wired up (via {@link ConnectionIdentity#panes()}), rather than a copy built for the test. + */ + public ConnectionIdentity identity() { + return identity; + } + /** Mark a connected spawned member available for the injector readiness gate. */ static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) { if (caller.isSpawnedMember()) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyConnectionIdentityTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyConnectionIdentityTest.java new file mode 100644 index 0000000..df5afdf --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyConnectionIdentityTest.java @@ -0,0 +1,222 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.config.ConfigRef; +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.herdr.PaneLocator; +import io.javalin.Javalin; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.function.LongSupplier; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; + +/** + * fleetd #612 step 2, unit B2 (CB-185, identity half). Replaces the deleted + * {@code FleetdConnectionIdentityConstructionTest}, which pinned this claim by reading {@code + * Fleetd.java}'s source text for {@code "new PaneLocator(herdr, memberHerdr)"}. That claim moved + * to {@code FleetdAssembly.java} (fleetd #612 Unit A) and is pinned here instead, by driving the + * real {@link ConnectionIdentity} — via {@code runtime.mcp().identity()}, not a copy — that {@link + * FleetdAssembly#assembleAndStart} built. + * + *

What this guards against (from the deleted test's own javadoc): pinning + * {@code PaneLocator} to {@code memberHerdr} alone leaves every LEAD's own MCP connection + * unresolvable ({@code callerTerminal == null}) the moment {@code memberHerdrSocket} names a + * second daemon, which breaks {@code fleet_reply}/{@code fleet_ask}/{@code fleet_whoami} for a + * lead. {@code PaneLocatorTest} already proves {@link PaneLocator} itself can search two clients + * given two — the gap this pins is that the assembly actually passes it two, and in the right + * order (lead first). + * + *

Why this cannot be driven through a real MCP/HTTP round trip. The natural + * way to observe {@code ConnectionIdentity} would be a real {@code fleet_whoami} call over the + * built {@code FleetMcp}, the way {@code FleetMcpContextExtractorTest} drives its own + * hand-built one. That does not work for the REAL assembly, because {@code FleetdAssembly} wires + * {@code ConnectionIdentity} with a hardcoded {@code new LsofPeerPidLookup()} (see {@code + * FleetdAssembly.java:444}), and {@code LsofPeerPidLookup} explicitly excludes its own PID — see + * its javadoc: "we exclude our own PID and take the other end". In a JUnit test the HTTP client + * and the daemon under test run in the very same JVM, so the "client" and "server" ends of the + * loopback connection ARE the same PID, and {@code pidForLocalPort} always returns {@code -1} + * before {@link PaneLocator} is ever reached — proving nothing about which daemon(s) got searched. + * This test instead reaches the real {@link PaneLocator} the assembly built (through {@link + * ConnectionIdentity#panes()}, added for exactly this) and drives it with a chosen pid directly, + * bypassing the OS-dependent PID lookup entirely — a legitimate substitute, since the pid lookup + * is not what CB-185 is about. + */ +class FleetdAssemblyConnectionIdentityTest { + + /** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, but keys {@code connectHerdr} by + * socket path so the lead and member daemons can be two DIFFERENT {@link FakeHerdr}s. */ + private static final class TwoHerdrResourcePorts implements ResourcePorts { + + final Map herdrsBySocket = new LinkedHashMap<>(); + final CopyOnWriteArrayList schedulers = new CopyOnWriteArrayList<>(); + Runnable shutdownHook; + + @Override + public Map environment() { + return Map.of(); + } + + @Override + public HerdrClient connectHerdr(Path socketPath) { + HerdrClient client = herdrsBySocket.get(socketPath); + if (client == null) { + throw new IllegalStateException("no fake herdr registered for socket " + socketPath); + } + return client; + } + + @Override + public Fleetd.AmqpOpener replyInboxOpener() { + return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox(); + } + + @Override + public Fleetd.LeadMailboxOpener leadMailboxOpener() { + return (uri, selfCoordId, prefetch) -> { + throw new UnsupportedOperationException( + "leadMailboxOpener must not be called — no coordinator: block is configured"); + }; + } + + @Override + public LongSupplier nanoClock() { + return System::nanoTime; + } + + @Override + public LongSupplier wallClockNanos() { + return System::nanoTime; + } + + @Override + public ScheduledExecutorService newScheduler(String purpose) { + ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); + schedulers.add(scheduler); + return scheduler; + } + + @Override + public void addShutdownHook(Runnable hook) { + this.shutdownHook = hook; + } + + @Override + public void startHttp(Javalin app, String host, int port) { + // Deliberately never bind — this test never issues a real HTTP request. + } + } + + private FleetdRuntime runtime; + private TwoHerdrResourcePorts ports; + + @AfterEach + void tearDown() { + if (ports != null && ports.shutdownHook != null) { + ports.shutdownHook.run(); + } + } + + private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock"); + private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock"); + + private static FleetConfig writeConfig(Path dir) throws Exception { + Path f = dir.resolve("fleetd.yaml"); + Files.writeString(f, """ + bind: + host: 127.0.0.1 + port: 8765 + herdrSocket: "%s" + memberHerdrSocket: "%s" + lifecycle: + idleTtlSeconds: 600 + health: + enabled: false + broker: + uri: "amqp://fake-test-broker/vh" + """.formatted(LEAD_SOCKET, MEMBER_SOCKET)); + return FleetConfig.load(f); + } + + private FleetdRuntime assemble(Path dir, FakeHerdr lead, FakeHerdr member) throws Exception { + FleetConfig cfg = writeConfig(dir); + ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg); + SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet()); + ports = new TwoHerdrResourcePorts(); + ports.herdrsBySocket.put(LEAD_SOCKET, lead); + ports.herdrsBySocket.put(MEMBER_SOCKET, member); + runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports); + return runtime; + } + + /** + * The pin. {@code lead} carries the one pane {@link FakeHerdr}'s canned {@code + * pane.process_info} ties to {@link FakeHerdr#WORKER_PID} (pane {@code w2:p7}); {@code member} + * reports NO panes at all ({@link FakeHerdr#withNoPanes()}) — modelling a second daemon that + * simply does not host the caller's pane, exactly the CB-185 javadoc's scenario for a lead's + * own connection. If {@code PaneLocator} only ever searches the member daemon (the bug), this + * pid resolves to nothing, because the pane that owns it lives on the LEAD daemon the bug + * skips. + */ + @Test + void connectionIdentitySearchesTheLeadDaemonNotJustTheMemberOne(@TempDir Path dir) throws Exception { + FakeHerdr lead = new FakeHerdr(); + FakeHerdr member = new FakeHerdr().withNoPanes(); + + assemble(dir, lead, member); + + PaneLocator panes = runtime.mcp().identity().panes(); + PaneLocator.Lookup lookup = panes.terminalForPid(FakeHerdr.WORKER_PID); + + assertEquals("term_a", lookup.terminal(), + "the pane owning WORKER_PID lives on the LEAD daemon only (the member fake reports " + + "no panes) — PaneLocator must still find it, which is only possible if it " + + "searches the lead client and not just the member one"); + } + + /** + * The mirror control: when the pane instead lives ONLY on the member daemon (the lead reports + * no panes), the lookup must still find it — proving the member client is genuinely searched + * too, not merely tolerated as a second, always-losing argument. + */ + @Test + void connectionIdentityAlsoSearchesTheMemberDaemon(@TempDir Path dir) throws Exception { + FakeHerdr lead = new FakeHerdr().withNoPanes(); + FakeHerdr member = new FakeHerdr(); + + assemble(dir, lead, member); + + PaneLocator panes = runtime.mcp().identity().panes(); + PaneLocator.Lookup lookup = panes.terminalForPid(FakeHerdr.WORKER_PID); + + assertEquals("term_a", lookup.terminal(), + "the pane owning WORKER_PID lives on the MEMBER daemon only — PaneLocator must " + + "find it there too"); + } + + /** Sanity control: a pid nobody owns resolves to nothing on either daemon. */ + @Test + void aPidNoPaneOwnsResolvesToNoTerminalOnEitherDaemon(@TempDir Path dir) throws Exception { + FakeHerdr lead = new FakeHerdr(); + FakeHerdr member = new FakeHerdr(); + + assemble(dir, lead, member); + + PaneLocator panes = runtime.mcp().identity().panes(); + PaneLocator.Lookup lookup = panes.terminalForPid(999_999L); + + assertNull(lookup.terminal()); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyFleetAppTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyFleetAppTest.java new file mode 100644 index 0000000..ba2b8c3 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyFleetAppTest.java @@ -0,0 +1,237 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.config.ConfigRef; +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import io.javalin.Javalin; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.function.LongSupplier; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * fleetd #612 step 2, unit B2 (CB-185, {@code FleetApp} half). Replaces the deleted {@code + * FleetdFleetAppConstructionTest}, which pinned this claim by reading {@code Fleetd.java}'s + * source text for {@code "new FleetApp(herdr, memberHerdr, workers,"}. That claim moved to {@code + * FleetdAssembly.java} (fleetd #612 Unit A) and is pinned here instead, by driving the real {@code + * Javalin} app — via {@code runtime.app()}, not a copy — that {@link + * FleetdAssembly#assembleAndStart} built and handed to {@link FleetdRuntime}. + * + *

What this guards against (from the deleted test's own javadoc): constructing + * {@code FleetApp} with the lead-only {@code herdr} client (dropping {@code memberHerdr}) makes + * {@code GET /healthz} report green while the MEMBER daemon is down — so every spawn fails + * invisibly — and silently drops every member workspace from {@code GET /sessions}. {@code + * FleetAppTwoDaemonTest} already proves {@code FleetApp} itself merges/gates correctly given two + * clients; the gap this pins is that the assembly actually passes it two. + * + *

Both directions, not just one (fleetd #612 issue comment 17525): the deleted + * guard's positive assertion required the exact pair {@code "new FleetApp(herdr, memberHerdr, + * workers,"}, which does not survive EITHER daemon being dropped. An earlier version of this class + * only proved the member-dropped direction, which left {@code new FleetApp(memberHerdr, + * memberHerdr, ...)} — the symmetric bug, {@code /healthz} green while the LEAD daemon is down — + * an undetected regression. {@link #healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp} + * closes that. + * + *

Unlike the {@code ConnectionIdentity} half of CB-185 ({@code + * FleetdAssemblyConnectionIdentityTest}), {@code /healthz} needs no caller identity at all, so + * this test can bind {@link FleetdRuntime#app()} to a REAL ephemeral port (exactly {@code + * FleetAppTwoDaemonTest} does for its own hand-built {@code FleetApp}) and drive it with a real + * {@code HttpClient} — no accessor needed for this half. + * + *

{@code GET /sessions} could not be driven the same way, so this class does + * not pin the merge half of the deleted test's javadoc. {@code /sessions} requires + * {@code Authz.Action.READ}, which — through the REAL assembly's real {@code + * CallerResolver}/{@code ConnectionIdentity} (built with a hardcoded {@code + * new LsofPeerPidLookup()}) — needs {@code Caller.resolved()}, i.e. a real positive pid from + * {@code lsof}. {@code LsofPeerPidLookup} excludes its own pid (see its javadoc), and a JUnit + * test's HTTP client and the daemon under test share one JVM pid, so the resolved pid is always + * {@code -1} and every such request is refused as {@code ANONYMOUS} (fleetd #317's fail-closed + * rule) before the route handler — and its {@code memberHerdr} merge — is ever reached. Verified + * directly: driving {@code GET /sessions} here returns {@code 401 unauthenticated}, not the + * merged body. {@code FleetAppTwoDaemonTest} avoids this because it builds {@code FleetApp} with + * {@code callers: null}, which is not what the real assembly passes. The {@code /healthz} pin + * below is what this class relies on for CB-185's {@code FleetApp} half; {@code + * FleetAppTwoDaemonTest} remains the full behavioural proof that {@code FleetApp} itself merges + * {@code /sessions} correctly once handed two clients. + */ +class FleetdAssemblyFleetAppTest { + + private static final class TwoHerdrResourcePorts implements ResourcePorts { + + final Map herdrsBySocket = new LinkedHashMap<>(); + final CopyOnWriteArrayList schedulers = new CopyOnWriteArrayList<>(); + Runnable shutdownHook; + + @Override + public Map environment() { + return Map.of(); + } + + @Override + public HerdrClient connectHerdr(Path socketPath) { + HerdrClient client = herdrsBySocket.get(socketPath); + if (client == null) { + throw new IllegalStateException("no fake herdr registered for socket " + socketPath); + } + return client; + } + + @Override + public Fleetd.AmqpOpener replyInboxOpener() { + return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox(); + } + + @Override + public Fleetd.LeadMailboxOpener leadMailboxOpener() { + return (uri, selfCoordId, prefetch) -> { + throw new UnsupportedOperationException( + "leadMailboxOpener must not be called — no coordinator: block is configured"); + }; + } + + @Override + public LongSupplier nanoClock() { + return System::nanoTime; + } + + @Override + public LongSupplier wallClockNanos() { + return System::nanoTime; + } + + @Override + public ScheduledExecutorService newScheduler(String purpose) { + ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); + schedulers.add(scheduler); + return scheduler; + } + + @Override + public void addShutdownHook(Runnable hook) { + this.shutdownHook = hook; + } + + @Override + public void startHttp(Javalin app, String host, int port) { + // Deliberately never bind here — this test binds runtime.app() itself, for real, below. + } + } + + private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock"); + private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock"); + + private final HttpClient http = HttpClient.newHttpClient(); + private FleetdRuntime runtime; + private TwoHerdrResourcePorts ports; + private Javalin boundApp; + + @AfterEach + void tearDown() { + if (boundApp != null) { + boundApp.stop(); + } + if (ports != null && ports.shutdownHook != null) { + ports.shutdownHook.run(); + } + } + + private static FleetConfig writeConfig(Path dir) throws Exception { + Path f = dir.resolve("fleetd.yaml"); + Files.writeString(f, """ + bind: + host: 127.0.0.1 + port: 8765 + herdrSocket: "%s" + memberHerdrSocket: "%s" + lifecycle: + idleTtlSeconds: 600 + health: + enabled: false + broker: + uri: "amqp://fake-test-broker/vh" + """.formatted(LEAD_SOCKET, MEMBER_SOCKET)); + return FleetConfig.load(f); + } + + /** Assembles the real graph, then binds the real {@code Javalin app} to an ephemeral port. */ + private int assembleAndBind(Path dir, FakeHerdr lead, FakeHerdr member) throws Exception { + FleetConfig cfg = writeConfig(dir); + ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg); + SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet()); + ports = new TwoHerdrResourcePorts(); + ports.herdrsBySocket.put(LEAD_SOCKET, lead); + ports.herdrsBySocket.put(MEMBER_SOCKET, member); + runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports); + boundApp = runtime.app().start("127.0.0.1", 0); + return boundApp.port(); + } + + private HttpResponse get(int port, String path) throws Exception { + HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path)).GET().build(); + return http.send(req, HttpResponse.BodyHandlers.ofString()); + } + + /** + * The pin. The MEMBER daemon is down; the LEAD daemon is healthy. If the assembly built + * {@code FleetApp} with only the lead client (the bug: passing {@code herdr} where {@code + * memberHerdr} is expected), the down member is invisible and {@code /healthz} stays 200. + */ + @Test + void healthzGoesRedWhenTheMemberDaemonIsDownEvenThoughTheLeadIsUp(@TempDir Path dir) throws Exception { + FakeHerdr lead = new FakeHerdr(); + FakeHerdr member = new FakeHerdr().healthy(false); + + int port = assembleAndBind(dir, lead, member); + + HttpResponse res = get(port, "/healthz"); + assertEquals(503, res.statusCode(), + "a down MEMBER daemon must not be masked by a healthy lead: " + res.body()); + } + + /** Sanity control: both daemons healthy must still be green through the real assembly. */ + @Test + void healthzIsGreenWhenBothDaemonsAreUp(@TempDir Path dir) throws Exception { + FakeHerdr lead = new FakeHerdr(); + FakeHerdr member = new FakeHerdr(); + + int port = assembleAndBind(dir, lead, member); + + assertEquals(200, get(port, "/healthz").statusCode()); + } + + /** + * The symmetric pin (fleetd #612 issue comment 17525): the LEAD daemon is down; the MEMBER + * daemon is healthy. If the assembly built {@code FleetApp} with only the member client + * (dropping {@code herdr} — the mirror of the bug above, {@code new FleetApp(memberHerdr, + * memberHerdr, ...)}), the down LEAD is invisible and {@code /healthz} stays 200. Without this + * case the pair above is one-directional and does not cover the deleted guard's positive + * assertion (it required BOTH {@code herdr,} and {@code memberHerdr,} in that order). + */ + @Test + void healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp(@TempDir Path dir) throws Exception { + FakeHerdr lead = new FakeHerdr().healthy(false); + FakeHerdr member = new FakeHerdr(); + + int port = assembleAndBind(dir, lead, member); + + HttpResponse res = get(port, "/healthz"); + assertEquals(503, res.statusCode(), + "a down LEAD daemon must not be masked by a healthy member: " + res.body()); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdCompletionResolverAssemblyTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdCompletionResolverAssemblyTest.java new file mode 100644 index 0000000..664af23 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdCompletionResolverAssemblyTest.java @@ -0,0 +1,307 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.config.ConfigRef; +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.inject.CompletionResolver; +import dev.ltms.fleet.msg.Rendezvous; +import dev.ltms.fleet.msg.TurnToken; +import dev.ltms.fleet.placement.PlacementException; +import dev.ltms.fleet.session.MemberSession; +import dev.ltms.fleet.session.WorktreeRequest; +import io.javalin.Javalin; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #612 step 2, Unit B1 — replaces {@code FleetdCompletionResolverWiringTest} (deleted in + * this same commit), whose four tests read {@code Fleetd.java}'s source text and asserted the + * {@code CompletionResolver} construction call still named the right arguments. That proved the + * call site's spelling, never that the assembled resolver actually behaves differently when an + * argument is dropped. + * + *

These tests drive {@link FleetdAssembly#assembleAndStart} — the real boot composition, + * fleetd #612 Unit A — and read {@link FleetdRuntime#completion()}: the exact {@link + * CompletionResolver} instance the assembled daemon uses, never a copy built alongside it for the + * test's benefit. Two behaviours are pinned, matching the ticket's own two measured mutations: + * + *

+ * + *

Both tests bypass {@link dev.ltms.fleet.inject.StatusPoller} and drive {@link + * CompletionResolver#onDelivered} / {@link CompletionResolver#resolveBeforePostAction} directly — + * the same public, synchronous entry points {@code CompletionResolverTest} uses — with a + * hand-built {@link CompletableFuture} waiter, so no real poller loop or herdr status poll is + * needed. The pane scrape comes from {@link FakeHerdr#readText}; the elapsed-time floor + * ({@code CompletionResolver.MIN_TURN_NANOS}) is controlled via a fake, advanceable {@link + * ResourcePorts#nanoClock()} rather than a real sleep. + */ +class FleetdCompletionResolverAssemblyTest { + + /** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, plus a nanoClock this test can advance. */ + private static final class ControllableResourcePorts implements ResourcePorts { + + final FakeHerdr herdr; + final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start + Runnable shutdownHook; + + ControllableResourcePorts(FakeHerdr herdr) { + this.herdr = herdr; + } + + void advanceSeconds(long seconds) { + nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds)); + } + + @Override + public Map environment() { + return Map.of(); + } + + @Override + public HerdrClient connectHerdr(Path socketPath) { + return herdr; + } + + @Override + public Fleetd.AmqpOpener replyInboxOpener() { + // Never invoked: this test's config has no `broker:` block, so Fleetd.selectReplyInbox + // returns the in-memory inbox before calling the opener at all. + return (uri, prefetch) -> { + throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block"); + }; + } + + @Override + public Fleetd.LeadMailboxOpener leadMailboxOpener() { + // Never invoked: no `coordinator:` block configured either. + return (uri, selfCoordId, prefetch) -> { + throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block"); + }; + } + + @Override + public LongSupplier nanoClock() { + return nowNanos::get; + } + + @Override + public LongSupplier wallClockNanos() { + return nowNanos::get; + } + + @Override + public ScheduledExecutorService newScheduler(String purpose) { + return Executors.newSingleThreadScheduledExecutor(); + } + + @Override + public void addShutdownHook(Runnable hook) { + this.shutdownHook = hook; + } + + @Override + public void startHttp(Javalin app, String host, int port) { + // Deliberately never bind a real port. + } + } + + private static FleetConfig writeConfig(Path dir, String profilesYaml, String extraGuardHost, + String worktreeRootYamlLine) throws Exception { + Path f = dir.resolve("fleetd.yaml"); + Files.writeString(f, """ + bind: + host: 127.0.0.1 + port: 8765 + idleSleepGuard: + enabled: false + %s + profiles: + %s + guard: + offSubscriptionHosts: + - %s + """.formatted(worktreeRootYamlLine == null ? "" : worktreeRootYamlLine, profilesYaml, extraGuardHost)); + return FleetConfig.load(f); + } + + private static void gitQuiet(Path cwd, String... args) throws Exception { + List cmd = new java.util.ArrayList<>(List.of("git")); + cmd.addAll(List.of(args)); + Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start(); + String out = new String(p.getInputStream().readAllBytes()); + assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args)); + assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out); + } + + private static Path initRepo(Path dir) throws Exception { + Files.createDirectories(dir); + gitQuiet(dir, "init", "-q", "-b", "main"); + gitQuiet(dir, "config", "user.email", "test@example.invalid"); + gitQuiet(dir, "config", "user.name", "Test"); + Files.writeString(dir.resolve("README.md"), "seed\n"); + gitQuiet(dir, "add", "README.md"); + gitQuiet(dir, "commit", "-q", "-m", "seed"); + return dir; + } + + /** + * fleetd #248's first measured mutation: replacing {@code CompletionResolver}'s 8th constructor + * argument with the inert {@code _ -> null} compiles clean and leaves every existing test green + * — it silently drops fleetd #241's fallback-report location. This drives the real assembled + * resolver through a member echoing its own injected brief back (no {@code fleet_reply}), which + * resolves via {@code noReportMessage(target)}, and proves the real member's {@code branch} — + * only obtainable via {@code Fleetd.worktreeBranchLookup(sessions::roster)} reading the real, + * worktree-provisioned {@link MemberSession} — appears in the reported text. + */ + @Test + void assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport(@TempDir Path dir) throws Exception { + Path repo = initRepo(dir.resolve("repo")); + FleetConfig cfg = writeConfig(dir, """ + wtprofile: + baseUrl: http://wthost.local:8000 + model: sonnet + """, "wthost.local", "worktreeRoot: " + dir.resolve("wts")); + ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg); + SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet()); + ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr()); + + FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports); + try { + MemberSession session = runtime.sessions().acquire("wtprofile", repo.toString(), repo.toString(), + null, new WorktreeRequest("fleetd-612-b1", null)); + String target = session.terminalId(); + String branch = session.branch(); + assertTrue(branch != null && branch.startsWith("worker/"), + "sanity: a worktree-provisioned session must carry a real branch, got: " + branch); + + CompletionResolver completion = runtime.completion(); + CompletableFuture waiter = new CompletableFuture<>(); + String echoedBrief = "z".repeat(450); // >= CompletionResolver.ECHO_MIN_CHARS normalised chars + + ports.herdr.readText("idle, nothing yet"); + completion.onDelivered(target, new TurnToken(target, waiter, echoedBrief)); + + ports.herdr.readText(echoedBrief); // the pane just echoes the injected brief back — no real report + ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s) without a real sleep + completion.resolveBeforePostAction(target); + + Rendezvous.Resolution resolution = waiter.getNow(null); + assertTrue(resolution != null, "the waiter must have resolved synchronously"); + assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind()); + assertTrue(resolution.text().contains(CompletionResolver.NO_REPORT_PREFIX), + "sanity: must have gone down the noReportMessage sub-path: " + resolution.text()); + assertTrue(resolution.text().contains("branch=" + branch), + "the assembled resolver must report the member's real branch (fleetd #241 via " + + "fleetd #248's worktreeBranchLookup wiring); got: " + resolution.text()); + } finally { + if (ports.shutdownHook != null) ports.shutdownHook.run(); + } + } + + /** + * fleetd #248's second measured mutation, and fleetd#201 Unit 5's own gap: replacing {@code + * backendErrorPatterns}/{@code backendErrorSink} with {@code BackendErrorPatternLookup.legacy()} + * / {@code BackendErrorSink.none()} compiles clean and leaves every existing behavioural test + * green. + * + *

Classification proof: this test's profile configures {@code errorPattern: "credential + * outage"} — text the built-in {@code (?i)\bAPI Error\s*:} fallback ({@code legacy()}'s only + * behaviour) never matches. So a real {@code Fleetd.backendErrorPatternLookup(...)} wiring + * classifies the send as {@code FAILED}; {@code legacy()} would fall through to the plain + * completion path instead ({@code Kind.COMPLETION}). + * + *

Cool-off proof: two distinct targets on the same profile/credential each classified as a + * backend error inside the 60s window must cool the credential off ({@link + * dev.ltms.fleet.placement.BackendOutagePolicy}, fleetd#201 Unit 5) — observable two ways: (1) + * the real {@code Fleetd.backendErrorSink(...)} marks each session {@code BACKEND_ERROR} (only + * the real sink calls {@code sessions.onBackendError}; {@code BackendErrorSink.none()} never + * does), and (2) a third explicit-profile spawn attempt is refused with a {@link + * PlacementException} naming the cool-off — only reachable because the real sink's {@code + * outagePolicy.record(...)} call actually ran. + */ + @Test + void assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern(@TempDir Path dir) throws Exception { + FleetConfig cfg = writeConfig(dir, """ + coolprofile: + baseUrl: http://coolhost.local:8000 + model: sonnet + errorPattern: "credential outage" + """, "coolhost.local", null); + ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg); + SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet()); + ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr()); + + FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports); + try { + MemberSession session1 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null); + MemberSession session2 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null); + String target1 = session1.terminalId(); + String target2 = session2.terminalId(); + assertTrue(!target1.equals(target2), "sanity: the two spawns must be distinct targets"); + + CompletionResolver completion = runtime.completion(); + + CompletableFuture waiter1 = new CompletableFuture<>(); + ports.herdr.readText("idle 1"); + completion.onDelivered(target1, new TurnToken(target1, waiter1, null)); + ports.herdr.readText("credential outage: upstream 503"); + ports.advanceSeconds(3); + completion.resolveBeforePostAction(target1); + Rendezvous.Resolution resolution1 = waiter1.getNow(null); + assertTrue(resolution1 != null, "target1's waiter must have resolved synchronously"); + assertEquals(Rendezvous.Kind.FAILED, resolution1.kind(), + "a configured errorPattern the built-in fallback never matches must classify as " + + "a backend error, not a plain completion; got: " + resolution1); + assertTrue(resolution1.text().contains("credential outage: upstream 503"), resolution1.text()); + + CompletableFuture waiter2 = new CompletableFuture<>(); + ports.herdr.readText("idle 2"); + completion.onDelivered(target2, new TurnToken(target2, waiter2, null)); + ports.herdr.readText("credential outage: upstream 503 again"); + ports.advanceSeconds(3); + completion.resolveBeforePostAction(target2); + Rendezvous.Resolution resolution2 = waiter2.getNow(null); + assertTrue(resolution2 != null, "target2's waiter must have resolved synchronously"); + assertEquals(Rendezvous.Kind.FAILED, resolution2.kind()); + + List roster = runtime.sessions().roster(); + assertTrue(roster.stream().anyMatch(s -> target1.equals(s.terminalId()) + && s.state() == MemberSession.State.BACKEND_ERROR), + "the real backendErrorSink must have transitioned target1 to BACKEND_ERROR: " + roster); + assertTrue(roster.stream().anyMatch(s -> target2.equals(s.terminalId()) + && s.state() == MemberSession.State.BACKEND_ERROR), + "the real backendErrorSink must have transitioned target2 to BACKEND_ERROR: " + roster); + + PlacementException coolOff = assertThrows(PlacementException.class, + () -> runtime.sessions().acquire("coolprofile", null, dir.toString(), null), + "two distinct targets classified within the 60s window must cool the credential " + + "off (BackendOutagePolicy), refusing a third explicit-profile spawn"); + assertTrue(coolOff.getMessage().contains("cooling off"), coolOff.getMessage()); + } finally { + if (ports.shutdownHook != null) ports.shutdownHook.run(); + } + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdCompletionResolverWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdCompletionResolverWiringTest.java deleted file mode 100644 index f8a1386..0000000 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdCompletionResolverWiringTest.java +++ /dev/null @@ -1,89 +0,0 @@ -package dev.ltms.fleet; - -import java.nio.file.Files; -import java.nio.file.Path; -import org.junit.jupiter.api.DisplayName; -import org.junit.jupiter.api.Test; - -import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertTrue; - -/** - * fleetd #248: this is the test that was actually missing. {@code Fleetd.main} builds its {@code - * CompletionResolver} from an 8-argument constructor, and the ticket's own measurement proved two - * ways to silently unwire it — both compiled with 0 errors and left every existing test green: - * - *

- * - *

Neither mutation could be caught by any test that constructs its own {@code - * CompletionResolver} (every test before this one did exactly that) or by a test of {@link - * Fleetd#worktreeBranchLookup}, {@link Fleetd#backendErrorPatternLookup}, or {@link - * Fleetd#backendErrorSink} in isolation (see {@code FleetdWorktreeBranchLookupTest}, {@code - * FleetdBackendErrorPatternLookupTest}, {@code FleetdBackendErrorSinkTest}) — those prove the - * factories work, never that {@code main} still calls them. This class is a plain source-text - * assertion on {@code Fleetd.java} — crude, but honest about what it checks, and it turns red the - * instant the wiring is dropped, mirroring the same fallback shape {@link - * FleetdFleetAppConstructionTest} already uses for a different constructor argument. - * - *

This test checks source text, not runtime behaviour. It never constructs a {@code - * CompletionResolver} and never runs {@code main}. - */ -class FleetdCompletionResolverWiringTest { - - private static String fleetdSource() throws Exception { - return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java")); - } - - @Test - @DisplayName("[SOURCE TEXT] CompletionResolver's construction call still names backendErrorPatterns and backendErrorSink") - void backendErrorArgumentsAreStillNamedAtTheCallSite() throws Exception { - String source = fleetdSource(); - assertTrue(source.contains( - "exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,"), - "CompletionResolver's construction call must still pass backendErrorPatterns and " - + "backendErrorSink as its 5th/6th arguments. Replacing them with " - + "BackendErrorPatternLookup.legacy()/BackendErrorSink.none() (fleetd #248's measured " - + "mutation) compiles with 0 errors and leaves every behavioural test green — this " - + "source check is what must go red instead."); - } - - @Test - @DisplayName("[SOURCE TEXT] CompletionResolver's construction call still passes worktreeBranchLookup(sessions::roster)") - void worktreeBranchLookupIsStillPassedAtTheCallSite() throws Exception { - String source = fleetdSource(); - assertTrue(source.contains("worktreeBranchLookup(sessions::roster)"), - "CompletionResolver's construction call must still pass worktreeBranchLookup(sessions::roster) " - + "as its 8th (last) argument. Replacing it with the inert `_ -> null` (fleetd #248's " - + "other measured mutation) compiles with 0 errors and leaves every behavioural test " - + "green — this source check is what must go red instead."); - assertFalse(source.contains("System::nanoTime,\n _ -> null"), - "the worktree/branch argument must never regress to the inert `_ -> null` literal"); - } - - @Test - @DisplayName("[SOURCE TEXT] backendErrorPatterns is assigned from the extracted backendErrorPatternLookup(...) factory") - void backendErrorPatternsComesFromTheFactory() throws Exception { - String source = fleetdSource(); - assertTrue(source.contains( - "BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,"), - "backendErrorPatterns must be assigned from Fleetd.backendErrorPatternLookup(...), not an " - + "inline lambda that a source check on the CompletionResolver call alone cannot see " - + "through"); - } - - @Test - @DisplayName("[SOURCE TEXT] backendErrorSink is assigned from the extracted backendErrorSink(...) factory") - void backendErrorSinkComesFromTheFactory() throws Exception { - String source = fleetdSource(); - assertTrue(source.contains( - "BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),"), - "backendErrorSink must be assigned from Fleetd.backendErrorSink(...), not an inline lambda " - + "that a source check on the CompletionResolver call alone cannot see through"); - } -} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdConnectionIdentityConstructionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdConnectionIdentityConstructionTest.java deleted file mode 100644 index 7274264..0000000 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdConnectionIdentityConstructionTest.java +++ /dev/null @@ -1,31 +0,0 @@ -package dev.ltms.fleet; - -import java.nio.file.Files; -import java.nio.file.Path; -import org.junit.jupiter.api.Test; - -import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertTrue; - -/** - * CB-185: {@code ConnectionIdentity} must resolve a caller's pane on EITHER herdr daemon (a - * lead's MCP connection resolves against the lead daemon; a member's against the member daemon). - * Pinning {@code PaneLocator} to {@code memberHerdr} alone — the bug this guards against — leaves - * every lead's own connection unresolvable ({@code callerTerminal == null}) the moment - * {@code memberHerdrSocket} names a second daemon, which breaks {@code fleet_reply}/{@code - * fleet_ask} and {@code fleet_whoami} for a lead. A unit test on {@link - * dev.ltms.fleet.herdr.PaneLocator} alone (see {@code PaneLocatorTest}) proves the class CAN - * search two clients, but not that {@code Fleetd.main} actually wires it that way — hence this - * source-level assertion, the same technique {@code FleetdHerdrControlConstructionTest} uses. - */ -class FleetdConnectionIdentityConstructionTest { - @Test - void connectionIdentitySearchesBothDaemonsNotJustTheMemberOne() throws Exception { - String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java")); - assertFalse(source.contains("new PaneLocator(memberHerdr)"), - "PaneLocator must not be pinned to the member daemon alone — a lead's own " - + "connection resolves against the LEAD daemon and would never be found"); - assertTrue(source.contains("new PaneLocator(herdr, memberHerdr)"), - "PaneLocator must search the lead daemon first, then the member daemon"); - } -} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdFleetAppConstructionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdFleetAppConstructionTest.java deleted file mode 100644 index a7e8482..0000000 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdFleetAppConstructionTest.java +++ /dev/null @@ -1,29 +0,0 @@ -package dev.ltms.fleet; - -import java.nio.file.Files; -import java.nio.file.Path; -import org.junit.jupiter.api.Test; - -import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertTrue; - -/** - * CB-185: {@code FleetApp} must be constructed with BOTH herdr clients (the lead's and the - * member's), never the raw lead-only {@code herdr}. Passing only {@code herdr} — the bug this - * guards against — makes {@code GET /healthz} green while the member daemon is down (so every - * spawn fails invisibly) and silently drops every member workspace from {@code GET /sessions}. - * A behavioural test on {@code FleetApp} alone (see {@code FleetAppTwoDaemonTest}) proves the - * class merges/gates correctly when given two clients, but not that {@code Fleetd.main} actually - * passes it two — hence this source-level assertion, mirroring - * {@code FleetdHerdrControlConstructionTest}. - */ -class FleetdFleetAppConstructionTest { - @Test - void fleetAppIsConstructedWithBothHerdrDaemons() throws Exception { - String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java")); - assertFalse(source.contains("new FleetApp(herdr, workers,"), - "FleetApp must not be constructed with the lead-only herdr client"); - assertTrue(source.contains("new FleetApp(herdr, memberHerdr, workers,"), - "FleetApp must be constructed with both the lead and the member herdr client"); - } -}