Compare commits
15 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 09159f2857 | |||
| 29cd1194c2 | |||
| 815e8f8b23 | |||
| 1e60ac0745 | |||
| 650a4c146b | |||
| 23f299e105 | |||
| dbf6fef0e9 | |||
| d292522d00 | |||
| 0241e0d3a8 | |||
| e4973eb8a4 | |||
| b6b88c5f1c | |||
| 86dddfe240 | |||
| 4a5030a5c6 | |||
| c1c8794c48 | |||
| f429ca1a50 |
@@ -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.
|
||||
*
|
||||
|
||||
@@ -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);
|
||||
|
||||
+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);
|
||||
}
|
||||
}
|
||||
@@ -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 —
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -393,6 +393,56 @@ class MessageServiceTest {
|
||||
* 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");
|
||||
|
||||
@@ -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