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 @@ -24,6 +24,7 @@
import io.vertx.core.net.SocketAddress;
import io.vertx.core.net.impl.NetSocketInternal;
import io.vertx.core.spi.metrics.ClientMetrics;
import io.vertx.core.spi.metrics.VertxMetrics;
import io.vertx.db2client.DB2ConnectOptions;
import io.vertx.sqlclient.SqlConnectOptions;
import io.vertx.sqlclient.SqlConnection;
Expand Down Expand Up @@ -54,7 +55,8 @@ protected Future<Connection> doConnectInternal(SqlConnectOptions options, Contex
int pipeliningLimit = db2Options.getPipeliningLimit();
NetClient netClient = netClient(options);
return netClient.connect(server).flatMap(so -> {
ClientMetrics metrics = clientMetricsProvider != null ? clientMetricsProvider.metricsFor(options) : null;
VertxMetrics vertxMetrics = vertx.metricsSPI();
ClientMetrics metrics = vertxMetrics != null ? vertxMetrics.createClientMetrics(db2Options.getSocketAddress(), "sql", db2Options.getMetricsName()) : null;
DB2SocketConnection conn = new DB2SocketConnection((NetSocketInternal) so, metrics, db2Options, cachePreparedStatements,
preparedStatementCacheSize, preparedStatementCacheSqlFilter, pipeliningLimit, context);
conn.init();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,8 +55,8 @@ public Pool newPool(Vertx vertx, Supplier<? extends Future<? extends SqlConnectO

private PoolImpl newPoolImpl(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> databases, PoolOptions options, CloseFuture closeFuture) {
boolean pipelinedPool = options instanceof Db2PoolOptions && ((Db2PoolOptions) options).isPipelined();
PoolImpl pool = new PoolImpl(vertx, this, pipelinedPool, options, null, null, closeFuture);
ConnectionFactory factory = createConnectionFactory(vertx, databases);
PoolImpl pool = new PoolImpl(vertx, this, pipelinedPool, options, factory.metricsProvider(), null, null, closeFuture);
pool.connectionProvider(factory::connect);
pool.init();
closeFuture.add(factory);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import io.vertx.core.net.*;
import io.vertx.core.net.impl.NetSocketInternal;
import io.vertx.core.spi.metrics.ClientMetrics;
import io.vertx.core.spi.metrics.VertxMetrics;
import io.vertx.mssqlclient.MSSQLConnectOptions;
import io.vertx.sqlclient.SqlConnectOptions;
import io.vertx.sqlclient.SqlConnection;
Expand Down Expand Up @@ -72,7 +73,8 @@ private Future<Connection> connectOrRedirect(MSSQLConnectOptions options, Contex
}

private MSSQLSocketConnection createSocketConnection(NetSocket so, MSSQLConnectOptions options, ContextInternal context) {
ClientMetrics metrics = clientMetricsProvider != null ? clientMetricsProvider.metricsFor(options) : null;
VertxMetrics vertxMetrics = vertx.metricsSPI();
ClientMetrics metrics = vertxMetrics != null ? vertxMetrics.createClientMetrics(options.getSocketAddress(), "sql", options.getMetricsName()) : null;
MSSQLSocketConnection conn = new MSSQLSocketConnection((NetSocketInternal) so, metrics, options, false, 0, sql -> true, 1, context);
conn.init();
return conn;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,8 @@ public Pool newPool(Vertx vertx, Supplier<? extends Future<? extends SqlConnectO
}

private PoolImpl newPoolImpl(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> databases, PoolOptions options, CloseFuture closeFuture) {
PoolImpl pool = new PoolImpl(vertx, this, false, options, null, null, closeFuture);
ConnectionFactory factory = createConnectionFactory(vertx, databases);
PoolImpl pool = new PoolImpl(vertx, this, false, options, factory.metricsProvider(), null, null, closeFuture);
pool.connectionProvider(context -> factory.connect(context, databases.get()));
pool.init();
closeFuture.add(factory);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import io.vertx.core.net.TrustOptions;
import io.vertx.core.net.impl.NetSocketInternal;
import io.vertx.core.spi.metrics.ClientMetrics;
import io.vertx.core.spi.metrics.VertxMetrics;
import io.vertx.mysqlclient.MySQLAuthenticationPlugin;
import io.vertx.mysqlclient.MySQLConnectOptions;
import io.vertx.mysqlclient.SslMode;
Expand Down Expand Up @@ -126,7 +127,8 @@ private Future<Connection> doConnect(MySQLConnectOptions options, SslMode sslMod
MySQLAuthenticationPlugin authenticationPlugin = options.getAuthenticationPlugin();
Future<NetSocket> fut = netClient(new NetClientOptions(options).setSsl(false)).connect(server);
return fut.flatMap(so -> {
ClientMetrics metrics = clientMetricsProvider != null ? clientMetricsProvider.metricsFor(options) : null;
VertxMetrics vertxMetrics = vertx.metricsSPI();
ClientMetrics metrics = vertxMetrics != null ? vertxMetrics.createClientMetrics(options.getSocketAddress(), "sql", options.getMetricsName()) : null;
MySQLSocketConnection conn = new MySQLSocketConnection((NetSocketInternal) so, metrics, options, cachePreparedStatements, preparedStatementCacheMaxSize, preparedStatementCacheSqlFilter, pipeliningLimit, context);
conn.init();
return Future.future(promise -> conn.sendStartupMessage(username, password, database, collation, serverRsaPublicKey, properties, sslMode, initialCapabilitiesFlags, charsetEncoding, authenticationPlugin, promise));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,8 +55,8 @@ public Pool newPool(Vertx vertx, Supplier<? extends Future<? extends SqlConnectO

private PoolImpl newPoolImpl(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> databases, PoolOptions options, CloseFuture closeFuture) {
boolean pipelinedPool = options instanceof MySQLPoolOptions && ((MySQLPoolOptions) options).isPipelined();
PoolImpl pool = new PoolImpl(vertx, this, pipelinedPool, options, null, null, closeFuture);
ConnectionFactory factory = createConnectionFactory(vertx, databases);
PoolImpl pool = new PoolImpl(vertx, this, pipelinedPool, options, factory.metricsProvider(), null, null, closeFuture);
pool.connectionProvider(context -> factory.connect(context, databases.get()));
pool.init();
closeFuture.add(factory);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,6 @@
import io.vertx.oracleclient.OracleConnectOptions;
import io.vertx.sqlclient.SqlConnectOptions;
import io.vertx.sqlclient.SqlConnection;
import io.vertx.sqlclient.impl.SingletonSupplier;
import io.vertx.sqlclient.impl.metrics.ClientMetricsProvider;
import io.vertx.sqlclient.impl.metrics.DynamicClientMetricsProvider;
import io.vertx.sqlclient.impl.metrics.SingleServerClientMetricsProvider;
import io.vertx.sqlclient.spi.ConnectionFactory;
import oracle.jdbc.OracleConnection;
import oracle.jdbc.datasource.OracleDataSource;
Expand All @@ -40,39 +36,15 @@ public class OracleConnectionFactory implements ConnectionFactory {

private final Supplier<? extends Future<? extends SqlConnectOptions>> options;
private final Map<JsonObject, OracleDataSource> datasources;
private final ClientMetricsProvider clientMetricsProvider;

public OracleConnectionFactory(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> options) {
VertxMetrics metrics = vertx.metricsSPI();
ClientMetricsProvider clientMetricsProvider;
if (metrics != null) {
if (options instanceof SingletonSupplier) {
SqlConnectOptions option = (SqlConnectOptions) ((SingletonSupplier) options).unwrap();
ClientMetrics<?, ?, ?, ?> clientMetrics = metrics.createClientMetrics(option.getSocketAddress(), "sql", option.getMetricsName());
clientMetricsProvider = new SingleServerClientMetricsProvider(clientMetrics);
} else {
clientMetricsProvider = new DynamicClientMetricsProvider(metrics);
}
} else {
clientMetricsProvider = null;
}
this.clientMetricsProvider = clientMetricsProvider;
this.options = options;
this.datasources = new HashMap<>();
}

@Override
public ClientMetricsProvider metricsProvider() {
return clientMetricsProvider;
}

@Override
public void close(Promise<Void> promise) {
if (clientMetricsProvider != null) {
clientMetricsProvider.close(promise);
} else {
promise.complete();
}
promise.complete();
}

@Override
Expand All @@ -96,7 +68,8 @@ private OracleDataSource getDatasource(SqlConnectOptions options) {
@Override
public Future<SqlConnection> connect(Context context, SqlConnectOptions options) {
OracleDataSource datasource = getDatasource(options);
ClientMetrics metrics = clientMetricsProvider != null ? clientMetricsProvider.metricsFor(options) : null;
VertxMetrics vertxMetrics = ((VertxInternal)context.owner()).metricsSPI();
ClientMetrics metrics = vertxMetrics != null ? vertxMetrics.createClientMetrics(options.getSocketAddress(), "sql", options.getMetricsName()) : null;
ContextInternal ctx = (ContextInternal) context;
return executeBlocking(context, () -> {
OracleConnection orac = datasource.createConnectionBuilder().build();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,8 @@ public Pool newPool(Vertx vertx, Supplier<? extends Future<? extends SqlConnectO
private PoolImpl newPoolImpl(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> databases, PoolOptions options, CloseFuture closeFuture) {
Function<Connection, Future<Void>> afterAcquire = conn -> ((OracleJdbcConnection) conn).afterAcquire();
Function<Connection, Future<Void>> beforeRecycle = conn -> ((OracleJdbcConnection) conn).beforeRecycle();
PoolImpl pool = new PoolImpl(vertx, this, false, options, afterAcquire, beforeRecycle, closeFuture);
ConnectionFactory factory = createConnectionFactory(vertx, databases);
PoolImpl pool = new PoolImpl(vertx, this, false, options, factory.metricsProvider(), afterAcquire, beforeRecycle, closeFuture);
pool.connectionProvider(context -> factory.connect(context, databases.get()));
pool.init();
closeFuture.add(factory);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -179,7 +179,9 @@ private PgSocketConnection newSocketConnection(ContextInternal context, NetSocke
Predicate<String> preparedStatementCacheSqlFilter = options.getPreparedStatementCacheSqlFilter();
int pipeliningLimit = options.getPipeliningLimit();
boolean useLayer7Proxy = options.getUseLayer7Proxy();
ClientMetrics metrics = clientMetricsProvider != null ? clientMetricsProvider.metricsFor(options) : null;
return new PgSocketConnection(socket, metrics, options, cachePreparedStatements, preparedStatementCacheMaxSize, preparedStatementCacheSqlFilter, pipeliningLimit, useLayer7Proxy, context);
VertxMetrics vertxMetrics = vertx.metricsSPI();
ClientMetrics metrics = vertxMetrics != null ? vertxMetrics.createClientMetrics(options.getSocketAddress(), "sql", options.getMetricsName()) : null;
PgSocketConnection conn = new PgSocketConnection(socket, metrics, options, cachePreparedStatements, preparedStatementCacheMaxSize, preparedStatementCacheSqlFilter, pipeliningLimit, useLayer7Proxy, context);
return conn;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,8 @@ public Pool newPool(Vertx vertx, Supplier<? extends Future<? extends SqlConnectO

private PoolImpl newPoolImpl(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> databases, PoolOptions options, CloseFuture closeFuture) {
boolean pipelinedPool = options instanceof PgPoolOptions && ((PgPoolOptions) options).isPipelined();
PoolImpl pool = new PoolImpl(vertx, this, pipelinedPool, options, null, null, closeFuture);
ConnectionFactory factory = createConnectionFactory(vertx, databases);
PoolImpl pool = new PoolImpl(vertx, this, pipelinedPool, options, factory.metricsProvider(), null, null, closeFuture);
pool.connectionProvider(context -> factory.connect(context)); // BEWARE!!!!
pool.init();
closeFuture.add(factory);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,8 @@
import io.vertx.core.net.NetClient;
import io.vertx.core.net.NetClientOptions;
import io.vertx.core.net.impl.NetClientBuilder;
import io.vertx.core.spi.metrics.ClientMetrics;
import io.vertx.core.spi.metrics.VertxMetrics;
import io.vertx.sqlclient.SqlConnectOptions;
import io.vertx.sqlclient.SqlConnection;
import io.vertx.sqlclient.impl.metrics.ClientMetricsProvider;
import io.vertx.sqlclient.impl.metrics.DynamicClientMetricsProvider;
import io.vertx.sqlclient.impl.metrics.SingleServerClientMetricsProvider;
import io.vertx.sqlclient.spi.ConnectionFactory;

import java.util.HashMap;
Expand All @@ -44,30 +39,14 @@ public abstract class ConnectionFactoryBase implements ConnectionFactory {
protected final VertxInternal vertx;
private final Map<JsonObject, NetClient> clients;
protected final Supplier<? extends Future<? extends SqlConnectOptions>> options;
protected final ClientMetricsProvider clientMetricsProvider;

// close hook
protected final CloseFuture clientCloseFuture = new CloseFuture();

protected ConnectionFactoryBase(VertxInternal vertx, Supplier<? extends Future<? extends SqlConnectOptions>> options) {
VertxMetrics metrics = vertx.metricsSPI();
ClientMetricsProvider clientMetricsProvider;
if (metrics != null) {
if (options instanceof SingletonSupplier) {
SqlConnectOptions option = (SqlConnectOptions) ((SingletonSupplier) options).unwrap();
ClientMetrics<?, ?, ?, ?> clientMetrics = metrics.createClientMetrics(option.getSocketAddress(), "sql", option.getMetricsName());
clientMetricsProvider = new SingleServerClientMetricsProvider(clientMetrics);
} else {
clientMetricsProvider = new DynamicClientMetricsProvider(metrics);
}
clientCloseFuture.add(clientMetricsProvider);
} else {
clientMetricsProvider = null;
}
this.vertx = vertx;
this.options = options;
this.clients = new HashMap<>();
this.clientMetricsProvider = clientMetricsProvider;
}

private NetClient createNetClient(NetClientOptions options) {
Expand Down Expand Up @@ -142,8 +121,4 @@ private void doConnectWithRetry(SqlConnectOptions options, PromiseInternal<Conne
*/
protected abstract Future<Connection> doConnectInternal(SqlConnectOptions options, ContextInternal context);

@Override
public ClientMetricsProvider metricsProvider() {
return clientMetricsProvider;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,10 @@
import io.vertx.core.impl.ContextInternal;
import io.vertx.core.impl.VertxInternal;
import io.vertx.core.impl.future.PromiseInternal;
import io.vertx.core.spi.metrics.PoolMetrics;
import io.vertx.core.spi.metrics.VertxMetrics;
import io.vertx.sqlclient.*;
import io.vertx.sqlclient.impl.command.CommandBase;
import io.vertx.sqlclient.impl.metrics.ClientMetricsProvider;
import io.vertx.sqlclient.impl.pool.SqlConnectionPool;
import io.vertx.sqlclient.spi.Driver;

Expand Down Expand Up @@ -57,12 +58,19 @@ public PoolImpl(VertxInternal vertx,
Driver driver,
boolean pipelined,
PoolOptions poolOptions,
ClientMetricsProvider clientMetricsProvider,
Function<Connection, Future<Void>> afterAcquire,
Function<Connection, Future<Void>> beforeRecycle,
CloseFuture closeFuture) {
super(driver);

VertxMetrics metrics = vertx.metricsSPI();
PoolMetrics<?> poolMetrics;
if (metrics != null) {
poolMetrics = metrics.createPoolMetrics("sql", poolOptions.getName(), poolOptions.getMaxSize());
} else {
poolMetrics = null;
}

this.idleTimeout = MILLISECONDS.convert(poolOptions.getIdleTimeout(), poolOptions.getIdleTimeoutUnit());
this.connectionTimeout = MILLISECONDS.convert(poolOptions.getConnectionTimeout(), poolOptions.getConnectionTimeoutUnit());
this.maxLifetime = MILLISECONDS.convert(poolOptions.getMaxLifetime(), poolOptions.getMaxLifetimeUnit());
Expand All @@ -71,8 +79,8 @@ public PoolImpl(VertxInternal vertx,
this.pipelined = pipelined;
this.vertx = vertx;
this.pool = new SqlConnectionPool(ctx -> connectionProvider.apply(ctx), () -> connectionInitializer,
clientMetricsProvider, afterAcquire, beforeRecycle, vertx, idleTimeout, maxLifetime, poolOptions.getMaxSize(),
pipelined, poolOptions.getMaxWaitQueueSize(), poolOptions.getEventLoopSize());
poolMetrics, afterAcquire, beforeRecycle, vertx, idleTimeout, maxLifetime, poolOptions.getMaxSize(), pipelined,
poolOptions.getMaxWaitQueueSize(), poolOptions.getEventLoopSize());
this.closeFuture = closeFuture;
}

Expand Down

This file was deleted.

Loading