Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -238,15 +238,19 @@ 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;
}

Connection conn = connections.get(connectionID);

if (conn != null) {
conn.disconnect();
if (failed) {
conn.disconnect();
} else {
conn.close();
}
}
}

Expand Down Expand Up @@ -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);
});
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -62,8 +62,6 @@ public class InVMConnection implements Connection {

private final ArtemisExecutor executor;

private volatile boolean closing;

private final ActiveMQPrincipal defaultActiveMQPrincipal;

private RemotingConnection protocolConnection;
Expand Down Expand Up @@ -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;
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,15 +214,19 @@ public BufferHandler getHandler() {
return handler;
}

public void disconnect(final String connectionID) {
public void disconnect(final String connectionID, final boolean failed) {
if (!started) {
return;
}

Connection conn = connections.get(connectionID);

if (conn != null) {
conn.close();
if (failed) {
conn.disconnect();
} else {
conn.close();
}
}
}

Expand Down Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
* <p>
* Two independent defects can cause this and are guarded here:
* <ol>
* <li>A race in {@code InVMConnection.internalClose} where a concurrent close could return before
* {@code connectionDestroyed} fired (fixed by making the close fully atomic).</li>
* <li>{@code InVMAcceptor}/{@code InVMConnector} reporting <em>every</em> 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.</li>
* </ol>
* 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<Callable<Void>> 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<Void> 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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand Down Expand Up @@ -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<ProtocolManager> {

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) {
}
}
}