diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java index ebdb068..95b984a 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java @@ -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()); + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java index b99046d..78ee0da 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -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); diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java new file mode 100644 index 0000000..e6d59e9 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java @@ -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 inboxEvents = attach(AmqpReplyInbox.class); + ListAppender 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 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 attach(Class owner) { + Logger logger = (Logger) LoggerFactory.getLogger(owner); + logger.setLevel(Level.DEBUG); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + return appender; + } + + private static void detach(Class owner, ListAppender 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 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); + } +}