Merge pull request 'fleetd #612 A-gaps: coordinator path + reportRoleFallbackGaps inside the assembly boundary' (#624) from worker/612-agaps-73a926-2 into worker/fleetd-612-unita-87807e-1
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m34s
CI / build (pull_request) Failing after 1m42s

This commit was merged in pull request #624.
This commit is contained in:
2026-09-22 06:41:42 +02:00
8 changed files with 458 additions and 41 deletions
+34 -31
View File
@@ -42,6 +42,7 @@ import dev.ltms.fleet.mcp.LsofProcessCwdLookup;
import dev.ltms.fleet.msg.AmqpReplyInbox; import dev.ltms.fleet.msg.AmqpReplyInbox;
import dev.ltms.fleet.msg.InMemoryReplyInbox; import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.LeadChannel; import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop; import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadMailbox; import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.MessageService;
@@ -177,27 +178,17 @@ public final class Fleetd {
// and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for // and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for
// why a name-by-name list here would have the same defect it replaces. // why a name-by-name list here would have the same defect it replaces.
cfg.validateAll(); cfg.validateAll();
// fleetd #613: validateAll() (validateMembers() inside it) only refuses a slot that names a
// bad role or profile — it says nothing about a role that has NO pool or NO charter at all,
// because both are legitimate ("unconstrained") states, not errors. Report them here, right
// after validation passes, so an operator sees the gap once per restart instead of finding
// it later in a roster row (see reportRoleFallbackGaps' javadoc for the measured cause).
reportRoleFallbackGaps(cfg);
// fleetd #469, follow-up to #464: validateAll() (and validateCharters() inside it) only
// checks that a charter's KEY is a role wire name and its text is non-blank — it never
// looks at what the text actually names. This is the separate check that does: it asks
// dev.ltms.fleet.mcp.FleetTool (the canonical registered-tool set) whether every fleet_*/
// bridge_* token a charter names is a tool this server actually registers. It cannot live
// inside FleetConfig#validateCharters() — config loads before the MCP server exists, and
// must not gain a dependency on the mcp package — so it runs here instead, at the one seam
// that already holds both a loaded FleetConfig and the mcp package, before anything below
// opens a socket or spawns a member. fleetd #474: the same check is also wired into `config`
// above as ConfigRef's extraValidation, so a reload refuses what this line refuses at startup.
assertChartersNameOnlyRegisteredTools(cfg);
// fleetd #612 Unit A: everything from here on used to run inline in this method. It now // fleetd #612 A-gaps (gap 2): everything from here on — including the two post-validation
// lives in FleetdAssembly.assembleAndStart, built against a real ResourcePorts — see that // reports that used to run inline right here (reportRoleFallbackGaps,
// class's javadoc for the full boot-order contract this preserves exactly. // assertChartersNameOnlyRegisteredTools) — now lives in FleetdAssembly.assembleAndStart,
// built against a real ResourcePorts. Unit A originally moved only the socket/broker/HTTP
// composition and left those two calls here, between validateAll() and the assembly call —
// outside the boundary FleetdAssemblyLifecycleTest drives, so deleting either call
// compiled clean and left the whole suite green. Moving the boundary to start immediately
// after validateAll() (this line) puts both back under test, in the same relative order,
// before either one does any I/O — see FleetdAssembly's javadoc for the full boot-order
// contract this preserves exactly.
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ResourcePorts.system()); FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ResourcePorts.system());
} }
@@ -1224,10 +1215,18 @@ public final class Fleetd {
ReplyInbox open(String uri, int prefetch); ReplyInbox open(String uri, int prefetch);
} }
/** Injection seam for {@link #openLeadMailbox}: production binds {@link LeadMailbox#open}. */ /**
* Injection seam for {@link #openLeadMailbox}: production binds {@link LeadMailbox#open}.
*
* <p>fleetd #612 A-gaps (gap 1): returns {@link LeadChannelHandle}, not the concrete {@link
* LeadMailbox}, so a test can supply a fake closeable channel instead of a real broker
* connection — see {@link LeadChannelHandle}'s own javadoc for why the narrower type existed
* and what widening it to add {@code close()} costs (nothing: {@code LeadMailbox} already
* implements it).
*/
@FunctionalInterface @FunctionalInterface
interface LeadMailboxOpener { interface LeadMailboxOpener {
LeadMailbox open(String uri, String selfCoordId, int prefetch); LeadChannelHandle open(String uri, String selfCoordId, int prefetch);
} }
/** /**
@@ -1250,7 +1249,7 @@ public final class Fleetd {
* fleet still works exactly as it did before this feature existed.</li> * fleet still works exactly as it did before this feature existed.</li>
* </ul> * </ul>
*/ */
static LeadMailbox openLeadMailbox(FleetConfig.Coordinator coordinator, Map<String, String> env, static LeadChannelHandle openLeadMailbox(FleetConfig.Coordinator coordinator, Map<String, String> env,
LeadMailboxOpener opener) { LeadMailboxOpener opener) {
if (coordinator == null) { if (coordinator == null) {
return null; // opt-in: nothing configured, nothing to say return null; // opt-in: nothing configured, nothing to say
@@ -1270,7 +1269,7 @@ public final class Fleetd {
return null; return null;
} }
try { try {
LeadMailbox mailbox = opener.open(uri, coordinator.selfId(), coordinator.prefetchOrDefault()); LeadChannelHandle mailbox = opener.open(uri, coordinator.selfId(), coordinator.prefetchOrDefault());
log.info("lead coordination: ON as coord-id {} (prefetch={})", log.info("lead coordination: ON as coord-id {} (prefetch={})",
coordinator.selfId(), coordinator.prefetchOrDefault()); coordinator.selfId(), coordinator.prefetchOrDefault());
return mailbox; return mailbox;
@@ -1650,16 +1649,20 @@ public final class Fleetd {
} }
/** /**
* fleetd #474: the one place both the startup call (right after {@code cfg.validateAll()} in * fleetd #474: the one place both the startup call (fleetd #612 A-gaps gap 2: right after
* {@link #main}) and the reload call (wired into {@code config}'s {@code extraValidation} above, * {@code cfg.validateAll()} in {@link FleetdAssembly#assembleAndStart}, immediately after
* via a method reference to this method) go through, so the two can never drift into checking * {@code main} hands off to it — see that method's javadoc) and the reload call (wired into
* different things. Extracted only to give {@link ConfigRef}'s {@code Consumer<FleetConfig>} * {@code config}'s {@code extraValidation} above, via a method reference to this method) go
* hook a {@code FleetConfig -> void} shape to bind to — {@link CharterToolSurface} itself still * through, so the two can never drift into checking different things. Extracted only to give
* takes the raw charter map and knows nothing about {@code ConfigRef} or {@code Fleetd}. * {@link ConfigRef}'s {@code Consumer<FleetConfig>} hook a {@code FleetConfig -> void} shape to
* bind to — {@link CharterToolSurface} itself still takes the raw charter map and knows nothing
* about {@code ConfigRef} or {@code Fleetd}.
* *
* <p>Package-private so a test can call it directly the same way the other startup-report * <p>Package-private so a test can call it directly the same way the other startup-report
* helpers above are tested, without needing to drive {@link #main} for a unit-level check; * helpers above are tested, without needing to drive {@link #main} for a unit-level check;
* {@code FleetdStartupValidationTest} proves the startup call site, and {@code * {@code FleetdStartupValidationTest} proves the startup call site by driving {@code main}
* itself end to end (the throw still happens before {@code main} reaches any real socket or
* broker work, since the assembly runs this before either), and {@code
* FleetdConfigRefCharterToolSurfaceWiringTest} — by constructing {@code ConfigRef} with this * FleetdConfigRefCharterToolSurfaceWiringTest} — by constructing {@code ConfigRef} with this
* exact method reference, the same way {@code main} does above — proves the reload call site. * exact method reference, the same way {@code main} does above — proves the reload call site.
* {@code dev.ltms.fleet.config.ConfigRefTest} pins the same reload behaviour too, through an * {@code dev.ltms.fleet.config.ConfigRefTest} pins the same reload behaviour too, through an
@@ -36,9 +36,9 @@ import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.member.MemberCredentialPolicyView; import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics; import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop; import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadHeartbeatLoop; import dev.ltms.fleet.msg.LeadHeartbeatLoop;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.ReplyInbox; import dev.ltms.fleet.msg.ReplyInbox;
@@ -77,8 +77,14 @@ import java.util.stream.Collectors;
* fleetd #612 Unit A: the real boot assembly, extracted out of {@code Fleetd.main} so a test can * fleetd #612 Unit A: the real boot assembly, extracted out of {@code Fleetd.main} so a test can
* drive it directly. {@link #assembleAndStart} is <em>the same statements {@code main} used to run * drive it directly. {@link #assembleAndStart} is <em>the same statements {@code main} used to run
* inline</em>, in the same order, against a real {@link ResourcePorts} in production and a fake one * inline</em>, in the same order, against a real {@link ResourcePorts} in production and a fake one
* in a test — see {@code FleetdAssemblyLifecycleTest}. {@code Fleetd.main} keeps config loading, * in a test — see {@code FleetdAssemblyLifecycleTest}. {@code Fleetd.main} keeps config loading and
* the startup reports and validation; everything from the herdr socket onward moved here. * {@code cfg.validateAll()}; everything from immediately after that call onward moved here —
* including, since fleetd #612 A-gaps (gap 2), the two post-validation reports ({@code
* reportRoleFallbackGaps}, {@code assertChartersNameOnlyRegisteredTools}) that Unit A originally
* left behind in {@code main}. Those two calls do no I/O themselves, but leaving them outside this
* method meant deleting either one compiled clean and left the whole suite green — nothing drove
* {@code main} itself, so nothing could notice. They run first here, in the same relative order,
* before the herdr socket or anything else that touches the outside world.
* *
* <p><strong>Construction and start order is preserved exactly, on purpose.</strong> This is not * <p><strong>Construction and start order is preserved exactly, on purpose.</strong> This is not
* rebuilt into "construct everything, then start everything" — that would change boot timing. The * rebuilt into "construct everything, then start everything" — that would change boot timing. The
@@ -119,6 +125,13 @@ final class FleetdAssembly {
ConfigRef config = inputs.config(); ConfigRef config = inputs.config();
SubscriptionGuard guard = inputs.guard(); SubscriptionGuard guard = inputs.guard();
// fleetd #612 A-gaps (gap 2): moved in from Fleetd.main, immediately after cfg.validateAll()
// there — the exact point main used to call these two, and still the first thing this
// method does, before any socket or broker work below. See this class's javadoc and each
// method's own for why they run here rather than in FleetConfig#validateAll() itself.
Fleetd.reportRoleFallbackGaps(cfg);
Fleetd.assertChartersNameOnlyRegisteredTools(cfg);
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank() Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
? Path.of(cfg.herdrSocket()) ? Path.of(cfg.herdrSocket())
: UnixSocketHerdrClient.defaultSocketPath(); : UnixSocketHerdrClient.defaultSocketPath();
@@ -346,7 +359,7 @@ final class FleetdAssembly {
// broker from the reply inbox by design. Absent a coordinator: block this is null and every // broker from the reply inbox by design. Absent a coordinator: block this is null and every
// lead path below is simply not wired, exactly the behaviour before this ticket. It owns a // lead path below is simply not wired, exactly the behaviour before this ticket. It owns a
// broker connection, so keep the reference for the ordered shutdown hook. // broker connection, so keep the reference for the ordered shutdown hook.
final LeadMailbox leadMailbox = Fleetd.openLeadMailbox(cfg.coordinator(), ports.environment(), final LeadChannelHandle leadMailbox = Fleetd.openLeadMailbox(cfg.coordinator(), ports.environment(),
ports.leadMailboxOpener()); ports.leadMailboxOpener());
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config). // CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null; String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
@@ -8,9 +8,9 @@ import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.Injector; import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller; import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.mcp.FleetMcp; import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop; import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadHeartbeatLoop; import dev.ltms.fleet.msg.LeadHeartbeatLoop;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.ReplyInbox; import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop; import dev.ltms.fleet.msg.ReplyPushLoop;
@@ -55,7 +55,7 @@ final class FleetdRuntime implements AutoCloseable {
private final SessionReaper reaper; // nullable — lifecycle.idleTtlSeconds opt-in private final SessionReaper reaper; // nullable — lifecycle.idleTtlSeconds opt-in
private final IdleSleepGuard idleSleepGuard; // nullable — idleSleepGuard.enabled: false private final IdleSleepGuard idleSleepGuard; // nullable — idleSleepGuard.enabled: false
private final ReplyInbox replyInbox; private final ReplyInbox replyInbox;
private final LeadMailbox leadMailbox; // nullable — coordinator: opt-in private final LeadChannelHandle leadMailbox; // nullable — coordinator: opt-in
private final CompletionResolver completion; private final CompletionResolver completion;
private final Injector injector; private final Injector injector;
/** /**
@@ -72,7 +72,7 @@ final class FleetdRuntime implements AutoCloseable {
LeadCoordLoop leadCoordLoop, ScheduledExecutorService leadCoordScheduler, LeadCoordLoop leadCoordLoop, ScheduledExecutorService leadCoordScheduler,
FleetHealthMonitor healthMonitor, ConfigWatcher configWatcher, FleetMcp mcp, FleetHealthMonitor healthMonitor, ConfigWatcher configWatcher, FleetMcp mcp,
SessionReaper reaper, IdleSleepGuard idleSleepGuard, ReplyInbox replyInbox, SessionReaper reaper, IdleSleepGuard idleSleepGuard, ReplyInbox replyInbox,
LeadMailbox leadMailbox, CompletionResolver completion, Injector injector) { LeadChannelHandle leadMailbox, CompletionResolver completion, Injector injector) {
this.cfg = cfg; this.cfg = cfg;
this.sessions = sessions; this.sessions = sessions;
this.router = router; this.router = router;
@@ -113,7 +113,7 @@ final class FleetdRuntime implements AutoCloseable {
SessionReaper reaper() { return reaper; } SessionReaper reaper() { return reaper; }
IdleSleepGuard idleSleepGuard() { return idleSleepGuard; } IdleSleepGuard idleSleepGuard() { return idleSleepGuard; }
ReplyInbox replyInbox() { return replyInbox; } ReplyInbox replyInbox() { return replyInbox; }
LeadMailbox leadMailbox() { return leadMailbox; } LeadChannelHandle leadMailbox() { return leadMailbox; }
CompletionResolver completion() { return completion; } CompletionResolver completion() { return completion; }
Injector injector() { return injector; } Injector injector() { return injector; }
Javalin app() { return app; } Javalin app() { return app; }
@@ -0,0 +1,24 @@
package dev.ltms.fleet.msg;
/**
* fleetd #612 A-gaps (gap 1): a {@link LeadChannel} that its owner can also close.
*
* <p>{@link LeadChannel}'s own javadoc says plainly that {@code close()} is deliberately left out
* of that interface — draining is a caller convenience nobody uses, and closing is the
* <em>owner's</em> job. This interface is that owner's own, wider view: whoever opens the
* coordination mailbox (the assembly that builds the daemon) also needs to close it from the
* shutdown path, and a test standing in for a real broker connection needs a fake it can mark
* closed, without ever holding a live connection. Every ordinary consumer ({@code FleetMcp},
* {@link LeadCoordLoop}) keeps taking the narrower {@link LeadChannel} exactly as before — only
* the owner speaks this wider one.
*
* <p>{@link LeadMailbox} is still the only production implementation. This only generalises the
* TYPE its owner holds it as (previously the concrete class), so a test can substitute a fake
* closeable channel instead of a real AMQP connection.
*/
public interface LeadChannelHandle extends LeadChannel, AutoCloseable {
/** Release the underlying connection. Declared with no checked exception, unlike the plain {@link AutoCloseable#close()}. */
@Override
void close();
}
@@ -63,7 +63,7 @@ import java.util.concurrent.TimeoutException;
* which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out * which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out
* the confirm timeout against a sequence number that means nothing on the new channel. * the confirm timeout against a sequence number that means nothing on the new channel.
*/ */
public final class LeadMailbox implements LeadChannel, AutoCloseable { public final class LeadMailbox implements LeadChannelHandle {
private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class); private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class);
@@ -0,0 +1,211 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 A-gaps (gap 1): {@code FleetdAssemblyLifecycleTest}'s own class javadoc says plainly
* that it leaves {@code coordinator:} unset, so {@code leadMailbox} and {@code leadCoordLoop} stay
* {@code null} throughout — the configured-coordinator path is never exercised by Unit A's own
* test. This class drives that path instead: a real {@code coordinator:} block, a fake {@link
* Fleetd.LeadMailboxOpener} returning a fake closeable channel (never a real broker connection),
* and proof that {@link FleetdAssembly#assembleAndStart} both builds it and, on shutdown, closes it.
*
* <p>Made possible by generalising {@code Fleetd.LeadMailboxOpener}'s return type (and {@code
* FleetdRuntime}'s field) from the concrete {@code LeadMailbox} to {@link LeadChannelHandle} — a
* {@link LeadChannel} its owner can also close. {@code FleetMcp} and {@code LeadCoordLoop} already
* consumed the narrower {@link LeadChannel}; this only widens the one seam that owns and closes it.
*/
class FleetdAssemblyCoordinatorLifecycleTest {
private static final String SELF_COORD_ID = "test-lead";
/** A fake {@link LeadChannelHandle}: never touches a broker, and records whether it was closed. */
private static final class FakeLeadChannel implements LeadChannelHandle {
volatile boolean closed = false;
@Override
public void publish(String toCoordId, LeadMessage m) {
}
@Override
public List<LeadMessage> peek() {
return List.of();
}
@Override
public void ack(String msgId) {
}
@Override
public String selfCoordId() {
return SELF_COORD_ID;
}
@Override
public boolean heldDurable() {
return true;
}
@Override
public MailboxState inspect(String coordId) {
return MailboxState.unknown(coordId);
}
@Override
public void close() {
closed = true;
}
}
/** Minimal fake {@link ResourcePorts}: a real herdr fake, a fake reply inbox, and a real, offered fake lead channel. */
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
final FakeLeadChannel leadChannel = new FakeLeadChannel();
String offeredUri;
String offeredSelfId;
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> replyInbox;
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
offeredUri = uri;
offeredSelfId = selfCoordId;
return leadChannel;
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// No real HTTP bind in a unit test.
}
}
private static final class SentinelReplyInbox implements ReplyInbox {
@Override
public void own(String target) {
}
@Override
public void release(String target) {
}
@Override
public void publish(String target, String msgId, String content) {
}
@Override
public List<InboxMessage> peek(String target) {
return List.of();
}
@Override
public boolean ack(String target, String msgId) {
return false;
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
coordinator:
uri: "amqp://fake-lead-broker/vh"
selfId: "test-lead"
""");
return FleetConfig.load(f);
}
@Test
void configuredCoordinatorIsBuiltByTheAssemblyAndClosedOnShutdown(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
// --- the assembly actually calls the configured opener and builds the coordinator path ---
assertEquals("amqp://fake-lead-broker/vh", ports.offeredUri,
"the assembly must open the mailbox at the configured broker uri");
assertEquals(SELF_COORD_ID, ports.offeredSelfId,
"the assembly must open the mailbox under the configured selfId");
assertSame(ports.leadChannel, runtime.leadMailbox(),
"FleetdRuntime must own the exact LeadChannelHandle the opener returned, not a copy");
assertNotNull(runtime.leadCoordLoop(),
"a configured coordinator: block must build the receiving LeadCoordLoop too");
assertFalse(ports.leadChannel.closed, "the channel must still be open while the daemon is running");
// --- shutting the assembly down closes it -------------------------------------------------
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
ports.shutdownHook.run();
assertTrue(ports.leadChannel.closed,
"FleetdRuntime.close() must close the configured LeadChannelHandle");
}
}
@@ -0,0 +1,165 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.testing.CapturedLog;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 A-gaps (gap 2): {@code Fleetd.java:185} used to call {@code
* reportRoleFallbackGaps(cfg)} <em>outside</em> the boundary {@code FleetdAssembly.assembleAndStart}
* — the assembly call itself sat at line 201, after it — so nothing that drives the assembly (the
* seam {@code FleetdAssemblyLifecycleTest} exercises) could ever notice the call being deleted.
* {@code Fleetd.main} still refuses a bad config end to end (see {@code
* FleetdStartupValidationTest}), but that only pins {@code validateAll()} and {@code
* assertChartersNameOnlyRegisteredTools} — a validator that <em>throws</em>. {@code
* reportRoleFallbackGaps} only logs; nothing about {@code main} throwing or not throwing can
* observe whether that particular call ran.
*
* <p>This drives {@link FleetdAssembly#assembleAndStart} directly — never a copy of its logic —
* with a config that has no {@code fleet:} pools or charters configured for any role, so every role
* trips both of {@code reportRoleFallbackGaps}' log branches, and asserts on the real log line a
* {@code ListAppender} attached to the shared {@code Fleetd}/{@code FleetdAssembly} logger
* captures. Deleting the call from {@code FleetdAssembly} (verified by hand, see the ticket) turns
* this test red; deleting it from {@code Fleetd.main} instead (its old location) would not, which
* is exactly the gap this test closes.
*/
class FleetdAssemblyRoleFallbackBoundaryTest {
/** Minimal fake {@link ResourcePorts}: enough for {@code assembleAndStart} to run with no real I/O. */
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> replyInbox;
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// No real HTTP bind in a unit test.
}
}
private static final class SentinelReplyInbox implements ReplyInbox {
@Override
public void own(String target) {
}
@Override
public void release(String target) {
}
@Override
public void publish(String target, String msgId, String content) {
}
@Override
public List<InboxMessage> peek(String target) {
return List.of();
}
@Override
public boolean ack(String target, String msgId) {
return false;
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
// Deliberately no `fleet:` block at all: every MemberRole has neither a pool nor a
// charter, so reportRoleFallbackGaps' "role fallback: no fleet.<role>s: pool for ..."
// branch is guaranteed to log something to assert on.
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
""");
return FleetConfig.load(f);
}
@Test
void assembleAndStartItselfReportsTheRoleFallbackGaps(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
try (var captured = CapturedLog.at(Fleetd.class, Level.INFO)) {
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
String infoLines = captured.events().stream()
.filter(e -> e.getLevel() == Level.INFO)
.map(ILoggingEvent::getFormattedMessage)
.reduce("", (a, b) -> a + "\n" + b);
assertTrue(infoLines.contains("role fallback: no fleet.<role>s: pool for"),
() -> "FleetdAssembly.assembleAndStart itself must call reportRoleFallbackGaps "
+ "(fleetd #612 A-gaps gap 2) — captured INFO lines: " + infoLines);
} finally {
if (ports.shutdownHook != null) {
ports.shutdownHook.run();
}
}
}
}
@@ -3,6 +3,7 @@ package dev.ltms.fleet;
import ch.qos.logback.classic.Level; import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent; import ch.qos.logback.classic.spi.ILoggingEvent;
import dev.ltms.fleet.config.FleetConfig; import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadMailbox; import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.testing.CapturedLog; import dev.ltms.fleet.testing.CapturedLog;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
@@ -35,7 +36,7 @@ class FleetdLeadMailboxSelectionTest {
boolean unreachable; boolean unreachable;
@Override @Override
public LeadMailbox open(String uri, String selfCoordId, int prefetch) { public LeadChannelHandle open(String uri, String selfCoordId, int prefetch) {
this.offeredUri = uri; this.offeredUri = uri;
this.offeredSelfId = selfCoordId; this.offeredSelfId = selfCoordId;
this.offeredPrefetch = prefetch; this.offeredPrefetch = prefetch;