From 6cccd458d4bc3f84324c8d25c857d63ebd97c81f Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 19 Sep 2026 15:24:11 +0700 Subject: [PATCH] #589 Group 3: wiring-test the 5 sites below line 500 in Fleetd.main() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Extracts the inline lambdas/method references at the 5 assigned wiring sites into named package-private factories on Fleetd, following the FleetdLoopHealthSourceWiringTest pattern from #584: - turnRegistrar(CompletionResolver) — was completion::register (Injector) - healthFailTarget(MessageService) — was messages::abandon (FleetHealthMonitor) - releaseCleanup(MessageService, ReplyInbox, PrimaryRegistry) — was the inline sessions.onRelease(detail -> {...}) cleanup lambda - replyInboxOpener() — was AmqpReplyInbox::open passed to selectReplyInbox - leadMailboxOpener() — was LeadMailbox::open passed to openLeadMailbox Each factory has a new runtime test (not source-text) that drives real collaborators through public APIs: MessageService.poll(ticket).phase(), InMemoryReplyInbox.peek(), PrimaryRegistry.nudgeTargetFor(), and the opener tests connect to a guaranteed-closed local port to prove a real network attempt vs. an inert stub. releaseCleanup was done first per the brief: MessageService.abandon's javadoc documents that losing this cleanup leaves a torn-down worker's rendezvous waiter open forever. Tests: 1789 -> 1794 (+5), 0 failures, 0 errors. mvn -q -o test exit 0, no BUILD FAILURE, no piped exit status. Each new test verified RED on the inert form named in the ticket, and GREEN after reformatting the call across lines and extracting the argument into a local/factory. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 142 +++++++++++++++--- .../FleetdHealthFailTargetWiringTest.java | 80 ++++++++++ .../FleetdLeadMailboxOpenerWiringTest.java | 48 ++++++ .../fleet/FleetdReleaseCleanupWiringTest.java | 107 +++++++++++++ .../FleetdReplyInboxOpenerWiringTest.java | 49 ++++++ .../fleet/FleetdTurnRegistrarWiringTest.java | 72 +++++++++ 6 files changed, 473 insertions(+), 25 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdHealthFailTargetWiringTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxOpenerWiringTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdReleaseCleanupWiringTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdReplyInboxOpenerWiringTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdTurnRegistrarWiringTest.java 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"); + } +}