Merge #598: fleetd #589 group 3 — wiring-test 5 sites in Fleetd.main()
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 50s
CI / build (push) Successful in 2m9s

Extracts 5 inline lambdas/method-refs in Fleetd.main() to package-private
factories and pins each with a wiring test. Production behaviour unchanged;
the releaseCleanup body moved verbatim.

Gate: test-merged onto current main in a scratch worktree, mvn exit=0,
1795 tests / 0 failures / 0 errors from 137 surefire reports. Diff read in
full. Worker's self-disclosed bare `git stash push` verified as recovered —
all 3 surviving stash entries predate today, so no other worktree lost work.
This commit was merged in pull request #598.
This commit is contained in:
2026-09-19 10:33:07 +02:00
6 changed files with 473 additions and 25 deletions
+117 -25
View File
@@ -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<String, String> 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.
*
* <p>{@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<SessionManager.ReleaseDetail> 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).
@@ -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.
*
* <p>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<String, String> 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");
}
}
@@ -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.
*
* <p>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");
}
}
@@ -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}.
*
* <p>{@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.
*
* <p>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<SessionManager.ReleaseDetail> 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");
}
}
@@ -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.
*
* <p>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");
}
}
@@ -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).
*
* <p>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<Rendezvous.Resolution> 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");
}
}