Merge PR #671: fleetd #670 — pin excludedWorkspaceLabels at FleetdAssembly.java:265
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 58s
CI / build (push) Failing after 1m49s

This commit is contained in:
Dai Ha
2026-10-03 20:27:15 +02:00
@@ -0,0 +1,201 @@
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.herdr.LeadTabScanner;
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 io.javalin.Javalin;
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.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #670 — pins the {@code excludedWorkspaceLabels} argument {@link FleetdAssembly}'s
* production boot path passes to {@link LeadTabScanner} at {@code FleetdAssembly.java:265}
* ({@code Set.of()}).
*
* <p>{@code LeadTabScannerTest} already covers this constructor parameter, but it builds its own
* {@link LeadTabScanner} with its own set, so it tests the seam and proves nothing about the
* producer. This test instead reaches the exact object {@link FleetdAssembly#assembleAndStart}
* builds: a {@code fleet.leaders:} block makes the assembly construct a real
* {@link LeadTabScanner} for its local {@code leads} supplier, and a {@code coordinator:} block
* makes it hand that same supplier instance to {@link LeadCoordLoop} (fleetd #637), which stores
* it as a field. Reflection recovers it from there, and then from the scanner itself, so the
* assertion is against the real production argument rather than a copy built for this test.
*/
class FleetdAssemblyLeadTabScannerExclusionTest {
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) {
// Do not bind a real port in this assembly test.
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
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"
fleet:
leaders:
primary:
tab: "lead: primary"
profile: sonnet
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
""");
return FleetConfig.load(file);
}
@Test
void productionBootPathPassesNoExcludedWorkspaceLabels(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
TestResourcePorts ports = new TestResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
try {
LeadCoordLoop coordLoop = runtime.leadCoordLoop();
assertNotNull(coordLoop, "control: a configured coordinator: block must build LeadCoordLoop");
Field leadsField = LeadCoordLoop.class.getDeclaredField("leads");
leadsField.setAccessible(true);
@SuppressWarnings("unchecked")
Supplier<Map<String, String>> leads = (Supplier<Map<String, String>>) leadsField.get(coordLoop);
assertInstanceOf(LeadTabScanner.class, leads,
"control: a non-empty fleet.leaders: block must make FleetdAssembly build a real "
+ "LeadTabScanner for its `leads` supplier, not the Map::of fallback — "
+ "otherwise this test would pass for the wrong reason");
Field excludedField = LeadTabScanner.class.getDeclaredField("excludedWorkspaceLabels");
excludedField.setAccessible(true);
Set<?> excluded = (Set<?>) excludedField.get(leads);
assertTrue(excluded.isEmpty(),
"FleetdAssembly.java:265 must pass an empty excludedWorkspaceLabels to "
+ "LeadTabScanner — scanning member tabs would demote the lead to a worker");
} finally {
assertNotNull(ports.shutdownHook, "control: assembly must capture its shutdown hook");
ports.shutdownHook.run();
}
}
}