From d292522d00273e72ae36b6a04995b1e5613ee092 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 5 Sep 2026 05:45:37 +0700 Subject: [PATCH 1/3] Name AMQP connection failure logs --- .../dev/ltms/fleet/msg/AmqpReplyInbox.java | 58 +++++++++- .../java/dev/ltms/fleet/msg/LeadMailbox.java | 17 ++- .../msg/AmqpConnectionFailureLoggerTest.java | 105 ++++++++++++++++++ 3 files changed, 168 insertions(+), 12 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java 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..67c5d1b 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.ForgivingExceptionHandler; 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,43 @@ 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 { + + 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) { + 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..a054a7f --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpConnectionFailureLoggerTest.java @@ -0,0 +1,105 @@ +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.ConnectionFactory; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; + +import java.io.IOException; + +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 = assertInstanceOf(AmqpConnectionFailureLogger.class, + inboxFactory.getExceptionHandler()); + AmqpConnectionFailureLogger mailboxHandler = assertInstanceOf(AmqpConnectionFailureLogger.class, + mailboxFactory.getExceptionHandler()); + 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 handlerOnlyChangesForgivingHandlerLogging() { + assertEquals(com.rabbitmq.client.impl.ForgivingExceptionHandler.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"); + } + + 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 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); + } +} From dbf6fef0e9f9625f6bf6117dd48f37794c5c5eda Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 5 Sep 2026 05:51:59 +0700 Subject: [PATCH 2/3] config: guard FleetConfig.withDefaults() against silently dropping a component Adding a component to FleetConfig follows an established pattern: the record grows by one arg, and a back-compat constructor is added at the OLD arity so existing callers keep compiling. That back-compat constructor also silently captures withDefaults()'s own literal-arity 'return new FleetConfig(...)' call the next time this happens, since that call is now a legal overload match too. It compiles, every other test passes, and the new component is defaulted away on every load(). This is not hypothetical - it happened live while building the (now parked) idle-sleep-guard PR, caught only because that branch's own new tests asserted on the new field. Add a reflective test that builds a FleetConfig through the true canonical constructor (resolved by record-component types, not arg count - the same pattern ConfigRefTopLevelReportingCoverageTest already uses in this file) with a real, non-null value in every component, runs the real withDefaults(), and asserts every value survives unchanged. Never hardcodes the arity - it enumerates FleetConfig.class.getRecordComponents() - so it keeps working as the record grows. No back-compat constructor is touched or removed. --- ...thDefaultsPreservesEveryComponentTest.java | 192 ++++++++++++++++++ 1 file changed, 192 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java new file mode 100644 index 0000000..1f710a0 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java @@ -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. + * + *

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. + * + *

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. + * + *

{@link #EXCLUDED_FROM_SURVIVAL_CHECK} exists anyway, kept deliberately empty and size-pinned + * by {@link #exclusionListSizeIsPinned()}: a future component that {@code withDefaults()} is + * documented 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 EXCLUDED_FROM_SURVIVAL_CHECK = Set.of(); + + /** One real, distinctive, non-null (non-blank where blankness would mean "unset") value per component. */ + private static Map baseValues() { + Map 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 values) { + Set 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 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 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 base = baseValues(); + FleetConfig config = configOf(base); + FleetConfig defaulted = config.withDefaults(); + + List 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); + } +} From 23f299e105faf2d4d769572f40f63ef67fb10e74 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 5 Sep 2026 05:53:44 +0700 Subject: [PATCH 3/3] Preserve strict AMQP exception handling --- .../dev/ltms/fleet/msg/AmqpReplyInbox.java | 5 ++- .../msg/AmqpConnectionFailureLoggerTest.java | 43 ++++++++++++++++--- 2 files changed, 39 insertions(+), 9 deletions(-) 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();