diff --git a/spring-jdbc/src/main/java/org/springframework/jdbc/datasource/SingleConnectionDataSource.java b/spring-jdbc/src/main/java/org/springframework/jdbc/datasource/SingleConnectionDataSource.java index c978e37aa0..b62fd4b728 100644 --- a/spring-jdbc/src/main/java/org/springframework/jdbc/datasource/SingleConnectionDataSource.java +++ b/spring-jdbc/src/main/java/org/springframework/jdbc/datasource/SingleConnectionDataSource.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2023 the original author or authors. + * Copyright 2002-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -22,6 +22,8 @@ import java.lang.reflect.Method; import java.lang.reflect.Proxy; import java.sql.Connection; import java.sql.SQLException; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.springframework.beans.factory.DisposableBean; import org.springframework.lang.Nullable; @@ -73,8 +75,8 @@ public class SingleConnectionDataSource extends DriverManagerDataSource @Nullable private Connection connection; - /** Synchronization monitor for the shared Connection. */ - private final Object connectionMonitor = new Object(); + /** Lifecycle lock for the shared Connection. */ + private final Lock connectionLock = new ReentrantLock(); /** @@ -181,7 +183,8 @@ public class SingleConnectionDataSource extends DriverManagerDataSource @Override public Connection getConnection() throws SQLException { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.connection == null) { // No underlying Connection -> lazy init via DriverManager. initConnection(); @@ -193,6 +196,9 @@ public class SingleConnectionDataSource extends DriverManagerDataSource } return this.connection; } + finally { + this.connectionLock.unlock(); + } } /** @@ -216,9 +222,13 @@ public class SingleConnectionDataSource extends DriverManagerDataSource */ @Override public boolean shouldClose(Connection con) { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { return (con != this.connection && con != this.target); } + finally { + this.connectionLock.unlock(); + } } /** @@ -241,11 +251,15 @@ public class SingleConnectionDataSource extends DriverManagerDataSource */ @Override public void destroy() { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.target != null) { closeConnection(this.target); } } + finally { + this.connectionLock.unlock(); + } } @@ -256,7 +270,8 @@ public class SingleConnectionDataSource extends DriverManagerDataSource if (getUrl() == null) { throw new IllegalStateException("'url' property is required for lazily initializing a Connection"); } - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.target != null) { closeConnection(this.target); } @@ -267,19 +282,26 @@ public class SingleConnectionDataSource extends DriverManagerDataSource } this.connection = (isSuppressClose() ? getCloseSuppressingConnectionProxy(this.target) : this.target); } + finally { + this.connectionLock.unlock(); + } } /** * Reset the underlying shared Connection, to be reinitialized on next access. */ public void resetConnection() { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.target != null) { closeConnection(this.target); } this.target = null; this.connection = null; } + finally { + this.connectionLock.unlock(); + } } /** diff --git a/spring-jms/src/main/java/org/springframework/jms/connection/SingleConnectionFactory.java b/spring-jms/src/main/java/org/springframework/jms/connection/SingleConnectionFactory.java index bdaeee82ef..f3ffd057bf 100644 --- a/spring-jms/src/main/java/org/springframework/jms/connection/SingleConnectionFactory.java +++ b/spring-jms/src/main/java/org/springframework/jms/connection/SingleConnectionFactory.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2023 the original author or authors. + * Copyright 2002-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -24,6 +24,8 @@ import java.util.ArrayList; import java.util.LinkedHashSet; import java.util.List; import java.util.Set; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import jakarta.jms.Connection; import jakarta.jms.ConnectionFactory; @@ -116,8 +118,8 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti /** Whether the shared Connection has been started. */ private int startedCount = 0; - /** Synchronization monitor for the shared Connection. */ - private final Object connectionMonitor = new Object(); + /** Lifecycle lock for the shared Connection. */ + private final Lock connectionLock = new ReentrantLock(); /** @@ -252,10 +254,14 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti @Override public QueueConnection createQueueConnection() throws JMSException { Connection con; - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { this.pubSubMode = Boolean.FALSE; con = createConnection(); } + finally { + this.connectionLock.unlock(); + } if (!(con instanceof QueueConnection queueConnection)) { throw new jakarta.jms.IllegalStateException( "This SingleConnectionFactory does not hold a QueueConnection but rather: " + con); @@ -272,10 +278,14 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti @Override public TopicConnection createTopicConnection() throws JMSException { Connection con; - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { this.pubSubMode = Boolean.TRUE; con = createConnection(); } + finally { + this.connectionLock.unlock(); + } if (!(con instanceof TopicConnection topicConnection)) { throw new jakarta.jms.IllegalStateException( "This SingleConnectionFactory does not hold a TopicConnection but rather: " + con); @@ -323,12 +333,16 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti * @see #initConnection() */ protected Connection getConnection() throws JMSException { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.connection == null) { initConnection(); } return this.connection; } + finally { + this.connectionLock.unlock(); + } } /** @@ -386,9 +400,13 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti */ @Override public boolean isRunning() { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { return (this.connection != null); } + finally { + this.connectionLock.unlock(); + } } @@ -404,7 +422,8 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti throw new IllegalStateException( "'targetConnectionFactory' is required for lazily initializing a Connection"); } - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.connection != null) { closeConnection(this.connection); } @@ -433,6 +452,9 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti logger.debug("Established shared JMS Connection: " + this.connection); } } + finally { + this.connectionLock.unlock(); + } } /** @@ -531,12 +553,16 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti * @see #closeConnection */ public void resetConnection() { - synchronized (this.connectionMonitor) { + this.connectionLock.lock(); + try { if (this.connection != null) { closeConnection(this.connection); } this.connection = null; } + finally { + this.connectionLock.unlock(); + } } /** @@ -634,7 +660,8 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti } case "setExceptionListener" -> { // Handle setExceptionListener method: add to the chain. - synchronized (connectionMonitor) { + connectionLock.lock(); + try { if (aggregatedExceptionListener != null) { ExceptionListener listener = (ExceptionListener) args[0]; if (listener != this.localExceptionListener) { @@ -656,9 +683,13 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti "which will allow for registering further ExceptionListeners to the recovery chain."); } } + finally { + connectionLock.unlock(); + } } case "getExceptionListener" -> { - synchronized (connectionMonitor) { + connectionLock.lock(); + try { if (this.localExceptionListener != null) { return this.localExceptionListener; } @@ -666,6 +697,9 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti return getExceptionListener(); } } + finally { + connectionLock.unlock(); + } } case "start" -> { localStart(); @@ -677,7 +711,8 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti } case "close" -> { localStop(); - synchronized (connectionMonitor) { + connectionLock.lock(); + try { if (this.localExceptionListener != null) { if (aggregatedExceptionListener != null) { aggregatedExceptionListener.delegates.remove(this.localExceptionListener); @@ -685,6 +720,9 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti this.localExceptionListener = null; } } + finally { + connectionLock.unlock(); + } return null; } case "createSession", "createQueueSession", "createTopicSession" -> { @@ -727,7 +765,8 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti } private void localStart() throws JMSException { - synchronized (connectionMonitor) { + connectionLock.lock(); + try { if (!this.locallyStarted) { this.locallyStarted = true; if (startedCount == 0 && connection != null) { @@ -736,10 +775,14 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti startedCount++; } } + finally { + connectionLock.unlock(); + } } private void localStop() throws JMSException { - synchronized (connectionMonitor) { + connectionLock.lock(); + try { if (this.locallyStarted) { this.locallyStarted = false; if (startedCount == 1 && connection != null) { @@ -750,6 +793,9 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti } } } + finally { + connectionLock.unlock(); + } } private SingleConnectionFactory factory() { @@ -771,9 +817,13 @@ public class SingleConnectionFactory implements ConnectionFactory, QueueConnecti // Iterate over temporary copy in order to avoid ConcurrentModificationException, // since listener invocations may in turn trigger registration of listeners... Set copy; - synchronized (connectionMonitor) { + connectionLock.lock(); + try { copy = new LinkedHashSet<>(this.delegates); } + finally { + connectionLock.unlock(); + } for (ExceptionListener listener : copy) { listener.onException(ex); }