From 7780e64dcae1681e3ad8e2328c9597238ecdb247 Mon Sep 17 00:00:00 2001 From: Goutam Adwant Date: Sat, 26 Sep 2026 23:06:07 -0700 Subject: [PATCH 1/2] Cancel pool connection timeout timer when acquisition fails See #1706 When a pooled query fails to acquire a connection, the connection timeout timer is not cancelled, so it fails the query handler a second time when it fires. When the timer fires first, the acquisition failure also fails the handler again. Handle the acquisition failure like a successful acquisition: cancel the timer, and if it already fired, do not complete the handler again. Signed-off-by: Goutam Adwant --- .../impl/pool/SqlConnectionPool.java | 4 + .../sqlclient/spi/backend/DriverBaseTest.java | 82 ++++++++++++++++++- 2 files changed, 85 insertions(+), 1 deletion(-) diff --git a/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java b/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java index c1a91bebe..aa9f9243d 100644 --- a/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java +++ b/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java @@ -248,6 +248,10 @@ public void execute(CommandBase cmd, Completable handler, long timeout }); }, t -> { dequeueMetric(metric); + if (timerId != -1 && !vertx.cancelTimer(timerId)) { + // The timer already failed the handler + return Future.failedFuture(POOL_QUERY_TIMEOUT_EXCEPTION); + } return Future.failedFuture(t); }).onComplete(ar -> { if (ar.succeeded()) { diff --git a/vertx-sql-client/src/test/java/io/vertx/tests/sqlclient/spi/backend/DriverBaseTest.java b/vertx-sql-client/src/test/java/io/vertx/tests/sqlclient/spi/backend/DriverBaseTest.java index d307df884..5512d2d77 100644 --- a/vertx-sql-client/src/test/java/io/vertx/tests/sqlclient/spi/backend/DriverBaseTest.java +++ b/vertx-sql-client/src/test/java/io/vertx/tests/sqlclient/spi/backend/DriverBaseTest.java @@ -3,6 +3,7 @@ import io.vertx.core.Completable; import io.vertx.core.Context; import io.vertx.core.Future; +import io.vertx.core.Promise; import io.vertx.core.Vertx; import io.vertx.core.net.NetClientOptions; import io.vertx.core.net.SocketAddress; @@ -18,10 +19,12 @@ import io.vertx.sqlclient.desc.ColumnDescriptor; import io.vertx.sqlclient.impl.RowBase; import io.vertx.sqlclient.spi.connection.Connection; +import io.vertx.sqlclient.internal.PoolInternal; import io.vertx.sqlclient.internal.QueryResultHandler; import io.vertx.sqlclient.internal.RowDescriptorBase; import io.vertx.sqlclient.spi.connection.ConnectionContext; import io.vertx.sqlclient.spi.protocol.CommandBase; +import io.vertx.sqlclient.spi.protocol.CommandScheduler; import io.vertx.sqlclient.spi.protocol.SimpleQueryCommand; import io.vertx.sqlclient.spi.connection.ConnectionFactory; import io.vertx.sqlclient.spi.DatabaseMetadata; @@ -29,11 +32,18 @@ import org.junit.Test; import java.sql.JDBCType; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; import java.util.stream.Collector; +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; import static org.junit.Assert.*; public class DriverBaseTest { @@ -136,6 +146,13 @@ private void scheduleQueryCommand(SimpleQueryCommand simpleQuery, Comp } private static DriverBase createDriver( + java.util.function.Function> afterAcquire, + java.util.function.Function> beforeRecycle) { + return createDriver(context -> Future.succeededFuture(fakeConnection()), afterAcquire, beforeRecycle); + } + + private static DriverBase createDriver( + java.util.function.Function> connector, java.util.function.Function> afterAcquire, java.util.function.Function> beforeRecycle) { return new DriverBase<>("generic", afterAcquire, beforeRecycle) { @@ -144,7 +161,7 @@ public ConnectionFactory createConnectionFactory(Vertx vertx, return new ConnectionFactory<>() { @Override public Future connect(Context context, SqlConnectOptions options) { - return Future.succeededFuture(fakeConnection()); + return connector.apply(context); } @Override @@ -234,4 +251,67 @@ public void testAfterAcquireFailureReleasesConnection() { vertx.close().await(); } } + + @Test + public void testConnectFailureCancelsPoolConnectionTimeout() throws Exception { + RuntimeException connectError = new RuntimeException("connect failed"); + DriverBase driver = createDriver(context -> Future.failedFuture(connectError), null, null); + + Vertx vertx = Vertx.vertx(); + + try { + PoolInternal pool = driver.newPool(vertx, () -> Future.succeededFuture(new SqlConnectOptions()), + new PoolOptions().setMaxSize(1).setConnectionTimeout(1000).setConnectionTimeoutUnit(MILLISECONDS), new NetClientOptions(), null); + + List failures = new CopyOnWriteArrayList<>(); + CountDownLatch latch = new CountDownLatch(1); + // The public query path completes its promise with tryFail, which hides a double completion + ((CommandScheduler) pool).schedule(new CommandBase() {}, (res, err) -> { + failures.add(err); + latch.countDown(); + }); + assertTrue(latch.await(10, SECONDS)); + + // Let the connection timeout elapse + Thread.sleep(1500); + assertEquals(Collections.singletonList(connectError), failures); + } finally { + vertx.close().await(); + } + } + + @Test + public void testConnectFailureAfterPoolConnectionTimeout() throws Exception { + RuntimeException connectError = new RuntimeException("connect failed"); + Promise connectPromise = Promise.promise(); + AtomicReference connectContext = new AtomicReference<>(); + DriverBase driver = createDriver(context -> { + connectContext.set(context); + return connectPromise.future(); + }, null, null); + + Vertx vertx = Vertx.vertx(); + + try { + PoolInternal pool = driver.newPool(vertx, () -> Future.succeededFuture(new SqlConnectOptions()), + new PoolOptions().setMaxSize(1).setConnectionTimeout(100).setConnectionTimeoutUnit(MILLISECONDS), new NetClientOptions(), null); + + List failures = new CopyOnWriteArrayList<>(); + CountDownLatch latch = new CountDownLatch(1); + // The public query path completes its promise with tryFail, which hides a double completion + ((CommandScheduler) pool).schedule(new CommandBase() {}, (res, err) -> { + failures.add(err); + latch.countDown(); + }); + assertTrue(latch.await(10, SECONDS)); + assertEquals(1, failures.size()); + assertEquals("Timeout waiting for connection", failures.get(0).getMessage()); + + connectContext.get().runOnContext(v -> connectPromise.fail(connectError)); + Thread.sleep(300); + assertEquals(1, failures.size()); + } finally { + vertx.close().await(); + } + } } From 1edf65cff26d6f4810ae796a1f7905f14272fa0a Mon Sep 17 00:00:00 2001 From: Goutam Adwant Date: Mon, 28 Sep 2026 00:03:16 -0700 Subject: [PATCH 2/2] Use if/else in the acquisition failure branch Signed-off-by: Goutam Adwant --- .../java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java b/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java index aa9f9243d..7f82e3b42 100644 --- a/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java +++ b/vertx-sql-client/src/main/java/io/vertx/sqlclient/impl/pool/SqlConnectionPool.java @@ -251,8 +251,9 @@ public void execute(CommandBase cmd, Completable handler, long timeout if (timerId != -1 && !vertx.cancelTimer(timerId)) { // The timer already failed the handler return Future.failedFuture(POOL_QUERY_TIMEOUT_EXCEPTION); + } else { + return Future.failedFuture(t); } - return Future.failedFuture(t); }).onComplete(ar -> { if (ar.succeeded()) { handler.succeed(ar.result());