diff --git a/bitrepository-core/src/main/java/org/bitrepository/protocol/activemq/ActiveMQMessageBus.java b/bitrepository-core/src/main/java/org/bitrepository/protocol/activemq/ActiveMQMessageBus.java index cc9246264..42a0450c1 100644 --- a/bitrepository-core/src/main/java/org/bitrepository/protocol/activemq/ActiveMQMessageBus.java +++ b/bitrepository-core/src/main/java/org/bitrepository/protocol/activemq/ActiveMQMessageBus.java @@ -35,8 +35,6 @@ import jakarta.jms.TextMessage; import jakarta.jms.Topic; import org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory; - -import java.io.ByteArrayInputStream; import org.bitrepository.bitrepositorymessages.Message; import org.bitrepository.bitrepositorymessages.MessageRequest; import org.bitrepository.common.DefaultThreadFactory; @@ -67,6 +65,7 @@ import org.slf4j.LoggerFactory; import org.xml.sax.SAXException; +import java.io.ByteArrayInputStream; import java.nio.charset.StandardCharsets; import java.util.Collections; import java.util.HashMap; @@ -174,18 +173,21 @@ public ActiveMQMessageBus(Settings settings, SecurityManager securityManager) { jaxbHelper = new JaxbHelper("xsd/", schemaLocation); ActiveMQConnectionFactory connectionFactory = ArtemisConnectionFactoryProvider.create(configuration); registerCustomMessageLoggers(); + Connection newConnection = null; try { - connection = connectionFactory.createConnection(); - connection.setClientID(clientID); - connection.setExceptionListener(new MessageBusExceptionListener()); + newConnection = connectionFactory.createConnection(); + newConnection.setClientID(clientID); + newConnection.setExceptionListener(new MessageBusExceptionListener()); - producerSession = connection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); - consumerSession = connection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); + producerSession = newConnection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); + consumerSession = newConnection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); producer = producerSession.createProducer(null); producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT); + connection = newConnection; startListeningForMessages(); } catch (JMSException e) { + closeQuietly(newConnection); throw new CoordinationLayerException("Unable to initialise connection to message bus", e); } log.debug("ActiveMQConnection initialized for '{}'", configuration); @@ -212,6 +214,20 @@ private void startListeningForMessages() { connectionStarter.start(); } + /** + * Closes the connection, suppressing any JMSException, so it can be safely used for cleanup after a + * failed initialisation without masking the original error. + */ + private void closeQuietly(Connection connection) { + if (connection != null) { + try { + connection.close(); + } catch (JMSException e) { + log.warn("Failed to close connection after initialisation failure", e); + } + } + } + @Override public synchronized void addListener(String destinationID, final MessageListener listener) { addListener(destinationID, listener, false); @@ -251,12 +267,9 @@ public synchronized void removeListener(String destinationID, MessageListener li public void close() throws JMSException { receivedMessageHandler.close(); log.info("Closing message bus: {}", configuration); - producerSession.close(); - log.debug("Producer session closed."); - consumerSession.close(); - log.debug("Consumer session closed."); - connection.close(); - log.debug("Connection closed."); + try (connection; consumerSession; producerSession) { + log.debug("Closing producer session, consumer session and connection."); + } } @Override diff --git a/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/ActiveMQMessageBusTest.java b/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/ActiveMQMessageBusTest.java index 139c7dbf0..2965d8c04 100644 --- a/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/ActiveMQMessageBusTest.java +++ b/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/ActiveMQMessageBusTest.java @@ -26,6 +26,7 @@ import org.bitrepository.bitrepositorymessages.DeleteFileRequest; import org.bitrepository.bitrepositorymessages.IdentifyPillarsForDeleteFileRequest; import org.bitrepository.bitrepositorymessages.IdentifyPillarsForDeleteFileResponse; +import org.bitrepository.protocol.CoordinationLayerException; import org.bitrepository.protocol.ProtocolComponentFactory; import org.bitrepository.protocol.activemq.ActiveMQMessageBus; import org.bitrepository.protocol.message.ExampleMessageFactory; @@ -181,4 +182,40 @@ final void toFilterTest() throws Exception { rawMessagebus.sendMessage(settingsForTestClient.getCollectionDestination(), rq); collectionReceiver.waitForMessage(DeleteFileRequest.class); } -} \ No newline at end of file + + @Test + @Tag("regressiontest") + final void closeReleasesJmsResourcesTest() throws Exception { + addDescription("Test that closing a message bus releases its JMS resources, so it can no longer " + + "be used to send messages afterwards."); + addStep("Create a dedicated message bus instance and close it", + "No exception should be thrown while closing."); + ActiveMQMessageBus dedicatedMessageBus = new ActiveMQMessageBus(settingsForTestClient, securityManager); + dedicatedMessageBus.close(); + + addStep("Attempt to send a message on the closed message bus", + "The send should fail, since the underlying JMS session and connection have been closed."); + IdentifyPillarsForDeleteFileRequest message = + ExampleMessageFactory.createMessage(IdentifyPillarsForDeleteFileRequest.class); + message.setDestination(settingsForTestClient.getCollectionDestination()); + Assertions.assertThrows(CoordinationLayerException.class, () -> dedicatedMessageBus.sendMessage(message)); + } + + @Test + @Tag("regressiontest") + final void rawMessagebusCloseReleasesJmsResourcesTest() throws Exception { + addDescription("Test that closing a RawMessagebus releases its JMS resources, so it can no longer " + + "be used afterwards."); + addStep("Create a raw message bus instance and close it", + "No exception should be thrown while closing."); + RawMessagebus rawMessagebus = new RawMessagebus( + settingsForTestClient.getMessageBusConfiguration(), + securityManager); + rawMessagebus.close(); + + addStep("Attempt to create a producer on the closed raw message bus", + "The call should fail, since the underlying JMS session and connection have been closed."); + Assertions.assertThrows(CoordinationLayerException.class, + () -> rawMessagebus.getProducer(settingsForTestClient.getCollectionDestination())); + } +} diff --git a/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/RawMessagebus.java b/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/RawMessagebus.java index 9c9cfe45d..8e0b4bcf5 100644 --- a/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/RawMessagebus.java +++ b/bitrepository-core/src/test/java/org/bitrepository/protocol/bus/RawMessagebus.java @@ -30,10 +30,10 @@ import jakarta.jms.MessageConsumer; import jakarta.jms.MessageProducer; import jakarta.jms.Session; -import org.bitrepository.protocol.activemq.ArtemisConnectionFactoryProvider; import org.bitrepository.common.JaxbHelper; import org.bitrepository.protocol.CoordinationLayerException; import org.bitrepository.protocol.activemq.ActiveMQMessageBus; +import org.bitrepository.protocol.activemq.ArtemisConnectionFactoryProvider; import org.bitrepository.protocol.security.SecurityManager; import org.bitrepository.settings.repositorysettings.MessageBusConfiguration; import org.slf4j.Logger; @@ -57,19 +57,47 @@ public RawMessagebus(MessageBusConfiguration messageBusConfiguration, SecurityMa this.securityManager = securityManager; var connectionFactory = ArtemisConnectionFactoryProvider.create(messageBusConfiguration); + Connection newConnection = null; try { - connection = connectionFactory.createConnection(); - connection.setExceptionListener(new MessageBusExceptionListener()); + newConnection = connectionFactory.createConnection(); + newConnection.setExceptionListener(new MessageBusExceptionListener()); - producerSession = connection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); - consumerSession = connection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); + producerSession = newConnection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); + consumerSession = newConnection.createSession(TRANSACTED, Session.AUTO_ACKNOWLEDGE); + connection = newConnection; connection.start(); } catch (JMSException e) { + closeQuietly(newConnection); throw new CoordinationLayerException("Unable to initialise connection to message bus", e); } } + /** + * Closes the connection, suppressing any JMSException, so it can be safely used for cleanup after a + * failed initialisation without masking the original error. + */ + private void closeQuietly(Connection connection) { + if (connection != null) { + try { + connection.close(); + } catch (JMSException e) { + log.warn("Failed to close connection after initialisation failure", e); + } + } + } + + /** + * Closes the producer session, consumer session and connection. Declared in reverse close order so + * try-with-resources closes the sessions before the connection, guaranteeing all are attempted even + * if one throws. + */ + public void close() throws JMSException { + try (connection; consumerSession; producerSession) { + log.debug("Closing raw message bus connection."); + } + } + public void addHeader(Message msg, String messageClass, String replyTo,