Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 09159f2857 | |||
| 29cd1194c2 | |||
| 815e8f8b23 | |||
| 1e60ac0745 | |||
| 650a4c146b | |||
| 23f299e105 | |||
| dbf6fef0e9 | |||
| d292522d00 | |||
| 0241e0d3a8 | |||
| e4973eb8a4 |
@@ -110,14 +110,6 @@ 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,8 +55,6 @@ 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;
|
||||
@@ -254,26 +252,6 @@ 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
|
||||
@@ -720,11 +698,6 @@ 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,10 +37,6 @@ 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
|
||||
@@ -133,9 +129,8 @@ 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),
|
||||
* recounted again after {@code idleSleepGuard:} was added.</strong>
|
||||
* {@code FleetConfig} has 23 top-level record components: 5 cold, 12 deferred, 3 split, 3
|
||||
* <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
|
||||
* 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()}
|
||||
@@ -218,7 +213,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", "idleSleepGuard");
|
||||
"quarantineCooldownSeconds", "profiles");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
@@ -423,14 +418,6 @@ 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,11 +106,6 @@ 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(
|
||||
@@ -135,22 +130,7 @@ public record FleetConfig(
|
||||
MemberCredentials memberCredentials,
|
||||
Coordinator coordinator,
|
||||
String worktreeGroup,
|
||||
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);
|
||||
}
|
||||
String memberLoginShell) {
|
||||
|
||||
/** Back-compat form before the {@code memberLoginShell} key was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
@@ -161,7 +141,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, null);
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the {@code worktreeGroup} key was added. */
|
||||
@@ -1265,25 +1245,6 @@ 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).
|
||||
*
|
||||
@@ -1534,7 +1495,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", "idleSleepGuard");
|
||||
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -2211,14 +2172,9 @@ 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, idleSleepGuard);
|
||||
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -8,6 +8,7 @@ import com.rabbitmq.client.DeliverCallback;
|
||||
import com.rabbitmq.client.Recoverable;
|
||||
import com.rabbitmq.client.RecoveryListener;
|
||||
import com.rabbitmq.client.Return;
|
||||
import com.rabbitmq.client.impl.DefaultExceptionHandler;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -153,17 +154,22 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
|
||||
public static AmqpReplyInbox open(String uri, int prefetch) {
|
||||
try {
|
||||
ConnectionFactory factory = new ConnectionFactory();
|
||||
factory.setUri(uri);
|
||||
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
|
||||
factory.setAutomaticRecoveryEnabled(true);
|
||||
factory.setTopologyRecoveryEnabled(true);
|
||||
return new AmqpReplyInbox(factory.newConnection("fleetd-reply-inbox"), prefetch);
|
||||
return new AmqpReplyInbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.REPLY_INBOX), prefetch);
|
||||
} catch (Exception e) {
|
||||
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
|
||||
}
|
||||
}
|
||||
|
||||
static ConnectionFactory connectionFactory(String uri) throws Exception {
|
||||
ConnectionFactory factory = new ConnectionFactory();
|
||||
factory.setUri(uri);
|
||||
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
|
||||
factory.setAutomaticRecoveryEnabled(true);
|
||||
factory.setTopologyRecoveryEnabled(true);
|
||||
factory.setExceptionHandler(new AmqpConnectionFailureLogger(AmqpConnectionFailureLogger.REPLY_INBOX, log));
|
||||
return factory;
|
||||
}
|
||||
|
||||
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
|
||||
AmqpReplyInbox(Connection connection) {
|
||||
this(connection, DEFAULT_PREFETCH);
|
||||
@@ -571,3 +577,44 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Keeps RabbitMQ's forgiving exception behaviour while adding the connection identity that its
|
||||
* default logger drops. Package-private so both AMQP connections use the same two names.
|
||||
*/
|
||||
final class AmqpConnectionFailureLogger extends DefaultExceptionHandler {
|
||||
|
||||
static final String REPLY_INBOX = "fleetd-reply-inbox";
|
||||
static final String LEAD_MAILBOX = "fleetd-lead-mailbox";
|
||||
|
||||
private final String connectionName;
|
||||
private final Logger logger;
|
||||
|
||||
AmqpConnectionFailureLogger(String connectionName, Logger logger) {
|
||||
this.connectionName = connectionName;
|
||||
this.logger = logger;
|
||||
}
|
||||
|
||||
String connectionName() {
|
||||
return connectionName;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void log(String message, Throwable cause) {
|
||||
if (isSocketClosedOrConnectionReset(cause)) {
|
||||
logger.warn("AMQP connection {}: {} (Exception message: {})", connectionName, message, cause.getMessage());
|
||||
} else {
|
||||
logger.error("AMQP connection {}: {}", connectionName, message, cause);
|
||||
}
|
||||
}
|
||||
|
||||
private static boolean isSocketClosedOrConnectionReset(Throwable cause) {
|
||||
// Deliberate copy of ForgivingExceptionHandler's private static helper; check it on amqp-client upgrades.
|
||||
if (!(cause instanceof IOException)) {
|
||||
return false;
|
||||
}
|
||||
return "Connection reset".equals(cause.getMessage())
|
||||
|| "Socket closed".equals(cause.getMessage())
|
||||
|| "Connection reset by peer".equals(cause.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -122,17 +122,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
|
||||
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
|
||||
try {
|
||||
ConnectionFactory factory = new ConnectionFactory();
|
||||
factory.setUri(uri);
|
||||
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
|
||||
factory.setAutomaticRecoveryEnabled(true);
|
||||
factory.setTopologyRecoveryEnabled(true);
|
||||
return new LeadMailbox(factory.newConnection("fleetd-lead-mailbox"), selfCoordId, prefetch);
|
||||
return new LeadMailbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.LEAD_MAILBOX), selfCoordId, prefetch);
|
||||
} catch (Exception e) {
|
||||
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
|
||||
}
|
||||
}
|
||||
|
||||
static ConnectionFactory connectionFactory(String uri) throws Exception {
|
||||
ConnectionFactory factory = new ConnectionFactory();
|
||||
factory.setUri(uri);
|
||||
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
|
||||
factory.setAutomaticRecoveryEnabled(true);
|
||||
factory.setTopologyRecoveryEnabled(true);
|
||||
factory.setExceptionHandler(new AmqpConnectionFailureLogger(AmqpConnectionFailureLogger.LEAD_MAILBOX, log));
|
||||
return factory;
|
||||
}
|
||||
|
||||
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
|
||||
LeadMailbox(Connection connection, String selfCoordId) {
|
||||
this(connection, selfCoordId, DEFAULT_PREFETCH);
|
||||
|
||||
@@ -1,93 +0,0 @@
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,105 +0,0 @@
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,11 +0,0 @@
|
||||
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();
|
||||
}
|
||||
@@ -1,21 +0,0 @@
|
||||
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,7 +107,6 @@ 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;
|
||||
}
|
||||
@@ -148,7 +147,6 @@ 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,58 +2580,4 @@ 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");
|
||||
}
|
||||
}
|
||||
|
||||
+192
@@ -0,0 +1,192 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* Guards against a "defect factory" built into this file's own established pattern, found live
|
||||
* while building the (parked) idle-sleep-guard PR: every time a component is added to
|
||||
* {@link FleetConfig}, the record grows by one arg AND a new back-compat constructor is added at
|
||||
* the OLD arity, so existing callers keep compiling. That is correct and required — see the
|
||||
* constructor ladder just below the record header. But {@link #withDefaults()}'s own {@code return
|
||||
* new FleetConfig(...)} call sits in this same file, written at a literal argument count. The very
|
||||
* next time a component is added, the freshly-added back-compat constructor at the OLD arity
|
||||
* silently captures that stale call, because it is now a legal overload at that arg count too. It
|
||||
* compiles. Every other test passes, because nothing else exercises the new field. The new
|
||||
* component is defaulted away — {@code null}, or whatever that back-compat overload defaults it to
|
||||
* — on every {@link FleetConfig#load}. Measured, not theoretical: this exact sequence happened
|
||||
* live when the {@code idleSleepGuard} component was added on a sibling branch; it was caught only
|
||||
* because that branch's own new tests happened to assert on the new field's value.
|
||||
*
|
||||
* <p>This test proves the opposite property, and does it in a way that survives the next field
|
||||
* being added without being rewritten: reflectively enumerate {@link FleetConfig}'s own record
|
||||
* components (never a hardcoded count — the arity is exactly what changes over time), build one
|
||||
* config through the true canonical constructor with a real, distinctive, non-null value in EVERY
|
||||
* component (reusing the exact reflective-construction pattern
|
||||
* {@link ConfigRefTopLevelReportingCoverageTest} already established for this file:
|
||||
* {@code getDeclaredConstructor(exact record-component types)}, which resolves the canonical
|
||||
* constructor by its true shape, not by binding to whichever overload happens to match arg count —
|
||||
* the same way Jackson resolves it), call the real {@link FleetConfig#withDefaults()}, and assert
|
||||
* every one of those values survives unchanged.
|
||||
*
|
||||
* <p>Why this is a valid check for every component, not just some: {@link #withDefaults()}'s own
|
||||
* comments document that it only ever REPLACES a component when the incoming value is {@code null}
|
||||
* (or blank, for {@code placement}) — {@code broker}/{@code primary}/{@code leadHeartbeat}/
|
||||
* {@code configReload}/{@code coordinator}/{@code worktreeGroup}/{@code memberLoginShell} are left
|
||||
* as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
|
||||
* {@code fleet}/{@code quarantineCooldownSeconds}/{@code memberCredentials}/{@code placement} are
|
||||
* replaced only on null/blank input. A value that is never null or blank going in must therefore
|
||||
* never change coming out, for every current component. No exclusion is needed today.
|
||||
*
|
||||
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} exists anyway, kept deliberately empty and size-pinned
|
||||
* by {@link #exclusionListSizeIsPinned()}: a future component that {@code withDefaults()} is
|
||||
* <em>documented</em> to transform unconditionally (unlike every field today) would legitimately
|
||||
* need one. Pinning the size at 0 means growing that set to make a failure go away is itself a
|
||||
* visible diff to this test, not a silent one — a checker that can be silenced by adding to its
|
||||
* own escape hatch is not a checker.
|
||||
*/
|
||||
class FleetConfigWithDefaultsPreservesEveryComponentTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
|
||||
|
||||
/** See the class javadoc — deliberately empty today; grow it only with a matching justification. */
|
||||
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
|
||||
|
||||
/** One real, distinctive, non-null (non-blank where blankness would mean "unset") value per component. */
|
||||
private static Map<String, Object> baseValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("bind", new FleetConfig.Bind("127.0.0.1", 8765));
|
||||
v.put("herdrSocket", "~/.config/herdr/guard.sock");
|
||||
v.put("memberHerdrSocket", "~/.config/herdr/member-guard.sock");
|
||||
v.put("profiles", Map.of("sonnet", minimalProfile("sonnet")));
|
||||
v.put("guard", new FleetConfig.Guard(List.of("host-guard")));
|
||||
v.put("worktreeRoot", "/wt/guard");
|
||||
v.put("lifecycle", new FleetConfig.Lifecycle(300, 5, 30, true));
|
||||
v.put("spawnReadyTimeoutMs", 12_345);
|
||||
v.put("spawnReadyPollMs", 234);
|
||||
v.put("broker", new FleetConfig.Broker("amqp://guard", null, 7));
|
||||
v.put("primary", new FleetConfig.Primary("term-guard", 4, 4000));
|
||||
v.put("fleet", new FleetConfig.Fleet(
|
||||
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-guard", 1, null, 10,
|
||||
"claude", null, null, null)),
|
||||
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
|
||||
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(301, 61_000L, 4));
|
||||
v.put("health", new FleetConfig.Health(true, 31, 601, 61, null));
|
||||
v.put("placement", "round-robin");
|
||||
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
|
||||
v.put("configReload", new FleetConfig.ConfigReload(true, 11));
|
||||
v.put("quarantineCooldownSeconds", 1801);
|
||||
v.put("memberCredentials", new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
|
||||
List.of("git"), List.of("git", "ssh"), null));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3));
|
||||
v.put("worktreeGroup", "group-guard");
|
||||
v.put("memberLoginShell", "/bin/zsh");
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
/** A minimal, otherwise-null {@link FleetConfig.Profile} — just enough to name one in a map. */
|
||||
private static FleetConfig.Profile minimalProfile(String name) {
|
||||
return new FleetConfig.Profile(name, null, null, null, null, null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null, null, null, null,
|
||||
null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Guards {@link #baseValues()} itself against drifting from the record's real shape — the same
|
||||
* assurance {@link ConfigRefTopLevelReportingCoverageTest} already relies on. This is what makes
|
||||
* "no hardcoded arity" true in practice: forgetting to add a new component here fails this
|
||||
* assertion by name, rather than silently checking one component fewer than the record has.
|
||||
*/
|
||||
private static void assertNamesMatchComponents(Map<String, Object> values) {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
names.add(rc.getName());
|
||||
}
|
||||
assertEquals(names, new TreeSet<>(values.keySet()),
|
||||
"this test's value map has drifted from FleetConfig's actual top-level components — "
|
||||
+ "update baseValues() alongside the record");
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a {@link FleetConfig} through the TRUE canonical constructor — resolved by the record's
|
||||
* own component types, not by argument count — so this never accidentally exercises a
|
||||
* back-compat overload the way a literal {@code new FleetConfig(...)} call risks doing.
|
||||
*/
|
||||
private static FleetConfig configOf(Map<String, Object> values) throws ReflectiveOperationException {
|
||||
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
|
||||
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
|
||||
Constructor<FleetConfig> ctor = FleetConfig.class.getDeclaredConstructor(types);
|
||||
return ctor.newInstance(args);
|
||||
}
|
||||
|
||||
@Test
|
||||
void exclusionListSizeIsPinned() {
|
||||
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
|
||||
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
|
||||
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
|
||||
+ "growing exclusion list that silences failures on its own is not a guard");
|
||||
}
|
||||
|
||||
/**
|
||||
* The mutation this is built to catch: make {@code withDefaults()}'s final constructor call
|
||||
* literal at some arg count, add one more component to the record with a new back-compat
|
||||
* constructor at the old arity, and the stale call silently rebinds. Every component here is
|
||||
* real and non-null (non-blank for the one String — {@code placement} — where blank has
|
||||
* meaning), so none of it should be replaced by {@code withDefaults()}; any component that
|
||||
* comes back different was silently dropped.
|
||||
*/
|
||||
@Test
|
||||
void everyComponentGivenARealValueSurvivesWithDefaults() throws ReflectiveOperationException {
|
||||
Map<String, Object> base = baseValues();
|
||||
FleetConfig config = configOf(base);
|
||||
FleetConfig defaulted = config.withDefaults();
|
||||
|
||||
List<String> dropped = new ArrayList<>();
|
||||
int checked = 0;
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
String name = rc.getName();
|
||||
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
|
||||
continue;
|
||||
}
|
||||
checked++;
|
||||
Object expected = base.get(name);
|
||||
Object actual;
|
||||
try {
|
||||
actual = rc.getAccessor().invoke(defaulted);
|
||||
} catch (ReflectiveOperationException e) {
|
||||
throw new RuntimeException("failed to read FleetConfig." + name + "()", e);
|
||||
}
|
||||
if (!Objects.equals(expected, actual)) {
|
||||
dropped.add(String.format(Locale.ROOT,
|
||||
"%s: withDefaults() was given a real, non-null value (%s) for '%s' but "
|
||||
+ "returned %s — a component silently dropped by withDefaults(), the "
|
||||
+ "shape of the defect this test exists to catch (its final "
|
||||
+ "\"return new FleetConfig(...)\" call binding to a back-compat "
|
||||
+ "constructor instead of the true canonical one)",
|
||||
name, expected, name, actual));
|
||||
}
|
||||
}
|
||||
|
||||
System.out.printf(Locale.ROOT,
|
||||
"FleetConfig.withDefaults() component-survival coverage — %d components, %d checked, "
|
||||
+ "%d excluded, %d survived%n",
|
||||
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
|
||||
assertEquals(List.of(), dropped,
|
||||
"withDefaults() silently dropped these real, given components: " + dropped);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,134 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.ConnectionFactory;
|
||||
import com.rabbitmq.client.impl.DefaultExceptionHandler;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.Proxy;
|
||||
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.assertInstanceOf;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class AmqpConnectionFailureLoggerTest {
|
||||
|
||||
@Test
|
||||
void installedHandlersLogTheirOwnConnectionNamesAtErrorWithTheCause() throws Exception {
|
||||
ConnectionFactory inboxFactory = AmqpReplyInbox.connectionFactory("amqp://127.0.0.1");
|
||||
ConnectionFactory mailboxFactory = LeadMailbox.connectionFactory("amqp://127.0.0.1");
|
||||
|
||||
AmqpConnectionFailureLogger inboxHandler = installedStrictHandler(inboxFactory, "reply inbox");
|
||||
AmqpConnectionFailureLogger mailboxHandler = installedStrictHandler(mailboxFactory, "lead mailbox");
|
||||
assertEquals(AmqpConnectionFailureLogger.REPLY_INBOX, inboxHandler.connectionName());
|
||||
assertEquals(AmqpConnectionFailureLogger.LEAD_MAILBOX, mailboxHandler.connectionName());
|
||||
|
||||
ListAppender<ILoggingEvent> inboxEvents = attach(AmqpReplyInbox.class);
|
||||
ListAppender<ILoggingEvent> mailboxEvents = attach(LeadMailbox.class);
|
||||
IllegalStateException inboxFailure = new IllegalStateException("inbox failure");
|
||||
IllegalStateException mailboxFailure = new IllegalStateException("mailbox failure");
|
||||
try {
|
||||
inboxHandler.handleUnexpectedConnectionDriverException(null, inboxFailure);
|
||||
mailboxHandler.handleConnectionRecoveryException(null, mailboxFailure);
|
||||
|
||||
assertError(inboxEvents, "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
|
||||
inboxFailure, "inbox failure line");
|
||||
assertError(mailboxEvents, "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
|
||||
mailboxFailure, "mailbox recovery line");
|
||||
} finally {
|
||||
detach(AmqpReplyInbox.class, inboxEvents);
|
||||
detach(LeadMailbox.class, mailboxEvents);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void connectionResetKeepsForgivingHandlerWarningSemantics() {
|
||||
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
|
||||
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
|
||||
ListAppender<ILoggingEvent> events = attach(AmqpReplyInbox.class);
|
||||
try {
|
||||
handler.handleUnexpectedConnectionDriverException(null, new IOException("Connection reset"));
|
||||
assertEquals(1, events.list.size(), "the handler must still log a reset");
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel(), "ForgivingExceptionHandler logs connection resets at WARN");
|
||||
assertEquals("AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred "
|
||||
+ "(Exception message: Connection reset)", event.getFormattedMessage());
|
||||
assertTrue(event.getThrowableProxy() == null, "ForgivingExceptionHandler does not attach a reset stack trace");
|
||||
} finally {
|
||||
detach(AmqpReplyInbox.class, events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void connectionNamesStayDistinct() {
|
||||
assertNotEquals(AmqpConnectionFailureLogger.REPLY_INBOX, AmqpConnectionFailureLogger.LEAD_MAILBOX,
|
||||
"reply-inbox and lead-mailbox failures must be distinguishable");
|
||||
}
|
||||
|
||||
@Test
|
||||
void strictConsumerExceptionStillClosesItsChannel() {
|
||||
AtomicInteger closes = new AtomicInteger();
|
||||
Channel channel = (Channel) Proxy.newProxyInstance(getClass().getClassLoader(), new Class<?>[] {Channel.class},
|
||||
(_, method, _) -> switch (method.getName()) {
|
||||
case "close" -> {
|
||||
closes.incrementAndGet();
|
||||
yield null;
|
||||
}
|
||||
case "toString" -> "test-channel";
|
||||
default -> throw new UnsupportedOperationException(method.getName());
|
||||
});
|
||||
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
|
||||
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
|
||||
|
||||
handler.handleConsumerException(channel, new IllegalStateException("consumer failed"), null, "tag", "handleDelivery");
|
||||
|
||||
assertEquals(1, closes.get(), "DefaultExceptionHandler must close a channel after a consumer exception");
|
||||
}
|
||||
|
||||
@Test
|
||||
void handlerOnlyChangesDefaultHandlerLogging() {
|
||||
assertEquals(DefaultExceptionHandler.class,
|
||||
AmqpConnectionFailureLogger.class.getSuperclass());
|
||||
assertFalse(java.util.Arrays.stream(AmqpConnectionFailureLogger.class.getDeclaredMethods())
|
||||
.anyMatch(method -> method.getName().startsWith("handle")),
|
||||
"all exception-handling methods must remain inherited from DefaultExceptionHandler");
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach(Class<?> owner) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(owner);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(Class<?> owner, ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(owner)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static AmqpConnectionFailureLogger installedStrictHandler(ConnectionFactory factory, String connection) {
|
||||
assertInstanceOf(DefaultExceptionHandler.class, factory.getExceptionHandler(),
|
||||
connection + " must keep DefaultExceptionHandler: replacing the strict handler with a forgiving one "
|
||||
+ "changes when a channel is closed");
|
||||
return assertInstanceOf(AmqpConnectionFailureLogger.class, factory.getExceptionHandler());
|
||||
}
|
||||
|
||||
private static void assertError(ListAppender<ILoggingEvent> events, String message, Throwable cause, String name) {
|
||||
assertEquals(1, events.list.size(), name);
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
assertEquals(Level.ERROR, event.getLevel(), name);
|
||||
assertEquals(message, event.getFormattedMessage(), name);
|
||||
assertEquals(cause.toString(), event.getThrowableProxy().getClassName() + ": "
|
||||
+ event.getThrowableProxy().getMessage(), name);
|
||||
}
|
||||
}
|
||||
@@ -1,54 +0,0 @@
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -1,56 +0,0 @@
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,118 +0,0 @@
|
||||
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());
|
||||
}
|
||||
}
|
||||
@@ -1,72 +0,0 @@
|
||||
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());
|
||||
}
|
||||
}
|
||||
@@ -116,6 +116,50 @@ check_log_path_matches_plist() {
|
||||
ok "log path check: script and plist agree ($resolved_out)"
|
||||
}
|
||||
|
||||
# Classify ERROR lines in one fresh log region. AMQP failure messages now include the connection
|
||||
# name, so a recovery can clear only errors for its own connection. A candidate with neither name
|
||||
# remains unexplained: it must never be quieted by a recovery on the other connection.
|
||||
classify_amqp_connection_errors() {
|
||||
local log_file="$1" line pending_inbox=0 pending_lead_mailbox=0
|
||||
REDEPLOY_ERROR_COUNT=0
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=0
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=0
|
||||
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
case "$line" in
|
||||
*' ERROR '*|*' SEVERE '*)
|
||||
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
|
||||
case "$line" in
|
||||
*'AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred'*|*'AMQP connection fleetd-reply-inbox: Caught an exception during connection recovery!'*)
|
||||
pending_inbox=$((pending_inbox + 1))
|
||||
;;
|
||||
*'AMQP connection fleetd-lead-mailbox: An unexpected connection driver error occurred'*|*'AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!'*)
|
||||
pending_lead_mailbox=$((pending_lead_mailbox + 1))
|
||||
;;
|
||||
*'AMQP connection'*'An unexpected connection driver error occurred'*|*'AMQP connection'*'Caught an exception during connection recovery!'*)
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1))
|
||||
;;
|
||||
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
|
||||
esac
|
||||
;;
|
||||
*'AMQP connection recovered; cleared held replies for fresh redelivery'*)
|
||||
if [ "$pending_inbox" -gt 0 ]; then
|
||||
pending_inbox=$((pending_inbox - 1))
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
|
||||
fi
|
||||
;;
|
||||
*'AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery'*)
|
||||
if [ "$pending_lead_mailbox" -gt 0 ]; then
|
||||
pending_lead_mailbox=$((pending_lead_mailbox - 1))
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
|
||||
fi
|
||||
;;
|
||||
esac
|
||||
done < "$log_file"
|
||||
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + pending_inbox + pending_lead_mailbox))
|
||||
}
|
||||
|
||||
# CB-600: sourceable for testing. When this file is SOURCED (not executed) it stops here — nothing
|
||||
# below runs — so a test harness can `source` it to call check_log_path_matches_plist (or the
|
||||
# other pure helpers above) against a throwaway plist fixture without ever reaching the mutating
|
||||
@@ -361,15 +405,21 @@ tail -n "+$((RESTART_MARK + 1))" "$OUT" 2>/dev/null \
|
||||
| grep -iE 'deferred|classification:|fleet health:|coverage' | tail -8 | sed 's/^/ /' \
|
||||
|| echo " (nothing reported)"
|
||||
|
||||
# Errors since the restart, anchored to the marker so old noise cannot leak in.
|
||||
ERRS="$(tail -n "+$((RESTART_MARK + 1))" "$OUT" 2>/dev/null | grep -cE ' (ERROR|SEVERE) ' || true)"
|
||||
# Errors since the restart, anchored to the marker so old noise cannot leak in. Keep the fresh
|
||||
# region in a file because the classifier must preserve the order of errors and recoveries.
|
||||
FRESH_LOG="$(mktemp -t fleetd-fresh-log)"
|
||||
trap 'rm -f "$FRESH_LOG"' EXIT
|
||||
tail -n "+$((RESTART_MARK + 1))" "$OUT" > "$FRESH_LOG" 2>/dev/null || true
|
||||
classify_amqp_connection_errors "$FRESH_LOG"
|
||||
say "result"
|
||||
ok "pid $NEW_PID, jar $(jar_id)"
|
||||
if [ "${ERRS:-0}" -gt 0 ]; then
|
||||
warn "$ERRS ERROR lines since restart:"
|
||||
tail -n "+$((RESTART_MARK + 1))" "$OUT" | grep -E ' (ERROR|SEVERE) ' | tail -5 | sed 's/^/ /'
|
||||
else
|
||||
if [ "$REDEPLOY_ERROR_COUNT" -eq 0 ]; then
|
||||
ok "no ERROR lines since restart"
|
||||
elif [ "$REDEPLOY_UNEXPLAINED_ERRORS" -eq 0 ]; then
|
||||
ok "$REDEPLOY_RECOVERED_AMQP_ERRORS AMQP connection reset ERROR lines recovered since restart"
|
||||
else
|
||||
warn "$REDEPLOY_ERROR_COUNT ERROR lines since restart:"
|
||||
grep -E ' (ERROR|SEVERE) ' "$FRESH_LOG" | tail -5 | sed 's/^/ /'
|
||||
fi
|
||||
echo
|
||||
echo " Next: call fleet_whoami and confirm it still answers 'primary'. A lead whose tab label"
|
||||
|
||||
Executable
+238
@@ -0,0 +1,238 @@
|
||||
#!/usr/bin/env bash
|
||||
# Self-contained checks for the pure log classifier in redeploy-fleetd.sh.
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
|
||||
TMP="$(mktemp -d "$ROOT/.redeploy-log-test.XXXXXX")"
|
||||
trap 'rm -rf "$TMP"' EXIT
|
||||
|
||||
# Sourcing stops before redeploy-fleetd.sh can build, stop, or start the daemon.
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
|
||||
fail() {
|
||||
printf 'FAIL: %s\n' "$*" >&2
|
||||
return 1
|
||||
}
|
||||
|
||||
assert_equals() {
|
||||
local expected="$1" actual="$2" description="$3"
|
||||
[ "$expected" = "$actual" ] || fail "$description: expected $expected, got $actual"
|
||||
}
|
||||
|
||||
classify_fixture() {
|
||||
local name="$1"
|
||||
classify_amqp_connection_errors "$TMP/$name"
|
||||
}
|
||||
|
||||
test_no_errors() {
|
||||
cat > "$TMP/no-errors.log" <<'LOG'
|
||||
2026-09-05 12:00:00 INFO fleetd listening
|
||||
LOG
|
||||
classify_fixture no-errors.log
|
||||
assert_equals 0 "$REDEPLOY_ERROR_COUNT" "no-errors total"
|
||||
assert_equals 0 "$REDEPLOY_UNEXPLAINED_ERRORS" "no-errors unexplained"
|
||||
}
|
||||
|
||||
test_recovery_patterns_match_source() {
|
||||
grep -F 'AMQP connection {}: {}' "$ROOT/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java" > /dev/null \
|
||||
|| fail "AMQP failure pattern no longer matches source"
|
||||
grep -F 'AMQP connection recovered; cleared held replies for fresh redelivery' \
|
||||
"$ROOT/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java" > /dev/null \
|
||||
|| fail "reply-inbox recovery pattern no longer matches source"
|
||||
grep -F 'AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery' \
|
||||
"$ROOT/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java" > /dev/null \
|
||||
|| fail "lead-mailbox recovery pattern no longer matches source"
|
||||
}
|
||||
|
||||
test_attributed_recovered_connection_error() {
|
||||
cat > "$TMP/attributed-recovered.log" <<'LOG'
|
||||
2026-09-05 12:00:00 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
|
||||
2026-09-05 12:00:01 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
|
||||
LOG
|
||||
classify_fixture attributed-recovered.log
|
||||
assert_equals 1 "$REDEPLOY_ERROR_COUNT" "attributed-recovered total"
|
||||
assert_equals 1 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "attributed-recovered errors"
|
||||
assert_equals 0 "$REDEPLOY_UNEXPLAINED_ERRORS" "attributed-recovered unexplained"
|
||||
}
|
||||
|
||||
test_source_derived_error_shapes_recover_by_connection() {
|
||||
# These ERROR shapes come from AmqpConnectionFailureLogger on main. They need a live-log check
|
||||
# after redeploy because the new code has not yet written a production line.
|
||||
cat > "$TMP/source-derived.log" <<'LOG'
|
||||
17:37:53.537 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
|
||||
17:37:54.537 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: Caught an exception during connection recovery!
|
||||
17:37:55.537 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
|
||||
17:37:56.537 ERROR [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP connection fleetd-lead-mailbox: An unexpected connection driver error occurred
|
||||
17:37:57.537 ERROR [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!
|
||||
17:37:58.537 ERROR [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP connection fleetd-lead-mailbox: An unexpected connection driver error occurred
|
||||
17:38:00.000 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
|
||||
17:38:01.000 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
|
||||
17:38:02.000 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
|
||||
17:38:03.000 INFO [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
17:38:04.000 INFO [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
17:38:05.000 INFO [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
LOG
|
||||
classify_fixture source-derived.log
|
||||
assert_equals 6 "$REDEPLOY_ERROR_COUNT" "source-derived total"
|
||||
assert_equals 6 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "source-derived recovered"
|
||||
assert_equals 0 "$REDEPLOY_UNEXPLAINED_ERRORS" "source-derived unexplained"
|
||||
}
|
||||
|
||||
test_cross_connection_unattributable_errors_stay_loud() {
|
||||
# This candidate has neither stable connection name, so LeadMailbox recovery must not consume it.
|
||||
cat > "$TMP/cross-unattributable.log" <<'LOG'
|
||||
2026-09-05 12:00:00 ERROR [AMQP Connection broker:5672] unknown - AMQP connection: An unexpected connection driver error occurred
|
||||
2026-09-05 12:00:01 ERROR [AMQP Connection broker:5672] unknown - AMQP connection: An unexpected connection driver error occurred
|
||||
2026-09-05 12:00:02 INFO [AMQP Connection 10.10.20.13:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
2026-09-05 12:00:03 INFO [AMQP Connection 10.10.20.13:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
LOG
|
||||
classify_fixture cross-unattributable.log
|
||||
assert_equals 2 "$REDEPLOY_ERROR_COUNT" "cross-unattributable total"
|
||||
assert_equals 0 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "cross-unattributable recovered"
|
||||
assert_equals 2 "$REDEPLOY_UNEXPLAINED_ERRORS" "cross-unattributable unexplained"
|
||||
}
|
||||
|
||||
test_attributed_cross_connection_errors_stay_loud() {
|
||||
# LeadMailbox recovery cannot heal AmqpReplyInbox errors.
|
||||
cat > "$TMP/cross-attributed.log" <<'LOG'
|
||||
2026-09-05 12:00:00 ERROR AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
|
||||
2026-09-05 12:00:01 ERROR AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
|
||||
2026-09-05 12:00:02 INFO LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
2026-09-05 12:00:03 INFO LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
|
||||
LOG
|
||||
classify_fixture cross-attributed.log
|
||||
assert_equals 2 "$REDEPLOY_ERROR_COUNT" "cross-attributed total"
|
||||
assert_equals 0 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "cross-attributed recovered"
|
||||
assert_equals 2 "$REDEPLOY_UNEXPLAINED_ERRORS" "cross-attributed unexplained"
|
||||
}
|
||||
|
||||
test_attributed_unrecovered_connection_error() {
|
||||
cat > "$TMP/unrecovered.log" <<'LOG'
|
||||
2026-09-05 12:00:00 ERROR AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
|
||||
LOG
|
||||
classify_fixture unrecovered.log
|
||||
assert_equals 1 "$REDEPLOY_ERROR_COUNT" "unrecovered total"
|
||||
assert_equals 0 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "unrecovered AMQP errors"
|
||||
assert_equals 1 "$REDEPLOY_UNEXPLAINED_ERRORS" "unrecovered unexplained"
|
||||
}
|
||||
|
||||
test_other_error_is_unexplained() {
|
||||
cat > "$TMP/other-error.log" <<'LOG'
|
||||
2026-09-05 12:00:00 ERROR dev.ltms.fleet.Fleetd - startup failed
|
||||
2026-09-05 12:00:01 INFO dev.ltms.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
|
||||
LOG
|
||||
classify_fixture other-error.log
|
||||
assert_equals 1 "$REDEPLOY_ERROR_COUNT" "other-error total"
|
||||
assert_equals 1 "$REDEPLOY_UNEXPLAINED_ERRORS" "other-error unexplained"
|
||||
}
|
||||
|
||||
test_recovery_requirement_mutation_is_caught() {
|
||||
classify_amqp_connection_errors() {
|
||||
local log_file="$1" line
|
||||
REDEPLOY_ERROR_COUNT=0
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=0
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=0
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
case "$line" in
|
||||
*' ERROR '*|*' SEVERE '*)
|
||||
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
|
||||
case "$line" in
|
||||
*'AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred'*)
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
|
||||
;;
|
||||
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
|
||||
esac
|
||||
;;
|
||||
esac
|
||||
done < "$log_file"
|
||||
}
|
||||
|
||||
if test_attributed_unrecovered_connection_error > "$TMP/mutation-output" 2>&1; then
|
||||
fail "mutation accepted an unrecovered connection error"
|
||||
fi
|
||||
grep -F 'FAIL: unrecovered AMQP errors: expected 0, got 1' "$TMP/mutation-output" > /dev/null \
|
||||
|| fail "mutation failed without the expected assertion"
|
||||
printf 'Recovery mutation: FAIL: unrecovered AMQP errors: expected 0, got 1\n'
|
||||
}
|
||||
|
||||
test_shared_counter_mutation_is_caught() {
|
||||
classify_amqp_connection_errors() {
|
||||
local log_file="$1" line pending=0
|
||||
REDEPLOY_ERROR_COUNT=0
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=0
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=0
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
case "$line" in
|
||||
*' ERROR '*|*' SEVERE '*)
|
||||
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
|
||||
case "$line" in
|
||||
*'AMQP connection'*'An unexpected connection driver error occurred'*|*'AMQP connection'*'Caught an exception during connection recovery!'*)
|
||||
case "$line" in
|
||||
*'fleetd-reply-inbox'*|*'fleetd-lead-mailbox'*) pending=$((pending + 1)) ;;
|
||||
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
|
||||
esac
|
||||
;;
|
||||
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
|
||||
esac
|
||||
;;
|
||||
*'AMQP connection recovered; cleared held replies for fresh redelivery'*|*'AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery'*)
|
||||
if [ "$pending" -gt 0 ]; then
|
||||
pending=$((pending - 1))
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
|
||||
fi
|
||||
;;
|
||||
esac
|
||||
done < "$log_file"
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + pending))
|
||||
}
|
||||
|
||||
if test_attributed_cross_connection_errors_stay_loud > "$TMP/shared-mutation-output" 2>&1; then
|
||||
fail "shared counter mutation accepted cross-connection recovery"
|
||||
fi
|
||||
grep -F 'FAIL: cross-attributed recovered: expected 0, got 2' "$TMP/shared-mutation-output" > /dev/null \
|
||||
|| fail "shared counter mutation failed without the expected assertion"
|
||||
printf 'Shared-counter mutation: FAIL: cross-attributed recovered: expected 0, got 2\n'
|
||||
}
|
||||
|
||||
test_unattributable_quiet_mutation_is_caught() {
|
||||
classify_amqp_connection_errors() {
|
||||
local log_file="$1" line
|
||||
REDEPLOY_ERROR_COUNT=0
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=0
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=0
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
case "$line" in
|
||||
*' ERROR '*|*' SEVERE '*)
|
||||
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
|
||||
case "$line" in
|
||||
*'AMQP connection'*'An unexpected connection driver error occurred'*|*'AMQP connection'*'Caught an exception during connection recovery!'*)
|
||||
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
|
||||
;;
|
||||
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
|
||||
esac
|
||||
;;
|
||||
esac
|
||||
done < "$log_file"
|
||||
}
|
||||
|
||||
if test_cross_connection_unattributable_errors_stay_loud > "$TMP/unattributable-mutation-output" 2>&1; then
|
||||
fail "unattributable mutation accepted an unknown connection"
|
||||
fi
|
||||
grep -F 'FAIL: cross-unattributable recovered: expected 0, got 2' "$TMP/unattributable-mutation-output" > /dev/null \
|
||||
|| fail "unattributable mutation failed without the expected assertion"
|
||||
printf 'Unattributable mutation: FAIL: cross-unattributable recovered: expected 0, got 2\n'
|
||||
}
|
||||
|
||||
test_no_errors
|
||||
test_recovery_patterns_match_source
|
||||
test_attributed_recovered_connection_error
|
||||
test_source_derived_error_shapes_recover_by_connection
|
||||
test_cross_connection_unattributable_errors_stay_loud
|
||||
test_attributed_cross_connection_errors_stay_loud
|
||||
test_attributed_unrecovered_connection_error
|
||||
test_other_error_is_unexplained
|
||||
test_recovery_requirement_mutation_is_caught
|
||||
test_shared_counter_mutation_is_caught
|
||||
test_unattributable_quiet_mutation_is_caught
|
||||
printf 'PASS: redeploy log classifier\n'
|
||||
Reference in New Issue
Block a user