fleetd #538: recover polling loops after errors #543
@@ -64,34 +64,39 @@ public final class StatusPoller {
|
||||
}
|
||||
|
||||
private void loop() {
|
||||
while (running) {
|
||||
Set<String> 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<String> 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;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<Void> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user