fleetd #811: schedule the lead tab-label repair so it actually runs
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 2m5s

reassertLeadLabels was only called once, at boot, before launch() had ever
populated launchedTabByLead — a guaranteed no-op in production. Add
LeadLabelRepairLoop, a scheduleWithFixedDelay wrapper that drives the repair
on a recurring, configurable interval (fleet.leaders.<name>.repairIntervalSeconds,
default 30) and swallows any Throwable from its tick so one failure can never
permanently cancel the schedule. Also split ensureLeads()'s try/catch so a
repair failure and a count failure log their own distinct cause.
This commit is contained in:
Dai Ha
2026-10-07 11:23:52 +02:00
parent 916ca315aa
commit 4a88406bfd
13 changed files with 201 additions and 19 deletions
+2
View File
@@ -671,6 +671,8 @@ fleet:
# tabPrefix: "lead:" # only used to guard against a worker tabLabel colliding with
# # this convention at startup; plays no part in matching a lead
# scanIntervalSeconds: 10 # rescan cadence, and the worst case before a new tab is seen
# repairIntervalSeconds: 30 # how often a launched lead's tab label is re-asserted against
# # whatever Claude Code last retitled it to
# workspace: leads # where a launched lead's tab is created (default "fleet",
# # the same shared space the members use). Sharing that space
# # with members is the normal shipped shape: the scanner tells
@@ -24,6 +24,7 @@ import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.lead.LeadLabelRepairLoop;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.mcp.ConnectionIdentity;
@@ -564,12 +565,29 @@ final class FleetdAssembly {
configWatcher = null;
}
// fleetd #811: re-asserts a lead's tab label against the tab LeadLauncher launched it in.
// ensureLeads() above calls the same repair once, at boot, but launchedTabByLead is only
// ever filled by a launch that happens later in that same pass — so without this loop the
// repair never actually runs. SEVENTH of the recurring background loops to start (optional
// in the sense that an empty fleet.leaders gives it nothing to repair).
var leadLabelRepairScheduler = ports.newScheduler("bridge-lead-repair-");
final LeadLabelRepairLoop leadLabelRepair;
if (!leaders.isEmpty()) {
long repairIntervalSeconds = leaders.values().iterator().next().repairIntervalSeconds();
leadLabelRepair = new LeadLabelRepairLoop(() -> leadLauncher.reassertLeadLabels(leaders),
leadLabelRepairScheduler, repairIntervalSeconds);
leadLabelRepair.start();
} else {
leadLabelRepair = null;
leadLabelRepairScheduler.shutdownNow();
}
// CB-303 part 3: single ordered shutdown hook — see FleetdRuntime#close() for the statements
// this used to be. Registered here, at the exact point `main` used to register it: after
// configWatcher, before the Javalin app exists (see this class's own javadoc for why).
FleetdRuntime runtime = new FleetdRuntime(cfg, sessions, router, poller, messages, pushLoop, heartbeat,
leadCoordLoop, leadCoordSchedulerRef, healthMonitor, configWatcher, mcp, reaper, idleSleepGuard,
replyInbox, leadMailbox, completion, injector);
leadCoordLoop, leadCoordSchedulerRef, healthMonitor, configWatcher, leadLabelRepair, mcp, reaper,
idleSleepGuard, replyInbox, leadMailbox, completion, injector);
ports.addShutdownHook(runtime::close);
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
@@ -7,6 +7,7 @@ import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.lead.LeadLabelRepairLoop;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop;
@@ -51,6 +52,7 @@ final class FleetdRuntime implements AutoCloseable {
private final ScheduledExecutorService leadCoordScheduler; // nullable, paired with leadCoordLoop
private final FleetHealthMonitor healthMonitor; // nullable — health.enabled opt-in
private final ConfigWatcher configWatcher; // nullable — configReload.enabled opt-in
private final LeadLabelRepairLoop leadLabelRepair; // nullable — empty fleet.leaders opt-out
private final FleetMcp mcp;
private final SessionReaper reaper; // nullable — lifecycle.idleTtlSeconds opt-in
private final IdleSleepGuard idleSleepGuard; // nullable — idleSleepGuard.enabled: false
@@ -70,8 +72,8 @@ final class FleetdRuntime implements AutoCloseable {
FleetdRuntime(FleetConfig cfg, SessionManager sessions, HerdrRouter router, StatusPoller poller,
MessageService messages, ReplyPushLoop pushLoop, LeadHeartbeatLoop heartbeat,
LeadCoordLoop leadCoordLoop, ScheduledExecutorService leadCoordScheduler,
FleetHealthMonitor healthMonitor, ConfigWatcher configWatcher, FleetMcp mcp,
SessionReaper reaper, IdleSleepGuard idleSleepGuard, ReplyInbox replyInbox,
FleetHealthMonitor healthMonitor, ConfigWatcher configWatcher, LeadLabelRepairLoop leadLabelRepair,
FleetMcp mcp, SessionReaper reaper, IdleSleepGuard idleSleepGuard, ReplyInbox replyInbox,
LeadChannelHandle leadMailbox, CompletionResolver completion, Injector injector) {
this.cfg = cfg;
this.sessions = sessions;
@@ -84,6 +86,7 @@ final class FleetdRuntime implements AutoCloseable {
this.leadCoordScheduler = leadCoordScheduler;
this.healthMonitor = healthMonitor;
this.configWatcher = configWatcher;
this.leadLabelRepair = leadLabelRepair;
this.mcp = mcp;
this.reaper = reaper;
this.idleSleepGuard = idleSleepGuard;
@@ -109,6 +112,7 @@ final class FleetdRuntime implements AutoCloseable {
LeadCoordLoop leadCoordLoop() { return leadCoordLoop; }
FleetHealthMonitor healthMonitor() { return healthMonitor; }
ConfigWatcher configWatcher() { return configWatcher; }
LeadLabelRepairLoop leadLabelRepair() { return leadLabelRepair; }
FleetMcp mcp() { return mcp; }
SessionReaper reaper() { return reaper; }
IdleSleepGuard idleSleepGuard() { return idleSleepGuard; }
@@ -137,6 +141,7 @@ final class FleetdRuntime implements AutoCloseable {
if (leadCoordScheduler != null) leadCoordScheduler.shutdownNow();
if (healthMonitor != null) healthMonitor.stop();
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
if (leadLabelRepair != null) leadLabelRepair.stop(); // fleetd #811: stop the label repair loop
mcp.close();
if (reaper != null) reaper.stop();
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
@@ -1160,11 +1160,13 @@ public record FleetConfig(
* @param model the model or selector it runs, for operators reading the roster
* @param workspace the space this lead's tab lives in — the uniqueness boundary
* identity now depends on. Default {@link #DEFAULT_WORKSPACE}
* @param repairIntervalSeconds how often the daemon re-asserts this lead's tab label against
* the tab it launched. Default 30
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Leader(String profile, String tab, Integer instances, String tabPrefix,
Integer scanIntervalSeconds, String kind, String model,
String workspace, String cwd) {
String workspace, String cwd, Integer repairIntervalSeconds) {
/** The tab label every lead is found by, and an auto-launched instance is created with. */
public static final String LEAD_TAB_LABEL = "lead";
@@ -1185,12 +1187,14 @@ public record FleetConfig(
workspace = (workspace == null || workspace.isBlank())
? DEFAULT_WORKSPACE : workspace.strip();
tab = (tab == null || tab.isBlank()) ? null : tab.strip();
repairIntervalSeconds =
(repairIntervalSeconds == null || repairIntervalSeconds <= 0) ? 30 : repairIntervalSeconds;
}
/** Back-compat 7-arg form — no workspace or cwd, so both take their defaults. */
/** Back-compat 7-arg form — no workspace, cwd or repairIntervalSeconds, so each takes its default. */
public Leader(String profile, String tab, Integer instances, String tabPrefix,
Integer scanIntervalSeconds, String kind, String model) {
this(profile, tab, instances, tabPrefix, scanIntervalSeconds, kind, model, null, null);
this(profile, tab, instances, tabPrefix, scanIntervalSeconds, kind, model, null, null, null);
}
/** True when this lead may be launched by the daemon rather than only recognised. */
@@ -0,0 +1,57 @@
package dev.ltms.fleet.lead;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
* Drives {@link LeadLauncher#reassertLeadLabels} on a recurring schedule (fleetd #811) — the
* only caller that makes the repair reachable, since {@link LeadLauncher#ensureLeads()} runs
* once, at boot, before any lead has a remembered tab id to repair.
*/
public final class LeadLabelRepairLoop {
private static final Logger log = LoggerFactory.getLogger(LeadLabelRepairLoop.class);
private final Runnable repair;
private final ScheduledExecutorService scheduler;
private final long interval;
private final TimeUnit unit;
public LeadLabelRepairLoop(Runnable repair, ScheduledExecutorService scheduler, long intervalSeconds) {
this(repair, scheduler, intervalSeconds, TimeUnit.SECONDS);
}
/** Test seam: an injectable {@link TimeUnit} so a test can run this on a millisecond cadence. */
LeadLabelRepairLoop(Runnable repair, ScheduledExecutorService scheduler, long interval, TimeUnit unit) {
this.repair = repair;
this.scheduler = scheduler;
this.interval = interval;
this.unit = unit;
}
/** Begin the recurring repair. */
public void start() {
scheduler.scheduleWithFixedDelay(this::tick, interval, interval, unit);
}
/**
* One repair pass. Never throws: {@code scheduleWithFixedDelay} cancels its task for good the
* first time its body throws, which would silently retire the repair for the rest of the
* daemon's life.
*/
void tick() {
try {
repair.run();
} catch (Throwable e) {
log.warn("lead label repair tick failed, still scheduled: {}", e.getMessage());
}
}
/** Stop the recurring repair. */
public void stop() {
scheduler.shutdownNow();
}
}
@@ -151,9 +151,14 @@ public final class LeadLauncher {
return 0;
}
ReconcileScan scan;
try {
reassertLeadLabels(leaders);
} catch (HerdrException e) {
log.warn("lead label repair skipped this pass: {}", e.getMessage());
}
ReconcileScan scan;
try {
scan = countLeads(leaders);
} catch (HerdrException e) {
// Counting is the whole safety mechanism against double-spawning. If we cannot count, we
@@ -430,7 +435,7 @@ public final class LeadLauncher {
* <p>Keyed on the remembered tab id alone, never on a label search — another tab that happens
* to carry the same overwritten label is never touched.
*/
void reassertLeadLabels(Map<String, FleetConfig.Leader> leaders) {
public void reassertLeadLabels(Map<String, FleetConfig.Leader> leaders) {
if (launchedTabByLead.isEmpty()) {
return;
}
@@ -53,7 +53,10 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
* NOT exercised here; see the class-level caveat in the implementer's hand-off. {@code
* idleSleepGuard.enabled: false} is set for the same kind of reason: it would otherwise try to spawn
* a real {@code caffeinate} subprocess, which is not one of the resources the ticket's acceptance
* criteria names (scheduler/inbox/mailbox/loop/MCP server/herdr router).
* criteria names (scheduler/inbox/mailbox/loop/MCP server/herdr router). {@code fleet.leaders:} is
* left unset too, which is why the {@code bridge-lead-repair-} scheduler is the one "always-created"
* scheduler this config shuts down immediately rather than leaving running (fleetd #811) — an empty
* {@code fleet.leaders} gives the label repair nothing to do.
*/
class FleetdAssemblyLifecycleTest {
@@ -208,8 +211,9 @@ class FleetdAssemblyLifecycleTest {
// order, then the shutdown hook is registered, then (finally) HTTP "starts" -------------
assertTrue(ports.ledger.indexOf("connectHerdr") < ports.ledger.indexOf("newScheduler:bridge-push-"),
"herdr must connect before the push scheduler is created: " + ports.ledger);
assertEquals(List.of("bridge-push-", "bridge-heartbeat-", "bridge-health-"), ports.schedulerPurposes,
"the three always-created schedulers must be requested in exactly this order: "
assertEquals(List.of("bridge-push-", "bridge-heartbeat-", "bridge-health-", "bridge-lead-repair-"),
ports.schedulerPurposes,
"the four always-created schedulers must be requested in exactly this order: "
+ ports.schedulerPurposes);
assertTrue(ports.ledger.indexOf("addShutdownHook") < ports.ledger.indexOf("startHttp"),
"the shutdown hook must be registered before HTTP starts — the one statement that "
@@ -239,8 +243,13 @@ class FleetdAssemblyLifecycleTest {
"FleetdRuntime must own the exact ReplyInbox instance this fake's AmqpOpener returned");
assertFalse(ports.herdr.closed, "herdr must still be open while the daemon is running");
assertFalse(ports.replyInbox.closed, "the reply inbox must still be open while the daemon is running");
for (ScheduledExecutorService scheduler : ports.schedulers) {
assertFalse(scheduler.isShutdown(), "a scheduler must still be running while the daemon is up");
// fleet.leaders: is unset in this config (see class javadoc), so bridge-lead-repair- is
// shut down immediately rather than left running — every OTHER scheduler must still be up.
int leadRepairIndex = ports.schedulerPurposes.indexOf("bridge-lead-repair-");
for (int i = 0; i < ports.schedulers.size(); i++) {
boolean expectRunning = i != leadRepairIndex;
assertEquals(expectRunning, !ports.schedulers.get(i).isShutdown(),
"scheduler for '" + ports.schedulerPurposes.get(i) + "' running state while the daemon is up");
}
// --- close order: invoke the captured shutdown-hook Runnable directly (no real JVM shutdown
@@ -95,7 +95,7 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("primary", new FleetConfig.Primary("term-a", 1, 1000));
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-a", 1, null, 10,
"claude", null, null, null)),
"claude", null, null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(300, 60_000L, 3));
v.put("health", new FleetConfig.Health(true, 30, 600, null, null));
@@ -143,7 +143,7 @@ class ConfigRefTopLevelReportingCoverageTest {
// test — that is by design, not a gap: this map exists to prove fleet.leaders is covered.
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-b", 1, null, 10,
"claude", null, null, null)),
"claude", null, null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(600, 120_000L, 5));
v.put("health", new FleetConfig.Health(false, 90, 900, null, null));
@@ -485,6 +485,8 @@ class FleetConfigTest {
assertEquals("lead:", lead.tabPrefix());
assertEquals(10, lead.scanIntervalSeconds());
assertEquals(1, lead.instances(), "one of a lead is the assumption worth defaulting to");
assertEquals(30, lead.repairIntervalSeconds(),
"no fleetd.yaml edit is required for the label repair loop to work");
}
@Test
@@ -499,11 +501,13 @@ class FleetConfigTest {
tab: "drive: opus"
tabPrefix: "drive:"
scanIntervalSeconds: 30
repairIntervalSeconds: 45
""");
FleetConfig.Leader lead = FleetConfig.load(f).fleet().leaders().get("opus");
assertEquals("drive:", lead.tabPrefix());
assertEquals(30, lead.scanIntervalSeconds());
assertEquals(45, lead.repairIntervalSeconds());
}
/**
@@ -81,7 +81,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
v.put("primary", new FleetConfig.Primary("term-guard", 4, 4000));
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-guard", 1, null, 10,
"claude", null, null, null)),
"claude", null, null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(301, 61_000L, 4));
v.put("health", new FleetConfig.Health(true, 31, 601, 61, null));
@@ -0,0 +1,78 @@
package dev.ltms.fleet.lead;
import org.junit.jupiter.api.Test;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #811 follow-up: {@code LeadLauncher.reassertLeadLabels} was unit-tested only via a direct
* call, a scenario production never reaches — nothing scheduled it. These tests drive the real
* {@code scheduleWithFixedDelay} path on a short, real interval rather than calling {@link
* LeadLabelRepairLoop#tick()} directly, so a future regression that removes the scheduling call
* again fails here too.
*/
class LeadLabelRepairLoopTest {
/** Poll up to this long for a background tick to land, rather than betting on one fixed sleep. */
private static void awaitAtLeast(AtomicInteger counter, int target, long timeoutMs) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeoutMs);
while (counter.get() < target && System.nanoTime() < deadline) {
Thread.sleep(10);
}
}
@Test
void theRepairActuallyRunsFromTheScheduledPathNotOnlyADirectCall() throws Exception {
AtomicInteger runs = new AtomicInteger();
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
LeadLabelRepairLoop loop = new LeadLabelRepairLoop(runs::incrementAndGet, scheduler, 10, TimeUnit.MILLISECONDS);
try {
loop.start();
awaitAtLeast(runs, 1, 2000);
assertTrue(runs.get() > 0,
"scheduleWithFixedDelay must fire the repair on its own — nobody here calls tick() directly");
} finally {
loop.stop();
}
}
@Test
void aTickThatThrowsDoesNotStopLaterTicks() throws Exception {
AtomicInteger calls = new AtomicInteger();
Runnable repair = () -> {
if (calls.incrementAndGet() == 1) {
throw new RuntimeException("boom");
}
};
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
LeadLabelRepairLoop loop = new LeadLabelRepairLoop(repair, scheduler, 10, TimeUnit.MILLISECONDS);
try {
loop.start();
awaitAtLeast(calls, 2, 2000);
assertTrue(calls.get() >= 2,
"a tick whose first call throws must still run again on a later interval");
} finally {
loop.stop();
}
}
@Test
void tickItselfSwallowsAThrowingRepairAndReturnsNormally() {
Runnable repair = () -> {
throw new IllegalStateException("boom");
};
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
LeadLabelRepairLoop loop = new LeadLabelRepairLoop(repair, scheduler, 10, TimeUnit.MILLISECONDS);
try {
assertDoesNotThrow(loop::tick);
} finally {
loop.stop();
}
}
}
@@ -63,7 +63,7 @@ class LeadLauncherTest {
private static FleetConfig.Leader lead(String profile, String tab, int instances) {
return new FleetConfig.Leader(profile, tab, instances, "lead:", 10, null, null,
"fleet", "/repo");
"fleet", "/repo", null);
}
private static LeadLauncher launcher(FakeHerdr herdr, FleetConfig cfg) {
@@ -114,7 +114,7 @@ class LeadRolloverTest {
null, null, null,
Map.of(), null, null, true, null);
FleetConfig.Leader leader = new FleetConfig.Leader("opus", "lead: opus", 1, "lead:", 10,
null, null, "fleet", null);
null, null, "fleet", null, null);
Map<String, FleetConfig.Leader> leaders = new LinkedHashMap<>();
leaders.put(LEAD_NAME, leader);
FleetConfig.Fleet fleet = new FleetConfig.Fleet(leaders, Map.of(), Map.of(), Map.of(), Map.of(), null);