Merge PR #631: fleetd #612 ranks 3+8 — behavioural pins for the AMQP assembly openers
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 59s
CI / build (push) Failing after 3m14s

Pins FleetdAssembly's replyInboxOpener (:357) and leadMailboxOpener (:363) so an
inert opener can no longer pass a green suite. Both directions covered: a durable
opener must leave the runtime owning the exact object it returned, and a failing
one must leave the in-memory fallback with coordination off — with the startup
report agreeing with the real state in each case.

Verified by the lead beyond the worker's own proof: a mutation that still CALLS
ports.replyInboxOpener() and discards the result is caught by assertSame, which is
rank 3's live defect (opener called, result thrown away, log still printing
'reply inbox: AMQP broker (durable)').

Test-only; no production change.
This commit is contained in:
Dai Ha
2026-10-01 16:38:16 +02:00
@@ -0,0 +1,176 @@
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<InboxMessage> 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<LeadMessage> 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<String, String> 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<ILoggingEvent> 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<ILoggingEvent> 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();
}
reports.list.clear();
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);
}
}
}