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
This commit was merged in pull request #624.
This commit is contained in:
@@ -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;
|
||||||
|
|||||||
Reference in New Issue
Block a user