Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 09159f2857 | |||
| 29cd1194c2 | |||
| 815e8f8b23 | |||
| 1e60ac0745 | |||
| 650a4c146b | |||
| 23f299e105 | |||
| d292522d00 | |||
| 0241e0d3a8 | |||
| e4973eb8a4 |
@@ -8,6 +8,7 @@ import com.rabbitmq.client.DeliverCallback;
|
|||||||
import com.rabbitmq.client.Recoverable;
|
import com.rabbitmq.client.Recoverable;
|
||||||
import com.rabbitmq.client.RecoveryListener;
|
import com.rabbitmq.client.RecoveryListener;
|
||||||
import com.rabbitmq.client.Return;
|
import com.rabbitmq.client.Return;
|
||||||
|
import com.rabbitmq.client.impl.DefaultExceptionHandler;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
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). */
|
/** 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) {
|
public static AmqpReplyInbox open(String uri, int prefetch) {
|
||||||
try {
|
try {
|
||||||
ConnectionFactory factory = new ConnectionFactory();
|
return new AmqpReplyInbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.REPLY_INBOX), prefetch);
|
||||||
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);
|
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, 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). */
|
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
|
||||||
AmqpReplyInbox(Connection connection) {
|
AmqpReplyInbox(Connection connection) {
|
||||||
this(connection, DEFAULT_PREFETCH);
|
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. */
|
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
|
||||||
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
|
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
|
||||||
try {
|
try {
|
||||||
ConnectionFactory factory = new ConnectionFactory();
|
return new LeadMailbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.LEAD_MAILBOX), selfCoordId, prefetch);
|
||||||
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);
|
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, 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). */
|
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
|
||||||
LeadMailbox(Connection connection, String selfCoordId) {
|
LeadMailbox(Connection connection, String selfCoordId) {
|
||||||
this(connection, selfCoordId, DEFAULT_PREFETCH);
|
this(connection, selfCoordId, DEFAULT_PREFETCH);
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -116,6 +116,50 @@ check_log_path_matches_plist() {
|
|||||||
ok "log path check: script and plist agree ($resolved_out)"
|
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
|
# 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
|
# 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
|
# 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/^/ /' \
|
| grep -iE 'deferred|classification:|fleet health:|coverage' | tail -8 | sed 's/^/ /' \
|
||||||
|| echo " (nothing reported)"
|
|| echo " (nothing reported)"
|
||||||
|
|
||||||
# Errors since the restart, anchored to the marker so old noise cannot leak in.
|
# Errors since the restart, anchored to the marker so old noise cannot leak in. Keep the fresh
|
||||||
ERRS="$(tail -n "+$((RESTART_MARK + 1))" "$OUT" 2>/dev/null | grep -cE ' (ERROR|SEVERE) ' || true)"
|
# 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"
|
say "result"
|
||||||
ok "pid $NEW_PID, jar $(jar_id)"
|
ok "pid $NEW_PID, jar $(jar_id)"
|
||||||
if [ "${ERRS:-0}" -gt 0 ]; then
|
if [ "$REDEPLOY_ERROR_COUNT" -eq 0 ]; then
|
||||||
warn "$ERRS ERROR lines since restart:"
|
|
||||||
tail -n "+$((RESTART_MARK + 1))" "$OUT" | grep -E ' (ERROR|SEVERE) ' | tail -5 | sed 's/^/ /'
|
|
||||||
else
|
|
||||||
ok "no ERROR lines since restart"
|
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
|
fi
|
||||||
echo
|
echo
|
||||||
echo " Next: call fleet_whoami and confirm it still answers 'primary'. A lead whose tab label"
|
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