From 3f5175dfd65432cd036a554c4d65c705b14dfebf Mon Sep 17 00:00:00 2001 From: Rico Neubauer Date: Thu, 2 Jul 2026 14:29:56 +0200 Subject: [PATCH] ARTEMIS-6142 Ensure InVM connections are removed on graceful close - Fix race condition in InVMConnection - InVMAcceptor/InVMConnector gracefully closing (not reported as disconnect anymore) - Adds InVMConnectionLeakStressTest --- .../core/remoting/impl/invm/InVMAcceptor.java | 10 +- .../remoting/impl/invm/InVMConnection.java | 21 ++-- .../remoting/impl/invm/InVMConnector.java | 10 +- .../server/impl/RemotingServiceImpl.java | 3 +- .../InVMConnectionLeakStressTest.java | 107 ++++++++++++++++++ .../impl/invm/InVMConnectionTest.java | 88 ++++++++++++++ 6 files changed, 219 insertions(+), 20 deletions(-) create mode 100644 tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/jms/connection/InVMConnectionLeakStressTest.java diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMAcceptor.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMAcceptor.java index 0d258081af7..193e247c2b5 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMAcceptor.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMAcceptor.java @@ -238,7 +238,7 @@ public void connect(final String connectionID, connectionListener.connectionCreated(this, inVMConnection, protocolMap.get(ActiveMQClient.DEFAULT_CORE_PROTOCOL)); } - public void disconnect(final String connectionID) { + public void disconnect(final String connectionID, final boolean failed) { if (!started) { return; } @@ -246,7 +246,11 @@ public void disconnect(final String connectionID) { Connection conn = connections.get(connectionID); if (conn != null) { - conn.disconnect(); + if (failed) { + conn.disconnect(); + } else { + conn.close(); + } } } @@ -301,7 +305,7 @@ public void connectionDestroyed(final Object connectionID, boolean failed) { // Execute on different thread after all the packets are sent, to avoid deadlocks connection.getExecutor().execute(() -> { // Remove on the other side too - connector.disconnect((String) connectionID); + connector.disconnect((String) connectionID, failed); }); } } diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnection.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnection.java index 699f34e7b95..957770fcc7b 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnection.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnection.java @@ -53,7 +53,7 @@ public class InVMConnection implements Connection { private final String id; - private boolean closed; + private volatile boolean closed; // Used on tests private static boolean flushEnabled = true; @@ -62,8 +62,6 @@ public class InVMConnection implements Connection { private final ArtemisExecutor executor; - private volatile boolean closing; - private final ActiveMQPrincipal defaultActiveMQPrincipal; private RemotingConnection protocolConnection; @@ -146,18 +144,15 @@ public void close() { } private void internalClose(boolean failed) { - if (closing) { - return; - } - - closing = true; - + // guarantee connectionDestroyed is fired exactly once synchronized (this) { - if (!closed) { - listener.connectionDestroyed(id, failed); - - closed = true; + if (closed) { + return; } + + listener.connectionDestroyed(id, failed); + + closed = true; } } diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnector.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnector.java index 5fe5c860667..db44466c53f 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnector.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/impl/invm/InVMConnector.java @@ -214,7 +214,7 @@ public BufferHandler getHandler() { return handler; } - public void disconnect(final String connectionID) { + public void disconnect(final String connectionID, final boolean failed) { if (!started) { return; } @@ -222,7 +222,11 @@ public void disconnect(final String connectionID) { Connection conn = connections.get(connectionID); if (conn != null) { - conn.close(); + if (failed) { + conn.disconnect(); + } else { + conn.close(); + } } } @@ -267,7 +271,7 @@ public void connectionCreated(final ActiveMQComponent component, public void connectionDestroyed(final Object connectionID, boolean failed) { if (connections.remove(connectionID) != null) { // Close the corresponding connection on the other side - acceptor.disconnect((String) connectionID); + acceptor.disconnect((String) connectionID, failed); // Execute on different thread to avoid deadlocks closeExecutor.execute(() -> listener.connectionDestroyed(connectionID, failed)); diff --git a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/server/impl/RemotingServiceImpl.java b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/server/impl/RemotingServiceImpl.java index f367ddf2d42..fa0b08d02db 100644 --- a/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/server/impl/RemotingServiceImpl.java +++ b/artemis-server/src/main/java/org/apache/activemq/artemis/core/remoting/server/impl/RemotingServiceImpl.java @@ -811,7 +811,8 @@ private void issueFailure(Object connectionID, ActiveMQException e) { private void issueClose(Object connectionID) { ConnectionEntry conn = connections.get(connectionID); - if (conn != null && !conn.connection.isSupportReconnect()) { + // always remove connection on graceful close + if (conn != null) { RemotingConnection removedConnection = removeConnection(connectionID); if (removedConnection != null) { try { diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/jms/connection/InVMConnectionLeakStressTest.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/jms/connection/InVMConnectionLeakStressTest.java new file mode 100644 index 00000000000..a68a1273524 --- /dev/null +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/jms/connection/InVMConnectionLeakStressTest.java @@ -0,0 +1,107 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.activemq.artemis.tests.integration.jms.connection; + +import javax.jms.Connection; +import javax.jms.Session; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import org.apache.activemq.artemis.api.core.TransportConfiguration; +import org.apache.activemq.artemis.api.jms.ActiveMQJMSClient; +import org.apache.activemq.artemis.api.jms.JMSFactoryType; +import org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory; +import org.apache.activemq.artemis.tests.util.JMSTestBase; +import org.apache.activemq.artemis.tests.util.Wait; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Attempts to reproduce a server-side InVM connection leak where {@code RemotingServiceImpl.removeConnection} is not + * invoked even though {@code ActiveMQConnection.close()} was called. + *

+ * Two independent defects can cause this and are guarded here: + *

    + *
  1. A race in {@code InVMConnection.internalClose} where a concurrent close could return before + * {@code connectionDestroyed} fired (fixed by making the close fully atomic).
  2. + *
  3. {@code InVMAcceptor}/{@code InVMConnector} reporting every close (including a graceful + * {@code close()}) to the server as a failure, which routed it through {@code issueFailure} where the + * {@code isSupportReconnect()} guard could skip removal. When the client enables a confirmation window the + * server-side connection reports {@code isSupportReconnect() == true}; if the session channel is still present at + * transport-teardown time the connection was never removed, and because the InVM connection-ttl is -1 the + * failure-check reaper never removes it either - a permanent leak.
  4. + *
+ * The bug is timing dependent (it depends on the ordering of {@code SESS_CLOSE} processing versus transport teardown), + * so this test floods the broker with many short-lived reconnect-capable connections to widen the race window. It may + * not fail on every machine, but on hardware/timing where the race is hit it will leave server connections behind. + */ +public class InVMConnectionLeakStressTest extends JMSTestBase { + + private ActiveMQConnectionFactory floodCf; + + @Override + @BeforeEach + public void setUp() throws Exception { + super.setUp(); + + floodCf = ActiveMQJMSClient.createConnectionFactoryWithoutHA(JMSFactoryType.CF, new TransportConfiguration(INVM_CONNECTOR_FACTORY)); + // A positive confirmation window makes the server-side connection report isSupportReconnect() == true, which is + // what used to make issueFailure()/issueClose() skip removeConnection(). + floodCf.setConfirmationWindowSize(1024 * 1024); + floodCf.setReconnectAttempts(-1); + } + + @Test + public void testConcurrentGracefulCloseRemovesAllConnections() throws Exception { + final int numConnections = 20_000; + final int threads = 100; + + List> tasks = new ArrayList<>(numConnections); + for (int i = 0; i < numConnections; i++) { + tasks.add(() -> { + Connection connection = floodCf.createConnection(); + Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); + session.createProducer(ActiveMQJMSClient.createQueue("stress-queue")); + // Graceful close + connection.close(); + return null; + }); + } + + ExecutorService executor = Executors.newFixedThreadPool(threads); + try { + for (Future future : executor.invokeAll(tasks)) { + future.get(); + } + } finally { + executor.shutdown(); + assertTrue(executor.awaitTermination(2, TimeUnit.MINUTES)); + } + + // Every gracefully-closed connection must be removed from the server. InVM connection-ttl is -1 so the + // failure-check reaper never removes them; if this never reaches 0 the connections have leaked. + Wait.assertEquals(0, () -> server.getRemotingService().getConnectionCount(), 10_000); + } +} diff --git a/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/remoting/impl/invm/InVMConnectionTest.java b/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/remoting/impl/invm/InVMConnectionTest.java index dc710389686..48b97d0f573 100644 --- a/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/remoting/impl/invm/InVMConnectionTest.java +++ b/tests/unit-tests/src/test/java/org/apache/activemq/artemis/tests/unit/core/remoting/impl/invm/InVMConnectionTest.java @@ -16,17 +16,25 @@ */ package org.apache.activemq.artemis.tests.unit.core.remoting.impl.invm; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; +import org.apache.activemq.artemis.api.core.ActiveMQException; import org.apache.activemq.artemis.api.core.TransportConfiguration; import org.apache.activemq.artemis.core.remoting.impl.invm.InVMConnection; import org.apache.activemq.artemis.core.remoting.impl.invm.InVMConnectorFactory; import org.apache.activemq.artemis.core.remoting.impl.invm.TransportConstants; +import org.apache.activemq.artemis.core.server.ActiveMQComponent; +import org.apache.activemq.artemis.spi.core.protocol.ProtocolManager; +import org.apache.activemq.artemis.spi.core.remoting.BaseConnectionLifeCycleListener; +import org.apache.activemq.artemis.spi.core.remoting.Connection; import org.junit.jupiter.api.Test; import java.util.HashMap; import java.util.Map; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.atomic.AtomicInteger; public class InVMConnectionTest { @@ -55,4 +63,84 @@ public void testIsTargetNode() throws Exception { assertTrue(conn.isSameTarget(tf2, tf0)); assertFalse(conn.isSameTarget(tf2, tf1)); } + + @Test + public void testConcurrentCloseFiresConnectionDestroyedExactlyOnce() throws Exception { + final int threads = 16; + // Repeat several rounds to widen the window for catching the race. + for (int round = 0; round < 50; round++) { + final CountingLifeCycleListener listener = new CountingLifeCycleListener(); + final InVMConnection conn = new InVMConnection(0, null, listener, null); + + final CyclicBarrier barrier = new CyclicBarrier(threads); + final Thread[] workers = new Thread[threads]; + final AtomicInteger prematureReturns = new AtomicInteger(); + + for (int i = 0; i < threads; i++) { + final boolean disconnect = (i % 2 == 0); + workers[i] = new Thread(() -> { + try { + // Line up all threads so they hit close()/disconnect() together. + barrier.await(); + } catch (Exception e) { + throw new RuntimeException(e); + } + if (disconnect) { + conn.disconnect(); + } else { + conn.close(); + } + // by the time any close()/disconnect() call returns, connectionDestroyed must already have fired. + if (!listener.destroyFired) { + prematureReturns.incrementAndGet(); + } + }); + } + + for (Thread worker : workers) { + worker.start(); + } + for (Thread worker : workers) { + worker.join(); + } + + assertEquals(1, listener.destroyedCount.get(), + "connectionDestroyed must be fired exactly once per connection (round " + round + ")"); + assertEquals(0, prematureReturns.get(), + "close()/disconnect() must not return before connectionDestroyed has fired (round " + round + ")"); + } + } + + private static final class CountingLifeCycleListener implements BaseConnectionLifeCycleListener { + + private final AtomicInteger destroyedCount = new AtomicInteger(); + + private volatile boolean destroyFired; + + @Override + public void connectionCreated(ActiveMQComponent component, Connection connection, ProtocolManager protocol) { + } + + @Override + public void connectionDestroyed(Object connectionID, boolean failed) { + // Hold inside the callback for a moment to widen the window in which a + // concurrent, non-atomic close() could wrongly observe the connection + // as "closing" and return early before this callback completes. + try { + Thread.sleep(5); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + destroyedCount.incrementAndGet(); + destroyFired = true; + } + + @Override + public void connectionException(Object connectionID, ActiveMQException me) { + } + + @Override + public void connectionReadyForWrites(Object connectionID, boolean ready) { + } + } }