Merge PR #688: fleetd #675 — pin three unpinned FleetdAssembly constructor args
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 1m35s

This commit is contained in:
Dai Ha
2026-10-03 22:27:32 +02:00
@@ -0,0 +1,190 @@
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.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
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.assertNotNull;
/**
* Asserts that the assembled loops use the production reminder, coordination, and delivery timing
* defaults when no {@code primary:} block configures the reply-push values.
*/
class FleetdAssemblyTimingDefaultsTest {
private static final class FakeLeadChannel 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 "test-lead";
}
@Override
public boolean heldDurable() {
return true;
}
@Override
public MailboxState inspect(String coordId) {
return MailboxState.unknown(coordId);
}
@Override
public void close() {
}
}
private static final class TestResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
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) -> new 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; }
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> new FakeLeadChannel();
}
@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) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
private TestResourcePorts ports;
@AfterEach
void tearDown() {
if (ports != null && ports.shutdownHook != null) {
ports.shutdownHook.run();
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
coordinator:
uri: "amqp://fake-lead-broker/vh"
selfId: "test-lead"
""");
return FleetConfig.load(file);
}
private FleetdRuntime assemble(Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ports = new TestResourcePorts();
return FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
}
private static long longField(Object target, String name) throws Exception {
Field field = target.getClass().getDeclaredField(name);
field.setAccessible(true);
return field.getLong(target);
}
@Test
void productionBootPathUsesTheExpectedLoopTimingDefaults(@TempDir Path dir) throws Exception {
FleetdRuntime runtime = assemble(dir);
ReplyPushLoop pushLoop = runtime.pushLoop();
assertEquals(5, longField(pushLoop, "maxReminders"),
"without primary:, ReplyPushLoop must stop after five reminder attempts");
assertEquals(15_000L, longField(pushLoop, "backoffMs"),
"without primary:, ReplyPushLoop must wait fifteen seconds before the next reminder");
LeadCoordLoop leadCoordLoop = runtime.leadCoordLoop();
assertNotNull(leadCoordLoop, "control: coordinator: must build LeadCoordLoop");
assertEquals(3_000L, longField(leadCoordLoop, "intervalMs"),
"LeadCoordLoop must poll for peer-lead mail every three seconds");
StatusPoller poller = runtime.poller();
assertEquals(Injector.POLL_INTERVAL_MILLIS, longField(poller, "intervalMillis"),
"StatusPoller must use Injector's delivery poll interval");
}
}