Make client single queued
Shouldn't need to protoct ack generation now that client is single queued
This commit is contained in:
parent
ed049e888d
commit
84dd3078d8
@ -22,6 +22,7 @@
|
|||||||
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
|
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
|
||||||
// THE SOFTWARE.
|
// THE SOFTWARE.
|
||||||
|
|
||||||
|
import Dispatch
|
||||||
import Foundation
|
import Foundation
|
||||||
|
|
||||||
public final class SocketAckEmitter : NSObject {
|
public final class SocketAckEmitter : NSObject {
|
||||||
@ -31,57 +32,52 @@ public final class SocketAckEmitter : NSObject {
|
|||||||
public var expected: Bool {
|
public var expected: Bool {
|
||||||
return ackNum != -1
|
return ackNum != -1
|
||||||
}
|
}
|
||||||
|
|
||||||
init(socket: SocketIOClient, ackNum: Int) {
|
init(socket: SocketIOClient, ackNum: Int) {
|
||||||
self.socket = socket
|
self.socket = socket
|
||||||
self.ackNum = ackNum
|
self.ackNum = ackNum
|
||||||
}
|
}
|
||||||
|
|
||||||
public func with(_ items: SocketData...) {
|
public func with(_ items: SocketData...) {
|
||||||
guard ackNum != -1 else { return }
|
guard ackNum != -1 else { return }
|
||||||
|
|
||||||
socket.emitAck(ackNum, with: items)
|
socket.emitAck(ackNum, with: items)
|
||||||
}
|
}
|
||||||
|
|
||||||
public func with(_ items: [Any]) {
|
public func with(_ items: [Any]) {
|
||||||
guard ackNum != -1 else { return }
|
guard ackNum != -1 else { return }
|
||||||
|
|
||||||
socket.emitAck(ackNum, with: items)
|
socket.emitAck(ackNum, with: items)
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public final class OnAckCallback : NSObject {
|
public final class OnAckCallback : NSObject {
|
||||||
private let ackNumber: Int
|
private let ackNumber: Int
|
||||||
private let items: [Any]
|
private let items: [Any]
|
||||||
private weak var socket: SocketIOClient?
|
private weak var socket: SocketIOClient?
|
||||||
|
|
||||||
init(ackNumber: Int, items: [Any], socket: SocketIOClient) {
|
init(ackNumber: Int, items: [Any], socket: SocketIOClient) {
|
||||||
self.ackNumber = ackNumber
|
self.ackNumber = ackNumber
|
||||||
self.items = items
|
self.items = items
|
||||||
self.socket = socket
|
self.socket = socket
|
||||||
}
|
}
|
||||||
|
|
||||||
deinit {
|
deinit {
|
||||||
DefaultSocketLogger.Logger.log("OnAckCallback for \(ackNumber) being released", type: "OnAckCallback")
|
DefaultSocketLogger.Logger.log("OnAckCallback for \(ackNumber) being released", type: "OnAckCallback")
|
||||||
}
|
}
|
||||||
|
|
||||||
public func timingOut(after seconds: Int, callback: @escaping AckCallback) {
|
public func timingOut(after seconds: Int, callback: @escaping AckCallback) {
|
||||||
guard let socket = self.socket else { return }
|
guard let socket = self.socket else { return }
|
||||||
|
|
||||||
socket.ackQueue.sync() {
|
socket.ackHandlers.addAck(ackNumber, callback: callback)
|
||||||
socket.ackHandlers.addAck(ackNumber, callback: callback)
|
|
||||||
}
|
|
||||||
|
|
||||||
socket._emit(items, ack: ackNumber)
|
socket._emit(items, ack: ackNumber)
|
||||||
|
|
||||||
guard seconds != 0 else { return }
|
guard seconds != 0 else { return }
|
||||||
|
|
||||||
let time = DispatchTime.now() + Double(UInt64(seconds) * NSEC_PER_SEC) / Double(NSEC_PER_SEC)
|
socket.handleQueue.asyncAfter(deadline: DispatchTime.now() + Double(seconds)) {
|
||||||
|
|
||||||
socket.handleQueue.asyncAfter(deadline: time) {
|
|
||||||
socket.ackHandlers.timeoutAck(self.ackNumber, onQueue: socket.handleQueue)
|
socket.ackHandlers.timeoutAck(self.ackNumber, onQueue: socket.handleQueue)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@ -22,6 +22,7 @@
|
|||||||
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
|
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
|
||||||
// THE SOFTWARE.
|
// THE SOFTWARE.
|
||||||
|
|
||||||
|
import Dispatch
|
||||||
import Foundation
|
import Foundation
|
||||||
|
|
||||||
public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineClient, SocketParsable {
|
public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineClient, SocketParsable {
|
||||||
@ -41,43 +42,38 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
}
|
}
|
||||||
|
|
||||||
public var forceNew = false
|
public var forceNew = false
|
||||||
|
public var handleQueue = DispatchQueue.main
|
||||||
public var nsp = "/"
|
public var nsp = "/"
|
||||||
public var config: SocketIOClientConfiguration
|
public var config: SocketIOClientConfiguration
|
||||||
public var reconnects = true
|
public var reconnects = true
|
||||||
public var reconnectWait = 10
|
public var reconnectWait = 10
|
||||||
|
|
||||||
private let logType = "SocketIOClient"
|
private let logType = "SocketIOClient"
|
||||||
private let parseQueue = DispatchQueue(label: "com.socketio.parseQueue")
|
|
||||||
|
|
||||||
private var anyHandler: ((SocketAnyEvent) -> Void)?
|
private var anyHandler: ((SocketAnyEvent) -> Void)?
|
||||||
private var currentReconnectAttempt = 0
|
private var currentReconnectAttempt = 0
|
||||||
private var handlers = [SocketEventHandler]()
|
private var handlers = [SocketEventHandler]()
|
||||||
private var reconnecting = false
|
private var reconnecting = false
|
||||||
|
|
||||||
private let ackSemaphore = DispatchSemaphore(value: 1)
|
|
||||||
private(set) var currentAck = -1
|
private(set) var currentAck = -1
|
||||||
private(set) var handleQueue = DispatchQueue.main
|
|
||||||
private(set) var reconnectAttempts = -1
|
private(set) var reconnectAttempts = -1
|
||||||
|
|
||||||
let ackQueue = DispatchQueue(label: "com.socketio.ackQueue")
|
|
||||||
let emitQueue = DispatchQueue(label: "com.socketio.emitQueue")
|
|
||||||
|
|
||||||
var ackHandlers = SocketAckManager()
|
var ackHandlers = SocketAckManager()
|
||||||
var waitingPackets = [SocketPacket]()
|
var waitingPackets = [SocketPacket]()
|
||||||
|
|
||||||
public var sid: String? {
|
public var sid: String? {
|
||||||
return engine?.sid
|
return engine?.sid
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Type safe way to create a new SocketIOClient. opts can be omitted
|
/// Type safe way to create a new SocketIOClient. opts can be omitted
|
||||||
public init(socketURL: URL, config: SocketIOClientConfiguration = []) {
|
public init(socketURL: URL, config: SocketIOClientConfiguration = []) {
|
||||||
self.config = config
|
self.config = config
|
||||||
self.socketURL = socketURL
|
self.socketURL = socketURL
|
||||||
|
|
||||||
if socketURL.absoluteString.hasPrefix("https://") {
|
if socketURL.absoluteString.hasPrefix("https://") {
|
||||||
self.config.insert(.secure(true))
|
self.config.insert(.secure(true))
|
||||||
}
|
}
|
||||||
|
|
||||||
for option in config {
|
for option in config {
|
||||||
switch option {
|
switch option {
|
||||||
case let .reconnects(reconnects):
|
case let .reconnects(reconnects):
|
||||||
@ -102,10 +98,10 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
}
|
}
|
||||||
|
|
||||||
self.config.insert(.path("/socket.io/"), replacing: false)
|
self.config.insert(.path("/socket.io/"), replacing: false)
|
||||||
|
|
||||||
super.init()
|
super.init()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Not so type safe way to create a SocketIOClient, meant for Objective-C compatiblity.
|
/// Not so type safe way to create a SocketIOClient, meant for Objective-C compatiblity.
|
||||||
/// If using Swift it's recommended to use `init(socketURL: NSURL, options: Set<SocketIOClientOption>)`
|
/// If using Swift it's recommended to use `init(socketURL: NSURL, options: Set<SocketIOClientOption>)`
|
||||||
public convenience init(socketURL: NSURL, config: NSDictionary?) {
|
public convenience init(socketURL: NSURL, config: NSDictionary?) {
|
||||||
@ -148,30 +144,25 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
} else {
|
} else {
|
||||||
engine?.connect()
|
engine?.connect()
|
||||||
}
|
}
|
||||||
|
|
||||||
guard timeoutAfter != 0 else { return }
|
guard timeoutAfter != 0 else { return }
|
||||||
|
|
||||||
let time = DispatchTime.now() + Double(UInt64(timeoutAfter) * NSEC_PER_SEC) / Double(NSEC_PER_SEC)
|
let time = DispatchTime.now() + Double(UInt64(timeoutAfter) * NSEC_PER_SEC) / Double(NSEC_PER_SEC)
|
||||||
|
|
||||||
handleQueue.asyncAfter(deadline: time) {[weak self] in
|
handleQueue.asyncAfter(deadline: time) {[weak self] in
|
||||||
guard let this = self, this.status != .connected && this.status != .disconnected else { return }
|
guard let this = self, this.status != .connected && this.status != .disconnected else { return }
|
||||||
|
|
||||||
this.status = .disconnected
|
this.status = .disconnected
|
||||||
this.engine?.disconnect(reason: "Connect timeout")
|
this.engine?.disconnect(reason: "Connect timeout")
|
||||||
|
|
||||||
handler?()
|
handler?()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private func nextAck() -> Int {
|
|
||||||
ackSemaphore.wait()
|
|
||||||
defer { ackSemaphore.signal() }
|
|
||||||
currentAck += 1
|
|
||||||
return currentAck
|
|
||||||
}
|
|
||||||
|
|
||||||
private func createOnAck(_ items: [Any]) -> OnAckCallback {
|
private func createOnAck(_ items: [Any]) -> OnAckCallback {
|
||||||
return OnAckCallback(ackNumber: nextAck(), items: items, socket: self)
|
currentAck += 1
|
||||||
|
|
||||||
|
return OnAckCallback(ackNumber: currentAck, items: items, socket: self)
|
||||||
}
|
}
|
||||||
|
|
||||||
func didConnect() {
|
func didConnect() {
|
||||||
@ -214,7 +205,7 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
handleEvent("error", data: ["Tried emitting \(event) when not connected"], isInternalMessage: true)
|
handleEvent("error", data: ["Tried emitting \(event) when not connected"], isInternalMessage: true)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
_emit([event] + items)
|
_emit([event] + items)
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -230,40 +221,40 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
}
|
}
|
||||||
|
|
||||||
func _emit(_ data: [Any], ack: Int? = nil) {
|
func _emit(_ data: [Any], ack: Int? = nil) {
|
||||||
emitQueue.async {
|
guard status == .connected else {
|
||||||
guard self.status == .connected else {
|
handleEvent("error", data: ["Tried emitting when not connected"], isInternalMessage: true)
|
||||||
self.handleEvent("error", data: ["Tried emitting when not connected"], isInternalMessage: true)
|
return
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
let packet = SocketPacket.packetFromEmit(data, id: ack ?? -1, nsp: self.nsp, ack: false)
|
|
||||||
let str = packet.packetString
|
|
||||||
|
|
||||||
DefaultSocketLogger.Logger.log("Emitting: %@", type: self.logType, args: str)
|
|
||||||
|
|
||||||
self.engine?.send(str, withData: packet.binary)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let packet = SocketPacket.packetFromEmit(data, id: ack ?? -1, nsp: nsp, ack: false)
|
||||||
|
let str = packet.packetString
|
||||||
|
|
||||||
|
DefaultSocketLogger.Logger.log("Emitting: %@", type: logType, args: str)
|
||||||
|
|
||||||
|
engine?.send(str, withData: packet.binary)
|
||||||
}
|
}
|
||||||
|
|
||||||
// If the server wants to know that the client received data
|
// If the server wants to know that the client received data
|
||||||
func emitAck(_ ack: Int, with items: [Any]) {
|
func emitAck(_ ack: Int, with items: [Any]) {
|
||||||
emitQueue.async {
|
guard status == .connected else { return }
|
||||||
guard self.status == .connected else { return }
|
|
||||||
|
let packet = SocketPacket.packetFromEmit(items, id: ack, nsp: nsp, ack: true)
|
||||||
let packet = SocketPacket.packetFromEmit(items, id: ack, nsp: self.nsp, ack: true)
|
let str = packet.packetString
|
||||||
let str = packet.packetString
|
|
||||||
|
DefaultSocketLogger.Logger.log("Emitting Ack: %@", type: logType, args: str)
|
||||||
DefaultSocketLogger.Logger.log("Emitting Ack: %@", type: self.logType, args: str)
|
|
||||||
|
engine?.send(str, withData: packet.binary)
|
||||||
self.engine?.send(str, withData: packet.binary)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public func engineDidClose(reason: String) {
|
public func engineDidClose(reason: String) {
|
||||||
parseQueue.async {
|
handleQueue.async {
|
||||||
self.waitingPackets.removeAll()
|
self._engineDidClose(reason: reason)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private func _engineDidClose(reason: String) {
|
||||||
|
waitingPackets.removeAll()
|
||||||
|
|
||||||
if status != .disconnected {
|
if status != .disconnected {
|
||||||
status = .notConnected
|
status = .notConnected
|
||||||
}
|
}
|
||||||
@ -276,13 +267,19 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// error
|
|
||||||
public func engineDidError(reason: String) {
|
public func engineDidError(reason: String) {
|
||||||
|
handleQueue.async {
|
||||||
|
self._engineDidError(reason: reason)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// error
|
||||||
|
private func _engineDidError(reason: String) {
|
||||||
DefaultSocketLogger.Logger.error("%@", type: logType, args: reason)
|
DefaultSocketLogger.Logger.error("%@", type: logType, args: reason)
|
||||||
|
|
||||||
handleEvent("error", data: [reason], isInternalMessage: true)
|
handleEvent("error", data: [reason], isInternalMessage: true)
|
||||||
}
|
}
|
||||||
|
|
||||||
public func engineDidOpen(reason: String) {
|
public func engineDidOpen(reason: String) {
|
||||||
DefaultSocketLogger.Logger.log(reason, type: "SocketEngineClient")
|
DefaultSocketLogger.Logger.log(reason, type: "SocketEngineClient")
|
||||||
}
|
}
|
||||||
@ -293,9 +290,7 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
|
|
||||||
DefaultSocketLogger.Logger.log("Handling ack: %@ with data: %@", type: logType, args: ack, data)
|
DefaultSocketLogger.Logger.log("Handling ack: %@ with data: %@", type: logType, args: ack, data)
|
||||||
|
|
||||||
handleQueue.async() {
|
ackHandlers.executeAck(ack, with: data, onQueue: handleQueue)
|
||||||
self.ackHandlers.executeAck(ack, with: data, onQueue: self.handleQueue)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Causes an event to be handled. Only use if you know what you're doing.
|
/// Causes an event to be handled. Only use if you know what you're doing.
|
||||||
@ -304,12 +299,10 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
|
|
||||||
DefaultSocketLogger.Logger.log("Handling event: %@ with data: %@", type: logType, args: event, data)
|
DefaultSocketLogger.Logger.log("Handling event: %@ with data: %@", type: logType, args: event, data)
|
||||||
|
|
||||||
handleQueue.async {
|
anyHandler?(SocketAnyEvent(event: event, items: data))
|
||||||
self.anyHandler?(SocketAnyEvent(event: event, items: data))
|
|
||||||
|
|
||||||
for handler in self.handlers where handler.event == event {
|
for handler in handlers where handler.event == event {
|
||||||
handler.executeCallback(with: data, withAck: ack, withSocket: self)
|
handler.executeCallback(with: data, withAck: ack, withSocket: self)
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -384,17 +377,17 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
public func parseEngineMessage(_ msg: String) {
|
public func parseEngineMessage(_ msg: String) {
|
||||||
DefaultSocketLogger.Logger.log("Should parse message: %@", type: "SocketIOClient", args: msg)
|
DefaultSocketLogger.Logger.log("Should parse message: %@", type: "SocketIOClient", args: msg)
|
||||||
|
|
||||||
parseQueue.async { self.parseSocketMessage(msg) }
|
handleQueue.async { self.parseSocketMessage(msg) }
|
||||||
}
|
}
|
||||||
|
|
||||||
public func parseEngineBinaryData(_ data: Data) {
|
public func parseEngineBinaryData(_ data: Data) {
|
||||||
parseQueue.async { self.parseBinaryData(data) }
|
handleQueue.async { self.parseBinaryData(data) }
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Tries to reconnect to the server.
|
/// Tries to reconnect to the server.
|
||||||
public func reconnect() {
|
public func reconnect() {
|
||||||
guard !reconnecting else { return }
|
guard !reconnecting else { return }
|
||||||
|
|
||||||
engine?.disconnect(reason: "manual reconnect")
|
engine?.disconnect(reason: "manual reconnect")
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -409,7 +402,7 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
|
|
||||||
DefaultSocketLogger.Logger.log("Starting reconnect", type: logType)
|
DefaultSocketLogger.Logger.log("Starting reconnect", type: logType)
|
||||||
handleEvent("reconnect", data: [reason], isInternalMessage: true)
|
handleEvent("reconnect", data: [reason], isInternalMessage: true)
|
||||||
|
|
||||||
_tryReconnect()
|
_tryReconnect()
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -425,26 +418,24 @@ public final class SocketIOClient : NSObject, SocketIOClientSpec, SocketEngineCl
|
|||||||
|
|
||||||
currentReconnectAttempt += 1
|
currentReconnectAttempt += 1
|
||||||
connect()
|
connect()
|
||||||
|
|
||||||
let deadline = DispatchTime.now() + Double(Int64(UInt64(reconnectWait) * NSEC_PER_SEC)) / Double(NSEC_PER_SEC)
|
handleQueue.asyncAfter(deadline: DispatchTime.now() + Double(reconnectWait), execute: _tryReconnect)
|
||||||
|
|
||||||
DispatchQueue.main.asyncAfter(deadline: deadline, execute: _tryReconnect)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Test properties
|
// Test properties
|
||||||
|
|
||||||
var testHandlers: [SocketEventHandler] {
|
var testHandlers: [SocketEventHandler] {
|
||||||
return handlers
|
return handlers
|
||||||
}
|
}
|
||||||
|
|
||||||
func setTestable() {
|
func setTestable() {
|
||||||
status = .connected
|
status = .connected
|
||||||
}
|
}
|
||||||
|
|
||||||
func setTestEngine(_ engine: SocketEngineSpec?) {
|
func setTestEngine(_ engine: SocketEngineSpec?) {
|
||||||
self.engine = engine
|
self.engine = engine
|
||||||
}
|
}
|
||||||
|
|
||||||
func emitTest(event: String, _ data: Any...) {
|
func emitTest(event: String, _ data: Any...) {
|
||||||
_emit([event] + data)
|
_emit([event] + data)
|
||||||
}
|
}
|
||||||
|
|||||||
@ -22,10 +22,13 @@
|
|||||||
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
|
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
|
||||||
// THE SOFTWARE.
|
// THE SOFTWARE.
|
||||||
|
|
||||||
|
import Dispatch
|
||||||
|
|
||||||
protocol SocketIOClientSpec : class {
|
protocol SocketIOClientSpec : class {
|
||||||
|
var handleQueue: DispatchQueue { get set }
|
||||||
var nsp: String { get set }
|
var nsp: String { get set }
|
||||||
var waitingPackets: [SocketPacket] { get set }
|
var waitingPackets: [SocketPacket] { get set }
|
||||||
|
|
||||||
func didConnect()
|
func didConnect()
|
||||||
func didDisconnect(reason: String)
|
func didDisconnect(reason: String)
|
||||||
func didError(reason: String)
|
func didError(reason: String)
|
||||||
@ -37,7 +40,7 @@ protocol SocketIOClientSpec : class {
|
|||||||
extension SocketIOClientSpec {
|
extension SocketIOClientSpec {
|
||||||
func didError(reason: String) {
|
func didError(reason: String) {
|
||||||
DefaultSocketLogger.Logger.error("%@", type: "SocketIOClient", args: reason)
|
DefaultSocketLogger.Logger.error("%@", type: "SocketIOClient", args: reason)
|
||||||
|
|
||||||
handleEvent("error", data: [reason], isInternalMessage: true, withAck: -1)
|
handleEvent("error", data: [reason], isInternalMessage: true, withAck: -1)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user