fleetd #675: pin assembly loop timing defaults
This commit is contained in:
@@ -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");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user