mirror of
https://github.com/swift-server/RediStack.git
synced 2026-06-02 07:37:33 +00:00
Significantly Improve the Configuration API for Pools and Connections
## Motivation The API for establishing the configuration of a connection pool had a lot of jargon and properties that developers had issues keeping straight and understanding what each does. This commit provides first-class API support for concepts such as retry strategies, and how the pool handles connection counts. ## Changes - Add: New ConnectionCountBehavior for determining leaky / non-leaky behavior - Add: New ConnectionRetryStrategy for allowing customization of retry behavior - Change: RedisConnection.defaultPort to be a computed property - Change: The logging keys of pool connection retry metadata - Rename: Several configuration properties to drop prefixes or to be combined into new structures ## Result Developers should have a much better experience exploring the available configuration options for pools and connections, being able to understand how each piece works with the underlying system.
This commit is contained in:
@@ -48,34 +48,22 @@ internal final class ConnectionPool {
|
||||
/// The event loop we're on.
|
||||
private let loop: EventLoop
|
||||
|
||||
/// The exponential backoff factor for connection attempts.
|
||||
internal let backoffFactor: Float32
|
||||
|
||||
/// The initial delay for backing off a reconnection attempt.
|
||||
internal let initialBackoffDelay: TimeAmount
|
||||
|
||||
/// The maximum number of connections the pool will preserve. Additional connections will be made available
|
||||
/// past this limit if `leaky` is set to `true`, but they will not be persisted in the pool once used.
|
||||
internal let maximumConnectionCount: Int
|
||||
/// The strategy to use for finding and returning connections when requested.
|
||||
internal let connectionRetryStrategy: RedisConnectionPool.PoolConnectionRetryStrategy
|
||||
|
||||
/// The minimum number of connections the pool will keep alive. If a connection is disconnected while in the
|
||||
/// pool such that the number of connections drops below this number, the connection will be re-established.
|
||||
internal let minimumConnectionCount: Int
|
||||
/// The maximum number of connections the pool will preserve.
|
||||
internal let maximumConnectionCount: Int
|
||||
/// The behavior to use for allowing or denying additional connections past the max connection count.
|
||||
internal let maxConnectionCountBehavior: RedisConnectionPool.ConnectionCountBehavior.MaxConnectionBehavior
|
||||
|
||||
/// The number of connection attempts currently outstanding.
|
||||
private var pendingConnectionCount: Int
|
||||
|
||||
/// The number of connections that have been handed out to users and are in active use.
|
||||
private(set) var leasedConnectionCount: Int
|
||||
|
||||
/// Whether this connection pool is "leaky".
|
||||
///
|
||||
/// The difference between a leaky and non-leaky connection pool is their behaviour when the pool is currently
|
||||
/// entirely in-use. For a leaky pool, if a connection is requested and none are available, a new connection attempt
|
||||
/// will be made and the connection will be passed to the user. For a non-leaky pool, the user will wait for a connection
|
||||
/// to be returned to the pool.
|
||||
internal let leaky: Bool
|
||||
|
||||
/// The current state of this connection pool.
|
||||
private var state: State
|
||||
|
||||
@@ -85,49 +73,48 @@ internal final class ConnectionPool {
|
||||
return self.availableConnections.count + self.pendingConnectionCount + self.leasedConnectionCount
|
||||
}
|
||||
|
||||
/// Whether a connection can be added into the availableConnections pool when it's returned. This is true
|
||||
/// for non-leaky pools if the sum of availableConnections and leased connections is less than max connections,
|
||||
/// and for leaky pools if the number of availableConnections is less than max connections (as we went to all
|
||||
/// the effort to create the connection, we may as well keep it).
|
||||
/// Note that this means connection attempts in flight may not be used for anything. This is ok!
|
||||
/// Whether a connection can be added into the availableConnections pool when it's returned.
|
||||
private var canAddConnectionToPool: Bool {
|
||||
if self.leaky {
|
||||
switch self.maxConnectionCountBehavior {
|
||||
// only if the current available count is less than the max
|
||||
case .elastic:
|
||||
return self.availableConnections.count < self.maximumConnectionCount
|
||||
} else {
|
||||
|
||||
// only if the total connections count is less than the max
|
||||
case .strict:
|
||||
return (self.availableConnections.count + self.leasedConnectionCount) < self.maximumConnectionCount
|
||||
}
|
||||
}
|
||||
|
||||
internal init(
|
||||
maximumConnectionCount: Int,
|
||||
minimumConnectionCount: Int,
|
||||
leaky: Bool,
|
||||
maximumConnectionCount: Int,
|
||||
maxConnectionCountBehavior: RedisConnectionPool.ConnectionCountBehavior.MaxConnectionBehavior,
|
||||
connectionRetryStrategy: RedisConnectionPool.PoolConnectionRetryStrategy,
|
||||
loop: EventLoop,
|
||||
poolLogger: Logger,
|
||||
connectionBackoffFactor: Float32 = 2,
|
||||
initialConnectionBackoffDelay: TimeAmount = .milliseconds(100),
|
||||
connectionFactory: @escaping (EventLoop) -> EventLoopFuture<RedisConnection>
|
||||
) {
|
||||
guard minimumConnectionCount <= maximumConnectionCount else {
|
||||
self.minimumConnectionCount = minimumConnectionCount
|
||||
self.maximumConnectionCount = maximumConnectionCount
|
||||
self.maxConnectionCountBehavior = maxConnectionCountBehavior
|
||||
|
||||
guard self.minimumConnectionCount <= self.maximumConnectionCount else {
|
||||
poolLogger.critical("pool's minimum connection count is higher than the maximum")
|
||||
preconditionFailure("Minimum connection count must not exceed maximum")
|
||||
preconditionFailure("minimum connection count must not exceed maximum")
|
||||
}
|
||||
|
||||
self.connectionFactory = connectionFactory
|
||||
self.pendingConnectionCount = 0
|
||||
self.leasedConnectionCount = 0
|
||||
self.availableConnections = []
|
||||
self.availableConnections.reserveCapacity(maximumConnectionCount)
|
||||
self.availableConnections.reserveCapacity(self.maximumConnectionCount)
|
||||
|
||||
// 8 is a good number to skip the first few buffer resizings
|
||||
self.connectionWaiters = CircularBuffer(initialCapacity: 8)
|
||||
self.loop = loop
|
||||
self.backoffFactor = connectionBackoffFactor
|
||||
self.initialBackoffDelay = initialConnectionBackoffDelay
|
||||
self.connectionFactory = connectionFactory
|
||||
self.connectionRetryStrategy = connectionRetryStrategy
|
||||
|
||||
self.maximumConnectionCount = maximumConnectionCount
|
||||
self.minimumConnectionCount = minimumConnectionCount
|
||||
self.pendingConnectionCount = 0
|
||||
self.leasedConnectionCount = 0
|
||||
self.leaky = leaky
|
||||
self.state = .active
|
||||
}
|
||||
|
||||
@@ -154,12 +141,13 @@ internal final class ConnectionPool {
|
||||
}
|
||||
}
|
||||
|
||||
func leaseConnection(deadline: NIODeadline, logger: Logger) -> EventLoopFuture<RedisConnection> {
|
||||
func leaseConnection(logger: Logger, deadline: NIODeadline? = nil) -> EventLoopFuture<RedisConnection> {
|
||||
let deadline = deadline ?? .now() + self.connectionRetryStrategy.timeout
|
||||
if self.loop.inEventLoop {
|
||||
return self._leaseConnection(deadline, logger: logger)
|
||||
return self._leaseConnection(logger: logger, deadline: deadline)
|
||||
} else {
|
||||
return self.loop.flatSubmit {
|
||||
return self._leaseConnection(deadline, logger: logger)
|
||||
return self._leaseConnection(logger: logger, deadline: deadline)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -191,12 +179,16 @@ extension ConnectionPool {
|
||||
RedisLogging.MetadataKeys.connectionCount: "\(neededConnections)"
|
||||
])
|
||||
while neededConnections > 0 {
|
||||
self._createConnection(backoff: self.initialBackoffDelay, startIn: .nanoseconds(0), logger: logger)
|
||||
self._createConnection(
|
||||
retryDelay: self.connectionRetryStrategy.initialDelay,
|
||||
startIn: .nanoseconds(0),
|
||||
logger: logger
|
||||
)
|
||||
neededConnections -= 1
|
||||
}
|
||||
}
|
||||
|
||||
private func _createConnection(backoff: TimeAmount, startIn delay: TimeAmount, logger: Logger) {
|
||||
private func _createConnection(retryDelay: TimeAmount, startIn delay: TimeAmount, logger: Logger) {
|
||||
self.loop.assertInEventLoop()
|
||||
self.pendingConnectionCount += 1
|
||||
|
||||
@@ -212,7 +204,7 @@ extension ConnectionPool {
|
||||
self.connectionCreationSucceeded(connection, logger: logger)
|
||||
|
||||
case .failure(let error):
|
||||
self.connectionCreationFailed(error, backoff: backoff, logger: logger)
|
||||
self.connectionCreationFailed(error, retryDelay: retryDelay, logger: logger)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -243,7 +235,7 @@ extension ConnectionPool {
|
||||
}
|
||||
}
|
||||
|
||||
private func connectionCreationFailed(_ error: Error, backoff: TimeAmount, logger: Logger) {
|
||||
private func connectionCreationFailed(_ error: Error, retryDelay: TimeAmount, logger: Logger) {
|
||||
self.loop.assertInEventLoop()
|
||||
|
||||
logger.warning("failed to create connection for pool", metadata: [
|
||||
@@ -260,15 +252,17 @@ extension ConnectionPool {
|
||||
// for this connection. Waiters can time out: if they do, we can just give up this connection.
|
||||
// We know folks need this in the following conditions:
|
||||
//
|
||||
// 1. For non-leaky buckets, we need this reconnection if there are any waiters AND the number of active connections (which includes
|
||||
// 1. For non-elastic buckets, we need this reconnection if there are any waiters AND the number of active connections (which includes
|
||||
// pending connection attempts) is less than max connections
|
||||
// 2. For leaky buckets, we need this reconnection if connectionWaiters.count is greater than the number of pending connection attempts.
|
||||
// 2. For elastic buckets, we need this reconnection if connectionWaiters.count is greater than the number of pending connection attempts.
|
||||
// 3. For either kind, if the number of active connections is less than the minimum.
|
||||
let shouldReconnect: Bool
|
||||
if self.leaky {
|
||||
switch self.maxConnectionCountBehavior {
|
||||
case .elastic:
|
||||
shouldReconnect = (self.connectionWaiters.count > self.pendingConnectionCount)
|
||||
|| (self.minimumConnectionCount > self.activeConnectionCount)
|
||||
} else {
|
||||
|
||||
case .strict:
|
||||
shouldReconnect = (!self.connectionWaiters.isEmpty && self.maximumConnectionCount > self.activeConnectionCount)
|
||||
|| (self.minimumConnectionCount > self.activeConnectionCount)
|
||||
}
|
||||
@@ -279,12 +273,12 @@ extension ConnectionPool {
|
||||
}
|
||||
|
||||
// Ok, we need the new connection.
|
||||
let newBackoff = TimeAmount.nanoseconds(Int64(Float32(backoff.nanoseconds) * self.backoffFactor))
|
||||
let nextRetryDelay = self.connectionRetryStrategy.determineNewDelay(currentDelay: retryDelay)
|
||||
logger.debug("reconnecting after failed connection attempt", metadata: [
|
||||
RedisLogging.MetadataKeys.poolConnectionRetryBackoff: "\(backoff)ns",
|
||||
RedisLogging.MetadataKeys.poolConnectionRetryNewBackoff: "\(newBackoff)ns"
|
||||
RedisLogging.MetadataKeys.poolConnectionRetryAmount: "\(retryDelay)ns",
|
||||
RedisLogging.MetadataKeys.poolConnectionRetryNewAmount: "\(nextRetryDelay)ns"
|
||||
])
|
||||
self._createConnection(backoff: newBackoff, startIn: backoff, logger: logger)
|
||||
self._createConnection(retryDelay: nextRetryDelay, startIn: retryDelay, logger: logger)
|
||||
}
|
||||
|
||||
/// A connection that was monitored by this pool has been closed.
|
||||
@@ -352,7 +346,7 @@ extension ConnectionPool {
|
||||
|
||||
/// This is the on-thread implementation for leasing connections out to users. Here we work out how to get a new
|
||||
/// connection, and attempt to do so.
|
||||
private func _leaseConnection(_ deadline: NIODeadline, logger: Logger) -> EventLoopFuture<RedisConnection> {
|
||||
private func _leaseConnection(logger: Logger, deadline: NIODeadline) -> EventLoopFuture<RedisConnection> {
|
||||
self.loop.assertInEventLoop()
|
||||
|
||||
guard case .active = self.state else {
|
||||
@@ -386,11 +380,22 @@ extension ConnectionPool {
|
||||
self.connectionWaiters.append(waiter)
|
||||
|
||||
// Ok, we have connection targets. If the number of active connections is
|
||||
// below the max, or the pool is leaky, we can create a new connection. Otherwise, we just have
|
||||
// below the max, or the pool is elastic, we can create a new connection. Otherwise, we just have
|
||||
// to wait for a connection to come back.
|
||||
if self.activeConnectionCount < self.maximumConnectionCount || self.leaky {
|
||||
|
||||
let shouldCreateConnection: Bool
|
||||
switch self.maxConnectionCountBehavior {
|
||||
case .elastic: shouldCreateConnection = true
|
||||
case .strict: shouldCreateConnection = false
|
||||
}
|
||||
|
||||
if self.activeConnectionCount < self.maximumConnectionCount || shouldCreateConnection {
|
||||
logger.trace("creating new connection")
|
||||
self._createConnection(backoff: self.initialBackoffDelay, startIn: .nanoseconds(0), logger: logger)
|
||||
self._createConnection(
|
||||
retryDelay: self.connectionRetryStrategy.initialDelay,
|
||||
startIn: .nanoseconds(0),
|
||||
logger: logger
|
||||
)
|
||||
}
|
||||
|
||||
return waiter.futureResult
|
||||
|
||||
+2
-112
@@ -2,7 +2,7 @@
|
||||
//
|
||||
// This source file is part of the RediStack open source project
|
||||
//
|
||||
// Copyright (c) 2020 RediStack project authors
|
||||
// Copyright (c) 2020-2022 RediStack project authors
|
||||
// Licensed under Apache License v2.0
|
||||
//
|
||||
// See LICENSE.txt for license information
|
||||
@@ -65,7 +65,7 @@ extension RedisConnection {
|
||||
/// The default port that Redis uses.
|
||||
///
|
||||
/// See [https://redis.io/topics/quickstart](https://redis.io/topics/quickstart)
|
||||
public static var defaultPort = 6379
|
||||
public static var defaultPort: Int { 6379 }
|
||||
|
||||
internal static let defaultLogger = Logger.redisBaseConnectionLogger
|
||||
|
||||
@@ -203,113 +203,3 @@ extension RedisConnection {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - RedisConnectionPool Config
|
||||
|
||||
extension RedisConnectionPool {
|
||||
/// A configuration object for creating Redis connections with a connection pool.
|
||||
/// - Warning: This type has **reference** semantics due to the `NIO.ClientBootstrap` reference.
|
||||
public struct ConnectionFactoryConfiguration {
|
||||
// this needs to be var so it can be updated by the pool with the pool id
|
||||
/// The logger prototype that will be used by connections by default when generating logs.
|
||||
public internal(set) var connectionDefaultLogger: Logger
|
||||
/// The password used to authenticate connections.
|
||||
public let connectionPassword: String?
|
||||
/// The initial database index that connections should use.
|
||||
public let connectionInitialDatabase: Int?
|
||||
/// The pre-configured TCP client for connections to use.
|
||||
public let tcpClient: ClientBootstrap?
|
||||
|
||||
/// Creates a new connection factory configuration with the provided options.
|
||||
/// - Parameters:
|
||||
/// - connectionInitialDatabase: The optional database index to initially connect to. The default is `nil`.
|
||||
/// Redis by default opens connections against index `0`, so only set this value if the desired default is not `0`.
|
||||
/// - connectionPassword: The optional password to authenticate connections with. The default is `nil`.
|
||||
/// - connectionDefaultLogger: The optional prototype logger to use as the default logger instance when generating logs from connections.
|
||||
/// If one is not provided, one will be generated. See `RedisLogging.baseConnectionLogger`.
|
||||
/// - tcpClient: If you have chosen to configure a `NIO.ClientBootstrap` yourself, this will be used instead of the `.makeRedisTCPClient` factory instance.
|
||||
public init(
|
||||
connectionInitialDatabase: Int? = nil,
|
||||
connectionPassword: String? = nil,
|
||||
connectionDefaultLogger: Logger? = nil,
|
||||
tcpClient: ClientBootstrap? = nil
|
||||
) {
|
||||
self.connectionInitialDatabase = connectionInitialDatabase
|
||||
self.connectionPassword = connectionPassword
|
||||
self.connectionDefaultLogger = connectionDefaultLogger ?? RedisConnection.Configuration.defaultLogger
|
||||
self.tcpClient = tcpClient
|
||||
}
|
||||
}
|
||||
|
||||
/// A configuration object for connection pools.
|
||||
/// - Warning: This type has **reference** semantics due to `ConnectionFactoryConfiguration`.
|
||||
public struct Configuration {
|
||||
/// The set of Redis servers to which this pool is initially willing to connect.
|
||||
public let initialConnectionAddresses: [SocketAddress]
|
||||
/// The minimum number of connections to preserve in the pool.
|
||||
///
|
||||
/// If the pool is mostly idle and the Redis servers close these idle connections,
|
||||
/// the `RedisConnectionPool` will initiate new outbound connections proactively to avoid the number of available connections dropping below this number.
|
||||
public let minimumConnectionCount: Int
|
||||
/// The maximum number of connections to for this pool, either to be preserved or as a hard limit.
|
||||
public let maximumConnectionCount: RedisConnectionPoolSize
|
||||
/// The configuration object that controls the connection retry behavior.
|
||||
public let connectionRetryConfiguration: (backoff: (initialDelay: TimeAmount, factor: Float32), timeout: TimeAmount)
|
||||
// these need to be var so they can be updated by the pool in some cases
|
||||
public internal(set) var factoryConfiguration: ConnectionFactoryConfiguration
|
||||
/// The logger prototype that will be used by the connection pool by default when generating logs.
|
||||
public internal(set) var poolDefaultLogger: Logger
|
||||
|
||||
/// Creates a new connection configuration with the provided options.
|
||||
/// - Parameters:
|
||||
/// - initialServerConnectionAddresses: The set of Redis servers to which this pool is initially willing to connect.
|
||||
/// This set can be updated over time directly on the connection pool.
|
||||
/// - maximumConnectionCount: The maximum number of connections to for this pool, either to be preserved or as a hard limit.
|
||||
/// - connectionFactoryConfiguration: The configuration to use while creating connections to fill the pool.
|
||||
/// - minimumConnectionCount: The minimum number of connections to preserve in the pool. If the pool is mostly idle
|
||||
/// and the Redis servers close these idle connections, the `RedisConnectionPool` will initiate new outbound
|
||||
/// connections proactively to avoid the number of available connections dropping below this number. Defaults to `1`.
|
||||
/// - connectionBackoffFactor: Used when connection attempts fail to control the exponential backoff. This is a multiplicative
|
||||
/// factor, each connection attempt will be delayed by this amount times the previous delay.
|
||||
/// - initialConnectionBackoffDelay: If a TCP connection attempt fails, this is the first backoff value on the reconnection attempt.
|
||||
/// Subsequent backoffs are computed by compounding this value by `connectionBackoffFactor`.
|
||||
/// - connectionRetryTimeout: The max time to wait for a connection to be available before failing a particular command or connection operation.
|
||||
/// The default is 60 seconds.
|
||||
/// - poolDefaultLogger: The `Logger` used by the connection pool itself.
|
||||
public init(
|
||||
initialServerConnectionAddresses: [SocketAddress],
|
||||
maximumConnectionCount: RedisConnectionPoolSize,
|
||||
connectionFactoryConfiguration: ConnectionFactoryConfiguration,
|
||||
minimumConnectionCount: Int = 1,
|
||||
connectionBackoffFactor: Float32 = 2,
|
||||
initialConnectionBackoffDelay: TimeAmount = .milliseconds(100),
|
||||
connectionRetryTimeout: TimeAmount? = .seconds(60),
|
||||
poolDefaultLogger: Logger? = nil
|
||||
) {
|
||||
self.initialConnectionAddresses = initialServerConnectionAddresses
|
||||
self.maximumConnectionCount = maximumConnectionCount
|
||||
self.factoryConfiguration = connectionFactoryConfiguration
|
||||
self.minimumConnectionCount = minimumConnectionCount
|
||||
self.connectionRetryConfiguration = (
|
||||
(initialConnectionBackoffDelay, connectionBackoffFactor),
|
||||
connectionRetryTimeout ?? .milliseconds(10) // always default to a baseline 10ms
|
||||
)
|
||||
self.poolDefaultLogger = poolDefaultLogger ?? .redisBaseConnectionPoolLogger
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// `RedisConnectionPoolSize` controls how the maximum number of connections in a pool are interpreted.
|
||||
public enum RedisConnectionPoolSize {
|
||||
/// The pool will allow no more than this number of connections to be "active" (that is, connecting, in-use,
|
||||
/// or pooled) at any one time. This will force possible future users of new connections to wait until a currently
|
||||
/// active connection becomes available by being returned to the pool, but provides a hard upper limit on concurrency.
|
||||
case maximumActiveConnections(Int)
|
||||
|
||||
/// The pool will only store up to this number of connections that are not currently in-use. However, if the pool is
|
||||
/// asked for more connections at one time than this number, it will create new connections to serve those waiting for
|
||||
/// connections. These "extra" connections will not be preserved: while they will be used to satisfy those waiting for new
|
||||
/// connections if needed, they will not be preserved in the pool if load drops low enough. This does not provide a hard
|
||||
/// upper bound on concurrency, but does provide an upper bound on low-level load.
|
||||
case maximumPreservedConnections(Int)
|
||||
}
|
||||
@@ -0,0 +1,248 @@
|
||||
//===----------------------------------------------------------------------===//
|
||||
//
|
||||
// This source file is part of the RediStack open source project
|
||||
//
|
||||
// Copyright (c) 2022 RediStack project authors
|
||||
// Licensed under Apache License v2.0
|
||||
//
|
||||
// See LICENSE.txt for license information
|
||||
// See CONTRIBUTORS.txt for the list of RediStack project authors
|
||||
//
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
//
|
||||
//===----------------------------------------------------------------------===//
|
||||
|
||||
import Logging
|
||||
import NIO
|
||||
|
||||
// MARK: - Pool Connection Congfiguration
|
||||
|
||||
extension RedisConnectionPool {
|
||||
/// A configuration object for a connection pool to use when creating Redis connections.
|
||||
/// - Warning: This type has **reference** semantics due to the `NIO.ClientBootstrap` reference.
|
||||
public struct PoolConnectionConfiguration {
|
||||
// this needs to be var so it can be updated by the pool with the pool id
|
||||
/// The logger that will be used by connections by default when generating logs.
|
||||
public internal(set) var defaultLogger: Logger
|
||||
/// The password used to authenticate connections.
|
||||
public let password: String?
|
||||
/// The initial database index that connections should use.
|
||||
public let initialDatabase: Int?
|
||||
/// The pre-configured TCP client for connections to use.
|
||||
public let tcpClient: ClientBootstrap?
|
||||
|
||||
/// Creates a new connection factory configuration with the provided options.
|
||||
/// - Parameters:
|
||||
/// - initialDatabase: The optional database index to initially connect to. The default is `nil`.
|
||||
/// Redis by default opens connections against index `0`, so only set this value if the desired default is not `0`.
|
||||
/// - password: The optional password to authenticate connections with. The default is `nil`.
|
||||
/// - defaultLogger: The optional prototype logger to use as the default logger instance when generating logs from connections.
|
||||
/// If one is not provided, one will be generated. See ``RedisLogging/baseConnectionLogger``.
|
||||
/// - tcpClient: If you have chosen to configure a `NIO.ClientBootstrap` yourself, this will be used instead of the `.makeRedisTCPClient` factory instance.
|
||||
public init(
|
||||
initialDatabase: Int? = nil,
|
||||
password: String? = nil,
|
||||
defaultLogger: Logger? = nil,
|
||||
tcpClient: ClientBootstrap? = nil
|
||||
) {
|
||||
self.initialDatabase = initialDatabase
|
||||
self.password = password
|
||||
self.defaultLogger = defaultLogger ?? RedisConnection.Configuration.defaultLogger
|
||||
self.tcpClient = tcpClient
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Connection Count Behavior
|
||||
|
||||
extension RedisConnectionPool {
|
||||
/// The desired behavior for a connection pool to maintain its pool of "active" (connecting, in-use, or pooled) connections.
|
||||
public struct ConnectionCountBehavior {
|
||||
/// The pool will allow no more than the specified maximum number of connections to be "active" at any given time.
|
||||
///
|
||||
/// This will force possible future users of connections to wait until an "active" connection becomes available
|
||||
/// by being returned to the pool.
|
||||
///
|
||||
/// In other words, this provides a hard upper limit on concurrency.
|
||||
/// - Parameters:
|
||||
/// - maximumConnectionCount: The maximum number of connections to preserve in the pool.
|
||||
/// - minimumConnectionCount: The minimum number of connections to preserve in the pool. The default is `1`.
|
||||
public static func strict(maximumConnectionCount: Int, minimumConnectionCount: Int = 1) -> Self {
|
||||
return .init(min: minimumConnectionCount, max: maximumConnectionCount, behavior: .strict)
|
||||
}
|
||||
|
||||
/// The pool will maintain the specified number of maxiumum connections,
|
||||
/// but will create more as needed based on demand.
|
||||
///
|
||||
/// Connections created to meet demaind are treated as "extra" connections,
|
||||
/// and will not be preserved after demand has reached below the specified ``maximumConnectionCount``.
|
||||
///
|
||||
/// In other words, this does not provide a hard upper bound on concurrency, but does provide an upper bound on low-level load.
|
||||
/// - Parameters:
|
||||
/// - maximumConnectionCount: The maximum number of connections to preserve in the pool.
|
||||
/// - minimumConnectionCount: The minimum number of connections to preserve in the pool. The default is `1`.
|
||||
public static func elastic(maximumConnectionCount: Int, minimumConnectionCount: Int = 1) -> Self {
|
||||
return .init(min: minimumConnectionCount, max: maximumConnectionCount, behavior: .elastic)
|
||||
}
|
||||
|
||||
/// Is the pool's maxiumum connection count elastic, allowing for additional on-demand connections?
|
||||
public var isElastic: Bool { self.maxConnectionBehavior == .elastic }
|
||||
|
||||
/// The minimum number of connections to preserve in the pool.
|
||||
///
|
||||
/// If the pool is mostly idle and the Redis servers close these idle connections,
|
||||
/// the ``RedisConnectionPool`` will initiate new outbound connections proactively
|
||||
/// to avoid the number of available connections dropping below this number.
|
||||
public let minimumConnectionCount: Int
|
||||
/// The maximum number of connections to preserve in the pool.
|
||||
///
|
||||
/// The actual maximum number of connections created by a pool could exceed this value,
|
||||
/// based on if the behavior ``isElastic``.
|
||||
public let maximumConnectionCount: Int
|
||||
|
||||
internal let maxConnectionBehavior: MaxConnectionBehavior
|
||||
internal enum MaxConnectionBehavior {
|
||||
case strict
|
||||
case elastic
|
||||
}
|
||||
|
||||
private init(min: Int, max: Int, behavior: MaxConnectionBehavior) {
|
||||
|
||||
self.minimumConnectionCount = min
|
||||
self.maximumConnectionCount = max
|
||||
self.maxConnectionBehavior = behavior
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Connection Retry Strategy
|
||||
|
||||
extension TimeAmount {
|
||||
fileprivate static var minimumTimeoutTolerance: Self { .milliseconds(10) }
|
||||
}
|
||||
|
||||
extension RedisConnectionPool {
|
||||
/// A definition of how a given connection pool will attempt to retry fulfilling requests for connections.
|
||||
///
|
||||
/// Each strategy defines an ``initialDelay`` that will be waited before asking again for a connection.
|
||||
///
|
||||
/// After that `initialDelay`, then the strategy's ``DeadlineProvider`` will be called
|
||||
/// to provide a new delay value to wait.
|
||||
///
|
||||
/// The strategy will continue to execute to fulfill a connection request until either
|
||||
/// the ``timeout`` is reached or a connection is made available.
|
||||
/// - Important: All `timeout` values are clamped to a minimum tolerance level to avoid false negative timeouts,
|
||||
/// as there is a slight overhead to the connection pool's logic for finding available connections.
|
||||
public struct PoolConnectionRetryStrategy {
|
||||
/// A closure that receives the current retry delay and returns a new delay value to wait.
|
||||
public typealias DeadlineProvider = (TimeAmount) -> TimeAmount
|
||||
|
||||
/// The default timeout strategies will use. The value is `.seconds(60)`.
|
||||
public static var defaultTimeout: TimeAmount { .seconds(60) }
|
||||
|
||||
/// Requests for a connection from a pool will exponentially backoff polling the pool, or timeout.
|
||||
/// - Parameters:
|
||||
/// - initialDelay: The initial delay of further retries. The default is `.milliseconds(100)`.
|
||||
/// - backoffFactor: The factor to multiply the current backoff amount by when additional retries are made. The default is `2`.
|
||||
/// - timeout: The maximum amount of time to wait before retrying ends. The default is ``defaultTimeout``.
|
||||
public static func exponentialBackoff(
|
||||
initialDelay: TimeAmount = .milliseconds(100),
|
||||
backoffFactor: Float32 = 2,
|
||||
timeout: TimeAmount = Self.defaultTimeout
|
||||
) -> Self {
|
||||
return .init(
|
||||
initialDelay: initialDelay,
|
||||
timeout: timeout,
|
||||
{ return .nanoseconds(Int64(Float32($0.nanoseconds) * backoffFactor)) }
|
||||
)
|
||||
}
|
||||
|
||||
/// No retrying will occur. Requests for a connection will fail immediately if a connection is not available.
|
||||
public static var none: Self { return .none(timeout: .minimumTimeoutTolerance) }
|
||||
|
||||
/// No retrying will occur. Requests for a connection will fail if a connection
|
||||
/// is not available after the specified timeout.
|
||||
/// - Parameter timeout: The maximum amount of time to wait before failing requests for a connection.
|
||||
public static func none(timeout: TimeAmount) -> Self {
|
||||
return .init(initialDelay: .zero, timeout: timeout, { _ in .zero })
|
||||
}
|
||||
|
||||
public let initialDelay: TimeAmount
|
||||
public let timeout: TimeAmount
|
||||
private let deadlineProvider: DeadlineProvider
|
||||
|
||||
/// Creates a strategy with a given initial delay value, a closure to calculate future delay values,
|
||||
/// and a timeout to provide an upper limit on waiting.
|
||||
/// - Parameters:
|
||||
/// - initialDelay: The initial time to wait before retrying.
|
||||
/// - timeout: The total time to wait before failing the request for a connection and cancelling retrying.
|
||||
///
|
||||
/// The value provided is clamped to a minimum tolerance level to avoid false negative timeouts,
|
||||
/// as there is a slight overhead to the connection pool's logic for finding available connections.
|
||||
/// - deadlineProvider: A method of calulating new delay values, given the current delay.
|
||||
public init(
|
||||
initialDelay: TimeAmount,
|
||||
timeout: TimeAmount = Self.defaultTimeout,
|
||||
_ deadlineProvider: @escaping DeadlineProvider
|
||||
) {
|
||||
self.initialDelay = initialDelay
|
||||
// because there's some overhead in the connection pooling logic,
|
||||
// we want a baseline minimum tolerance so we don't always have false immediate timeouts
|
||||
self.timeout = timeout <= .minimumTimeoutTolerance ? .minimumTimeoutTolerance : timeout
|
||||
self.deadlineProvider = deadlineProvider
|
||||
}
|
||||
|
||||
/// Determines the new delay amount to wait before the next retry attempt.
|
||||
/// - Parameter currentDelay: The current delay value that was waited before the current retry attempt.
|
||||
/// - Returns: A new delay value to wait before the next retry attempt.
|
||||
public func determineNewDelay(currentDelay: TimeAmount) -> TimeAmount {
|
||||
return self.deadlineProvider(currentDelay)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Pool Configuration
|
||||
|
||||
extension RedisConnectionPool {
|
||||
/// A configuration object for connection pools.
|
||||
/// - Warning: This type has **reference** semantics due to ``PoolConnectionConfiguration``.
|
||||
public struct Configuration {
|
||||
/// The set of Redis servers to which this pool is initially willing to connect.
|
||||
public let initialConnectionAddresses: [SocketAddress]
|
||||
/// The behavior the pool should use for maintaining its pool of "active" connections and providing connections upon request.
|
||||
public let connectionCountBehavior: ConnectionCountBehavior
|
||||
/// The strategy used by the connection pool to handle retrying to find an available "active" connection to use.
|
||||
public let retryStrategy: PoolConnectionRetryStrategy
|
||||
|
||||
// these need to be var so they can be updated by the pool in some cases
|
||||
|
||||
/// The configuration used when creating connections.
|
||||
public internal(set) var connectionConfiguration: PoolConnectionConfiguration
|
||||
/// The logger prototype that will be used by the connection pool by default when generating logs.
|
||||
public internal(set) var poolDefaultLogger: Logger
|
||||
|
||||
/// Creates a new connection configuration with the provided options.
|
||||
/// - Parameters:
|
||||
/// - initialServerConnectionAddresses: The set of Redis servers to which this pool is initially willing to connect.
|
||||
/// This set can be updated over time directly on the connection pool.
|
||||
/// - connectionCountBehavior: The behavior used by the pool for maintaining it's count of connections.
|
||||
/// - connectionConfiguration: The configuration to use when creating connections to fill the pool.
|
||||
/// - retryStrategy: The retry strategy to apply while waiting for connections to become available for a particular command or connection operation.
|
||||
///
|
||||
/// The default is ``RedisConnectionPool/PoolConnectionRetryStrategy/exponentialBackoff(initialDelay:backoffFactor:timeout:)`` with default values.
|
||||
/// - poolDefaultLogger: The `Logger` used by the connection pool itself.
|
||||
public init(
|
||||
initialServerConnectionAddresses: [SocketAddress],
|
||||
connectionCountBehavior: ConnectionCountBehavior,
|
||||
connectionConfiguration: PoolConnectionConfiguration,
|
||||
retryStrategy: PoolConnectionRetryStrategy = .exponentialBackoff(),
|
||||
poolDefaultLogger: Logger? = nil
|
||||
) {
|
||||
self.initialConnectionAddresses = initialServerConnectionAddresses
|
||||
self.connectionCountBehavior = connectionCountBehavior
|
||||
self.connectionConfiguration = connectionConfiguration
|
||||
self.retryStrategy = retryStrategy
|
||||
self.poolDefaultLogger = poolDefaultLogger ?? .redisBaseConnectionPoolLogger
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -73,13 +73,10 @@ public class RedisConnectionPool {
|
||||
|
||||
self.loop = boundEventLoop
|
||||
self.serverConnectionAddresses = ConnectionAddresses(initialAddresses: config.initialConnectionAddresses)
|
||||
|
||||
// mix of terminology here with the loggers
|
||||
// as we're being "forward thinking" in terms of the 'baggage context' future type
|
||||
|
||||
var taggedConnectionLogger = config.factoryConfiguration.connectionDefaultLogger
|
||||
|
||||
var taggedConnectionLogger = config.connectionConfiguration.defaultLogger
|
||||
taggedConnectionLogger[metadataKey: RedisLogging.MetadataKeys.connectionPoolID] = "\(self.id)"
|
||||
config.factoryConfiguration.connectionDefaultLogger = taggedConnectionLogger
|
||||
config.connectionConfiguration.defaultLogger = taggedConnectionLogger
|
||||
|
||||
var taggedPoolLogger = config.poolDefaultLogger
|
||||
taggedPoolLogger[metadataKey: RedisLogging.MetadataKeys.connectionPoolID] = "\(self.id)"
|
||||
@@ -88,13 +85,12 @@ public class RedisConnectionPool {
|
||||
self.configuration = config
|
||||
|
||||
self.pool = ConnectionPool(
|
||||
maximumConnectionCount: config.maximumConnectionCount.size,
|
||||
minimumConnectionCount: config.minimumConnectionCount,
|
||||
leaky: config.maximumConnectionCount.leaky,
|
||||
minimumConnectionCount: self.configuration.connectionCountBehavior.minimumConnectionCount,
|
||||
maximumConnectionCount: self.configuration.connectionCountBehavior.maximumConnectionCount,
|
||||
maxConnectionCountBehavior: self.configuration.connectionCountBehavior.maxConnectionBehavior,
|
||||
connectionRetryStrategy: self.configuration.retryStrategy,
|
||||
loop: boundEventLoop,
|
||||
poolLogger: config.poolDefaultLogger,
|
||||
connectionBackoffFactor: config.connectionRetryConfiguration.backoff.factor,
|
||||
initialConnectionBackoffDelay: config.connectionRetryConfiguration.backoff.initialDelay,
|
||||
connectionFactory: self.connectionFactory(_:)
|
||||
)
|
||||
}
|
||||
@@ -249,8 +245,6 @@ extension RedisConnectionPool {
|
||||
// Validate the loop invariants.
|
||||
self.loop.preconditionInEventLoop()
|
||||
targetLoop.preconditionInEventLoop()
|
||||
|
||||
let factoryConfig = self.configuration.factoryConfiguration
|
||||
|
||||
guard let nextTarget = self.serverConnectionAddresses.nextTarget() else {
|
||||
// No valid connection target, we'll keep track of the request and attempt to satisfy it later.
|
||||
@@ -270,9 +264,7 @@ extension RedisConnectionPool {
|
||||
do {
|
||||
connectionConfig = try .init(
|
||||
address: nextTarget,
|
||||
password: factoryConfig.connectionPassword,
|
||||
initialDatabase: factoryConfig.connectionInitialDatabase,
|
||||
defaultLogger: factoryConfig.connectionDefaultLogger
|
||||
prototypeConfiguration: self.configuration.connectionConfiguration
|
||||
)
|
||||
} catch {
|
||||
// config validation failed, return the error
|
||||
@@ -283,7 +275,7 @@ extension RedisConnectionPool {
|
||||
.make(
|
||||
configuration: connectionConfig,
|
||||
boundEventLoop: targetLoop,
|
||||
configuredTCPClient: factoryConfig.tcpClient
|
||||
configuredTCPClient: self.configuration.connectionConfiguration.tcpClient
|
||||
)
|
||||
.map { connection in
|
||||
// disallow subscriptions on all connections by default to enforce our management of PubSub state
|
||||
@@ -533,10 +525,7 @@ extension RedisConnectionPool: RedisClient {
|
||||
|
||||
guard let connection = preferredConnection else {
|
||||
return pool
|
||||
.leaseConnection(
|
||||
deadline: .now() + self.configuration.connectionRetryConfiguration.timeout,
|
||||
logger: logger
|
||||
)
|
||||
.leaseConnection(logger: logger)
|
||||
.flatMap { operation($0, pool.returnConnection(_:logger:), logger) }
|
||||
}
|
||||
|
||||
@@ -544,6 +533,19 @@ extension RedisConnectionPool: RedisClient {
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: Helper for creating connection configs
|
||||
|
||||
extension RedisConnection.Configuration {
|
||||
fileprivate init(address: SocketAddress, prototypeConfiguration: RedisConnectionPool.PoolConnectionConfiguration) throws {
|
||||
try self.init(
|
||||
address: address,
|
||||
password: prototypeConfiguration.password,
|
||||
initialDatabase: prototypeConfiguration.initialDatabase,
|
||||
defaultLogger: prototypeConfiguration.defaultLogger
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: Helper for round-robin connection establishment
|
||||
extension RedisConnectionPool {
|
||||
/// A helper structure for valid connection addresses. This structure implements round-robin connection establishment.
|
||||
@@ -581,22 +583,3 @@ extension RedisConnectionPool {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: RedisConnectionPoolSize helpers
|
||||
extension RedisConnectionPoolSize {
|
||||
fileprivate var size: Int {
|
||||
switch self {
|
||||
case .maximumActiveConnections(let size), .maximumPreservedConnections(let size):
|
||||
return size
|
||||
}
|
||||
}
|
||||
|
||||
fileprivate var leaky: Bool {
|
||||
switch self {
|
||||
case .maximumActiveConnections:
|
||||
return false
|
||||
case .maximumPreservedConnections:
|
||||
return true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,8 +45,8 @@ public enum RedisLogging {
|
||||
internal static var command: String { "rdstk_command" }
|
||||
internal static var commandResult: String { "rdstk_result" }
|
||||
internal static var connectionCount: String { "rdstk_conn_count" }
|
||||
internal static var poolConnectionRetryBackoff: String { "rdstk_conn_retry_prev_backoff" }
|
||||
internal static var poolConnectionRetryNewBackoff: String { "rdstk_conn_retry_new_backoff" }
|
||||
internal static var poolConnectionRetryAmount: String { "rdstk_conn_retry_prev_amount" }
|
||||
internal static var poolConnectionRetryNewAmount: String { "rdstk_conn_retry_new_amount" }
|
||||
internal static var poolConnectionCount: String { "rdstk_pool_active_connection_count" }
|
||||
internal static let pubsubTarget = "rdstk_ps_target"
|
||||
internal static let subscriptionCount = "rdstk_sub_count"
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
//
|
||||
// This source file is part of the RediStack open source project
|
||||
//
|
||||
// Copyright (c) 2020 RediStack project authors
|
||||
// Copyright (c) 2020-2022 RediStack project authors
|
||||
// Licensed under Apache License v2.0
|
||||
//
|
||||
// See LICENSE.txt for license information
|
||||
@@ -78,18 +78,16 @@ open class RedisConnectionPoolIntegrationTestCase: XCTestCase {
|
||||
public func makeNewPool(
|
||||
initialAddresses: [SocketAddress]? = nil,
|
||||
initialConnectionBackoffDelay: TimeAmount = .milliseconds(100),
|
||||
connectionRetryTimeout: TimeAmount? = .seconds(5),
|
||||
connectionRetryTimeout: TimeAmount = .seconds(5),
|
||||
minimumConnectionCount: Int = 0
|
||||
) throws -> RedisConnectionPool {
|
||||
let addresses = try initialAddresses ?? [SocketAddress.makeAddressResolvingHost(self.redisHostname, port: self.redisPort)]
|
||||
let pool = RedisConnectionPool(
|
||||
configuration: .init(
|
||||
initialServerConnectionAddresses: addresses,
|
||||
maximumConnectionCount: .maximumActiveConnections(4),
|
||||
connectionFactoryConfiguration: .init(connectionPassword: self.redisPassword),
|
||||
minimumConnectionCount: minimumConnectionCount,
|
||||
initialConnectionBackoffDelay: initialConnectionBackoffDelay,
|
||||
connectionRetryTimeout: connectionRetryTimeout
|
||||
connectionCountBehavior: .strict(maximumConnectionCount: 4, minimumConnectionCount: minimumConnectionCount),
|
||||
connectionConfiguration: .init(password: self.redisPassword),
|
||||
retryStrategy: .exponentialBackoff(initialDelay: initialConnectionBackoffDelay, timeout: connectionRetryTimeout)
|
||||
),
|
||||
boundEventLoop: self.eventLoopGroup.next()
|
||||
)
|
||||
|
||||
@@ -45,7 +45,7 @@ extension RedisConnectionPoolTests {
|
||||
}
|
||||
|
||||
func test_nilConnectionRetryTimeoutStillWorks() throws {
|
||||
let pool = try self.makeNewPool(connectionRetryTimeout: nil)
|
||||
let pool = try self.makeNewPool(connectionRetryTimeout: .zero)
|
||||
defer { pool.close() }
|
||||
XCTAssertNoThrow(try pool.get(#function).wait())
|
||||
}
|
||||
|
||||
@@ -67,8 +67,8 @@ final class RedisLoggingTests: RediStackIntegrationTestCase {
|
||||
let pool = RedisConnectionPool(
|
||||
configuration: .init(
|
||||
initialServerConnectionAddresses: [try .makeAddressResolvingHost(self.redisHostname, port: self.redisPort)],
|
||||
maximumConnectionCount: .maximumActiveConnections(1),
|
||||
connectionFactoryConfiguration: .init(connectionPassword: self.redisPassword)
|
||||
connectionCountBehavior: .strict(maximumConnectionCount: 1),
|
||||
connectionConfiguration: .init(password: self.redisPassword)
|
||||
),
|
||||
boundEventLoop: self.connection.eventLoop
|
||||
)
|
||||
@@ -92,8 +92,8 @@ final class RedisLoggingTests: RediStackIntegrationTestCase {
|
||||
let hosts = InMemoryServiceDiscovery<String, SocketAddress>(configuration: .init())
|
||||
let config = RedisConnectionPool.Configuration(
|
||||
initialServerConnectionAddresses: [],
|
||||
maximumConnectionCount: .maximumActiveConnections(1),
|
||||
connectionFactoryConfiguration: .init(connectionPassword: self.redisPassword)
|
||||
connectionCountBehavior: .strict(maximumConnectionCount: 1),
|
||||
connectionConfiguration: .init(password: self.redisPassword)
|
||||
)
|
||||
let client = RedisConnectionPool.activatedServiceDiscoveryPool(
|
||||
service: "default.local",
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
//
|
||||
// This source file is part of the RediStack open source project
|
||||
//
|
||||
// Copyright (c) 2020 RediStack project authors
|
||||
// Copyright (c) 2020-2022 RediStack project authors
|
||||
// Licensed under Apache License v2.0
|
||||
//
|
||||
// See LICENSE.txt for license information
|
||||
@@ -24,9 +24,8 @@ final class RedisServiceDiscoveryTests: RediStackConnectionPoolIntegrationTestCa
|
||||
let hosts = InMemoryServiceDiscovery<String, SocketAddress>(configuration: .init())
|
||||
let config = RedisConnectionPool.Configuration(
|
||||
initialServerConnectionAddresses: [],
|
||||
maximumConnectionCount: .maximumActiveConnections(5),
|
||||
connectionFactoryConfiguration: .init(connectionPassword: self.redisPassword),
|
||||
minimumConnectionCount: 1
|
||||
connectionCountBehavior: .strict(maximumConnectionCount: 5),
|
||||
connectionConfiguration: .init(password: self.redisPassword)
|
||||
)
|
||||
let client = RedisConnectionPool.activatedServiceDiscoveryPool(
|
||||
service: "default.local",
|
||||
|
||||
@@ -39,23 +39,33 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
return RedisConnection(configuredRESPChannel: channel, defaultLogger: .redisBaseConnectionLogger)
|
||||
}
|
||||
|
||||
func createPool(maximumConnectionCount: Int, minimumConnectionCount: Int, leaky: Bool) -> ConnectionPool {
|
||||
func createPool(
|
||||
maximumConnectionCount: Int,
|
||||
minimumConnectionCount: Int,
|
||||
behavior: RedisConnectionPool.ConnectionCountBehavior.MaxConnectionBehavior
|
||||
) -> ConnectionPool {
|
||||
return ConnectionPool(
|
||||
maximumConnectionCount: maximumConnectionCount,
|
||||
minimumConnectionCount: minimumConnectionCount,
|
||||
leaky: leaky,
|
||||
maximumConnectionCount: maximumConnectionCount,
|
||||
maxConnectionCountBehavior: behavior,
|
||||
connectionRetryStrategy: .exponentialBackoff(),
|
||||
loop: self.server.loop,
|
||||
poolLogger: .redisBaseConnectionPoolLogger
|
||||
) { loop in
|
||||
return loop.makeSucceededFuture(self.createAConnection())
|
||||
}
|
||||
poolLogger: .redisBaseConnectionPoolLogger,
|
||||
connectionFactory: { return $0.makeSucceededFuture(self.createAConnection()) }
|
||||
)
|
||||
}
|
||||
|
||||
func createPool(maximumConnectionCount: Int, minimumConnectionCount: Int, leaky: Bool, connectionFactory: @escaping (EventLoop) -> EventLoopFuture<RedisConnection>) -> ConnectionPool {
|
||||
func createPool(
|
||||
maximumConnectionCount: Int,
|
||||
minimumConnectionCount: Int,
|
||||
behavior: RedisConnectionPool.ConnectionCountBehavior.MaxConnectionBehavior,
|
||||
connectionFactory: @escaping (EventLoop) -> EventLoopFuture<RedisConnection>
|
||||
) -> ConnectionPool {
|
||||
return ConnectionPool(
|
||||
maximumConnectionCount: maximumConnectionCount,
|
||||
minimumConnectionCount: minimumConnectionCount,
|
||||
leaky: leaky,
|
||||
maximumConnectionCount: maximumConnectionCount,
|
||||
maxConnectionCountBehavior: behavior,
|
||||
connectionRetryStrategy: .exponentialBackoff(),
|
||||
loop: self.server.loop,
|
||||
poolLogger: .redisBaseConnectionPoolLogger,
|
||||
connectionFactory: connectionFactory
|
||||
@@ -63,7 +73,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testPoolMaintainsMinimumConnections() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 4, leaky: true)
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 4, behavior: .elastic)
|
||||
XCTAssertNoThrow(try self.server.runWhileActive())
|
||||
XCTAssertEqual(self.server.channels.count, 0)
|
||||
|
||||
@@ -93,7 +103,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testConnectionPoolCanLeaseConnections() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 4, leaky: true)
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 4, behavior: .elastic)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -120,7 +130,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testNonLeakyParallelLease() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 1, leaky: false)
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 1, behavior: .strict)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -179,7 +189,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testLeakyParallelLease() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 1, leaky: true)
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 1, behavior: .elastic)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -231,7 +241,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testReturningClosedConnectionsGetReopened() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false)
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -266,7 +276,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testLeasingFromClosedPoolsFails() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false)
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict)
|
||||
pool.activate()
|
||||
pool.close()
|
||||
|
||||
@@ -276,7 +286,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testNothingBadHappensWhenYouRepeatedlyCloseAPool() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false)
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict)
|
||||
pool.activate()
|
||||
|
||||
// Just spam close
|
||||
@@ -286,7 +296,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testPendingWaitersAreFailedOnPoolClose() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false)
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -322,7 +332,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testConnectionsThatCompleteAfterCloseAreClosed() throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -345,7 +355,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testConnectionsCanFailAfterCloseWithoutIncident() throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -372,7 +382,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testExponentialConnectionBackoff() throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, leaky: false) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 1, behavior: .strict) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -382,7 +392,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
XCTAssertEqual(self.server.channels.count, 0)
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
var delay = pool.initialBackoffDelay
|
||||
var delay = pool.connectionRetryStrategy.initialDelay
|
||||
let oneNanosecond = TimeAmount.nanoseconds(1)
|
||||
for _ in 0..<10 {
|
||||
let promise = connectionPromise
|
||||
@@ -394,7 +404,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
self.server.loop.advanceTime(by: oneNanosecond)
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
delay = .nanoseconds(Int64(Float32(delay.nanoseconds) * pool.backoffFactor))
|
||||
delay = pool.connectionRetryStrategy.determineNewDelay(currentDelay: delay)
|
||||
}
|
||||
|
||||
pool.close()
|
||||
@@ -403,7 +413,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testNonLeakyBucketWillKeepConnectingIfThereIsSpaceAndWaiters() throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, leaky: false) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, behavior: .strict) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -418,7 +428,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
XCTAssertNoThrow(try self.server.runWhileActive())
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
var delay = pool.initialBackoffDelay
|
||||
var delay = pool.connectionRetryStrategy.initialDelay
|
||||
let oneNanosecond = TimeAmount.nanoseconds(1)
|
||||
for _ in 0..<10 {
|
||||
let promise = connectionPromise
|
||||
@@ -430,7 +440,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
self.server.loop.advanceTime(by: oneNanosecond)
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
delay = .nanoseconds(Int64(Float32(delay.nanoseconds) * pool.backoffFactor))
|
||||
delay = pool.connectionRetryStrategy.determineNewDelay(currentDelay: delay)
|
||||
}
|
||||
|
||||
pool.close()
|
||||
@@ -442,7 +452,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testLeakyBucketWillKeepConnectingIfThereAreWaitersEvenIfTheresNoSpace() throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, leaky: true) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, behavior: .elastic) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -469,7 +479,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
XCTAssertNoThrow(try self.server.runWhileActive())
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
var delay = pool.initialBackoffDelay
|
||||
var delay = pool.connectionRetryStrategy.initialDelay
|
||||
let oneNanosecond = TimeAmount.nanoseconds(1)
|
||||
for _ in 0..<10 {
|
||||
let promise = connectionPromise
|
||||
@@ -481,7 +491,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
self.server.loop.advanceTime(by: oneNanosecond)
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
delay = .nanoseconds(Int64(Float32(delay.nanoseconds) * pool.backoffFactor))
|
||||
delay = pool.connectionRetryStrategy.determineNewDelay(currentDelay: delay)
|
||||
}
|
||||
|
||||
pool.close()
|
||||
@@ -493,7 +503,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testDeadlinesWork() throws {
|
||||
var promises: [EventLoopPromise<RedisConnection>] = []
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 0, leaky: true) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 8, minimumConnectionCount: 0, behavior: .elastic) { loop in
|
||||
let connectionPromise = self.server.loop.makePromise(of: RedisConnection.self)
|
||||
promises.append(connectionPromise)
|
||||
return connectionPromise.futureResult
|
||||
@@ -550,7 +560,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
|
||||
func testPoolWillStoreConnectionIfWaiterGoesAway() throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, leaky: true) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, behavior: .elastic) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -586,7 +596,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testPoolCorrectlyClosesItselfWhenLeasedConnectionsAreReturned() throws {
|
||||
let pool = self.createPool(maximumConnectionCount: 2, minimumConnectionCount: 1, leaky: false)
|
||||
let pool = self.createPool(maximumConnectionCount: 2, minimumConnectionCount: 1, behavior: .strict)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -615,7 +625,7 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
func testLeasedConnectionsInExcessOfMaxReplacePooledOnes() throws {
|
||||
// This test validates that if a leaky pool has allowed extra connections, and all those connections are
|
||||
// returned back, the active connections are the ones that were returned to the pool last.
|
||||
let pool = self.createPool(maximumConnectionCount: 4, minimumConnectionCount: 0, leaky: true)
|
||||
let pool = self.createPool(maximumConnectionCount: 4, minimumConnectionCount: 0, behavior: .elastic)
|
||||
defer {
|
||||
pool.close()
|
||||
}
|
||||
@@ -651,9 +661,9 @@ final class ConnectionPoolTests: XCTestCase {
|
||||
}
|
||||
|
||||
extension ConnectionPoolTests {
|
||||
private func stopReconnectingIfThereAreNoWaiters(leaky: Bool) throws {
|
||||
private func stopReconnectingIfThereAreNoWaiters(behavior: RedisConnectionPool.ConnectionCountBehavior.MaxConnectionBehavior) throws {
|
||||
var connectionPromise: EventLoopPromise<RedisConnection>? = nil
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, leaky: leaky) { loop in
|
||||
let pool = self.createPool(maximumConnectionCount: 1, minimumConnectionCount: 0, behavior: behavior) { loop in
|
||||
XCTAssertTrue(loop === self.server.loop)
|
||||
connectionPromise = self.server.loop.makePromise()
|
||||
return connectionPromise!.futureResult
|
||||
@@ -683,7 +693,7 @@ extension ConnectionPoolTests {
|
||||
XCTAssertNil(connectionPromise)
|
||||
|
||||
// Now advance time the remaining amount.
|
||||
XCTAssertNoThrow(self.server.loop.advanceTime(by: pool.initialBackoffDelay))
|
||||
XCTAssertNoThrow(self.server.loop.advanceTime(by: pool.connectionRetryStrategy.initialDelay))
|
||||
XCTAssertNoThrow(try self.server.runWhileActive())
|
||||
XCTAssertNotNil(connectionPromise)
|
||||
|
||||
@@ -699,12 +709,12 @@ extension ConnectionPoolTests {
|
||||
}
|
||||
|
||||
func testLeakyPoolStopsReconnecting() throws {
|
||||
try self.stopReconnectingIfThereAreNoWaiters(leaky: true)
|
||||
try self.stopReconnectingIfThereAreNoWaiters(behavior: .elastic)
|
||||
}
|
||||
|
||||
func testNonLeakyPoolStopsReconnectingIfThereAreNoWaiters() throws {
|
||||
// This is the same as the test above, but the pool isn't leaky.
|
||||
try self.stopReconnectingIfThereAreNoWaiters(leaky: false)
|
||||
try self.stopReconnectingIfThereAreNoWaiters(behavior: .strict)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -714,7 +724,7 @@ extension ConnectionPool {
|
||||
func activate() { self.activate(logger: .redisBaseConnectionPoolLogger) }
|
||||
|
||||
func leaseConnection(deadline: NIODeadline) -> EventLoopFuture<RedisConnection> {
|
||||
return self.leaseConnection(deadline: deadline, logger: .redisBaseConnectionPoolLogger)
|
||||
return self.leaseConnection(logger: .redisBaseConnectionPoolLogger, deadline: deadline)
|
||||
}
|
||||
|
||||
func returnConnection(_ connection: RedisConnection) {
|
||||
|
||||
Reference in New Issue
Block a user