mirror of
https://github.com/swift-server/RediStack.git
synced 2026-06-02 07:37:33 +00:00
Merge branch 'composeability' into rev-2
* composeability: Update unit test references for new type names and to cover more test cases Update `DispatchRedis` module with new classes Remove deprecated classes Update command extensions Add `RedisDriver` to eventually replace `NIORedis` that works more deeply with NIO Add `RedisPipeline` to eventually replace `NIORedisPipeline` that works more deeply with NIO Add `RedisConnection` to eventually replace `NIORedisConnection` that works more deeply with NIO Add `RedisCommandHandler` to eventually replace `RedisMessenger` that works more deeply with NIO
This commit is contained in:
@@ -4,12 +4,12 @@ import NIORedis
|
||||
|
||||
/// A factory that handles all necessary details for creating `RedisConnection` instances.
|
||||
public final class Redis {
|
||||
private let driver: NIORedis
|
||||
private let driver: RedisDriver
|
||||
|
||||
deinit { try? driver.terminate() }
|
||||
|
||||
public init(threadCount: Int = 1) {
|
||||
self.driver = NIORedis(executionModel: .spawnThreads(threadCount))
|
||||
self.driver = RedisDriver(executionModel: .spawnThreads(threadCount))
|
||||
}
|
||||
|
||||
public func makeConnection(
|
||||
|
||||
@@ -2,20 +2,20 @@ import Foundation
|
||||
import NIORedis
|
||||
|
||||
public final class RedisConnection {
|
||||
let _driverConnection: NIORedisConnection
|
||||
let _driverConnection: NIORedis.RedisConnection
|
||||
|
||||
private let queue: DispatchQueue
|
||||
|
||||
deinit { _driverConnection.close() }
|
||||
|
||||
init(driver: NIORedisConnection, callbackQueue: DispatchQueue) {
|
||||
init(driver: NIORedis.RedisConnection, callbackQueue: DispatchQueue) {
|
||||
self._driverConnection = driver
|
||||
self.queue = callbackQueue
|
||||
}
|
||||
|
||||
/// Creates a `RedisPipeline` for executing a batch of commands.
|
||||
public func makePipeline(callbackQueue: DispatchQueue? = nil) -> RedisPipeline {
|
||||
return .init(using: self, callbackQueue: callbackQueue ?? queue)
|
||||
return .init(connection: self, callbackQueue: callbackQueue ?? queue)
|
||||
}
|
||||
|
||||
public func get(
|
||||
|
||||
@@ -14,15 +14,15 @@ import NIORedis
|
||||
/// - Important: The larger the pipeline queue, the more memory both the Redis driver and Redis server will use.
|
||||
/// See https://redis.io/topics/pipelining#redis-pipelining
|
||||
public final class RedisPipeline {
|
||||
private let _driverPipeline: NIORedisPipeline
|
||||
private let _driverPipeline: NIORedis.RedisPipeline
|
||||
private let queue: DispatchQueue
|
||||
|
||||
/// Creates a new pipeline queue using the provided `RedisConnection`, executing callbacks on the provided `DispatchQueue`.
|
||||
/// - Parameters:
|
||||
/// - using: The connection to execute the commands on.
|
||||
/// - callbackQueue: The queue to execute all callbacks on.
|
||||
public init(using connection: RedisConnection, callbackQueue: DispatchQueue) {
|
||||
self._driverPipeline = NIORedisPipeline(using: connection._driverConnection)
|
||||
public init(connection: RedisConnection, callbackQueue: DispatchQueue) {
|
||||
self._driverPipeline = NIORedis.RedisPipeline(channel: connection._driverConnection.channel)
|
||||
self.queue = callbackQueue
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
import Foundation
|
||||
import NIO
|
||||
|
||||
/// A context for `RedisCommandHandler` to operate within.
|
||||
public struct RedisCommandContext {
|
||||
/// A full command keyword and arguments stored as a single `RESPValue`.
|
||||
public let command: RESPValue
|
||||
/// A promise expected to be fulfilled with the `RESPValue` response to the command from Redis.
|
||||
public let promise: EventLoopPromise<RESPValue>
|
||||
}
|
||||
|
||||
/// A `ChannelDuplexHandler` that works with `RedisCommandContext`s to send commands and forward responses.
|
||||
open class RedisCommandHandler {
|
||||
/// Queue of promises waiting to receive a response value from a sent command.
|
||||
private var commandResponseQueue: [EventLoopPromise<RESPValue>]
|
||||
|
||||
public init() {
|
||||
self.commandResponseQueue = []
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: ChannelInboundHandler
|
||||
|
||||
extension RedisCommandHandler: ChannelInboundHandler {
|
||||
/// See `ChannelInboundHandler.InboundIn`
|
||||
public typealias InboundIn = RESPValue
|
||||
|
||||
/// Invoked by NIO when an error has been thrown. The command response promise at the front of the queue will be
|
||||
/// failed with the error.
|
||||
///
|
||||
/// See `ChannelInboundHandler.errorCaught(ctx:error:)`
|
||||
public func errorCaught(ctx: ChannelHandlerContext, error: Error) {
|
||||
guard let leadPromise = commandResponseQueue.last else {
|
||||
return assertionFailure("Received unexpected error while idle: \(error.localizedDescription)")
|
||||
}
|
||||
leadPromise.fail(error: error)
|
||||
}
|
||||
|
||||
/// Invoked by NIO when a read has been fired from earlier in the response chain. This forwards the unwrapped
|
||||
/// `RESPValue` to the promise awaiting a response at the front of the queue.
|
||||
///
|
||||
/// See `ChannelInboundHandler.channelRead(ctx:data:)`
|
||||
public func channelRead(ctx: ChannelHandlerContext, data: NIOAny) {
|
||||
let value = unwrapInboundIn(data)
|
||||
|
||||
guard let leadPromise = commandResponseQueue.last else {
|
||||
return assertionFailure("Read triggered with an empty promise queue! Ignoring: \(value)")
|
||||
}
|
||||
|
||||
let popped = commandResponseQueue.popLast()
|
||||
assert(popped != nil)
|
||||
|
||||
switch value {
|
||||
case .error(let e): leadPromise.fail(error: e)
|
||||
default: leadPromise.succeed(result: value)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: ChannelOutboundHandler
|
||||
|
||||
extension RedisCommandHandler: ChannelOutboundHandler {
|
||||
/// See `ChannelOutboundHandler.OutboundIn`
|
||||
public typealias OutboundIn = RedisCommandContext
|
||||
/// See `ChannelOutboundHandler.OutboundOut`
|
||||
public typealias OutboundOut = RESPValue
|
||||
|
||||
/// Invoked by NIO when a `write` has been requested on the `Channel`.
|
||||
/// This unwraps a `RedisCommandContext`, retaining a callback to forward a response to later, and forwards
|
||||
/// the underlying command data further into the pipeline.
|
||||
///
|
||||
/// See `ChannelOutboundHandler.write(ctx:data:promise:)`
|
||||
public func write(ctx: ChannelHandlerContext, data: NIOAny, promise: EventLoopPromise<Void>?) {
|
||||
let context = unwrapOutboundIn(data)
|
||||
commandResponseQueue.insert(context.promise, at: 0)
|
||||
ctx.write(wrapOutboundOut(context.command), promise: promise)
|
||||
}
|
||||
}
|
||||
@@ -1,102 +0,0 @@
|
||||
import NIO
|
||||
|
||||
/// `ChannelInboundHandler` that is responsible for coordinating incoming and outgoing messages on a particular
|
||||
/// connection to Redis.
|
||||
internal final class RedisMessenger {
|
||||
private let eventLoop: EventLoop
|
||||
|
||||
/// Context to be used for writing outgoing messages with.
|
||||
private var channelContext: ChannelHandlerContext?
|
||||
|
||||
/// Queue of promises waiting to receive an incoming response value from an outgoing message.
|
||||
private var waitingResponseQueue: [EventLoopPromise<RESPValue>]
|
||||
/// Queue of unset outgoing messages, with the oldest messages at the end of the array.
|
||||
private var outgoingMessageQueue: [RESPValue]
|
||||
|
||||
/// Creates a new handler that works on the specified `EventLoop`.
|
||||
init(on eventLoop: EventLoop) {
|
||||
self.waitingResponseQueue = []
|
||||
self.outgoingMessageQueue = []
|
||||
self.eventLoop = eventLoop
|
||||
}
|
||||
|
||||
/// Adds a complete message encoded as `RESPValue` to the queue and returns an `EventLoopFuture` that resolves
|
||||
/// the response from Redis.
|
||||
func enqueue(_ output: RESPValue) -> EventLoopFuture<RESPValue> {
|
||||
// ensure that we are on the event loop before modifying our data
|
||||
guard eventLoop.inEventLoop else {
|
||||
return eventLoop.submit({}).then { return self.enqueue(output) }
|
||||
}
|
||||
|
||||
// add the new output to the writing queue at the front
|
||||
outgoingMessageQueue.insert(output, at: 0)
|
||||
|
||||
// every outgoing message is expected to receive some form of response, so create a promise that we'll resolve
|
||||
// with the response
|
||||
let promise = eventLoop.makePromise(of: RESPValue.self)
|
||||
waitingResponseQueue.insert(promise, at: 0)
|
||||
|
||||
// if we have a context for writing, flush the outgoing queue
|
||||
channelContext?.eventLoop.execute {
|
||||
self._flushOutgoingQueue()
|
||||
}
|
||||
|
||||
return promise.futureResult
|
||||
}
|
||||
|
||||
/// Writes all queued outgoing messages to the channel.
|
||||
func _flushOutgoingQueue() {
|
||||
guard let context = channelContext else { return }
|
||||
|
||||
while let output = outgoingMessageQueue.popLast() {
|
||||
context.write(wrapOutboundOut(output), promise: nil)
|
||||
context.flush()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: ChannelInboundHandler
|
||||
|
||||
extension RedisMessenger: ChannelInboundHandler {
|
||||
/// See `ChannelInboundHandler.InboundIn`
|
||||
public typealias InboundIn = RESPValue
|
||||
|
||||
/// See `ChannelInboundHandler.OutboundOut`
|
||||
public typealias OutboundOut = RESPValue
|
||||
|
||||
/// Invoked by NIO when the channel for this handler has become active, receiving a context that is ready to
|
||||
/// send messages.
|
||||
///
|
||||
/// Any queued messages will be flushed at this point.
|
||||
/// See `ChannelInboundHandler.channelActive(ctx:)`
|
||||
public func channelActive(ctx: ChannelHandlerContext) {
|
||||
channelContext = ctx
|
||||
_flushOutgoingQueue()
|
||||
}
|
||||
|
||||
/// Invoked by NIO when an error was thrown earlier in the response chain. The waiting promise at the front
|
||||
/// of the queue will be failed with the error.
|
||||
/// See `ChannelInboundHandler.errorCaught(ctx:error:)`
|
||||
public func errorCaught(ctx: ChannelHandlerContext, error: Error) {
|
||||
guard let leadPromise = waitingResponseQueue.last else {
|
||||
return assertionFailure("Received unexpected error while idle: \(error.localizedDescription)")
|
||||
}
|
||||
leadPromise.fail(error: error)
|
||||
}
|
||||
|
||||
/// Invoked by NIO when a read has been fired from earlier in the response chain. This forwards the unwrapped
|
||||
/// `RESPValue` to the response at the front of the queue.
|
||||
/// See `ChannelInboundHandler.channelRead(ctx:data:)`
|
||||
public func channelRead(ctx: ChannelHandlerContext, data: NIOAny) {
|
||||
let input = unwrapInboundIn(data)
|
||||
|
||||
guard let leadPromise = waitingResponseQueue.last else {
|
||||
return assertionFailure("Read triggered with an empty input queue! Ignoring: \(input)")
|
||||
}
|
||||
|
||||
let popped = waitingResponseQueue.popLast()
|
||||
assert(popped != nil)
|
||||
|
||||
leadPromise.succeed(result: input)
|
||||
}
|
||||
}
|
||||
@@ -1,13 +1,13 @@
|
||||
import Foundation
|
||||
import NIO
|
||||
|
||||
extension NIORedisConnection {
|
||||
extension RedisConnection {
|
||||
/// Select the Redis logical database having the specified zero-based numeric index.
|
||||
/// New connections always use the database 0.
|
||||
///
|
||||
/// https://redis.io/commands/select
|
||||
public func select(_ id: Int) -> EventLoopFuture<Void> {
|
||||
return command("SELECT", [RESPValue(bulk: id.description)])
|
||||
return command("SELECT", arguments: [RESPValue(bulk: id.description)])
|
||||
.map { _ in return () }
|
||||
}
|
||||
|
||||
@@ -15,7 +15,8 @@ extension NIORedisConnection {
|
||||
///
|
||||
/// https://redis.io/commands/auth
|
||||
public func authorize(with password: String) -> EventLoopFuture<Void> {
|
||||
return command("AUTH", [RESPValue(bulk: password)]).map { _ in return () }
|
||||
return command("AUTH", arguments: [RESPValue(bulk: password)])
|
||||
.map { _ in return () }
|
||||
}
|
||||
|
||||
/// Removes the specified keys. A key is ignored if it does not exist.
|
||||
@@ -24,7 +25,7 @@ extension NIORedisConnection {
|
||||
/// - Returns: A future number of keys that were removed.
|
||||
public func delete(_ keys: String...) -> EventLoopFuture<Int> {
|
||||
let keyArgs = keys.map { RESPValue(bulk: $0) }
|
||||
return command("DEL", keyArgs)
|
||||
return command("DEL", arguments: keyArgs)
|
||||
.thenThrowing { res in
|
||||
guard let count = res.int else {
|
||||
throw RedisError(identifier: "delete", reason: "Unexpected response: \(res)")
|
||||
@@ -41,7 +42,7 @@ extension NIORedisConnection {
|
||||
/// - after: The lifetime (in seconds) the key will expirate at.
|
||||
/// - Returns: A future bool indicating if the expiration was set or not.
|
||||
public func expire(_ key: String, after deadline: Int) -> EventLoopFuture<Bool> {
|
||||
return command("EXPIRE", [RESPValue(bulk: key), RESPValue(bulk: deadline.description)])
|
||||
return command("EXPIRE", arguments: [RESPValue(bulk: key), RESPValue(bulk: deadline.description)])
|
||||
.thenThrowing { res in
|
||||
guard let value = res.int else {
|
||||
throw RedisError(identifier: "expire", reason: "Unexpected response: \(res)")
|
||||
@@ -56,7 +57,7 @@ extension NIORedisConnection {
|
||||
///
|
||||
/// https://redis.io/commands/get
|
||||
public func get(_ key: String) -> EventLoopFuture<String?> {
|
||||
return command("GET", [RESPValue(bulk: key)])
|
||||
return command("GET", arguments: [RESPValue(bulk: key)])
|
||||
.map { return $0.string }
|
||||
}
|
||||
|
||||
@@ -66,7 +67,7 @@ extension NIORedisConnection {
|
||||
///
|
||||
/// https://redis.io/commands/set
|
||||
public func set(_ key: String, to value: String) -> EventLoopFuture<Void> {
|
||||
return command("SET", [RESPValue(bulk: key), RESPValue(bulk: value)])
|
||||
return command("SET", arguments: [RESPValue(bulk: key), RESPValue(bulk: value)])
|
||||
.map { _ in return () }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,60 +0,0 @@
|
||||
import NIO
|
||||
import NIOConcurrencyHelpers
|
||||
|
||||
public final class NIORedisConnection {
|
||||
/// The `EventLoop` this connection uses to execute commands on.
|
||||
public var eventLoop: EventLoop { return channel.eventLoop }
|
||||
|
||||
/// Has the connection been closed?
|
||||
public private(set) var isClosed = Atomic<Bool>(value: false)
|
||||
|
||||
internal let messenger: RedisMessenger
|
||||
|
||||
private let channel: Channel
|
||||
|
||||
deinit { assert(isClosed.load(), "Redis connection was not properly shut down!") }
|
||||
|
||||
/// Creates a new connection on the provided channel, using the handler for executing commands.
|
||||
/// - Important: Call `close()` before deinitializing to properly cleanup resources!
|
||||
init(channel: Channel, handler: RedisMessenger) {
|
||||
self.channel = channel
|
||||
self.messenger = handler
|
||||
}
|
||||
|
||||
/// Closes the connection to Redis.
|
||||
public func close() {
|
||||
guard isClosed.exchange(with: true) else { return }
|
||||
|
||||
channel.close(promise: nil)
|
||||
}
|
||||
|
||||
/// Executes the desired command with the specified arguments.
|
||||
/// - Important: All arguments should be in `.bulkString` format.
|
||||
public func command(_ command: String, _ arguments: [RESPValue] = []) -> EventLoopFuture<RESPValue> {
|
||||
return _send(.array([RESPValue(bulk: command)] + arguments))
|
||||
.thenThrowing { response in
|
||||
switch response {
|
||||
case let .error(error): throw error
|
||||
default: return response
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Creates a `NIORedisPipeline` for executing a batch of commands.
|
||||
public func makePipeline() -> NIORedisPipeline {
|
||||
return .init(using: self)
|
||||
}
|
||||
|
||||
func _send(_ message: RESPValue) -> EventLoopFuture<RESPValue> {
|
||||
// ensure the connection is still open
|
||||
guard !isClosed.load() else { return eventLoop.makeFailedFuture(error: RedisError.connectionClosed) }
|
||||
|
||||
// create a new promise to store
|
||||
let promise = eventLoop.makePromise(of: RESPValue.self)
|
||||
|
||||
// cascade this enqueue to the newly created promise
|
||||
messenger.enqueue(message).cascade(promise: promise)
|
||||
|
||||
return promise.futureResult
|
||||
}
|
||||
}
|
||||
@@ -1,85 +0,0 @@
|
||||
import Foundation
|
||||
import NIO
|
||||
|
||||
/// An object that provides a mechanism to "pipeline" multiple Redis commands in sequence, providing an aggregate response
|
||||
/// of all the Redis responses for each individual command.
|
||||
///
|
||||
/// let results = connection.makePipeline()
|
||||
/// .enqueue(command: "SET", arguments: ["my_key", 3])
|
||||
/// .enqueue(command: "INCR", arguments: ["my_key"])
|
||||
/// .execute()
|
||||
/// // results == Future<[RESPValue]>
|
||||
/// // results[0].string == Optional("OK")
|
||||
/// // results[1].int == Optional(4)
|
||||
/// - Important: The larger the pipeline queue, the more memory both NIORedis and Redis will use.
|
||||
/// See https://redis.io/topics/pipelining#redis-pipelining
|
||||
public final class NIORedisPipeline {
|
||||
/// The client to execute the commands on.
|
||||
private let connection: NIORedisConnection
|
||||
|
||||
/// The queue of complete, encoded commands to execute.
|
||||
private var queue: [RESPValue]
|
||||
private var messageCount: Int
|
||||
|
||||
/// Creates a new pipeline queue using the provided `NIORedisConnection`.
|
||||
/// - Parameter using: The connection to execute the commands on.
|
||||
public init(using connection: NIORedisConnection) {
|
||||
self.connection = connection
|
||||
self.queue = []
|
||||
self.messageCount = 0
|
||||
}
|
||||
|
||||
/// Queues the provided command and arguments to be executed when `execute()` is invoked.
|
||||
/// - Parameters:
|
||||
/// - command: The command to execute. See https://redis.io/commands
|
||||
/// - arguments: The arguments, if any, to send with the command.
|
||||
/// - Returns: A self-reference to this `NIORedisPipeline` instance for chaining commands.
|
||||
@discardableResult
|
||||
public func enqueue(command: String, arguments: [RESPConvertible] = []) throws -> NIORedisPipeline {
|
||||
let args = try arguments.map { try $0.convertToRESP() }
|
||||
|
||||
queue.append(.array([RESPValue(bulk: command)] + args))
|
||||
|
||||
return self
|
||||
}
|
||||
|
||||
/// Flushes the queue, sending all of the commands to Redis in the same order as they were enqueued.
|
||||
/// - Important: If any of the commands fail, the remaining commands will not execute and the `EventLoopFuture` will fail.
|
||||
/// - Returns: A `EventLoopFuture` that resolves the `RESPValue` responses, in the same order as the command queue.
|
||||
public func execute() -> EventLoopFuture<[RESPValue]> {
|
||||
let promise = connection.eventLoop.makePromise(of: [RESPValue].self)
|
||||
|
||||
var results = [RESPValue]()
|
||||
var iterator = queue.makeIterator()
|
||||
|
||||
// recursive internal method for chaining each request and
|
||||
// attaching callbacks for failing or ultimately succeeding
|
||||
func handle(_ command: RESPValue) {
|
||||
let future = connection._send(command)
|
||||
future.whenSuccess { response in
|
||||
switch response {
|
||||
case let .error(error): promise.fail(error: error)
|
||||
default:
|
||||
results.append(response)
|
||||
|
||||
if let next = iterator.next() {
|
||||
handle(next)
|
||||
} else {
|
||||
promise.succeed(result: results)
|
||||
}
|
||||
}
|
||||
}
|
||||
future.whenFailure { promise.fail(error: $0) }
|
||||
}
|
||||
|
||||
if let first = iterator.next() {
|
||||
handle(first)
|
||||
} else {
|
||||
promise.succeed(result: [])
|
||||
}
|
||||
|
||||
promise.futureResult.whenComplete { self.queue = [] }
|
||||
|
||||
return promise.futureResult
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
import NIO
|
||||
import NIOConcurrencyHelpers
|
||||
|
||||
/// A connection to a Redis database instance, with the ability to send and receive commands.
|
||||
///
|
||||
/// let result = connection.send(command: "GET", arguments: ["my_key"]
|
||||
/// // result == EventLoopFuture<RESPValue>
|
||||
///
|
||||
/// See https://redis.io/commands
|
||||
public final class RedisConnection {
|
||||
/// The `Channel` this connection is associated with.
|
||||
public let channel: Channel
|
||||
|
||||
/// Has the connection been closed?
|
||||
public private(set) var isClosed = Atomic<Bool>(value: false)
|
||||
|
||||
deinit { assert(isClosed.load(), "Redis connection was not properly shut down!") }
|
||||
|
||||
/// Creates a new connection on the provided channel.
|
||||
/// - Note: This connection will take ownership of the `Channel` object.
|
||||
/// - Important: Call `close()` before deinitializing to properly cleanup resources.
|
||||
public init(channel: Channel) {
|
||||
self.channel = channel
|
||||
}
|
||||
|
||||
/// Closes the connection to Redis.
|
||||
/// - Returns: An `EventLoopFuture` that resolves when the connection has been closed.
|
||||
@discardableResult
|
||||
public func close() -> EventLoopFuture<Void> {
|
||||
guard isClosed.exchange(with: true) else { return channel.eventLoop.makeSucceededFuture(result: ()) }
|
||||
|
||||
let promise = channel.eventLoop.makePromise(of: Void.self)
|
||||
|
||||
channel.close(promise: promise)
|
||||
|
||||
return promise.futureResult
|
||||
}
|
||||
|
||||
/// Sends the desired command with the specified arguments.
|
||||
/// - Parameters:
|
||||
/// - command: The command to execute.
|
||||
/// - arguments: The arguments to be sent with the command.
|
||||
/// - Returns: An `EventLoopFuture` that will resolve with the Redis command response.
|
||||
public func send(command: String, with arguments: [RESPConvertible] = []) throws -> EventLoopFuture<RESPValue> {
|
||||
let args = try arguments.map { try $0.convertToRESP() }
|
||||
return self.command(command, arguments: args)
|
||||
}
|
||||
|
||||
/// Invokes a command against Redis with the provided arguments.
|
||||
/// - Important: Arguments should be stored as `.bulkString`.
|
||||
/// - Parameters:
|
||||
/// - command: The command to execute.
|
||||
/// - arguments: The arguments to be sent with the command.
|
||||
/// - Returns: An `EventLoopFuture` that will resolve with the Redis command response.
|
||||
public func command(_ command: String, arguments: [RESPValue] = []) -> EventLoopFuture<RESPValue> {
|
||||
guard !isClosed.load() else {
|
||||
return channel.eventLoop.makeFailedFuture(error: RedisError.connectionClosed)
|
||||
}
|
||||
|
||||
let promise = channel.eventLoop.makePromise(of: RESPValue.self)
|
||||
let context = RedisCommandContext(
|
||||
command: .array([RESPValue(bulk: command)] + arguments),
|
||||
promise: promise
|
||||
)
|
||||
|
||||
_ = channel.writeAndFlush(context)
|
||||
|
||||
return promise.futureResult
|
||||
}
|
||||
|
||||
/// Creates a `RedisPipeline` for executing a batch of commands.
|
||||
public func makePipeline() -> RedisPipeline {
|
||||
return .init(channel: channel)
|
||||
}
|
||||
}
|
||||
@@ -1,8 +1,7 @@
|
||||
import NIO
|
||||
import NIOConcurrencyHelpers
|
||||
|
||||
/// A factory that handles all necessary details for creating connections to a Redis database instance.
|
||||
public final class NIORedis {
|
||||
public final class RedisDriver {
|
||||
/// The threading model to use for asynchronous tasks.
|
||||
///
|
||||
/// Using `.eventLoopGroup` will allow an external provider to handle the lifetime of the `EventLoopGroup`,
|
||||
@@ -13,62 +12,63 @@ public final class NIORedis {
|
||||
}
|
||||
|
||||
private let executionModel: ExecutionModel
|
||||
private let elg: EventLoopGroup
|
||||
private let eventLoopGroup: EventLoopGroup
|
||||
|
||||
private let isRunning = Atomic<Bool>(value: true)
|
||||
|
||||
deinit { assert(!isRunning.load(), "Redis driver was not properly shut down!") }
|
||||
|
||||
/// Creates a handle to create connections to a Redis instance using the `ExecutionModel` provided.
|
||||
/// - Parameter executionModel: The model to use for handling asynchronous scheduling.
|
||||
/// Creates a driver instance to create connections to a Redis.
|
||||
/// - Important: Call `terminate()` before deinitializing to properly cleanup resources.
|
||||
/// - Parameter executionModel: The model to use for handling connection resources.
|
||||
public init(executionModel model: ExecutionModel) {
|
||||
self.executionModel = model
|
||||
|
||||
switch model {
|
||||
case .spawnThreads(let count):
|
||||
self.elg = MultiThreadedEventLoopGroup(numberOfThreads: count)
|
||||
self.eventLoopGroup = MultiThreadedEventLoopGroup(numberOfThreads: count)
|
||||
case .eventLoopGroup(let group):
|
||||
self.elg = group
|
||||
self.eventLoopGroup = group
|
||||
}
|
||||
}
|
||||
|
||||
/// Creates a new `NIORedisConnection` with the connection parameters provided.
|
||||
/// Handles the proper shutdown of managed resources.
|
||||
/// - Important: This method should always be called, even when running in `.eventLoopGroup` mode.
|
||||
public func terminate() throws {
|
||||
guard isRunning.exchange(with: false) else { return }
|
||||
|
||||
switch executionModel {
|
||||
case .spawnThreads: try self.eventLoopGroup.syncShutdownGracefully()
|
||||
case .eventLoopGroup: return
|
||||
}
|
||||
}
|
||||
|
||||
/// Creates a new `RedisConnection` with the parameters provided.
|
||||
public func makeConnection(
|
||||
hostname: String = "localhost",
|
||||
port: Int = 6379,
|
||||
password: String? = nil
|
||||
) -> EventLoopFuture<NIORedisConnection> {
|
||||
let channelHandler = RedisMessenger(on: elg.next())
|
||||
let bootstrap = ClientBootstrap(group: self.elg)
|
||||
) -> EventLoopFuture<RedisConnection> {
|
||||
let bootstrap = ClientBootstrap(group: eventLoopGroup)
|
||||
.channelOption(ChannelOptions.socket(SocketOptionLevel(SOL_SOCKET), SO_REUSEADDR), value: 1)
|
||||
.channelInitializer { channel in
|
||||
return channel.pipeline.addHandlers(
|
||||
RESPEncoder(),
|
||||
ByteToMessageHandler(RESPDecoder()),
|
||||
channelHandler
|
||||
RedisCommandHandler()
|
||||
)
|
||||
}
|
||||
|
||||
return bootstrap.connect(host: hostname, port: port)
|
||||
.map { return NIORedisConnection(channel: $0, handler: channelHandler) }
|
||||
.map { return RedisConnection(channel: $0) }
|
||||
.then { connection in
|
||||
guard let pw = password else {
|
||||
return self.elg.next().makeSucceededFuture(result: connection)
|
||||
return self.eventLoopGroup.next().makeSucceededFuture(result: connection)
|
||||
}
|
||||
|
||||
return connection.authorize(with: pw).map { _ in return connection }
|
||||
}
|
||||
}
|
||||
|
||||
/// Handles the proper shutdown of managed resources.
|
||||
/// - Important: This method should always be called before deinit.
|
||||
public func terminate() throws {
|
||||
guard isRunning.exchange(with: false) else { return }
|
||||
|
||||
switch executionModel {
|
||||
case .spawnThreads: try self.elg.syncShutdownGracefully()
|
||||
case .eventLoopGroup: return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private extension ChannelPipeline {
|
||||
@@ -0,0 +1,70 @@
|
||||
import Foundation
|
||||
|
||||
/// An object that provides a mechanism to "pipeline" multiple Redis commands in sequence,
|
||||
/// providing an aggregate response of all the Redis responses for each individual command.
|
||||
///
|
||||
/// let results = connection.makePipeline()
|
||||
/// .enqueue(command: "SET", arguments: ["my_key", 3])
|
||||
/// .enqueue(command: "INCR", arguments: ["my_key"])
|
||||
/// .execute()
|
||||
/// // results == Future<[RESPValue]>
|
||||
/// // results[0].string == Optional("OK")
|
||||
/// // results[1].int == Optional(4)
|
||||
///
|
||||
/// See https://redis.io/topics/pipelining#redis-pipelining
|
||||
/// - Important: The larger the pipeline queue, the more memory both NIORedis and Redis will use.
|
||||
public final class RedisPipeline {
|
||||
/// The number of commands in the pipeline.
|
||||
public var count: Int {
|
||||
return queuedCommandResults.count
|
||||
}
|
||||
|
||||
/// The channel being used to send commands with.
|
||||
private let channel: Channel
|
||||
|
||||
/// The queue of response handlers that have been queued.
|
||||
private var queuedCommandResults: [EventLoopFuture<RESPValue>]
|
||||
|
||||
/// Creates a new pipeline queue that will write to the channel provided.
|
||||
/// - Parameter channel: The `Channel` to write to.
|
||||
public init(channel: Channel) {
|
||||
self.channel = channel
|
||||
self.queuedCommandResults = []
|
||||
}
|
||||
|
||||
/// Queues the provided command and arguments to be executed when `execute()` is invoked.
|
||||
/// - Parameters:
|
||||
/// - command: The command to execute. See https://redis.io/commands
|
||||
/// - arguments: The arguments, if any, to send with the command.
|
||||
/// - Returns: A self-reference for chaining commands.
|
||||
@discardableResult
|
||||
public func enqueue(command: String, arguments: [RESPConvertible] = []) throws -> RedisPipeline {
|
||||
let args = try arguments.map { try $0.convertToRESP() }
|
||||
|
||||
let promise = channel.eventLoop.makePromise(of: RESPValue.self)
|
||||
let context = RedisCommandContext(
|
||||
command: .array([RESPValue(bulk: command)] + args),
|
||||
promise: promise
|
||||
)
|
||||
|
||||
queuedCommandResults.append(promise.futureResult)
|
||||
|
||||
_ = channel.write(context)
|
||||
|
||||
return self
|
||||
}
|
||||
|
||||
/// Flushes the queue, sending all of the commands to Redis.
|
||||
/// - Important: If any of the commands fail, the remaining commands will not execute and the `EventLoopFuture` will fail.
|
||||
/// - Returns: An `EventLoopFuture` that resolves the `RESPValue` responses, in the same order as the command queue.
|
||||
public func execute() -> EventLoopFuture<[RESPValue]> {
|
||||
channel.flush()
|
||||
|
||||
return EventLoopFuture<[RESPValue]>.reduce(
|
||||
into: [],
|
||||
queuedCommandResults,
|
||||
eventLoop: channel.eventLoop,
|
||||
{ (results, response) in results.append(response) }
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -2,10 +2,10 @@
|
||||
import XCTest
|
||||
|
||||
final class BasicCommandsTests: XCTestCase {
|
||||
private let redis = NIORedis(executionModel: .spawnThreads(1))
|
||||
private let redis = RedisDriver(executionModel: .spawnThreads(1))
|
||||
deinit { try? redis.terminate() }
|
||||
|
||||
private var connection: NIORedisConnection?
|
||||
private var connection: RedisConnection?
|
||||
|
||||
override func setUp() {
|
||||
do {
|
||||
|
||||
@@ -1,27 +0,0 @@
|
||||
@testable import NIORedis
|
||||
import XCTest
|
||||
|
||||
final class NIORedisTests: XCTestCase {
|
||||
func test_makeConnection() {
|
||||
let redis = NIORedis(executionModel: .spawnThreads(1))
|
||||
defer { try? redis.terminate() }
|
||||
|
||||
XCTAssertNoThrow(try redis.makeConnection().wait().close())
|
||||
}
|
||||
|
||||
func test_command() throws {
|
||||
let redis = NIORedis(executionModel: .spawnThreads(1))
|
||||
defer { try? redis.terminate() }
|
||||
|
||||
let connection = try redis.makeConnection().wait()
|
||||
let result = try connection.command("SADD", [.bulkString("key".convertedToData()), try 3.convertToRESP()]).wait()
|
||||
XCTAssertNotNil(result.int)
|
||||
XCTAssertEqual(result.int, 1)
|
||||
try connection.command("DEL", [.bulkString("key".convertedToData())]).wait()
|
||||
connection.close()
|
||||
}
|
||||
|
||||
static var allTests = [
|
||||
("test_makeConnection", test_makeConnection),
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
@testable import NIORedis
|
||||
import XCTest
|
||||
|
||||
final class RedisDriverTests: XCTestCase {
|
||||
private var driver: RedisDriver!
|
||||
private var connection: RedisConnection!
|
||||
|
||||
override func setUp() {
|
||||
let driver = RedisDriver(executionModel: .spawnThreads(1))
|
||||
|
||||
guard let connection = try? driver.makeConnection().wait() else {
|
||||
return XCTFail("Failed to create a connection!")
|
||||
}
|
||||
|
||||
self.driver = driver
|
||||
self.connection = connection
|
||||
}
|
||||
|
||||
override func tearDown() {
|
||||
_ = connection.command("FLUSHALL")
|
||||
.then { _ in self.connection.close() }
|
||||
.map { _ in try? self.driver.terminate() }
|
||||
}
|
||||
|
||||
func test_makeConnection() {
|
||||
XCTAssertNoThrow(try driver.makeConnection().wait().close())
|
||||
}
|
||||
|
||||
func test_command_succeeds() throws {
|
||||
let result = try connection.command(
|
||||
"SADD",
|
||||
arguments: [.bulkString("key".convertedToData()), try 3.convertToRESP()
|
||||
]).wait()
|
||||
|
||||
XCTAssertNotNil(result.int)
|
||||
XCTAssertEqual(result.int, 1)
|
||||
}
|
||||
|
||||
func test_command_fails() {
|
||||
let command = connection.command("GET")
|
||||
|
||||
XCTAssertThrowsError(try command.wait())
|
||||
}
|
||||
|
||||
static var allTests = [
|
||||
("test_makeConnection", test_makeConnection),
|
||||
("test_command_succeeds", test_command_succeeds),
|
||||
("test_command_fails", test_command_fails),
|
||||
]
|
||||
}
|
||||
+7
-6
@@ -1,12 +1,12 @@
|
||||
@testable import NIORedis
|
||||
import XCTest
|
||||
|
||||
final class NIORedisPipelineTests: XCTestCase {
|
||||
private var redis: NIORedis!
|
||||
private var connection: NIORedisConnection!
|
||||
final class RedisPipelineTests: XCTestCase {
|
||||
private var redis: RedisDriver!
|
||||
private var connection: RedisConnection!
|
||||
|
||||
override func setUp() {
|
||||
let redis = NIORedis(executionModel: .spawnThreads(2))
|
||||
let redis = RedisDriver(executionModel: .spawnThreads(2))
|
||||
|
||||
guard let connection = try? redis.makeConnection().wait() else {
|
||||
return XCTFail("Failed to create connection!")
|
||||
@@ -30,10 +30,11 @@ final class NIORedisPipelineTests: XCTestCase {
|
||||
}
|
||||
|
||||
func test_executeFails() throws {
|
||||
let pipeline = try connection.makePipeline()
|
||||
let future = try connection.makePipeline()
|
||||
.enqueue(command: "GET")
|
||||
.execute()
|
||||
|
||||
XCTAssertThrowsError(try pipeline.execute().wait())
|
||||
XCTAssertThrowsError(try future.wait())
|
||||
}
|
||||
|
||||
func test_singleCommand() throws {
|
||||
@@ -3,14 +3,14 @@ import XCTest
|
||||
#if !os(macOS)
|
||||
public func allTests() -> [XCTestCaseEntry] {
|
||||
return [
|
||||
testCase(NIORedisTests.allTests),
|
||||
testCase(RedisDriverTests.allTests),
|
||||
testCase(RESPDecoderTests.allTests),
|
||||
testCase(RESPDecoderParsingTests.allTests),
|
||||
testCase(RESPDecoderByteToMessageDecoderTests.allTests),
|
||||
testCase(RESPEncoderTests.allTests),
|
||||
testCase(RESPEncoderParsingTests.allTests),
|
||||
testCase(BasicCommandsTests.allTests),
|
||||
testCase(NIORedisPipelineTests.allTests)
|
||||
testCase(RedisPipelineTests.allTests)
|
||||
]
|
||||
}
|
||||
#endif
|
||||
|
||||
Reference in New Issue
Block a user