Skip to content
Merged
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 @@ -248,7 +248,12 @@ public <R> Future<R> execute(ContextInternal context, CommandBase<R> cmd, long t
});
}, t -> {
dequeueAndReject(queueMetric);
return Future.failedFuture(t);
if (timerId != -1 && !vertx.cancelTimer(timerId)) {
// The timer already failed the result
return Future.failedFuture(POOL_QUERY_TIMEOUT_EXCEPTION);
Comment thread
tsegismont marked this conversation as resolved.
} else {
return Future.failedFuture(t);
}
}).onComplete(ar -> {
if (ar.succeeded()) {
res.complete(ar.result());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,9 @@
import io.vertx.sqlclient.spi.Driver;
import org.junit.Test;

import java.util.Collections;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
Expand Down Expand Up @@ -214,4 +217,88 @@ public void testAfterAcquireFailureReleasesConnection() throws Exception {
vertx.close();
}
}

@Test
public void testConnectFailureCancelsConnectionTimeout() throws Exception {
RuntimeException connectError = new RuntimeException("connect failed");

VertxInternal vertx = (VertxInternal) Vertx.vertx();
List<Throwable> uncaught = new CopyOnWriteArrayList<>();
vertx.exceptionHandler(uncaught::add);

try {
SqlConnectionPool pool = createPool(vertx, ctx -> Future.failedFuture(connectError));

ContextInternal ctx = vertx.getOrCreateContext();
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Throwable> failure = new AtomicReference<>();
ctx.runOnContext(v -> {
pool.execute(ctx, new CommandBase<Void>() {}, 1000).onComplete(ar -> {
failure.set(ar.cause());
latch.countDown();
});
});
assertTrue(latch.await(5, TimeUnit.SECONDS));
assertEquals(connectError, failure.get());

// Let the connection timeout elapse
Thread.sleep(1500);
assertEquals(Collections.emptyList(), uncaught);

pool.close();
} finally {
vertx.close();
}
}

@Test
public void testConnectFailureAfterConnectionTimeout() throws Exception {
RuntimeException connectError = new RuntimeException("connect failed");
Promise<SqlConnection> connectPromise = Promise.promise();

VertxInternal vertx = (VertxInternal) Vertx.vertx();
List<Throwable> uncaught = new CopyOnWriteArrayList<>();
vertx.exceptionHandler(uncaught::add);

try {
SqlConnectionPool pool = createPool(vertx, ctx -> connectPromise.future());

ContextInternal ctx = vertx.getOrCreateContext();
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Throwable> failure = new AtomicReference<>();
ctx.runOnContext(v -> {
pool.execute(ctx, new CommandBase<Void>() {}, 100).onComplete(ar -> {
failure.set(ar.cause());
latch.countDown();
});
});
assertTrue(latch.await(5, TimeUnit.SECONDS));
assertEquals("Timeout waiting for connection", failure.get().getMessage());

ctx.runOnContext(v -> connectPromise.fail(connectError));
Thread.sleep(300);
assertEquals(Collections.emptyList(), uncaught);

pool.close();
} finally {
vertx.close();
}
}

private static SqlConnectionPool createPool(VertxInternal vertx, java.util.function.Function<Context, Future<SqlConnection>> connectionProvider) {
return new SqlConnectionPool(
connectionProvider,
() -> null,
null,
null,
null,
vertx,
0,
0,
1,
false,
-1,
0
);
}
}
Loading