From 91be5079d6fc8574e58b2ac0ddf3af2b0a5adb73 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 1 Oct 2026 16:26:02 +0200 Subject: [PATCH 1/2] fleetd #612: pin AMQP assembly openers --- .../fleet/FleetdAssemblyAmqpOpenersTest.java | 175 ++++++++++++++++++ 1 file changed, 175 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java new file mode 100644 index 0000000..89e27e6 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java @@ -0,0 +1,175 @@ +package dev.ltms.fleet; + +import ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; +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.InMemoryReplyInbox; +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 org.slf4j.LoggerFactory; + +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.concurrent.atomic.AtomicInteger; +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.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #612 step 4 ranks 3 and 8: the assembled daemon must use the AMQP openers from {@link + * ResourcePorts}, and its startup report must describe the object the runtime actually owns. These + * fakes never open a socket. + */ +class FleetdAssemblyAmqpOpenersTest { + + private static final String COORD_ID = "assembly-test"; + + private static final class DurableReplyInbox 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 peek(String target) { return List.of(); } + @Override public boolean ack(String target, String msgId) { return false; } + } + + private static final class DurableLeadMailbox implements LeadChannelHandle { + @Override public void publish(String toCoordId, LeadMessage message) { } + @Override public List peek() { return List.of(); } + @Override public void ack(String msgId) { } + @Override public String selfCoordId() { return COORD_ID; } + @Override public boolean heldDurable() { return true; } + @Override public MailboxState inspect(String coordId) { return MailboxState.unknown(coordId); } + @Override public void close() { } + } + + private static final class RecordingPorts implements ResourcePorts { + final FakeHerdr herdr = new FakeHerdr(); + final DurableReplyInbox replyInbox = new DurableReplyInbox(); + final DurableLeadMailbox leadMailbox = new DurableLeadMailbox(); + final AtomicInteger replyOpenCalls = new AtomicInteger(); + final AtomicInteger mailboxOpenCalls = new AtomicInteger(); + final boolean openSucceeds; + + RecordingPorts(boolean openSucceeds) { + this.openSucceeds = openSucceeds; + } + + @Override public Map environment() { return Map.of(); } + @Override public HerdrClient connectHerdr(Path socketPath) { return herdr; } + @Override public Fleetd.AmqpOpener replyInboxOpener() { + return (uri, prefetch) -> { + replyOpenCalls.incrementAndGet(); + if (!openSucceeds) throw new IllegalStateException("fake reply broker is down"); + return replyInbox; + }; + } + @Override public Fleetd.LeadMailboxOpener leadMailboxOpener() { + return (uri, selfId, prefetch) -> { + mailboxOpenCalls.incrementAndGet(); + if (!openSucceeds) throw new IllegalStateException("fake coordination broker is down"); + return leadMailbox; + }; + } + @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) { } + @Override public void startHttp(Javalin app, String host, int port) { } + } + + private static FleetConfig writeConfig(Path dir) throws Exception { + Files.createDirectories(dir); + Path config = dir.resolve("fleetd.yaml"); + Files.writeString(config, """ + bind: + host: 127.0.0.1 + port: 8765 + idleSleepGuard: + enabled: false + broker: + uri: "amqp://fake-reply-broker/vh" + coordinator: + uri: "amqp://fake-coordination-broker/vh" + selfId: "assembly-test" + """); + return FleetConfig.load(config); + } + + private static FleetdRuntime assemble(Path dir, RecordingPorts ports) throws Exception { + FleetConfig cfg = writeConfig(dir); + return FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, new ConfigRef(dir.resolve("fleetd.yaml"), cfg), + new SubscriptionGuard(cfg.guard().hostSet())), ports); + } + + private static boolean reportContains(ListAppender appender, String text) { + return appender.list.stream().map(ILoggingEvent::getFormattedMessage).anyMatch(message -> message.contains(text)); + } + + @Test + void assembledAmqpOpenersAndTheirReportsAgreeOnDurableAndFallbackStates(@TempDir Path dir) throws Exception { + Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class); + Level oldLevel = logger.getLevel(); + ListAppender reports = new ListAppender<>(); + reports.start(); + logger.setLevel(Level.INFO); + logger.addAppender(reports); + try { + RecordingPorts durablePorts = new RecordingPorts(true); + FleetdRuntime durable = assemble(dir.resolve("durable"), durablePorts); + try { + // Control: this fails loudly if the assembly did not run or used an inert opener. + assertEquals(1, durablePorts.replyOpenCalls.get(), "assembly must call replyInboxOpener once"); + assertEquals(1, durablePorts.mailboxOpenCalls.get(), "assembly must call leadMailboxOpener once"); + assertSame(durablePorts.replyInbox, durable.replyInbox(), + "the durable reply report must describe the exact inbox the runtime owns"); + assertSame(durablePorts.leadMailbox, durable.leadMailbox(), + "the coordination-on report must describe the exact mailbox the runtime owns"); + assertNotNull(durable.leadCoordLoop(), "a durable mailbox must start lead coordination"); + assertTrue(reportContains(reports, "reply inbox: AMQP broker (durable)")); + assertTrue(reportContains(reports, "lead coordination: ON as coord-id " + COORD_ID)); + } finally { + durable.close(); + } + + RecordingPorts fallbackPorts = new RecordingPorts(false); + FleetdRuntime fallback = assemble(dir.resolve("fallback"), fallbackPorts); + try { + assertEquals(1, fallbackPorts.replyOpenCalls.get(), "assembly must call the failing reply opener once"); + assertEquals(1, fallbackPorts.mailboxOpenCalls.get(), "assembly must call the failing mailbox opener once"); + assertTrue(fallback.replyInbox() instanceof InMemoryReplyInbox, + "a failed reply opener must make the runtime own the in-memory fallback"); + assertNull(fallback.leadMailbox(), "a failed mailbox opener must leave coordination off"); + assertNull(fallback.leadCoordLoop(), "coordination must not start without a mailbox"); + assertTrue(reportContains(reports, "reply inbox: in-memory (soft-state)")); + assertTrue(reportContains(reports, "lead-to-lead messaging is OFF")); + } finally { + fallback.close(); + } + } finally { + logger.detachAppender(reports); + logger.setLevel(oldLevel); + } + } +} From c3b0406826315d2952ee7f543b194b979ab74597 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 1 Oct 2026 16:33:45 +0200 Subject: [PATCH 2/2] fleetd #612: isolate fallback AMQP reports --- .../test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java index 89e27e6..b5060f7 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyAmqpOpenersTest.java @@ -153,6 +153,7 @@ class FleetdAssemblyAmqpOpenersTest { durable.close(); } + reports.list.clear(); RecordingPorts fallbackPorts = new RecordingPorts(false); FleetdRuntime fallback = assemble(dir.resolve("fallback"), fallbackPorts); try {