From 7f575ec64adfc1f20af0f172c9428da86a824fb4 Mon Sep 17 00:00:00 2001 From: Goutam Adwant <8672451+goutamadwant@users.noreply.github.com> Date: Thu, 8 Oct 2026 06:19:12 -0700 Subject: [PATCH] Cancel pool connection timeout timer when acquisition fails (#1712) * 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 * Use if/else in the acquisition failure branch Signed-off-by: Goutam Adwant --------- Signed-off-by: Goutam Adwant --- .../impl/pool/SqlConnectionPool.java | 7 +- .../sqlclient/spi/backend/DriverBaseTest.java | 82 ++++++++++++++++++- 2 files changed, 87 insertions(+), 2 deletions(-) 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..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 @@ -248,7 +248,12 @@ public void execute(CommandBase cmd, Completable handler, long timeout }); }, t -> { dequeueMetric(metric); - return Future.failedFuture(t); + if (timerId != -1 && !vertx.cancelTimer(timerId)) { + // The timer already failed the handler + return Future.failedFuture(POOL_QUERY_TIMEOUT_EXCEPTION); + } else { + return Future.failedFuture(t); + } }).onComplete(ar -> { if (ar.succeeded()) { handler.succeed(ar.result()); 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(); + } + } }