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 67c5d1b..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,7 +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.ForgivingExceptionHandler; +import com.rabbitmq.client.impl.DefaultExceptionHandler; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -582,7 +582,7 @@ 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 ForgivingExceptionHandler { +final class AmqpConnectionFailureLogger extends DefaultExceptionHandler { static final String REPLY_INBOX = "fleetd-reply-inbox"; static final String LEAD_MAILBOX = "fleetd-lead-mailbox"; @@ -609,6 +609,7 @@ final class AmqpConnectionFailureLogger extends ForgivingExceptionHandler { } 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; } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java index a054a7f..e6d59e9 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java @@ -4,11 +4,15 @@ 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; @@ -23,10 +27,8 @@ class AmqpConnectionFailureLoggerTest { ConnectionFactory inboxFactory = AmqpReplyInbox.connectionFactory("amqp://127.0.0.1"); ConnectionFactory mailboxFactory = LeadMailbox.connectionFactory("amqp://127.0.0.1"); - AmqpConnectionFailureLogger inboxHandler = assertInstanceOf(AmqpConnectionFailureLogger.class, - inboxFactory.getExceptionHandler()); - AmqpConnectionFailureLogger mailboxHandler = assertInstanceOf(AmqpConnectionFailureLogger.class, - mailboxFactory.getExceptionHandler()); + AmqpConnectionFailureLogger inboxHandler = installedStrictHandler(inboxFactory, "reply inbox"); + AmqpConnectionFailureLogger mailboxHandler = installedStrictHandler(mailboxFactory, "lead mailbox"); assertEquals(AmqpConnectionFailureLogger.REPLY_INBOX, inboxHandler.connectionName()); assertEquals(AmqpConnectionFailureLogger.LEAD_MAILBOX, mailboxHandler.connectionName()); @@ -73,12 +75,32 @@ class AmqpConnectionFailureLoggerTest { } @Test - void handlerOnlyChangesForgivingHandlerLogging() { - assertEquals(com.rabbitmq.client.impl.ForgivingExceptionHandler.class, + 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 ForgivingExceptionHandler"); + "all exception-handling methods must remain inherited from DefaultExceptionHandler"); } private static ListAppender attach(Class owner) { @@ -94,6 +116,13 @@ class AmqpConnectionFailureLoggerTest { ((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();