diff --git a/Sources/DispatchRedis/Redis.swift b/Sources/DispatchRedis/Redis.swift index df6d329..bcfb056 100644 --- a/Sources/DispatchRedis/Redis.swift +++ b/Sources/DispatchRedis/Redis.swift @@ -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( diff --git a/Sources/DispatchRedis/RedisConnection.swift b/Sources/DispatchRedis/RedisConnection.swift index 24cf0dd..e165983 100644 --- a/Sources/DispatchRedis/RedisConnection.swift +++ b/Sources/DispatchRedis/RedisConnection.swift @@ -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( diff --git a/Sources/DispatchRedis/RedisPipeline.swift b/Sources/DispatchRedis/RedisPipeline.swift index 599f0e2..c8f6441 100644 --- a/Sources/DispatchRedis/RedisPipeline.swift +++ b/Sources/DispatchRedis/RedisPipeline.swift @@ -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 } diff --git a/Sources/NIORedis/ChannelHandlers/RedisCommandHandler.swift b/Sources/NIORedis/ChannelHandlers/RedisCommandHandler.swift new file mode 100644 index 0000000..ab6cfee --- /dev/null +++ b/Sources/NIORedis/ChannelHandlers/RedisCommandHandler.swift @@ -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 +} + +/// 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] + + 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?) { + let context = unwrapOutboundIn(data) + commandResponseQueue.insert(context.promise, at: 0) + ctx.write(wrapOutboundOut(context.command), promise: promise) + } +} diff --git a/Sources/NIORedis/ChannelHandlers/RedisMessenger.swift b/Sources/NIORedis/ChannelHandlers/RedisMessenger.swift deleted file mode 100644 index 52ff673..0000000 --- a/Sources/NIORedis/ChannelHandlers/RedisMessenger.swift +++ /dev/null @@ -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] - /// 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 { - // 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) - } -} diff --git a/Sources/NIORedis/Commands/BasicCommands.swift b/Sources/NIORedis/Commands/BasicCommands.swift index 9c27880..4e16e11 100644 --- a/Sources/NIORedis/Commands/BasicCommands.swift +++ b/Sources/NIORedis/Commands/BasicCommands.swift @@ -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 { - 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 { - 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 { 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 { - 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 { - 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 { - return command("SET", [RESPValue(bulk: key), RESPValue(bulk: value)]) + return command("SET", arguments: [RESPValue(bulk: key), RESPValue(bulk: value)]) .map { _ in return () } } } diff --git a/Sources/NIORedis/NIORedisConnection.swift b/Sources/NIORedis/NIORedisConnection.swift deleted file mode 100644 index 990ca36..0000000 --- a/Sources/NIORedis/NIORedisConnection.swift +++ /dev/null @@ -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(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 { - 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 { - // 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 - } -} diff --git a/Sources/NIORedis/NIORedisPipeline.swift b/Sources/NIORedis/NIORedisPipeline.swift deleted file mode 100644 index 03c5094..0000000 --- a/Sources/NIORedis/NIORedisPipeline.swift +++ /dev/null @@ -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 - } -} diff --git a/Sources/NIORedis/RedisConnection.swift b/Sources/NIORedis/RedisConnection.swift new file mode 100644 index 0000000..9cabc6d --- /dev/null +++ b/Sources/NIORedis/RedisConnection.swift @@ -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 +/// +/// 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(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 { + 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 { + 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 { + 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) + } +} diff --git a/Sources/NIORedis/NIORedis.swift b/Sources/NIORedis/RedisDriver.swift similarity index 68% rename from Sources/NIORedis/NIORedis.swift rename to Sources/NIORedis/RedisDriver.swift index 8c6adc2..2ff4b3f 100644 --- a/Sources/NIORedis/NIORedis.swift +++ b/Sources/NIORedis/RedisDriver.swift @@ -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(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 { - let channelHandler = RedisMessenger(on: elg.next()) - let bootstrap = ClientBootstrap(group: self.elg) + ) -> EventLoopFuture { + 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 { diff --git a/Sources/NIORedis/RedisPipeline.swift b/Sources/NIORedis/RedisPipeline.swift new file mode 100644 index 0000000..199500b --- /dev/null +++ b/Sources/NIORedis/RedisPipeline.swift @@ -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] + + /// 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) } + ) + } +} diff --git a/Tests/NIORedisTests/Commands/BasicCommandsTests.swift b/Tests/NIORedisTests/Commands/BasicCommandsTests.swift index e3a548b..ec49425 100644 --- a/Tests/NIORedisTests/Commands/BasicCommandsTests.swift +++ b/Tests/NIORedisTests/Commands/BasicCommandsTests.swift @@ -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 { diff --git a/Tests/NIORedisTests/NIORedisTests.swift b/Tests/NIORedisTests/NIORedisTests.swift deleted file mode 100644 index 8781859..0000000 --- a/Tests/NIORedisTests/NIORedisTests.swift +++ /dev/null @@ -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), - ] -} diff --git a/Tests/NIORedisTests/RedisDriverTests.swift b/Tests/NIORedisTests/RedisDriverTests.swift new file mode 100644 index 0000000..029ec8d --- /dev/null +++ b/Tests/NIORedisTests/RedisDriverTests.swift @@ -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), + ] +} diff --git a/Tests/NIORedisTests/NIORedisPipelineTests.swift b/Tests/NIORedisTests/RedisPipelineTests.swift similarity index 88% rename from Tests/NIORedisTests/NIORedisPipelineTests.swift rename to Tests/NIORedisTests/RedisPipelineTests.swift index 8d6d789..b81fbc6 100644 --- a/Tests/NIORedisTests/NIORedisPipelineTests.swift +++ b/Tests/NIORedisTests/RedisPipelineTests.swift @@ -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 { diff --git a/Tests/NIORedisTests/XCTestManifests.swift b/Tests/NIORedisTests/XCTestManifests.swift index 4633259..f8f3148 100644 --- a/Tests/NIORedisTests/XCTestManifests.swift +++ b/Tests/NIORedisTests/XCTestManifests.swift @@ -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