Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha e4eb3dbed4 fleetd #504 item 1: stop swallowing real launchctl/systemctl failures on the 'loaded but not running' path
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Successful in 1m55s
The two 'loaded but not currently running' branches in the stop step (launchd/systemd, reached
when $OLD_PID is empty) ran 'launchctl unload'/'systemctl --user stop' with '2>/dev/null || true'
and printed 'ok' unconditionally. That swallowed a real supervisor failure (e.g. launchd or the
systemd user bus unreachable) exactly like a harmless already-stopped answer, and let the script
proceed to start a new daemon believing nothing was loaded -- the two-daemons failure fleetd #492
exists to prevent.

Adds unload_launchd_if_loaded/stop_systemd_if_loaded, applying systemd_loaded's own pattern
(capture stderr separately; a non-zero exit WITH stderr is a real failure, a non-zero exit with
empty stderr is a clean already-stopped answer) to the write side. The two call sites now use
these functions instead of the bare '|| true'.

Adds 5 tests: dies-on-real-failure and tolerates-clean-negative for each function, plus a
source-text check that the main flow calls the new functions instead of the original bare
'2>/dev/null || true'. All 5 verified by mutation (reintroducing the swallow, and separately
over-correcting to die unconditionally) -- each goes red with its own message, restores
byte-identical (full sha256), and passes a green control.
2026-09-12 13:42:27 +07:00
6 changed files with 192 additions and 238 deletions
@@ -64,39 +64,34 @@ public final class StatusPoller {
}
private void loop() {
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);
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 (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();
}
} finally {
if (running) {
log.error("status poller loop exited unexpectedly; it can be restarted");
}
running = false;
sleep();
}
}
@@ -60,21 +60,14 @@ public final class SessionReaper {
}
private void loop() {
try {
while (running) {
try {
sessions.reapIdle(idleTtlNanos);
} catch (Throwable e) {
log.error("session reaper iteration failed; continuing", e);
}
maybeSweepWipRefs();
sleep();
while (running) {
try {
sessions.reapIdle(idleTtlNanos);
} catch (RuntimeException e) {
log.warn("session reaper iteration failed; continuing", e);
}
} finally {
if (running) {
log.error("session reaper loop exited unexpectedly; it can be restarted");
}
running = false;
maybeSweepWipRefs();
sleep();
}
}
@@ -97,8 +90,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 (Throwable e) {
log.error("refs/wip retention sweep failed; continuing", e);
} catch (RuntimeException e) {
log.warn("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.
@@ -1,99 +0,0 @@
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();
}
}
}
@@ -1,91 +0,0 @@
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);
}
}
+49 -2
View File
@@ -287,6 +287,49 @@ systemd_loaded() {
return "$rc"
}
# fleetd #504: the "loaded but not currently running" branches in the main stop step (case
# launchd/systemd, reached when $OLD_PID is empty) used to run `launchctl unload`/`systemctl --user
# stop` with `2>/dev/null || true` and then print `ok` unconditionally — the exact conflation
# systemd_loaded/systemd_installed above already fixed on the READ side (fleetd #492 follow-up): a
# genuine "already stopped" answer (nonzero exit, nothing on stderr) is harmless, but a real tool
# failure (nonzero exit WITH a stderr message — e.g. launchd or the systemd user bus is
# unreachable) is not, and reporting `ok` on THAT means the start step below can register a fresh
# load on top of a supervisor that never actually let go: the exact two-daemons failure fleetd #492
# exists to prevent, reached from the one state (already odd) where a false `ok` is least
# affordable. These two functions apply the same "capture stderr separately, flag only a nonzero
# exit WITH stderr as a real failure" pattern to the WRITE side. No ${VAR:-default} anywhere here —
# see the #492 follow-up constraints comment above detect_supervisor for why a default would hide a
# lost value instead of surfacing it (fleetd #497's defect class).
unload_launchd_if_loaded() {
local err_file rc=0
if ! err_file="$(mktemp -t launchd-unload-err)"; then
die "could not create a temp file to capture 'launchctl unload' stderr — cannot tell a real
failure from a clean already-unloaded answer, so refusing to guess. The daemon's
supervision state was NOT touched."
fi
launchctl unload -w "$LAUNCHD_PLIST" >/dev/null 2>"$err_file" || rc=$?
if [ "$rc" -ne 0 ] && [ -s "$err_file" ]; then
die "'launchctl unload -w $LAUNCHD_PLIST' failed: $(cat "$err_file")
The daemon may still be under supervision; investigate before retrying."
fi
rm -f "$err_file"
}
stop_systemd_if_loaded() {
local err_file rc=0
if ! err_file="$(mktemp -t systemd-stop-err)"; then
die "could not create a temp file to capture 'systemctl --user stop' stderr — cannot tell a
real failure from a clean already-stopped answer, so refusing to guess. The daemon's
supervision state was NOT touched."
fi
systemctl --user stop "$SYSTEMD_UNIT" >/dev/null 2>"$err_file" || rc=$?
if [ "$rc" -ne 0 ] && [ -s "$err_file" ]; then
die "'systemctl --user stop $SYSTEMD_UNIT' failed: $(cat "$err_file")
The daemon may still be under supervision; investigate before retrying."
fi
rm -f "$err_file"
}
# fleetd #492: three real answers, not two — launchd, systemd, or genuinely unsupervised — plus a
# fourth, "ambiguous", for the one case this script cannot tell apart: both signals firing at once.
# That is exactly "I cannot tell who supervises this process", and guessing wrong here is how two
@@ -881,16 +924,20 @@ if [ -n "$OLD_PID" ]; then
elif [ "$SUPERVISOR_KIND" = "launchd" ]; then
# Loaded but not currently running (e.g. throttled after a crash loop). Unload it anyway so the
# start step below does a clean load, never a load stacked on an already-loaded label.
# fleetd #504: unload_launchd_if_loaded (above) tolerates a genuine already-unloaded answer but
# dies on a real `launchctl` failure — never a bare `|| true` that would print `ok` either way.
say "stop"
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true
unload_launchd_if_loaded
ok "launchd agent unloaded (was already not running)"
elif [ "$SUPERVISOR_KIND" = "systemd" ]; then
# Same case for systemd: the unit is known/active-capable but not currently running. `stop` on an
# already-stopped unit is a harmless no-op — kept for symmetry with the launchd branch above.
# fleetd #504: stop_systemd_if_loaded (above) tolerates that genuine no-op but dies on a real
# `systemctl` failure — never a bare `|| true` that would print `ok` either way.
say "stop"
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
systemctl --user stop "$SYSTEMD_UNIT" 2>/dev/null || true
stop_systemd_if_loaded
ok "systemd --user unit stopped (was already not running)"
else
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
+109
View File
@@ -621,6 +621,110 @@ test_refuse_drain_gate_call_site_present() {
|| fail "could not find the main flow's refuse_drain_gate call site in redeploy-fleetd.sh"
}
# fleetd #504 — the "loaded but not currently running" branches for launchd/systemd used to run
# `launchctl unload`/`systemctl --user stop` with `2>/dev/null || true` and print `ok`
# unconditionally, so a real supervisor failure (e.g. it cannot reach launchd/the systemd user bus)
# read exactly like a harmless already-stopped answer. unload_launchd_if_loaded/
# stop_systemd_if_loaded (redeploy-fleetd.sh, right after systemd_loaded) apply systemd_loaded's own
# "capture stderr separately — only a non-zero exit WITH stderr is a real failure" pattern to the
# WRITE side. Both call the real `launchctl`/`systemctl` binaries directly (they are not overridable
# wrapper functions the way launchd_loaded/systemd_loaded are), so these tests put a stub binary
# first on PATH — the same technique test_detect_supervisor_systemd_probe_error_is_unclear above
# already uses for `systemctl`.
test_unload_launchd_if_loaded_dies_on_real_failure() {
local bin_dir output rc=0 saved_plist="$LAUNCHD_PLIST"
bin_dir="$TMP/stub-bin-launchctl-error"; mkdir -p "$bin_dir"
cat > "$bin_dir/launchctl" <<'STUB'
#!/usr/bin/env bash
echo "Could not find specified service" >&2
exit 1
STUB
chmod +x "$bin_dir/launchctl"
LAUNCHD_PLIST="$TMP/fake-fail.plist"
output="$(PATH="$bin_dir:$PATH" unload_launchd_if_loaded 2>&1)" || rc=$?
LAUNCHD_PLIST="$saved_plist"
[ "$rc" -ne 0 ] \
|| fail "unload_launchd_if_loaded must die when launchctl exits non-zero AND writes to stderr"
printf '%s' "$output" | grep -qF 'launchctl unload' \
|| fail "die message does not name the failing launchctl unload command"
}
# Captured via $(...) rather than called bare: unload_launchd_if_loaded's own die() does a hard
# `exit`, and calling it directly at this level would let a regression that makes it die on this
# clean-negative case kill the WHOLE suite before the `|| fail` below ever ran — printing die's own
# message instead of this test's. Inside a command substitution, that `exit` only ends the subshell
# (a-guard-is-defeated-by-its-calling-context: the same reason the *_dies_on_real_failure tests
# above capture this way), so this test's own message is what actually reaches the report.
test_unload_launchd_if_loaded_tolerates_clean_negative() {
local bin_dir saved_plist="$LAUNCHD_PLIST" output rc=0
bin_dir="$TMP/stub-bin-launchctl-noop"; mkdir -p "$bin_dir"
cat > "$bin_dir/launchctl" <<'STUB'
#!/usr/bin/env bash
exit 1
STUB
chmod +x "$bin_dir/launchctl"
LAUNCHD_PLIST="$TMP/fake-noop.plist"
output="$(PATH="$bin_dir:$PATH" unload_launchd_if_loaded 2>&1)" || rc=$?
LAUNCHD_PLIST="$saved_plist"
[ "$rc" -eq 0 ] \
|| fail "unload_launchd_if_loaded must tolerate a clean already-unloaded answer (non-zero exit, empty stderr): $output"
}
test_stop_systemd_if_loaded_dies_on_real_failure() {
local bin_dir output rc=0
bin_dir="$TMP/stub-bin-systemctl-stop-error"; mkdir -p "$bin_dir"
cat > "$bin_dir/systemctl" <<'STUB'
#!/usr/bin/env bash
echo "Failed to connect to bus: No such file or directory" >&2
exit 1
STUB
chmod +x "$bin_dir/systemctl"
output="$(PATH="$bin_dir:$PATH" stop_systemd_if_loaded 2>&1)" || rc=$?
[ "$rc" -ne 0 ] \
|| fail "stop_systemd_if_loaded must die when systemctl exits non-zero AND writes to stderr"
printf '%s' "$output" | grep -qF 'systemctl --user stop' \
|| fail "die message does not name the failing systemctl --user stop command"
}
# Same subshell-capture reasoning as test_unload_launchd_if_loaded_tolerates_clean_negative above:
# stop_systemd_if_loaded's own die() does a hard `exit`, so this must run inside $(...) or a
# regression here would kill the whole suite with die's message instead of this test's.
test_stop_systemd_if_loaded_tolerates_clean_negative() {
local bin_dir output rc=0
bin_dir="$TMP/stub-bin-systemctl-stop-noop"; mkdir -p "$bin_dir"
cat > "$bin_dir/systemctl" <<'STUB'
#!/usr/bin/env bash
exit 1
STUB
chmod +x "$bin_dir/systemctl"
output="$(PATH="$bin_dir:$PATH" stop_systemd_if_loaded 2>&1)" || rc=$?
[ "$rc" -eq 0 ] \
|| fail "stop_systemd_if_loaded must tolerate a clean already-stopped answer (non-zero exit, empty stderr): $output"
}
# Closes the same gap test_refuse_drain_gate_call_site_present closes for the drain gate: the four
# tests above call unload_launchd_if_loaded/stop_systemd_if_loaded directly, and sourcing stops
# before the main flow ever runs (the SOURCED guard), so none of them can prove the main flow still
# CALLS these two functions instead of the original bare `2>/dev/null || true`. A source-text check,
# like test_swap_ordered_after_wait_and_before_start. The call-site needle is anchored (`^ name$`)
# so it cannot be satisfied by the comment lines above each call site that merely mention the
# function by name.
test_stop_branches_call_tolerant_helpers_not_bare_or_true() {
local src="$ROOT/scripts/redeploy-fleetd.sh" unload_call_line stop_call_line
unload_call_line="$(grep -n '^ unload_launchd_if_loaded$' "$src" | head -1 | cut -d: -f1 || true)"
stop_call_line="$(grep -n '^ stop_systemd_if_loaded$' "$src" | head -1 | cut -d: -f1 || true)"
[ -n "$unload_call_line" ] \
|| fail "could not find the main flow's call to unload_launchd_if_loaded in redeploy-fleetd.sh"
[ -n "$stop_call_line" ] \
|| fail "could not find the main flow's call to stop_systemd_if_loaded in redeploy-fleetd.sh"
if grep -qF 'launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true' "$src"; then
fail "the bare 'launchctl unload ... 2>/dev/null || true' defect (fleetd #504) is back in redeploy-fleetd.sh"
fi
if grep -qF 'systemctl --user stop "$SYSTEMD_UNIT" 2>/dev/null || true' "$src"; then
fail "the bare 'systemctl --user stop ... 2>/dev/null || true' defect (fleetd #504) is back in redeploy-fleetd.sh"
fi
}
test_no_errors() {
cat > "$TMP/no-errors.log" <<'LOG'
2026-09-05 12:00:00 INFO fleetd listening
@@ -1016,6 +1120,11 @@ test_refuse_drain_gate_build_ran_staged_absent
test_refuse_drain_gate_no_build_staged_present
test_refuse_drain_gate_no_build_staged_absent
test_refuse_drain_gate_call_site_present
test_unload_launchd_if_loaded_dies_on_real_failure
test_unload_launchd_if_loaded_tolerates_clean_negative
test_stop_systemd_if_loaded_dies_on_real_failure
test_stop_systemd_if_loaded_tolerates_clean_negative
test_stop_branches_call_tolerant_helpers_not_bare_or_true
test_no_errors
test_recovery_patterns_match_source
test_attributed_recovered_connection_error