Compare commits

...

12 Commits

Author SHA1 Message Date
Dai Ha 130962e8e3 fleetd must not let the host idle-sleep while members are live
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 1m54s
Adds a small IdleSleepGuard (dev.ltms.fleet.power) that holds a macOS
caffeinate -i child while at least one fleet member is live, and
releases it once none are. It hangs off SessionManager's existing
onAcquire/onRelease hooks and SessionManager#size() rather than
tracking members a second way. New idleSleepGuard: config block,
on by default, following the FleetConfig.Health/ConfigReload pattern.
2026-09-05 05:37:38 +07:00
Dai Ha b6b88c5f1c #334: pin the fresh-owner gate on ask()'s timeout teardown
CI / build (push) Successful in 2m7s
CI / contract (push) Successful in 34m37s
2026-09-04 17:12:49 +07:00
Dai Ha 86dddfe240 Merge #334: ask()'s timeout closes the turn before it forgets the task mapping 2026-09-04 17:08:40 +07:00
Dai Ha 0d5944af63 fleetd #334: close ask()'s turn before forgetting its Task, closing the last stranding window
CI / build (pull_request) Successful in 2m14s
CI / contract (pull_request) Successful in 19m14s
ask()'s TimeoutException catch used to run clearAsyncQuestion(turnId, true) -- forgetting the
Task's asyncTasksByTurn mapping -- before rendezvous.closeAsk(turnId) ran in the shared finally.
Between those two calls the ask was still "answerable" (askSession(turnId) non-null) but the Task
mapping was already gone, so a racing answer() call found task == null, skipped
finishAsyncTask, and stranded the async ticket at PENDING even though answer() itself reported a
result. #329 fixed one step of this same race; this closes the remaining one.

The fix reorders the fresh owner's teardown: closeAsk runs first, then markAskTimedOut and
clearAsyncQuestion. A racing answer() call now either sees the ask still open (and the Task
mapping guaranteed intact) or sees it already closed (STALE_TURN, before it ever reaches
asyncTasksByTurn). It also gates the whole block by ticket.fresh(), matching the invariant the
finally block already states ("only the fresh owner tears down the shared turn") -- a duplicate
coalesced ask() timing out no longer forgets bookkeeping the fresh owner still needs.

Adds aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket, which pins the exact window
with a new test-only hook (askTimeoutRaceHookForTest) and proves both invariants: a late answer()
racing the timeout sees STALE_TURN, and the async ticket still resolves DONE from the worker's
real reply. Reverting the reorder (verified locally, not committed) makes this test fail with
"expected STALE_TURN but was TIMED_OUT_WORKING".
2026-09-04 17:06:15 +07:00
Dai Ha 4a5030a5c6 #348: drop a chrome skip that cannot fire, and pin the live pattern shape
CI / contract (push) Successful in 2m10s
CI / build (push) Successful in 2m12s
2026-09-04 16:57:36 +07:00
Dai Ha c1c8794c48 Merge #348: a member's prose about a usage limit no longer quarantines a credential 2026-09-04 16:53:19 +07:00
Dai Ha 65a78932c1 Merge #335: a per-task cleanup throw in abandon() no longer strands the tasks behind it
CI / contract (push) Successful in 1m24s
CI / build (push) Successful in 2m8s
2026-09-04 16:44:04 +07:00
Dai Ha f429ca1a50 Avoid exhaustion cooldown for member prose 2026-09-04 16:43:16 +07:00
Dai Ha 73aab3f83e Merge #342: teardown resolves the pane's real tab instead of trusting the delegate's placement config
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m47s
2026-09-04 16:38:31 +07:00
Dai Ha 887aca0183 fleetd#335: abandon()'s per-task cleanup and sendAsync's terminal hook must not swallow throws
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 2m21s
Site 1 (abandon()'s matching loop, reachable): the recovery/put-back branch calls
inbox.publish, which AmqpReplyInbox implements as a real broker round trip that
throws IllegalStateException on an unroutable/unconfirmed/interrupted publish.
An uncaught throw there aborted the loop, stranding every task after it in
`matching` PENDING forever. Fixed by recording each task's own future.complete()
result before any cleanup runs, then wrapping the cleanup in try/catch so one
task's failure cannot stop its siblings from getting their outcome. Reaching the
throwing branch by real timing needs a race the file's own #137 follow-up already
found unreachable through the public API, so the reproducing test uses a
test-only hook (same technique as the existing fleetd #324/#329 hooks) to inject
the throw at that exact point.

Site 2 (sendAsync's task.future.whenComplete, reachable): the returned stage is
discarded, so an uncaught throw from pushLoop.onTicketTerminal vanished with no
log line. Reproduced for real: Fleetd's shutdown hook runs messages.close()
(stops the async executor from taking new work, but does not cancel a send
already in flight) before pushLoop.close() (shuts its scheduler down
immediately) — a ticket completing in that window makes onTicketTerminal's own
scheduler.schedule(...) throw a genuine RejectedExecutionException. Fixed with a
try/catch(Throwable) plus log.error inside the whenComplete action.

Site 3 (the two `finally { asyncTasksByWaiter.remove(reply); rendezvous.close(...);
}` blocks in send() and answer()): read Rendezvous.close/closeAsk and the
ConcurrentHashMap operations behind them — both are plain map ops on a non-null
key with no user-overridable code, so neither can throw. Left unchanged; not a
defect.

Mutation-proven: reverting either fix reproduces the failure it exists to catch
— removing site 1's try/catch aborts abandon() with the injected exception
(MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks
errors); removing site 2's try/catch leaves the RejectedExecutionException
unlogged (MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently
fails its log assertion). Full suite: mvn clean install, Tests run: 1365,
Failures: 0, Errors: 0, BUILD SUCCESS.
2026-09-04 16:38:14 +07:00
Dai Ha f379847942 Merge #345: the timeout path's use of Cancellation.DELIVERED is now pinned
CI / contract (push) Successful in 47s
CI / build (push) Successful in 2m8s
2026-09-04 16:32:45 +07:00
Dai Ha ea12107497 fleetd #345: test timeout cancellation race
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m3s
2026-09-04 16:30:46 +07:00
18 changed files with 1201 additions and 49 deletions
+8
View File
@@ -110,6 +110,14 @@ bind:
# notifications:
# mode: disabled
# Idle-sleep guard: while at least one member is live, hold an OS-level assertion against idle
# sleep (macOS only — a `caffeinate -i` child; a no-op elsewhere or if caffeinate is missing), so
# an unattended host does not idle-sleep out from under a member's long turn. Unlike health/
# configReload above, this is ON BY DEFAULT — omitting the block entirely leaves it enabled, the
# same as `enabled: true`. Uncomment only to turn it off:
# idleSleepGuard:
# enabled: false
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
herdrSocket: ~/.config/herdr/herdr.sock
@@ -55,6 +55,8 @@ import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.member.OpenCodeLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
import dev.ltms.fleet.power.IdleSleepGuard;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -252,6 +254,26 @@ public final class Fleetd {
System::nanoTime, contextCap, clearAfterTurn);
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
// Idle-sleep guard: hold an OS-level assertion against idle sleep while at least one
// member is live, so an unattended host does not idle-sleep out from under a member's
// long turn (see FleetConfig.IdleSleepGuard / dev.ltms.fleet.power.IdleSleepGuard for the
// measurement that motivated this). Opt-out via idleSleepGuard.enabled: false; on by
// default. Hangs off SessionManager's own onAcquire/onRelease hooks (CB-520/CB-516,
// previously wired only to the reply inbox) and SessionManager#size() — the exact registry
// fleet_list's live/capacity numbers are themselves computed from — rather than tracking
// members a second way. No-op (never constructed) off macOS or when idleSleepGuard.enabled
// is explicitly false; the mechanism itself is additionally a no-op if 'caffeinate' cannot
// be started, so this can never fail a spawn, a release, or startup.
boolean idleSleepGuardEnabled = cfg.idleSleepGuard() == null || cfg.idleSleepGuard().isEnabled();
final IdleSleepGuard idleSleepGuard;
if (idleSleepGuardEnabled) {
idleSleepGuard = new IdleSleepGuard(new CaffeinateSleepAssertionMechanism(), sessions::size);
sessions.onAcquire(_ -> idleSleepGuard.recheck());
sessions.onRelease(_ -> idleSleepGuard.recheck());
} else {
idleSleepGuard = null;
}
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
final SessionReaper reaper;
if (cfg.lifecycle() != null
@@ -698,6 +720,11 @@ public final class Fleetd {
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
mcp.close();
if (reaper != null) reaper.stop();
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
// drained every session (and each release already drove the live count to 0, which
// releases the guard's assertion on its own) — this is the backstop for a drain that was
// itself interrupted or threw, so no caffeinate child ever outlives the daemon.
if (idleSleepGuard != null) idleSleepGuard.close();
// Release the broker connection last among message resources (no-op for the in-memory inbox).
if (replyInbox instanceof AutoCloseable closeable) {
try {
@@ -37,6 +37,10 @@ import java.util.function.Supplier;
* makes {@code fleet:} split rather than hot — see below.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code idleSleepGuard:} ({@code Fleetd.java} reads it once, at startup, to decide whether
* to construct an {@code IdleSleepGuard} and wire {@code SessionManager}'s
* {@code onAcquire}/{@code onRelease} hooks to it — neither is rebuilt on reload, so a
* running daemon keeps whatever this was at startup regardless of a later edit),
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
* {@code guard:}, {@code worktreeRoot:} and {@code worktreeGroup:} (both baked once into the
@@ -129,8 +133,9 @@ import java.util.function.Supplier;
* five of COLD_KEYS" rather than re-listing them, so prose and set cannot drift again.</li>
* </ul>
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333).</strong>
* {@code FleetConfig} has 22 top-level record components: 5 cold, 11 deferred, 3 split, 3
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333),
* recounted again after {@code idleSleepGuard:} was added.</strong>
* {@code FleetConfig} has 23 top-level record components: 5 cold, 12 deferred, 3 split, 3
* hot-excluded. Three of them are named nowhere in this file, and the reason is the same for all
* three: {@code placement}, {@code memberCredentials} and {@code memberLoginShell} are
* <strong>hot</strong> and correctly absent — all three are read live off {@code config.get()}
@@ -213,7 +218,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
static final Set<String> DEFERRED_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles");
"quarantineCooldownSeconds", "profiles", "idleSleepGuard");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -418,6 +423,14 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.configReload(), fresh.configReload())) {
changed.add("configReload");
}
// Fleetd.java reads cfg.idleSleepGuard() once, at startup, to decide whether to construct
// an IdleSleepGuard at all and wire SessionManager's onAcquire/onRelease hooks to it —
// neither is rebuilt on reload, so a running daemon keeps whatever this was at startup
// (armed or not) regardless of a later edit here. Not cold: nothing already-open goes
// inconsistent with the new value, an armed-or-not guard just keeps its original answer.
if (!Objects.equals(old.idleSleepGuard(), fresh.idleSleepGuard())) {
changed.add("idleSleepGuard");
}
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
@@ -106,6 +106,11 @@ import java.util.regex.PatternSyntaxException;
* When {@code memberHerdrSocket} is NOT configured this field is never
* consulted at all; fleetd keeps reading its own {@code $SHELL}, exactly as
* before this field existed.
* @param idleSleepGuard opt-in-by-default: hold an OS-level assertion against idle sleep while at
* least one member is live, so an unattended host does not sleep out from
* under a member's long turn. {@code null} (the block omitted) behaves the
* same as an explicit {@code enabled: true}; set {@code enabled: false} to
* turn it off. See {@link dev.ltms.fleet.power.IdleSleepGuard}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -130,7 +135,22 @@ public record FleetConfig(
MemberCredentials memberCredentials,
Coordinator coordinator,
String worktreeGroup,
String memberLoginShell) {
String memberLoginShell,
IdleSleepGuard idleSleepGuard) {
/** Back-compat form before the {@code idleSleepGuard:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
String memberLoginShell) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
memberLoginShell, null);
}
/** Back-compat form before the {@code memberLoginShell} key was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
@@ -141,7 +161,7 @@ public record FleetConfig(
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null);
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null, null);
}
/** Back-compat form before the {@code worktreeGroup} key was added. */
@@ -1245,6 +1265,25 @@ public record FleetConfig(
}
}
/**
* Hold an OS-level assertion against idle sleep while at least one member is live (see
* {@link dev.ltms.fleet.power.IdleSleepGuard}).
*
* <p>Unlike most opt-in blocks in this file, this one defaults to <em>on</em>: an unattended
* host idle-sleeping mid-turn is a correctness problem (a dropped AMQP link, a frozen member),
* not a convenience, so the safer default is armed. An operator who wants the previous
* behaviour (no assertion held, ever) sets {@code enabled: false} explicitly.
*
* @param enabled {@code false} turns the guard off; {@code null} (the block omitted
* entirely) or {@code true} leaves it on
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record IdleSleepGuard(Boolean enabled) {
public boolean isEnabled() {
return !Boolean.FALSE.equals(enabled);
}
}
/**
* The terminal → lead-name map seeded from the legacy singular {@code primary:} pin (CB-530).
*
@@ -1495,7 +1534,7 @@ public record FleetConfig(
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell");
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "idleSleepGuard");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -2172,9 +2211,14 @@ public record FleetConfig(
// memberLoginShell is left as-is (fleetd #213), like worktreeGroup: null/blank is "not
// configured", and there is no sane non-null default — a member's login shell is
// operator-specific and only meaningful when memberHerdrSocket is also set.
// idleSleepGuard is left as-is, like leadHeartbeat/configReload above, but for the opposite
// reason: it is on by default already (its own isEnabled() treats null the same as
// enabled: true — see its javadoc), so defaulting the block here would change nothing a
// reader observes and would only obscure that "block omitted" and "block present and
// enabled" are deliberately the same outcome.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell);
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, idleSleepGuard);
}
/**
@@ -380,7 +380,9 @@ public final class CompletionResolver implements TurnListener {
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
if (startsWithExhaustion(matchedLine, exhausted)) {
exhaustionSink.onExhausted(target, reason);
}
}
return;
}
@@ -478,7 +480,9 @@ public final class CompletionResolver implements TurnListener {
+ "usable assistant block; no fleet_reply): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
if (startsWithExhaustion(matchedLine, exhausted)) {
exhaustionSink.onExhausted(target, reason);
}
}
return true;
}
@@ -613,6 +617,46 @@ public final class CompletionResolver implements TurnListener {
return null;
}
/**
* True when nothing before the match on this pane line ends a sentence — that is, the match is
* still inside the line's first sentence rather than inside prose a member wrote about it.
* Used to decide whether an exhaustion match may quarantine a credential (fleetd #348).
*
* <p><strong>Why this is looser than {@link #startsWithBackendError}.</strong> An
* {@code exhaustedPattern} is written per profile and may name only the decisive words of a
* provider message — {@code "usage limit has been reached"} without its leading {@code "The"}.
* A start-of-line check would then reject the genuine refusal. That is the false negative
* fleetd #348's invariant 1 calls the worse direction: an unrecorded exhaustion leaves the
* fleet spawning into a credential with no capacity, and a quarantine runs 1800s against the
* backend-error cooldown's fixed 60s.
*
* <p>This rule accepts a superset of what a start-of-line check accepts: if the match begins
* right after the chrome, there is nothing in front of it, so there is no sentence ending
* either. So moving to it cannot add a false negative.
*
* <p><strong>No chrome skipping here, deliberately.</strong> The first version of this method
* copied {@code startsWithBackendError}'s leading-chrome loop. Measured on merge: deleting that
* loop left all 1369 tests green, and it must — the scan only looks for {@code . ! ?}, and no
* terminal chrome character is one of those. A step that cannot change the result is worse than
* no step, because the next reader takes it as evidence that chrome was handled.
*
* <p>It stays a heuristic. Prose whose <em>first</em> sentence carries the pattern still
* notifies the sink, and a genuine refusal behind an earlier full stop (a hostname, a version
* number) still does not. Both are known and neither is fixed here.
*/
private static boolean startsWithExhaustion(String line, Pattern pattern) {
var matcher = pattern.matcher(line);
if (!matcher.find()) {
return false;
}
for (int prefix = 0; prefix < matcher.start(); prefix++) {
if (".!?".indexOf(line.charAt(prefix)) >= 0) {
return false;
}
}
return true;
}
/**
* True when the error pattern begins the matched pane line, rather than appearing in prose.
*
@@ -709,27 +709,47 @@ public final class MessageService {
boolean isRecovery = task == recoveryTask && recovered != null;
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
String turnId = task.turnId;
if (task.future.complete(outcome)) {
if (outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
// fleetd #335: completing THIS task's future must not depend on any other task's
// cleanup succeeding — every task in `matching` is owed its own outcome regardless of
// what happens below, so decide and record that before doing anything that can throw.
boolean completedHere = task.future.complete(outcome);
if (completedHere && outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
}
try {
if (abandonCleanupHookForTest != null) {
// Test-only (fleetd #335 site 1): see the field's own javadoc.
abandonCleanupHookForTest.run();
}
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
if (completedHere) {
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
} catch (RuntimeException e) {
// fleetd #335: inbox.publish reaches a broker (AmqpReplyInbox throws
// IllegalStateException on an unroutable/unconfirmed/interrupted publish) and this
// loop has no other teardown path — a caller on the release path, or the health
// monitor's GONE/NEVER_READY sweep. Losing this exception uncaught would abort the
// loop and leave every task still to come in `matching` PENDING forever (fleetd
// #335 site 1). completedHere is already recorded above, so only this task's
// best-effort bookkeeping is lost — log it and let the loop reach the rest.
log.error("abandon: per-task cleanup failed for ticket {} (target {}, turnId {})",
task.ticket, target, turnId, e);
}
}
if (failed) {
@@ -871,6 +891,10 @@ public final class MessageService {
boolean wasDelivered = delivery.completion().isDone()
&& !delivery.completion().isCompletedExceptionally();
if (!wasDelivered) {
if (timeoutCancellationRaceHookForTest != null) {
// Test-only (fleetd #345): see the field's own javadoc.
timeoutCancellationRaceHookForTest.run();
}
// The target monitor makes cancellation atomic with onStatus picking this
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
@@ -941,13 +965,41 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
// what keeps the target from staying BUSY forever), but it would otherwise also erase
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
// Only the fresh owner tears down the shared turn (mirrors the finally block below).
// A duplicate's own timeoutMillis says nothing about whether the SHARED ask is actually
// done — it must leave the close/forget bookkeeping to the fresh owner, exactly as it
// already leaves closeAsk to it.
if (ticket.fresh()) {
// fleetd #334: close the ask turn BEFORE forgetting this task's turnId mapping below.
// Before this fix the order was reversed — the mapping was forgotten here first, and
// rendezvous.closeAsk only ran afterward, in the shared finally. A primary's answer()
// call racing this exact timeout could then find rendezvous.askSession(turnId) still
// non-null (the ask still "answerable") after the Task mapping was already gone:
// answer()'s own asyncTasksByTurn lookup returned null, its task != null guard skipped
// the completion, and the async ticket sat at PENDING forever even though answer()
// itself reported the worker's real reply. Closing here first removes that window:
// any answer() call that still observes askSession(turnId) != null is necessarily
// racing a point BEFORE the forgetting below runs (both happen on this one thread, in
// this order, with nothing that yields in between), so the Task mapping is still there
// for it to find; any call that observes askSession(turnId) == null now correctly
// bails out STALE_TURN (see answer()'s own top check) before ever reaching
// asyncTasksByTurn. rendezvous.closeAsk is idempotent — a no-op once the turn is
// already removed, see its own javadoc — so the shared finally below re-running it
// for this same fresh call is harmless.
rendezvous.closeAsk(ticket.turnId());
if (askTimeoutRaceHookForTest != null) {
// Test-only (fleetd #334): see the field's own javadoc.
askTimeoutRaceHookForTest.run();
}
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it
// out of asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and
// stays (it is what keeps the target from staying BUSY forever), but it would
// otherwise also erase askAnsweredAsyncTasks' only signal that the worker's eventual
// real fleet_reply still belongs to this task, stranding it in the inbox with a false
// "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
}
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -1046,19 +1098,28 @@ public final class MessageService {
// the chained ask deliberately left open.
//
// A null task is NOT only "this was never an async ticket". That reading was in this
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
// is still answerable but the map entry is already gone, so the lookup at :991
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
// firing only that first half before answer() runs printed
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
// method's own waiter is live. So #329 narrows this window rather than closing it.
// Open as fleetd #334 — do not read this guard as complete.
// comment when #329 merged and it is wrong; it is still not the whole story after
// #334. A genuine async ticket can still land here with task == null — a blocking
// (wait:true) send's fleet_ask never has a Task at all, so that case is expected and
// fine. What #334 fixed was a SECOND, unintended way to get here with task == null:
// ask()'s timeout path used to run clearAsyncQuestion(turnId, true) — which drops the
// asyncTasksByTurn entry — in its catch block, while rendezvous.closeAsk(turnId) ran
// later, in its finally. Between those two the ask was still answerable but the map
// entry was already gone, so the lookup at :1053 returned null and this ticket was
// never completed. Measured on 2026-09-04: a probe firing only that first half before
// answer() ran printed "answer=REPLIED phase=PENDING reply=null" — the same stranded
// ticket #329 set out to fix, one step earlier in the same race; the probe used
// forgetTurnForTest, which omits ask()'s markAskTimedOut, and that omission does not
// change the outcome, because askTimedOut is read only by askAnsweredAsyncTasks, and
// reply() never reaches it while this method's own waiter is live. #334's fix
// reorders ask()'s timeout catch to run closeAsk before the forgetting (see the
// fresh-owner block there), which removes this path entirely rather than narrowing
// it further: once closeAsk has run, rendezvous.askSession(turnId) is null and
// answer() returns STALE_TURN from its own top check, before it ever reaches this
// lookup — see aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket in
// MessageServiceTest, which pins the exact window with askTimeoutRaceHookForTest.
// So by the time this line runs, task == null means only the ordinary blocking-send
// case (or #329's own already-fixed race elsewhere) — not this one.
if (result.outcome() != Outcome.QUESTION && task != null) {
finishAsyncTask(task, result);
}
@@ -1118,7 +1179,18 @@ public final class MessageService {
// class javadoc on sendAsync/CB-107.
task.future.whenComplete((reply, ex) -> {
boolean failed = ex != null || reply == null || !reply.completed();
pushLoop.onTicketTerminal(ticket, target, failed);
try {
pushLoop.onTicketTerminal(ticket, target, failed);
} catch (Throwable t) {
// fleetd #335 (site 2): this stage's own CompletableFuture is discarded, so an
// uncaught throw here (e.g. a RejectedExecutionException from
// ReplyPushLoop's scheduler, already shut down while this in-flight send's
// whenComplete fires during the daemon's own shutdown sequence — messages.close()
// only stops accepting NEW async work, it does not cancel a delivery already
// running) vanishes with no log line and no metric, and the push loop never learns
// the ticket went terminal — the exact thing this hook exists to tell it.
log.error("push loop failed to learn ticket {} (target {}) went terminal", ticket, target, t);
}
});
}
asyncExecutor.submit(() -> {
@@ -1370,6 +1442,24 @@ public final class MessageService {
this.afterFinishAsyncTaskCompleteHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after
* {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel}
* takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued
* delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}.
* This deterministically covers the caller's need to use that result rather than relying on a
* timing-sensitive real race.
*/
private volatile Runnable timeoutCancellationRaceHookForTest;
/**
* Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setTimeoutCancellationRaceHookForTest(Runnable hook) {
this.timeoutCancellationRaceHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
@@ -1399,6 +1489,56 @@ public final class MessageService {
clearAsyncQuestion(turnId, true);
}
/**
* Null in production; test seam for fleetd #335 (site 1) — invoked from {@link #abandon(String,
* String, boolean)}'s per-task loop, once per task, right before that task's own cleanup
* (turnId bookkeeping, or the stranded-reply put-back) runs. A test installs this to inject a
* throw at that exact point deterministically.
*
* <p>The one production call there that can really throw is {@code inbox.publish} in the
* put-back branch — {@link AmqpReplyInbox#publish} reaches a broker and throws {@link
* IllegalStateException} on an unroutable, unconfirmed, or interrupted publish — but reaching
* that branch requires a second completion of the very task {@code abandon} is about to
* complete to win the race first (see the branch's own comment), and the #137 follow-up
* investigation above already found the combination this needs (a stranded reply coinciding
* with an open matching task) unreachable through the public API, not merely hard to time.
* This hook reproduces the resulting shape — a per-task cleanup throw — directly, the same
* technique {@link #finishAsyncTaskRaceHook} and {@link
* #afterFinishAsyncTaskCompleteHookForTest} already use for their own hard-to-time races.
*/
private volatile Runnable abandonCleanupHookForTest;
/**
* Test-only (fleetd #335, site 1): install {@link #abandonCleanupHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setAbandonCleanupHookForTest(Runnable hook) {
this.abandonCleanupHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #334 — invoked from {@link #ask}'s {@code
* TimeoutException} catch, only for the fresh owner, right after {@code rendezvous.closeAsk}
* has run and before {@link #markAskTimedOut} / {@link #clearAsyncQuestion} forget this task's
* turnId mapping. A test installs this to call {@link #answer} for the very same {@code turnId}
* synchronously from inside that exact window, deterministically reproducing the race a real
* concurrent {@code answer()} call could otherwise only win by timing luck: with the ask already
* closed, that call must see {@code rendezvous.askSession(turnId) == null} and return {@link
* Outcome#STALE_TURN} immediately, never reaching {@code asyncTasksByTurn} at all — proving the
* window fleetd #334 describes (mapping forgotten while the ask was still "answerable") is
* closed, rather than merely narrowed the way fleetd #329 narrowed the sibling race in {@link
* #finishAsyncTask}.
*/
private volatile Runnable askTimeoutRaceHookForTest;
/**
* Test-only (fleetd #334): install {@link #askTimeoutRaceHookForTest}. Package-private so the
* test, in the same package, can reach it without widening any production API.
*/
void setAskTimeoutRaceHookForTest(Runnable hook) {
this.askTimeoutRaceHookForTest = hook;
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -0,0 +1,93 @@
package dev.ltms.fleet.power;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.Locale;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* Holds macOS idle sleep off by keeping a {@code caffeinate -i} child process alive for the life
* of the returned {@link SleepAssertion}.
*
* <p>{@code -i} asserts only against <em>idle</em> sleep — it does not stop the lid closing or an
* operator-requested sleep from taking effect. That is deliberate: this class exists to stop an
* unattended host from sleeping out from under a member's long turn, never to override the
* operator. {@code -s}/{@code -d} (which also block system/display sleep on demand) are
* intentionally not used here.
*
* <p>{@link #acquire()} never throws. It returns {@code null} — a no-op — off macOS, and again if
* starting the {@code caffeinate} child fails for any reason (binary missing, process table full,
* …); either case is logged once at INFO, not on every occurrence, so a daemon that runs for
* weeks with the tool unavailable does not fill its log.
*/
public final class CaffeinateSleepAssertionMechanism implements SleepAssertionMechanism {
private static final Logger log = LoggerFactory.getLogger(CaffeinateSleepAssertionMechanism.class);
private final AtomicBoolean loggedOnce = new AtomicBoolean(false);
/** {@code true} when running on macOS, the only platform {@code caffeinate} ships on. */
public static boolean isSupportedPlatform() {
return isSupportedPlatform(System.getProperty("os.name"));
}
/** Package-visible so a test can drive the platform check without touching a real property. */
static boolean isSupportedPlatform(String osName) {
return osName != null && osName.toLowerCase(Locale.ROOT).contains("mac");
}
@Override
public SleepAssertion acquire() {
if (!isSupportedPlatform()) {
logOnce("not running on macOS (os.name={}); the idle-sleep guard is a no-op on this platform",
System.getProperty("os.name"));
return null;
}
try {
Process process = new ProcessBuilder("caffeinate", "-i")
.redirectOutput(ProcessBuilder.Redirect.DISCARD)
.redirectError(ProcessBuilder.Redirect.DISCARD)
.start();
return new CaffeinateAssertion(process);
} catch (IOException | RuntimeException e) {
logOnce("could not start 'caffeinate -i' ({}); the host may idle-sleep while members are live",
e.toString());
return null;
}
}
private void logOnce(String format, Object arg) {
if (loggedOnce.compareAndSet(false, true)) {
log.info("idle-sleep guard: " + format, arg);
}
}
/** Wraps the live {@code caffeinate} child; {@link #close} force-destroys it, idempotently. */
private static final class CaffeinateAssertion implements SleepAssertion {
private final Process process;
CaffeinateAssertion(Process process) {
this.process = process;
}
@Override
public void close() {
if (!process.isAlive()) {
return;
}
process.destroy();
try {
if (!process.waitFor(2, TimeUnit.SECONDS)) {
process.destroyForcibly();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
process.destroyForcibly();
}
}
}
}
@@ -0,0 +1,105 @@
package dev.ltms.fleet.power;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.function.IntSupplier;
/**
* Holds an OS-level assertion against idle sleep for exactly as long as at least one fleet
* member is live.
*
* <p><strong>Why this exists:</strong> a fleetd host was measured idle-sleeping after as little
* as one minute of inactivity (its {@code pmset -g custom} reports {@code sleep 1} on battery).
* Overnight the daemon's AMQP link to the broker dropped 13 times, and cross-checking every drop
* minute against {@code pmset -g log} found a sleep or wake event in the same minute or the one
* before, every time. The AMQP churn is only the visible symptom — the real problem is that a
* member mid-turn freezes with the host, and a long turn with nobody typing is exactly the case
* that goes idle.
*
* <p><strong>How it tracks "live":</strong> this is driven by {@code SessionManager}'s existing
* {@code onAcquire}/{@code onRelease} lifecycle hooks (added for CB-520/CB-516, previously wired
* to nothing but the reply inbox) rather than a second member count kept in parallel. Wire it as:
* <pre>{@code
* IdleSleepGuard guard = new IdleSleepGuard(mechanism, sessions::size);
* sessions.onAcquire(_ -> guard.recheck());
* sessions.onRelease(_ -> guard.recheck());
* }</pre>
* Every acquire/release event re-reads {@code SessionManager#size()} — the same registry {@code
* fleet_list}'s live/capacity numbers are themselves computed from — and only an actual 0→1 or
* 1→0 crossing touches the OS. A listener exception is already caught and logged by {@code
* SessionManager} itself (it must never let a listener failure block the acquire/release it is
* reacting to), so {@link #recheck()} does not need its own top-level try/catch to honor that.
*
* <p><strong>Failure posture:</strong> every method here is safe to call whether or not {@link
* SleepAssertionMechanism#acquire()} actually works. A mechanism that returns {@code null} (wrong
* platform, missing tool, spawn failure) simply means this guard never holds anything — it never
* throws and never blocks a spawn, a release, or shutdown.
*/
public final class IdleSleepGuard implements AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(IdleSleepGuard.class);
private final SleepAssertionMechanism mechanism;
private final IntSupplier liveCount;
private final Object lock = new Object();
private SleepAssertion held;
public IdleSleepGuard(SleepAssertionMechanism mechanism, IntSupplier liveCount) {
this.mechanism = mechanism;
this.liveCount = liveCount;
}
/**
* Re-read the live count and acquire or release the held assertion to match: nothing held and
* at least one member live ⇒ acquire; something held and no member live ⇒ release. A steady
* count (still zero, still positive) is a no-op either way, so a single spawn or release only
* ever touches the OS on the crossing, not on every call.
*/
public void recheck() {
synchronized (lock) {
int live = liveCount.getAsInt();
if (live > 0 && held == null) {
held = mechanism.acquire();
if (held != null) {
log.debug("idle-sleep guard armed: {} live member(s)", live);
}
} else if (live == 0 && held != null) {
releaseHeldLocked();
}
}
}
/** {@code true} while an assertion is actually held. Exposed for tests. */
boolean isHeld() {
synchronized (lock) {
return held != null;
}
}
/**
* Release whatever is held, if anything. Idempotent and safe to call at any time, including
* repeatedly — a daemon shutdown hook calls this unconditionally so no assertion (and no
* {@code caffeinate} child) survives the process, even if the drain that would otherwise have
* driven the live count to zero was itself interrupted or threw.
*/
@Override
public void close() {
synchronized (lock) {
if (held != null) {
releaseHeldLocked();
}
}
}
/** Caller must hold {@link #lock}. */
private void releaseHeldLocked() {
try {
held.close();
} catch (RuntimeException e) {
log.warn("idle-sleep guard: failed to release its assertion cleanly: {}", e.toString());
} finally {
held = null;
}
}
}
@@ -0,0 +1,11 @@
package dev.ltms.fleet.power;
/**
* A held OS-level assertion against idle sleep. {@link #close} must be idempotent — safe to call
* more than once — and must never throw, matching {@link IdleSleepGuard}'s "never break the
* fleet" contract.
*/
public interface SleepAssertion extends AutoCloseable {
@Override
void close();
}
@@ -0,0 +1,21 @@
package dev.ltms.fleet.power;
/**
* The OS mechanism {@link IdleSleepGuard} uses to hold and release an idle-sleep assertion. This
* is the seam a test exercises instead of the real effect (a live {@code caffeinate} child) — see
* {@code IdleSleepGuardTest}.
*
* <p>Implementations must never throw. Every failure — wrong platform, missing tool, a spawn
* error — must show up as {@link #acquire()} returning {@code null}, so a caller can treat "no
* assertion held" and "the mechanism could not be used" identically and the fleet keeps running
* either way.
*/
public interface SleepAssertionMechanism {
/**
* Acquire a fresh assertion against idle sleep, or {@code null} when this mechanism is not
* usable right now (wrong platform, the tool is missing, the child process could not start).
* Never throws.
*/
SleepAssertion acquire();
}
@@ -107,6 +107,7 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
v.put("worktreeGroup", "group-a");
v.put("memberLoginShell", null);
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
assertNamesMatchComponents(v);
return v;
}
@@ -147,6 +148,7 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
v.put("worktreeGroup", "group-b");
v.put("memberLoginShell", null);
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(false));
assertNamesMatchComponents(v);
return v;
}
@@ -2580,4 +2580,58 @@ class FleetConfigTest {
"with no pool to choose from, every configured profile is a candidate and the "
+ "first one wins");
}
// ── idle-sleep guard: default-on config block ───────────────────────────────────────────────
@Test
void idleSleepGuardIsOnByDefaultWhenTheBlockIsEntirelyAbsent(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8080
""");
FleetConfig cfg = FleetConfig.load(f);
assertNull(cfg.idleSleepGuard(), "an absent block parses to null, unlike most other blocks here");
// The block itself is absent, but the FEATURE stays on: Fleetd treats a null block the
// same as enabled: true (see FleetConfig.idleSleepGuard's javadoc) — this test only pins
// the parse result, the on-by-default behaviour is Fleetd's own null check.
}
@Test
void idleSleepGuardExplicitlyEnabledIsOn(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
idleSleepGuard:
enabled: true
""");
FleetConfig cfg = FleetConfig.load(f);
assertTrue(cfg.idleSleepGuard().isEnabled());
}
@Test
void idleSleepGuardExplicitlyDisabledIsOff(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
idleSleepGuard:
enabled: false
""");
FleetConfig cfg = FleetConfig.load(f);
assertFalse(cfg.idleSleepGuard().isEnabled());
}
@Test
void idleSleepGuardBlockPresentButEmptyDefaultsToEnabled(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
idleSleepGuard: {}
""");
FleetConfig cfg = FleetConfig.load(f);
assertTrue(cfg.idleSleepGuard().isEnabled(),
"unlike ConfigReload/Health, this block defaults to ON even when present but empty");
}
}
@@ -484,6 +484,27 @@ class CompletionResolverTest {
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
@Test
void aNormalMemberReportMentioningTheExhaustionPatternDoesNotNotifyTheSink() {
String block = "⏺ I reviewed capacity handling. The usage limit has been reached means no more work can start.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"a matching report still fails the send as exhausted");
assertTrue(waiter.getNow(null).text().contains("I reviewed capacity handling."),
"the exhausted result keeps the whole matched pane line");
assertTrue(notified.isEmpty(),
"a normal report mentioning an exhaustion pattern must not quarantine a credential");
}
@Test
void classifiesAMatchingScrapeAsBackendExhaustedInsteadOfACompletedReply() {
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
@@ -534,6 +555,53 @@ class CompletionResolverTest {
"the sink is told the matched reason: " + notified.get(0));
}
@Test
void aRealExhaustionBehindTerminalChromeStillNotifiesTheSink() {
String block = "⏺ │ The usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"a real exhaustion must still fail the send as exhausted");
assertEquals(1, notified.size(),
"a real exhaustion behind terminal chrome must reach the sink");
}
/**
* The live fleet configures {@code exhaustedPattern: "The usage limit has been reached"} — with
* the leading {@code "The"}. Every other test here uses a pattern without it, which is the shape
* that made fleetd #348 need a looser rule than a start-of-line check. This pins the deployed
* shape as well, so a later tightening of {@link CompletionResolver} cannot silently stop
* recording the exhaustion this fleet actually reports.
*
* <p>What it does not prove: that this is the only pattern shape an operator will write.
*/
@Test
void anExhaustionPatternCarryingItsLeadingWordsStillNotifiesTheSink() {
String block = "⏺ │ The usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("The usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"the live pattern shape must still fail the send as exhausted");
assertEquals(1, notified.size(),
"the live pattern shape must still reach the sink");
}
@Test
void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() {
// The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed —
@@ -377,6 +377,126 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
}
/**
* fleetd #334. {@code ask()}'s {@code TimeoutException} catch used to forget this task's
* {@code turnId} mapping ({@code clearAsyncQuestion(turnId, true)}) BEFORE closing the ask
* ({@code rendezvous.closeAsk}, in the shared {@code finally}). A primary's {@code answer()}
* call racing that exact window found the ask still "answerable" ({@code
* rendezvous.askSession(turnId)} still non-null) while the {@code Task} was already forgotten,
* so its {@code task != null} guard skipped the completion and the async ticket sat at
* {@code PENDING} forever even though {@code answer()} itself reported a result. The fix
* (closing the ask first) makes this window impossible: a racing {@code answer()} call either
* still finds the ask open (and the {@code Task} mapping guaranteed intact) or finds it already
* closed (and bails {@code STALE_TURN} before ever touching the {@code Task}). This test pins
* the exact window with {@code askTimeoutRaceHookForTest} and proves both invariants the ticket
* named: (1) a late/racing {@code answer()} sees the ask as already lapsed ({@code STALE_TURN}),
* never made answerable again, and (2) the async ticket still resolves {@code DONE} once the
* worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}.
*/
/**
* fleetd #334 gated the ask-timeout teardown on {@code ticket.fresh()}, matching the {@code
* finally} block that already did. This pins that gate. A coalesced duplicate passes its own
* {@code timeoutMillis}, which says nothing about whether the shared ask is done — so a
* duplicate timing out first must leave the fresh owner's still-open ask answerable.
*
* <p>Measured on merge: without this test, removing the {@code ticket.fresh()} gate left all
* 1371 tests green. The gate shipped with the reorder and nothing held it there.
*
* <p>What this does not prove: anything about the ordering inside the gate — that is
* {@code aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket}'s job.
*/
@Test
void aCoalescedDuplicateAskTimingOutLeavesTheFreshOwnersAskOpen() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
CompletableFuture<MessageService.AskResult> fresh =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = null;
long deadline = System.currentTimeMillis() + 2000;
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
&& System.currentTimeMillis() < deadline) {
asking = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(asking, "the fresh owner's question must surface before the duplicate asks");
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
// A coalesced duplicate on the same session, with its own much shorter timeout.
MessageService.AskResult duplicate = messages.ask(T, "which config file?", 100);
assertEquals(MessageService.AskOutcome.TIMED_OUT, duplicate.outcome(),
"the duplicate's own timeout elapses first");
CompletableFuture<MessageService.Reply> answered =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500));
MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(),
"a duplicate's timeout must not lapse the ask the fresh owner still holds");
assertEquals("fleetd.yaml", a.answer());
assertFalse(answered.get(5, TimeUnit.SECONDS).outcome() == MessageService.Outcome.STALE_TURN,
"the answer must not be rejected as stale");
}
@Test
void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 200));
// Wait for the question to actually surface (poll sees ASKING) before racing the timeout.
MessageService.TaskView asking = null;
long deadline = System.currentTimeMillis() + 2000;
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
&& System.currentTimeMillis() < deadline) {
asking = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(asking, "the question must surface before the ask times out");
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
messages.setAskTimeoutRaceHookForTest(() ->
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
try {
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
MessageService.Reply late = lateAnswer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.STALE_TURN, late.outcome(),
"a late answer racing the timeout teardown must see the ask as already lapsed");
// The worker resumes on its own (per the ask() contract) and eventually sends its real
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
} finally {
messages.setAskTimeoutRaceHookForTest(null);
}
MessageService.TaskView done = null;
deadline = System.currentTimeMillis() + 2000;
while ((done == null || done.phase() == MessageService.Phase.PENDING)
&& System.currentTimeMillis() < deadline) {
done = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(done);
assertEquals(MessageService.Phase.DONE, done.phase(), "the async ticket must not be stranded PENDING");
assertEquals("real result", done.reply());
}
@Test
void answeringAnUnknownTurnIsStale() {
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
@@ -410,6 +530,29 @@ class MessageServiceTest {
"a delivered send whose worker never replies times out as still working");
}
/**
* fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send}
* first observes its completion as incomplete, but before {@code cancel} takes the target monitor.
* The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report
* {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed.
*
* <p>What this does not prove: that this precise interleaving happens by itself under production
* timing. The test forces it through a test-only hook; it proves the timeout caller handles the
* injector result when the interleaving occurs.
*/
@Test
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
try {
MessageService.Reply reply = messages.send(T, "race delivery", 50);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
"cancel reporting DELIVERED means the worker received the timed-out message");
} finally {
messages.setTimeoutCancellationRaceHookForTest(null);
}
}
@Test
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
@@ -785,6 +928,49 @@ class MessageServiceTest {
assertFailedTicket(third, "agent target term_a not found");
}
// --- fleetd #335 (site 1): a per-task cleanup failure inside the abandon() loop must not -----
// strand the tasks that come after it. abandon()'s own comment on the loop documents the one
// real production call that can throw there (inbox.publish, in the stranded-reply put-back
// branch, reached when a concurrent reply() or a second abandon() races this one) — but
// reaching that branch requires the exact combination the #137 follow-up above already found
// unreachable through the public API. abandonCleanupHookForTest reproduces the resulting SHAPE
// (one task's cleanup throws) directly instead, the same technique this file already uses for
// fleetd #324/#329's own hard-to-time races.
@Test
void aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try {
String first = messages.sendAsync(T, "first task");
awaitWaiting(); // first task owns the target lock and rendezvous waiter
String second = messages.sendAsync(T, "second task"); // parked on the same lock
String third = messages.sendAsync(T, "third task"); // parked too — the whole sweep must survive
java.util.concurrent.atomic.AtomicInteger calls = new java.util.concurrent.atomic.AtomicInteger();
messages.setAbandonCleanupHookForTest(() -> {
if (calls.getAndIncrement() == 0) {
throw new RuntimeException("PROBE-335-SITE1");
}
});
assertTrue(messages.abandon(T, "agent target term_a not found"));
// Every task in the loop still gets its own outcome — the one whose cleanup threw
// included — even though the loop had no way to know in advance which one that would be.
assertFailedTicket(first, "agent target term_a not found");
assertFailedTicket(second, "agent target term_a not found");
assertFailedTicket(third, "agent target term_a not found");
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getThrowableProxy() != null
&& "PROBE-335-SITE1".equals(e.getThrowableProxy().getMessage())),
"a per-task cleanup failure must still reach the log, not vanish silently");
} finally {
messages.setAbandonCleanupHookForTest(null);
detachMessageServiceLog(appender);
}
}
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
//
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
@@ -1472,6 +1658,42 @@ class MessageServiceTest {
}
}
// --- fleetd #335 (site 2): task.future.whenComplete's own returned stage is discarded, so an --
// uncaught throw from ReplyPushLoop.onTicketTerminal used to vanish with no log line and no
// metric. The real production trigger is the daemon's own shutdown sequence (Fleetd's shutdown
// hook): messages.close() only stops the async executor from taking NEW work — it does not
// cancel a send already in flight — while pushLoop.close() shuts its scheduler down immediately
// right after, so a ticket that completes in that narrow window has onTicketTerminal's own
// scheduler.schedule(...) throw a real RejectedExecutionException. Reproduced here by shutting
// the very same scheduler down before the ticket resolves — no test-only hook needed, this
// reachable path throws for real.
@Test
void aTicketTerminalPushFailureDoesNotVanishSilently() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try (var wiring = wireWithPushLoop(1, 50)) {
String ticket = wiring.service().sendAsync(T, "long task");
awaitWaiting();
wiring.scheduler().shutdownNow(); // simulate pushLoop.close() racing an in-flight send
injectDelivery();
assertTrue(rendezvous.resolve(T, "async result"));
// The ticket's own outcome must be unaffected by the swallowed exception — finishAsyncTask
// completes task.future before whenComplete's action (and thus onTicketTerminal) ever runs.
MessageService.TaskView done = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE);
assertEquals("async result", done.reply());
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getFormattedMessage().contains(ticket)
&& e.getThrowableProxy() != null
&& "java.util.concurrent.RejectedExecutionException"
.equals(e.getThrowableProxy().getClassName())),
"onTicketTerminal throwing must still reach the log, not vanish silently");
} finally {
detachMessageServiceLog(appender);
}
}
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
@Test
@@ -0,0 +1,54 @@
package dev.ltms.fleet.power;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Platform-detection unit tests for {@link CaffeinateSleepAssertionMechanism}.
*
* <p>This deliberately never calls {@link CaffeinateSleepAssertionMechanism#acquire()} itself —
* doing so on a real macOS machine would actually start a live {@code caffeinate} child and hold
* a real idle-sleep assertion, which the ticket this class exists for explicitly forbids testing
* with. Instead this exercises the pure {@code isSupportedPlatform(String)} predicate that
* {@code acquire()} consults before ever touching {@link ProcessBuilder} — so it proves the
* platform check itself is correct on any CI OS, but it does <strong>not</strong> prove that a
* real {@code caffeinate -i} spawn succeeds or that its child is torn down correctly; that half is
* exercised indirectly by {@link IdleSleepGuardTest} against a {@link FakeSleepAssertionMechanism}
* instead, which is the seam invariant 2/3 in the ticket call for.
*/
class CaffeinateSleepAssertionMechanismTest {
@Test
void macOsNamesAreSupported() {
assertTrue(CaffeinateSleepAssertionMechanism.isSupportedPlatform("Mac OS X"));
assertTrue(CaffeinateSleepAssertionMechanism.isSupportedPlatform("macOS"));
assertTrue(CaffeinateSleepAssertionMechanism.isSupportedPlatform("MAC OS X"));
}
@Test
void nonMacNamesAreNotSupported() {
assertFalse(CaffeinateSleepAssertionMechanism.isSupportedPlatform("Linux"));
assertFalse(CaffeinateSleepAssertionMechanism.isSupportedPlatform("Windows 11"));
}
@Test
void nullOsNameIsNotSupported() {
assertFalse(CaffeinateSleepAssertionMechanism.isSupportedPlatform(null));
}
/**
* The overload {@code isSupportedPlatform()} (no args) reads the JVM's real {@code os.name} —
* proves the wiring is live, without asserting a specific answer (this suite itself must pass
* on both macOS and Linux CI).
*/
@Test
void noArgOverloadReadsRealSystemProperty() {
boolean expected = CaffeinateSleepAssertionMechanism
.isSupportedPlatform(System.getProperty("os.name"));
boolean actual = CaffeinateSleepAssertionMechanism.isSupportedPlatform();
assertEquals(expected, actual);
}
}
@@ -0,0 +1,56 @@
package dev.ltms.fleet.power;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicInteger;
/**
* Recording fake {@link SleepAssertionMechanism} — the seam behind the real OS effect (a live
* {@code caffeinate} child process). No test in this package ever spawns that real process; every
* assertion here is against this fake's own call log instead.
*
* <p>Each acquired {@link FakeAssertion} records its own {@code close()} calls, and every
* acquired instance is kept in {@link #acquired} so a test can inspect all of them, including
* ones {@link IdleSleepGuard} has already released.
*/
final class FakeSleepAssertionMechanism implements SleepAssertionMechanism {
/** Every {@link FakeAssertion} this mechanism has ever handed out, in order. */
final CopyOnWriteArrayList<FakeAssertion> acquired = new CopyOnWriteArrayList<>();
private final AtomicInteger acquireCalls = new AtomicInteger();
private volatile boolean unavailable = false;
/** Make the next (and every subsequent) {@link #acquire()} return {@code null}, like a missing tool. */
void makeUnavailable() {
unavailable = true;
}
int acquireCallCount() {
return acquireCalls.get();
}
@Override
public SleepAssertion acquire() {
acquireCalls.incrementAndGet();
if (unavailable) {
return null;
}
FakeAssertion a = new FakeAssertion();
acquired.add(a);
return a;
}
/** A held fake assertion; records how many times {@code close()} was actually called. */
static final class FakeAssertion implements SleepAssertion {
private final AtomicInteger closeCalls = new AtomicInteger();
int closeCallCount() {
return closeCalls.get();
}
@Override
public void close() {
closeCalls.incrementAndGet();
}
}
}
@@ -0,0 +1,118 @@
package dev.ltms.fleet.power;
import org.junit.jupiter.api.Test;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* {@link IdleSleepGuard} against a {@link FakeSleepAssertionMechanism} — the seam that stands in
* for a real {@code caffeinate} child process. No test in this class ever spawns a real OS
* process or asserts against real idle sleep; every assertion is against the fake's call log
* (how many times {@code acquire()}/{@code close()} were actually called). That proves the
* <em>orchestration</em> — when the guard decides to hold or release an assertion, and that it
* never throws — but it does <strong>not</strong> prove that {@code caffeinate -i} itself
* actually stops macOS from idle-sleeping; that half is outside what a unit test can safely
* exercise (see {@link CaffeinateSleepAssertionMechanismTest}'s class doc).
*/
class IdleSleepGuardTest {
@Test
void acquiresOnZeroToOneAndReleasesOnOneToZero() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(0);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
assertFalse(guard.isHeld(), "nothing held before any member is live");
liveCount.set(1);
guard.recheck();
assertTrue(guard.isHeld(), "an assertion must be held once a member is live");
assertEquals(1, mechanism.acquired.size());
assertEquals(0, mechanism.acquired.get(0).closeCallCount());
liveCount.set(0);
guard.recheck();
assertFalse(guard.isHeld(), "the assertion must be released once the last member goes");
assertEquals(1, mechanism.acquired.get(0).closeCallCount(), "the SAME held assertion must be closed");
}
@Test
void steadyLiveCountDoesNotReacquireOrRerelease() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(2);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck(); // 0 -> 2 crossing: acquires
guard.recheck(); // still 2: must be a no-op
guard.recheck(); // still 2: must be a no-op
assertEquals(1, mechanism.acquireCallCount(), "only the crossing touches the mechanism");
liveCount.set(1); // 2 -> 1: still > 0, still a no-op
guard.recheck();
assertTrue(guard.isHeld());
assertEquals(0, mechanism.acquired.get(0).closeCallCount());
assertEquals(1, mechanism.acquireCallCount());
}
/**
* Invariant 2: a missing/unavailable mechanism must never throw, and the guard must simply
* hold nothing. {@link FakeSleepAssertionMechanism#makeUnavailable()} makes {@code acquire()}
* return {@code null}, exactly like {@link CaffeinateSleepAssertionMechanism} does off macOS
* or when the {@code caffeinate} binary is missing.
*/
@Test
void unavailableMechanismNeverThrowsAndHoldsNothing() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
mechanism.makeUnavailable();
AtomicInteger liveCount = new AtomicInteger(1);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck(); // must not throw
assertFalse(guard.isHeld(), "acquire() returned null, so nothing is held");
assertEquals(1, mechanism.acquireCallCount());
// still must not throw or leak on release, even though nothing was ever actually held
liveCount.set(0);
guard.recheck();
assertFalse(guard.isHeld());
guard.close(); // teardown with nothing held must also be a safe no-op
}
/**
* Invariant 3 (teardown). This is the test the mutation testing step removes the production
* release call to fail: with {@code releaseHeldLocked()} not invoked from {@link
* IdleSleepGuard#close()}, the held fake assertion's {@code close()} would never be called and
* this assertion would fail.
*/
@Test
void closeReleasesAHeldAssertionEvenWithoutAZeroCrossing() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(1);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck();
assertTrue(guard.isHeld());
guard.close();
assertFalse(guard.isHeld(), "close() must release whatever is held, independent of live count");
assertEquals(1, mechanism.acquired.get(0).closeCallCount());
}
@Test
void closeIsIdempotent() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(1);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck();
guard.close();
guard.close(); // must not throw, must not double-release
assertEquals(1, mechanism.acquired.get(0).closeCallCount());
}
}
@@ -0,0 +1,72 @@
package dev.ltms.fleet.power;
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 dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Proves the wiring {@code Fleetd.main} actually performs — {@code
* sessions.onAcquire(_ -> guard.recheck())} / {@code sessions.onRelease(_ -> guard.recheck())} —
* not just {@link IdleSleepGuard}'s own orchestration logic in isolation
* ({@link IdleSleepGuardTest} already covers that in isolation, which on its own would not catch
* a wiring gap — e.g. an {@code onAcquire} call typo'd to a no-op lambda, or the listener wired to
* the wrong SessionManager instance — see fleetd's own "a test on the seam does not prove the
* caller" lesson). This test builds a real {@link SessionManager} exactly as
* {@code SessionManagerTest} does (a {@link FakeHerdr}-backed {@link ClaudeCodeLauncher}, no live
* herdr process), wires it to an {@link IdleSleepGuard} the same two lines {@code Fleetd.main}
* uses, and drives real {@link SessionManager#acquire} / {@link SessionManager#release} calls.
*/
class IdleSleepGuardWiringTest {
private SessionManager sessionManager(FakeHerdr herdr) {
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 workers = 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(workers);
}
@Test
void acquiringAndReleasingRealSessionsDrivesTheGuardThroughTheSameWiringFleetdUses() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
IdleSleepGuard guard = new IdleSleepGuard(mechanism, sessions::size);
// The exact two lines Fleetd.main wires up.
sessions.onAcquire(_ -> guard.recheck());
sessions.onRelease(_ -> guard.recheck());
assertFalse(guard.isHeld(), "no member yet: nothing held");
MemberSession a = sessions.acquire("ltms-local", "/a", "/caller", "ownerA");
assertTrue(guard.isHeld(), "0 -> 1: the first live member must arm the guard");
MemberSession b = sessions.acquire("ltms-local", "/b", "/caller", "ownerB");
assertEquals(1, mechanism.acquireCallCount(), "2nd member: still just 1 live-to-2 step, no new acquire");
sessions.release(a.paneId());
assertTrue(guard.isHeld(), "one member still live: the guard must stay armed");
assertEquals(0, mechanism.acquired.get(0).closeCallCount());
sessions.release(b.paneId());
assertFalse(guard.isHeld(), "1 -> 0: the last member releasing must disarm the guard");
assertEquals(1, mechanism.acquired.get(0).closeCallCount());
}
}