From 820820d8770debaa4eadf15bd4fb2138bb760ec9 Mon Sep 17 00:00:00 2001 From: Nathan Harris Date: Sun, 24 Apr 2022 23:37:22 -0500 Subject: [PATCH] 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. --- .../ConnectionPool/ConnectionPool.swift | 123 ++++----- ...ft => RedisConnection+Configuration.swift} | 114 +------- .../RedisConnectionPool+Configuration.swift | 248 ++++++++++++++++++ Sources/RediStack/RedisConnectionPool.swift | 63 ++--- Sources/RediStack/RedisLogging.swift | 4 +- ...disConnectionPoolIntegrationTestCase.swift | 12 +- .../RedisConnectionPoolTests.swift | 2 +- .../RedisLoggingTests.swift | 8 +- .../RedisServiceDiscoveryTests.swift | 7 +- .../RediStackTests/ConnectionPoolTests.swift | 88 ++++--- 10 files changed, 401 insertions(+), 268 deletions(-) rename Sources/RediStack/{Configuration.swift => RedisConnection+Configuration.swift} (56%) create mode 100644 Sources/RediStack/RedisConnectionPool+Configuration.swift diff --git a/Sources/RediStack/ConnectionPool/ConnectionPool.swift b/Sources/RediStack/ConnectionPool/ConnectionPool.swift index 23dc02d..5a1c21c 100644 --- a/Sources/RediStack/ConnectionPool/ConnectionPool.swift +++ b/Sources/RediStack/ConnectionPool/ConnectionPool.swift @@ -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 ) { - 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 { + func leaseConnection(logger: Logger, deadline: NIODeadline? = nil) -> EventLoopFuture { + 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 { + private func _leaseConnection(logger: Logger, deadline: NIODeadline) -> EventLoopFuture { 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 diff --git a/Sources/RediStack/Configuration.swift b/Sources/RediStack/RedisConnection+Configuration.swift similarity index 56% rename from Sources/RediStack/Configuration.swift rename to Sources/RediStack/RedisConnection+Configuration.swift index 573fff0..325e4ff 100644 --- a/Sources/RediStack/Configuration.swift +++ b/Sources/RediStack/RedisConnection+Configuration.swift @@ -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) -} diff --git a/Sources/RediStack/RedisConnectionPool+Configuration.swift b/Sources/RediStack/RedisConnectionPool+Configuration.swift new file mode 100644 index 0000000..0b49981 --- /dev/null +++ b/Sources/RediStack/RedisConnectionPool+Configuration.swift @@ -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 + } + } +} diff --git a/Sources/RediStack/RedisConnectionPool.swift b/Sources/RediStack/RedisConnectionPool.swift index 5a9f9c3..791e762 100644 --- a/Sources/RediStack/RedisConnectionPool.swift +++ b/Sources/RediStack/RedisConnectionPool.swift @@ -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 - } - } -} diff --git a/Sources/RediStack/RedisLogging.swift b/Sources/RediStack/RedisLogging.swift index ef7fa03..dd60b2d 100644 --- a/Sources/RediStack/RedisLogging.swift +++ b/Sources/RediStack/RedisLogging.swift @@ -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" diff --git a/Sources/RediStackTestUtils/RedisConnectionPoolIntegrationTestCase.swift b/Sources/RediStackTestUtils/RedisConnectionPoolIntegrationTestCase.swift index 4aab3a6..f680f17 100644 --- a/Sources/RediStackTestUtils/RedisConnectionPoolIntegrationTestCase.swift +++ b/Sources/RediStackTestUtils/RedisConnectionPoolIntegrationTestCase.swift @@ -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() ) diff --git a/Tests/RediStackIntegrationTests/RedisConnectionPoolTests.swift b/Tests/RediStackIntegrationTests/RedisConnectionPoolTests.swift index 179b139..9524ed5 100644 --- a/Tests/RediStackIntegrationTests/RedisConnectionPoolTests.swift +++ b/Tests/RediStackIntegrationTests/RedisConnectionPoolTests.swift @@ -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()) } diff --git a/Tests/RediStackIntegrationTests/RedisLoggingTests.swift b/Tests/RediStackIntegrationTests/RedisLoggingTests.swift index 5a49c96..d372677 100644 --- a/Tests/RediStackIntegrationTests/RedisLoggingTests.swift +++ b/Tests/RediStackIntegrationTests/RedisLoggingTests.swift @@ -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(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", diff --git a/Tests/RediStackIntegrationTests/RedisServiceDiscoveryTests.swift b/Tests/RediStackIntegrationTests/RedisServiceDiscoveryTests.swift index 3b769f3..31249d9 100644 --- a/Tests/RediStackIntegrationTests/RedisServiceDiscoveryTests.swift +++ b/Tests/RediStackIntegrationTests/RedisServiceDiscoveryTests.swift @@ -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(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", diff --git a/Tests/RediStackTests/ConnectionPoolTests.swift b/Tests/RediStackTests/ConnectionPoolTests.swift index 9d021d1..6222991 100644 --- a/Tests/RediStackTests/ConnectionPoolTests.swift +++ b/Tests/RediStackTests/ConnectionPoolTests.swift @@ -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) -> ConnectionPool { + func createPool( + maximumConnectionCount: Int, + minimumConnectionCount: Int, + behavior: RedisConnectionPool.ConnectionCountBehavior.MaxConnectionBehavior, + connectionFactory: @escaping (EventLoop) -> EventLoopFuture + ) -> 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? = 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? = 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? = 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? = 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? = 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] = [] - 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? = 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? = 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 { - return self.leaseConnection(deadline: deadline, logger: .redisBaseConnectionPoolLogger) + return self.leaseConnection(logger: .redisBaseConnectionPoolLogger, deadline: deadline) } func returnConnection(_ connection: RedisConnection) {