diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java index a2e7966..3c47cf9 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java @@ -64,34 +64,39 @@ public final class StatusPoller { } private void loop() { - while (running) { - Set active = injector.activeTargets(); - for (String target : active) { - if (!running) return; - try { - // herdr's agent_status can misreport a settled worker as `unknown`; refine it - // against the pane content before it drives delivery/completion (CB-115). - // CB-185: refine THROUGH the same control the raw status came from — a router - // splits lead/member targets across two herdr daemons, and reading a lead's pane - // through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN. - AgentControl control = router != null ? router.agentsFor(target) : agents; - AgentStatus status = refiner.refine(target, control.status(target), control); - injector.onStatus(target, status); - } catch (HerdrException e) { - // The worker's agent is gone — stop trying and unblock its waiters. - if (e.code() != null && e.code().endsWith("_not_found")) { - log.debug("target {} gone; dropping its queue", target); - injector.drop(target, e); - } else { - log.debug("status poll for {} failed (will retry): {}", target, e.getMessage()); + try { + while (running) { + Set active = injector.activeTargets(); + for (String target : active) { + if (!running) return; + try { + // herdr's agent_status can misreport a settled worker as `unknown`; refine it + // against the pane content before it drives delivery/completion (CB-115). + // CB-185: refine THROUGH the same control the raw status came from — a router + // splits lead/member targets across two herdr daemons, and reading a lead's pane + // through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN. + AgentControl control = router != null ? router.agentsFor(target) : agents; + AgentStatus status = refiner.refine(target, control.status(target), control); + injector.onStatus(target, status); + } catch (HerdrException e) { + // The worker's agent is gone — stop trying and unblock its waiters. + if (e.code() != null && e.code().endsWith("_not_found")) { + log.debug("target {} gone; dropping its queue", target); + injector.drop(target, e); + } else { + log.debug("status poll for {} failed (will retry): {}", target, e.getMessage()); + } + } catch (Throwable e) { + log.error("unexpected failure polling {}; skipping this round", target, e); } - } catch (RuntimeException e) { - // Never let one target's unexpected error (e.g. an odd agent.get shape) kill - // the single poller thread and stall injection for every worker. - log.warn("unexpected error polling {}; skipping this round", target, e); } + sleep(); } - sleep(); + } finally { + if (running) { + log.error("status poller loop exited unexpectedly; it can be restarted"); + } + running = false; } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java b/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java index 69687a0..4a85342 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java +++ b/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java @@ -60,14 +60,21 @@ public final class SessionReaper { } private void loop() { - while (running) { - try { - sessions.reapIdle(idleTtlNanos); - } catch (RuntimeException e) { - log.warn("session reaper iteration failed; continuing", e); + try { + while (running) { + try { + sessions.reapIdle(idleTtlNanos); + } catch (Throwable e) { + log.error("session reaper iteration failed; continuing", e); + } + maybeSweepWipRefs(); + sleep(); } - maybeSweepWipRefs(); - sleep(); + } finally { + if (running) { + log.error("session reaper loop exited unexpectedly; it can be restarted"); + } + running = false; } } @@ -90,8 +97,8 @@ public final class SessionReaper { log.info("refs/wip retention sweep deleted {} snapshot ref(s) older than 24h whose " + "content was already reachable from main", deleted); } - } catch (RuntimeException e) { - log.warn("refs/wip retention sweep failed; continuing", e); + } catch (Throwable e) { + log.error("refs/wip retention sweep failed; continuing", e); } // Set even when the sweep threw, so a broken repo is retried on the slow cadence rather // than hammering git on every 5-second iteration. diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerResilienceTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerResilienceTest.java new file mode 100644 index 0000000..d1f03b1 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerResilienceTest.java @@ -0,0 +1,99 @@ +package dev.ltms.fleet.inject; + +import com.fasterxml.jackson.databind.JsonNode; +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.msg.TestTurnTokens; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class StatusPollerResilienceTest { + + @Test + void anErrorForOneTargetDoesNotStopPollingTheNextTarget() throws Exception { + FakeHerdr fake = new FakeHerdr().withAgent("worker", "term_b", "w2:p8", "w2:t8"); + CountDownLatch errorThrown = new CountDownLatch(1); + AgentControl agents = new AgentControl(new ErrorOnceForFirstTarget(fake, errorThrown)); + Injector injector = new Injector(agents); + StatusPoller poller = new StatusPoller(agents, injector, 1); + poller.start(); + try { + injector.enqueue("term_a", "first", TestTurnTokens.inert("term_a")); + assertTrue(errorThrown.await(2, TimeUnit.SECONDS), + "the first target must throw its test Error"); + + CompletableFuture delivered = + injector.enqueue("term_b", "second", TestTurnTokens.inert("term_b")).completion(); + delivered.get(2, TimeUnit.SECONDS); + } finally { + poller.stop(); + } + } + + @Test + void anAbnormalExitClearsRunningSoStartCreatesANewLoop() throws Exception { + StatusPoller poller = new StatusPoller(new AgentControl(new FakeHerdr()), new Injector(new AgentControl(new FakeHerdr())), -1); + poller.start(); + Thread first = threadOf(poller); + first.join(2000); + assertFalse(runningOf(poller), "an abnormal loop exit must clear running"); + + poller.start(); + Thread restarted = threadOf(poller); + try { + assertNotSame(first, restarted, "start() must create a new loop after an abnormal exit"); + restarted.join(2000); + } finally { + poller.stop(); + } + } + + private static Thread threadOf(StatusPoller poller) throws ReflectiveOperationException { + Field field = StatusPoller.class.getDeclaredField("thread"); + field.setAccessible(true); + return (Thread) field.get(poller); + } + + private static boolean runningOf(StatusPoller poller) throws ReflectiveOperationException { + Field field = StatusPoller.class.getDeclaredField("running"); + field.setAccessible(true); + return field.getBoolean(poller); + } + + private static final class ErrorOnceForFirstTarget implements HerdrClient { + private final FakeHerdr delegate; + private final CountDownLatch errorThrown; + private final AtomicBoolean first = new AtomicBoolean(true); + + private ErrorOnceForFirstTarget(FakeHerdr delegate, CountDownLatch errorThrown) { + this.delegate = delegate; + this.errorThrown = errorThrown; + } + + @Override + public JsonNode call(String method, Object params) { + if (method.equals("agent.get") && params instanceof Map map + && "w2:p7".equals(map.get("target")) && first.compareAndSet(true, false)) { + errorThrown.countDown(); + throw new AssertionError("test Error from the first poll target"); + } + return delegate.call(method, params); + } + + @Override + public void close() { + delegate.close(); + } + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperResilienceTest.java b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperResilienceTest.java new file mode 100644 index 0000000..14da39d --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperResilienceTest.java @@ -0,0 +1,91 @@ +package dev.ltms.fleet.session; + +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.WorkspaceControl; +import dev.ltms.fleet.member.ClaudeCodeLauncher; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.LongSupplier; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class SessionReaperResilienceTest { + + @Test + void anErrorInOneIterationDoesNotStopTheNextIteration() throws Exception { + CountDownLatch errorThrown = new CountDownLatch(1); + CountDownLatch nextIteration = new CountDownLatch(1); + AtomicBoolean first = new AtomicBoolean(true); + LongSupplier clock = () -> { + if (first.compareAndSet(true, false)) { + errorThrown.countDown(); + throw new AssertionError("test Error from the first reap iteration"); + } + nextIteration.countDown(); + return System.nanoTime(); + }; + SessionReaper reaper = new SessionReaper(sessionManager(clock), 60, 1); + reaper.start(); + try { + assertTrue(errorThrown.await(2, TimeUnit.SECONDS), + "the first reap iteration must throw its test Error"); + assertTrue(nextIteration.await(2, TimeUnit.SECONDS), + "the reaper must continue to the next iteration after an Error"); + } finally { + reaper.stop(); + } + } + + @Test + void anAbnormalExitClearsRunningSoStartCreatesANewLoop() throws Exception { + SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, -1); + reaper.start(); + Thread first = threadOf(reaper); + first.join(2000); + assertFalse(runningOf(reaper), "an abnormal loop exit must clear running"); + + reaper.start(); + Thread restarted = threadOf(reaper); + try { + assertNotSame(first, restarted, "start() must create a new loop after an abnormal exit"); + restarted.join(2000); + } finally { + reaper.stop(); + } + } + + private static SessionManager sessionManager(LongSupplier clock) { + FakeHerdr herdr = new FakeHerdr(); + FleetConfig.Profile cfg = new FleetConfig.Profile( + "ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", + List.of("ccs", "ltms-local"), "tab", "fleetd-workers", + "worker: {profile} #{n}", null, null, null); + ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null); + return new SessionManager(launcher, new FakeWorktrees(), clock); + } + + private static Thread threadOf(SessionReaper reaper) throws ReflectiveOperationException { + Field field = SessionReaper.class.getDeclaredField("thread"); + field.setAccessible(true); + return (Thread) field.get(reaper); + } + + private static boolean runningOf(SessionReaper reaper) throws ReflectiveOperationException { + Field field = SessionReaper.class.getDeclaredField("running"); + field.setAccessible(true); + return field.getBoolean(reaper); + } +}