diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index a93637f..a6355bf 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -24,6 +24,7 @@ import dev.ltms.fleet.inject.Injector; import dev.ltms.fleet.inject.LoopWatchdog; import dev.ltms.fleet.inject.StatusPoller; import dev.ltms.fleet.inject.TurnListener; +import dev.ltms.fleet.inject.TurnRegistrar; import dev.ltms.fleet.inject.MemberPresence; import dev.ltms.fleet.auth.MemberRegistry; import dev.ltms.fleet.auth.CallerResolver; @@ -80,7 +81,9 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BiConsumer; import java.util.function.BooleanSupplier; +import java.util.function.Consumer; import java.util.function.Function; import java.util.function.LongSupplier; import java.util.function.Predicate; @@ -501,21 +504,26 @@ public final class Fleetd { // fleetd #556: registration is wired directly to `completion`, not folded into the // `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future // listener) throwing, regardless of call order. See TurnRegistrar's javadoc. + // fleetd #589 Group 3 (:505): extracted to turnRegistrar(...) — see FleetdTurnRegistrarWiringTest. Injector injector = new Injector(router, turnListener, deliverable, - presence::forget, completion::register); + presence::forget, turnRegistrar(completion)); StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS); poller.start(); // CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or // unusable), fleetd stays soft-state on the in-memory inbox. The AMQP inbox owns a broker // connection, so keep the reference to close it in the ordered shutdown hook. - final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open); + // fleetd #589 Group 3 (:512): opener extracted to replyInboxOpener() — see + // FleetdReplyInboxOpenerWiringTest. + final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), replyInboxOpener()); // CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate // broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator: // block this is null and every lead path below is simply not wired, which is exactly the // behaviour before this ticket. It owns a broker connection, so keep the reference for the // ordered shutdown hook. - final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open); + // fleetd #589 Group 3 (:518): opener extracted to leadMailboxOpener() — see + // FleetdLeadMailboxOpenerWiringTest. + final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), leadMailboxOpener()); // CB-307: learn the primary's terminal from orchestration tool calls (or pin from config). // The pin also feeds CallerResolver below: a primary running inside a herdr pane would // otherwise resolve as a worker and be refused every orchestration tool. @@ -585,10 +593,12 @@ public final class Fleetd { // fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also // gets a wall-clock source to detect and correct for that freeze. Every other decision // in FleetHealthMonitor stays on the monotonic clock, unchanged. + // fleetd #589 Group 3 (:591): failTarget extracted to healthFailTarget(...) — see + // FleetdHealthFailTargetWiringTest. healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler, System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()), cfg.health().intervalOrDefault(), - cfg.health().workingSuspectAfterOrDefault(), messages::abandon); + cfg.health().workingSuspectAfterOrDefault(), healthFailTarget(messages)); String coverage = FleetHealthMonitor.coverage(true, cfg.health().notifications() != null && cfg.health().notifications().configured()); if ("detection-only".equals(coverage)) { @@ -608,27 +618,11 @@ public final class Fleetd { // CB-516: releasing a worker must fail whatever send was waiting on it. Without this a // torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never // reached /metrics — the delegation was unresolvable and nothing said so. - sessions.onRelease(detail -> { - // CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead - // where to re-dispatch onto the same tree, not just that the worker vanished. - String reason = "the worker session was released before it replied"; - if (detail.worktreePath() != null) { - reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch() - + " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none"); - } - // CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the - // member's conversation instead of only re-dispatching a fresh one onto the same files. - if (detail.agentSessionId() != null) { - reason += " agentSessionId=" + detail.agentSessionId(); - } - // fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the - // worker's pane is being stopped right now, so an open fleet_ask has no turn left to - // resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call - // (see MessageService.abandon's javadoc for why those two must differ). - messages.abandon(detail.terminalId(), reason, true); - replyInbox.release(detail.terminalId()); - primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding - }); + // fleetd #589 Group 3 (:611-631): the whole cleanup lambda extracted to releaseCleanup(...) + // — see FleetdReleaseCleanupWiringTest, and MessageService.abandon's javadoc for the + // documented incident (a torn-down worker's rendezvous waiter left open) this lambda exists + // to prevent. + sessions.onRelease(releaseCleanup(messages, replyInbox, primaryRegistry)); // MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp. // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments. @@ -1066,6 +1060,104 @@ public final class Fleetd { () -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health()); } + /** + * fleetd #589 Group 3 (site {@code :505}): package-private factory for the {@link Injector}'s + * {@link TurnRegistrar}, extracted out of {@code main} for the same reason {@link + * #loopHealthSource} was — before this ticket {@code completion::register} was an inline + * argument to {@code new Injector(...)}, so nothing could pin it directly. Replacing it with + * {@link TurnRegistrar#NOOP} compiles clean and leaves every existing test green: {@code + * onDelivered}'s own {@code captureBaseline} does the identical {@code inFlight} check-and-put a + * moment later on the ordinary path, so the two are indistinguishable once a turn's delivery + * finishes normally. The gap {@link TurnRegistrar}'s own javadoc (fleetd #556) exists to close is + * a {@code turnListener} callback throwing between the two — {@link + * FleetdTurnRegistrarWiringTest} pins that {@code register} itself (not just {@code + * captureBaseline}) makes a delivered turn's waiter resolvable. + */ + static TurnRegistrar turnRegistrar(CompletionResolver completion) { + return completion::register; + } + + /** + * fleetd #589 Group 3 (site {@code :591}): package-private factory for {@link + * FleetHealthMonitor}'s {@code failTarget} callback, extracted out of {@code main} for the same + * reason {@link #loopHealthSource} was. Before this ticket {@code messages::abandon} was an + * inline argument to {@code new FleetHealthMonitor(...)}; replacing it with a no-op {@code + * BiConsumer} compiles clean and leaves every existing test green, and in production it means a + * member found {@code GONE}/{@code NEVER_READY} never fails the ticket waiting on it — the + * caller reports {@code PENDING} for the full 30-minute async timeout instead of the immediate, + * accurate failure CB-580 exists to give it. {@link FleetdHealthFailTargetWiringTest} pins that + * the returned callback actually reaches the real {@link MessageService#abandon}. + */ + static BiConsumer healthFailTarget(MessageService messages) { + return messages::abandon; + } + + /** + * fleetd #589 Group 3 (site {@code :611-631}): package-private factory for the whole {@link + * SessionManager#onRelease} cleanup callback, extracted out of {@code main} for the same reason + * {@link #loopHealthSource} was. Before this ticket this was an inline lambda built directly + * inside {@code main}; replacing its body with a no-op {@code detail -> { }} compiles clean and + * leaves every existing test green, and in production it is the exact incident {@link + * MessageService#abandon}'s own javadoc documents: a torn-down worker's rendezvous waiter is + * left open, so a blocking {@code fleet_send} keeps blocking and an async one reports {@code + * PENDING} for a hardcoded thirty minutes on every {@code fleet_stop} and every idle-reap. + * + *

{@link FleetdReleaseCleanupWiringTest} pins all three collaborator calls this lambda makes + * — {@code messages.abandon}, {@code replyInbox.release}, and {@code + * primaryRegistry.forgetDelegation} — each already tested on its own ({@code MessageServiceTest}, + * {@code PrimaryRegistryTest}), but never before proven to actually be reached from here. + */ + static Consumer releaseCleanup(MessageService messages, ReplyInbox replyInbox, + PrimaryRegistry primaryRegistry) { + return detail -> { + // CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead + // where to re-dispatch onto the same tree, not just that the worker vanished. + String reason = "the worker session was released before it replied"; + if (detail.worktreePath() != null) { + reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch() + + " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none"); + } + // CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the + // member's conversation instead of only re-dispatching a fresh one onto the same files. + if (detail.agentSessionId() != null) { + reason += " agentSessionId=" + detail.agentSessionId(); + } + // fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the + // worker's pane is being stopped right now, so an open fleet_ask has no turn left to + // resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call + // (see MessageService.abandon's javadoc for why those two must differ). + messages.abandon(detail.terminalId(), reason, true); + replyInbox.release(detail.terminalId()); + primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding + }; + } + + /** + * fleetd #589 Group 3 (site {@code :512}): package-private factory for {@code selectReplyInbox}'s + * production {@link AmqpOpener}, extracted out of {@code main} for the same reason {@link + * #loopHealthSource} was. Before this ticket {@code AmqpReplyInbox::open} was an inline argument + * to the {@code selectReplyInbox(...)} call; replacing it with {@code (uri, prefetch) -> new + * InMemoryReplyInbox()} compiles clean and leaves every existing test green — {@link + * FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own injected opener and + * never sees what {@code main} actually passes. {@link FleetdReplyInboxOpenerWiringTest} pins + * that this returns the real opener by pointing it at a guaranteed-closed local port and + * asserting the real network attempt throws — the inert stub never attempts a connection at all. + */ + static AmqpOpener replyInboxOpener() { + return AmqpReplyInbox::open; + } + + /** + * fleetd #589 Group 3 (site {@code :518}): package-private factory for {@code + * openLeadMailbox}'s production {@link LeadMailboxOpener}, extracted out of {@code main} for the + * same reason {@link #replyInboxOpener} was — same gap, same fix, the lead-coordination mailbox + * instead of the reply inbox. {@link FleetdLeadMailboxOpenerWiringTest} pins that this returns + * the real opener the same way. + */ + static LeadMailboxOpener leadMailboxOpener() { + return LeadMailbox::open; + } + /** * fleetd #248: package-private factory for the member worktree/branch lookup {@link * CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241). diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdHealthFailTargetWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdHealthFailTargetWiringTest.java new file mode 100644 index 0000000..dc835ad --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdHealthFailTargetWiringTest.java @@ -0,0 +1,80 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.msg.MessageService; +import dev.ltms.fleet.msg.Rendezvous; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.util.function.BiConsumer; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #589 Group 3, site {@code :591}. {@code Fleetd.main} wires {@link + * dev.ltms.fleet.health.FleetHealthMonitor}'s {@code failTarget} callback with {@code + * messages::abandon} — before this ticket that was an inline argument to {@code new + * FleetHealthMonitor(...)}. Measured: replacing it with a no-op {@code BiConsumer} at the call + * site compiles with 0 errors and leaves the full suite green, because nothing else in the tree + * ever drives that specific constructor argument. In production it means a member the monitor + * classifies {@code GONE}/{@code NEVER_READY} never has its pending ticket failed — the caller + * keeps reporting {@code PENDING} for the full 30-minute async timeout instead of the immediate, + * accurate failure CB-580 exists to give it. + * + *

This test calls {@link Fleetd#healthFailTarget} directly — never {@code FleetHealthMonitor} + * or {@code main} — against a real {@link MessageService}, using the same {@code sendAsync} + + * {@code poll} observable {@link MessageServiceTest} already relies on to pin {@code + * MessageService.abandon} itself. + */ +class FleetdHealthFailTargetWiringTest { + + private static final String T = "term_a"; + + @Test + @DisplayName("Fleetd.healthFailTarget delegates to the real MessageService.abandon, not a no-op") + void healthFailTargetDelegatesToMessagesAbandon() throws Exception { + FakeHerdr herdr = new FakeHerdr().readText("$ prompt"); + AgentControl agents = new AgentControl(herdr); + Rendezvous rendezvous = new Rendezvous(); + Injector injector = new Injector(agents); + MessageService messages = new MessageService(agents, injector, rendezvous); + + BiConsumer failTarget = Fleetd.healthFailTarget(messages); + + String ticket = messages.sendAsync(T, "long task"); + awaitWaiting(rendezvous); + assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase()); + + failTarget.accept(T, "member unreachable (health monitor)"); + + MessageService.TaskView view = null; + long deadline = System.currentTimeMillis() + 3000; + while (System.currentTimeMillis() < deadline) { + view = messages.poll(ticket); + if (view.phase() != MessageService.Phase.PENDING) { + break; + } + //noinspection BusyWait + Thread.sleep(10); + } + assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(), + "Fleetd.healthFailTarget(messages) must return messages::abandon — replacing it " + + "with a no-op BiConsumer at the Fleetd.healthFailTarget call site means " + + "this ticket is never failed and keeps polling as PENDING"); + assertTrue(view.detail() != null && view.detail().contains("member unreachable"), + "the failure reason passed to failTarget.accept must reach MessageService.abandon " + + "and end up in the ticket's detail"); + } + + private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException { + long deadline = System.currentTimeMillis() + 2000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + //noinspection BusyWait + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxOpenerWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxOpenerWiringTest.java new file mode 100644 index 0000000..5aa1aff --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxOpenerWiringTest.java @@ -0,0 +1,48 @@ +package dev.ltms.fleet; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.net.ServerSocket; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #589 Group 3, site {@code :518}. {@code Fleetd.main} passes {@code LeadMailbox::open} as + * the {@link Fleetd.LeadMailboxOpener} argument to {@code openLeadMailbox(...)} — before this + * ticket that method reference was inline at the call site. Measured: replacing it with the inert + * {@code (uri, selfCoordId, prefetch) -> null} compiles with 0 errors and leaves the full suite + * green — {@code FleetdLeadMailboxSelectionTest} drives {@code openLeadMailbox} with its own + * injected opener and never observes what {@code main} itself actually passes. + * + *

This test calls {@link Fleetd#leadMailboxOpener} directly and proves it is the real, + * network-attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must + * throw, exactly mirroring {@link FleetdReplyInboxOpenerWiringTest} for the reply-inbox opener. + * The inert form never attempts a connection and returns {@code null} without throwing, so it + * fails this assertion silently. + */ +class FleetdLeadMailboxOpenerWiringTest { + + @Test + @DisplayName("Fleetd.leadMailboxOpener is the real LeadMailbox::open, not a stub that never connects") + void leadMailboxOpenerAttemptsARealConnection() throws Exception { + int closedPort; + try (ServerSocket socket = new ServerSocket(0)) { + closedPort = socket.getLocalPort(); + } // released immediately — connecting to it now is a guaranteed refusal, not a fluke + + Fleetd.LeadMailboxOpener opener = Fleetd.leadMailboxOpener(); + + IllegalStateException thrown = assertThrows(IllegalStateException.class, + () -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/coord", "coord-1", 50), + "Fleetd.leadMailboxOpener() must be LeadMailbox::open — a real network attempt " + + "against a genuinely unreachable broker must throw. The inert form " + + "(uri, selfCoordId, prefetch) -> null never attempts a connection and " + + "returns null instead of throwing, so it would fail this assertion " + + "silently."); + assertTrue(thrown.getMessage().contains("cannot connect to AMQP coordination broker"), + "must be LeadMailbox.open's own real failure message, not a different exception " + + "shape standing in for it"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdReleaseCleanupWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdReleaseCleanupWiringTest.java new file mode 100644 index 0000000..55236bf --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdReleaseCleanupWiringTest.java @@ -0,0 +1,107 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.mcp.PrimaryRegistry; +import dev.ltms.fleet.msg.InMemoryReplyInbox; +import dev.ltms.fleet.msg.MessageService; +import dev.ltms.fleet.msg.Rendezvous; +import dev.ltms.fleet.session.SessionManager; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.util.function.Consumer; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #589 Group 3, site {@code :611-631}. {@code Fleetd.main} wires {@code + * sessions.onRelease(...)} with a lambda that calls three collaborators — {@code + * messages.abandon}, {@code replyInbox.release}, and {@code primaryRegistry.forgetDelegation} — + * before this ticket built inline inside {@code main}. Measured: replacing the whole lambda body + * with {@code detail -> { }} compiles with 0 errors and leaves the full suite green, because each + * collaborator is separately tested in isolation ({@code MessageServiceTest}, {@code + * PrimaryRegistryTest}) but nothing before this ticket drove the lambda that calls all three from + * {@code main}. + * + *

{@code MessageService.abandon}'s own javadoc documents the consequence: without this, + * tearing a worker down leaves its rendezvous waiter open, so a blocking {@code fleet_send} keeps + * blocking and an async one reports {@code PENDING} for a hardcoded thirty minutes on every {@code + * fleet_stop} and every idle-reap. + * + *

This test calls {@link Fleetd#releaseCleanup} directly — never {@code SessionManager} or + * {@code main} — against real {@link MessageService}, {@link InMemoryReplyInbox}, and {@link + * PrimaryRegistry} instances, and asserts each collaborator's own observable effect: the pending + * async ticket transitions to {@code FAILED} (abandon), the inbox no longer owns the target's + * queue (release), and the recorded delegation is forgotten (forgetDelegation). + */ +class FleetdReleaseCleanupWiringTest { + + private static final String T = "term_a"; + + @Test + @DisplayName("Fleetd.releaseCleanup reaches messages.abandon, replyInbox.release, and primaryRegistry.forgetDelegation") + void releaseCleanupReachesAllThreeCollaborators() throws Exception { + FakeHerdr herdr = new FakeHerdr().readText("$ prompt"); + AgentControl agents = new AgentControl(herdr); + Rendezvous rendezvous = new Rendezvous(); + Injector injector = new Injector(agents); + MessageService messages = new MessageService(agents, injector, rendezvous); + InMemoryReplyInbox replyInbox = new InMemoryReplyInbox(); + PrimaryRegistry primaryRegistry = new PrimaryRegistry(null); + + // Set up the "before" state each collaborator's own effect is measured against. + String ticket = messages.sendAsync(T, "long task"); + awaitWaiting(rendezvous); + assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(), + "sanity: the async ticket is pending before cleanup runs"); + + replyInbox.own(T); + replyInbox.publish(T, "msg-1", "hello"); + assertEquals(1, replyInbox.peek(T).size(), + "sanity: the inbox owns T and holds one message before cleanup runs"); + + primaryRegistry.recordDelegation(T, "lead-1"); + assertEquals("lead-1", primaryRegistry.nudgeTargetFor(T).orElse(null), + "sanity: the delegation is recorded before cleanup runs"); + + Consumer cleanup = + Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry); + cleanup.accept(new SessionManager.ReleaseDetail(T, null, null, null, null)); + + MessageService.TaskView view = null; + long deadline = System.currentTimeMillis() + 3000; + while (System.currentTimeMillis() < deadline) { + view = messages.poll(ticket); + if (view.phase() != MessageService.Phase.PENDING) { + break; + } + //noinspection BusyWait + Thread.sleep(10); + } + assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(), + "releaseCleanup must call messages.abandon(...) — an inert detail -> { } lambda " + + "leaves this ticket PENDING forever"); + assertTrue(view.detail() != null && view.detail().contains("released"), + "the abandon reason must say the worker session was released"); + + assertTrue(replyInbox.peek(T).isEmpty(), + "releaseCleanup must call replyInbox.release(...) — an inert lambda leaves the " + + "inbox still owning T with its message"); + + assertTrue(primaryRegistry.nudgeTargetFor(T).isEmpty(), + "releaseCleanup must call primaryRegistry.forgetDelegation(...) — an inert lambda " + + "leaves the stale delegation in place"); + } + + private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException { + long deadline = System.currentTimeMillis() + 2000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + //noinspection BusyWait + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdReplyInboxOpenerWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdReplyInboxOpenerWiringTest.java new file mode 100644 index 0000000..29cad31 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdReplyInboxOpenerWiringTest.java @@ -0,0 +1,49 @@ +package dev.ltms.fleet; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.net.ServerSocket; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #589 Group 3, site {@code :512}. {@code Fleetd.main} passes {@code AmqpReplyInbox::open} + * as the {@link Fleetd.AmqpOpener} argument to {@code selectReplyInbox(...)} — before this ticket + * that method reference was inline at the call site. Measured: replacing it with the inert {@code + * (uri, prefetch) -> new InMemoryReplyInbox()} compiles with 0 errors and leaves the full suite + * green — {@code FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own + * injected opener (including one test that passes the real {@code AmqpReplyInbox::open} + * explicitly) and never observes what {@code main} itself actually passes. + * + *

This test calls {@link Fleetd#replyInboxOpener} directly and proves it is the real, network- + * attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must throw — + * the same shape {@code FleetdReplyInboxSelectionTest.aRealUnreachableBrokerFallsBackViaTheRealOpener} + * already relies on for {@code AmqpReplyInbox.open} itself. The inert form never attempts a + * connection and never throws, so it fails this assertion silently (by returning normally). + */ +class FleetdReplyInboxOpenerWiringTest { + + @Test + @DisplayName("Fleetd.replyInboxOpener is the real AmqpReplyInbox::open, not a stub that never connects") + void replyInboxOpenerAttemptsARealConnection() throws Exception { + int closedPort; + try (ServerSocket socket = new ServerSocket(0)) { + closedPort = socket.getLocalPort(); + } // released immediately — connecting to it now is a guaranteed refusal, not a fluke + + Fleetd.AmqpOpener opener = Fleetd.replyInboxOpener(); + + IllegalStateException thrown = assertThrows(IllegalStateException.class, + () -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/vh", 50), + "Fleetd.replyInboxOpener() must be AmqpReplyInbox::open — a real network attempt " + + "against a genuinely unreachable broker must throw. The inert form " + + "(uri, prefetch) -> new InMemoryReplyInbox() never attempts a connection " + + "and never throws, so it would return normally here and fail this " + + "assertion silently."); + assertTrue(thrown.getMessage().contains("cannot connect to AMQP broker"), + "must be AmqpReplyInbox.open's own real failure message, not a different exception " + + "shape standing in for it"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdTurnRegistrarWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdTurnRegistrarWiringTest.java new file mode 100644 index 0000000..0d3c8a8 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdTurnRegistrarWiringTest.java @@ -0,0 +1,72 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.inject.CompletionResolver; +import dev.ltms.fleet.inject.ExhaustedPatternLookup; +import dev.ltms.fleet.inject.ExhaustionSink; +import dev.ltms.fleet.inject.TurnRegistrar; +import dev.ltms.fleet.msg.Rendezvous; +import dev.ltms.fleet.msg.TurnToken; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #589 Group 3, site {@code :505}. {@code Fleetd.main} wires the {@link + * dev.ltms.fleet.inject.Injector}'s {@link TurnRegistrar} with {@code completion::register} — + * before this ticket that was an inline argument to {@code new Injector(...)}, so nothing could + * pin it directly. Measured: replacing it with {@link TurnRegistrar#NOOP} at the call site + * compiles with 0 errors and leaves the full suite green, because {@code onDelivered}'s own {@code + * captureBaseline} performs the identical {@code inFlight} check-and-put a moment later on the + * ordinary delivery path — the two are indistinguishable unless something reads the resolver + * between {@code register} and {@code onDelivered}, or {@code onDelivered} never runs at all (the + * gap fleetd #556 introduced {@link TurnRegistrar} to close). + * + *

This test calls {@link Fleetd#turnRegistrar} directly — never {@code Injector} or {@code + * main} — and drives {@link CompletionResolver} entirely through its public API: {@link + * TurnRegistrar#register} followed by {@link CompletionResolver#resolveBeforePostAction}, which + * looks up the same {@code inFlight} entry {@code onTurnComplete} would. With the real registrar, + * that entry exists and the waiter opened by {@link Rendezvous#open} resolves; with {@link + * TurnRegistrar#NOOP} nothing was ever registered, {@code resolveBeforePostAction} finds no + * in-flight turn, and the waiter is left exactly as it started — never done. + */ +class FleetdTurnRegistrarWiringTest { + + private static final String T = "term_a"; + + @Test + @DisplayName("Fleetd.turnRegistrar delegates to the real CompletionResolver, not a no-op") + void turnRegistrarDelegatesToCompletionRegister() { + FakeHerdr herdr = new FakeHerdr().readText("⏺ BUILD GREEN: 391 files\n❯ "); + AgentControl agents = new AgentControl(herdr); + Rendezvous rendezvous = new Rendezvous(); + // fleetd#164: an ever-advancing fake clock stands in for the real time a turn would take + // between delivery and resolution, so the MIN_TURN_NANOS "too fast" floor never trips here — + // see MessageServiceTest's resolverClock for the same technique. + AtomicLong clock = new AtomicLong(); + CompletionResolver completion = new CompletionResolver(agents, rendezvous, + ExhaustedPatternLookup.none(), ExhaustionSink.none(), + () -> clock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1)); + + TurnRegistrar registrar = Fleetd.turnRegistrar(completion); + + CompletableFuture waiter = rendezvous.open(T); + registrar.register(T, new TurnToken(T, waiter)); + + // Mirrors what Injector.onStatus's confirmed working->idle boundary would trigger via + // CompletionResolver.onTurnComplete — resolveBeforePostAction is the public, synchronous + // twin of that path and reads the exact same inFlight entry register() must have written. + completion.resolveBeforePostAction(T); + + assertTrue(waiter.isDone(), + "Fleetd.turnRegistrar(completion) must return completion::register — replacing it " + + "with TurnRegistrar.NOOP at the Fleetd.turnRegistrar call site means this " + + "turn is never registered with CompletionResolver, so resolveBeforePostAction " + + "finds no in-flight turn and this waiter is never resolved"); + } +}