Make engine single queued

Fix polling

allow building of refactor branch

remove self reference

Some refactoring
This commit is contained in:
Erik 2017-05-03 22:48:25 -04:00
parent 7e494f4bcb
commit ed049e888d
No known key found for this signature in database
GPG Key ID: 4930B7C5FBC1A69D
4 changed files with 119 additions and 118 deletions

View File

@ -6,6 +6,7 @@ branches:
only: only:
- master - master
- development - development
- refactor-engine
before_install: before_install:
- brew update - brew update
- brew outdated xctool || brew upgrade xctool - brew outdated xctool || brew upgrade xctool

View File

@ -22,12 +22,11 @@
// 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 SocketEngine : NSObject, URLSessionDelegate, SocketEnginePollable, SocketEngineWebsocket { public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePollable, SocketEngineWebsocket {
public let emitQueue = DispatchQueue(label: "com.socketio.engineEmitQueue", attributes: []) public let engineQueue = DispatchQueue(label: "com.socketio.engineHandleQueue", attributes: [])
public let handleQueue = DispatchQueue(label: "com.socketio.engineHandleQueue", attributes: [])
public let parseQueue = DispatchQueue(label: "com.socketio.engineParseQueue", attributes: [])
public var connectParams: [String: Any]? { public var connectParams: [String: Any]? {
didSet { didSet {
@ -174,6 +173,12 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
/// Starts the connection to the server /// Starts the connection to the server
public func connect() { public func connect() {
engineQueue.async {
self._connect()
}
}
private func _connect() {
if connected { if connected {
DefaultSocketLogger.Logger.error("Engine tried opening while connected. Assuming this was a reconnect", type: logType) DefaultSocketLogger.Logger.error("Engine tried opening while connected. Assuming this was a reconnect", type: logType)
disconnect(reason: "reconnect") disconnect(reason: "reconnect")
@ -191,8 +196,7 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
return return
} }
var reqPolling = URLRequest(url: urlPolling, cachePolicy: .reloadIgnoringLocalCacheData, var reqPolling = URLRequest(url: urlPolling, cachePolicy: .reloadIgnoringLocalCacheData, timeoutInterval: 60.0)
timeoutInterval: 60.0)
if cookies != nil { if cookies != nil {
let headers = HTTPCookie.requestHeaderFields(with: cookies!) let headers = HTTPCookie.requestHeaderFields(with: cookies!)
@ -244,7 +248,6 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
} }
private func createWebsocketAndConnect() { private func createWebsocketAndConnect() {
ws?.delegate = nil ws?.delegate = nil
ws = WebSocket(url: urlWebSocketWithSid as URL) ws = WebSocket(url: urlWebSocketWithSid as URL)
@ -261,7 +264,7 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
} }
} }
ws?.callbackQueue = handleQueue ws?.callbackQueue = engineQueue
ws?.voipEnabled = voipEnabled ws?.voipEnabled = voipEnabled
ws?.delegate = self ws?.delegate = self
ws?.disableSSLCertValidation = selfSigned ws?.disableSSLCertValidation = selfSigned
@ -277,6 +280,12 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
} }
public func disconnect(reason: String) { public func disconnect(reason: String) {
engineQueue.async {
self._disconnect(reason: reason)
}
}
private func _disconnect(reason: String) {
guard connected else { return closeOutEngine(reason: reason) } guard connected else { return closeOutEngine(reason: reason) }
DefaultSocketLogger.Logger.log("Engine is being closed.", type: logType) DefaultSocketLogger.Logger.log("Engine is being closed.", type: logType)
@ -296,12 +305,10 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
// We need to take special care when we're polling that we send it ASAP // We need to take special care when we're polling that we send it ASAP
// Also make sure we're on the emitQueue since we're touching postWait // Also make sure we're on the emitQueue since we're touching postWait
private func disconnectPolling(reason: String) { private func disconnectPolling(reason: String) {
emitQueue.sync { postWait.append(String(SocketEnginePacketType.close.rawValue))
self.postWait.append(String(SocketEnginePacketType.close.rawValue))
let req = self.createRequestForPostWithPostWait() doRequest(for: createRequestForPostWithPostWait()) {_, _, _ in }
self.doRequest(for: req) {_, _, _ in } closeOutEngine(reason: reason)
self.closeOutEngine(reason: reason)
}
} }
public func doFastUpgrade() { public func doFastUpgrade() {
@ -321,16 +328,14 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
private func flushProbeWait() { private func flushProbeWait() {
DefaultSocketLogger.Logger.log("Flushing probe wait", type: logType) DefaultSocketLogger.Logger.log("Flushing probe wait", type: logType)
emitQueue.async { for waiter in probeWait {
for waiter in self.probeWait { write(waiter.msg, withType: waiter.type, withData: waiter.data)
self.write(waiter.msg, withType: waiter.type, withData: waiter.data) }
}
self.probeWait.removeAll(keepingCapacity: false) probeWait.removeAll(keepingCapacity: false)
if self.postWait.count != 0 { if postWait.count != 0 {
self.flushWaitingForPostToWebSocket() flushWaitingForPostToWebSocket()
}
} }
} }
@ -456,13 +461,16 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
// Puts the engine back in its default state // Puts the engine back in its default state
private func resetEngine() { private func resetEngine() {
let queue = OperationQueue()
queue.underlyingQueue = engineQueue
closed = false closed = false
connected = false connected = false
fastUpgrade = false fastUpgrade = false
polling = true polling = true
probing = false probing = false
invalidated = false invalidated = false
session = Foundation.URLSession(configuration: .default, delegate: sessionDelegate, delegateQueue: OperationQueue.main) session = Foundation.URLSession(configuration: .default, delegate: sessionDelegate, delegateQueue: queue)
sid = "" sid = ""
waitingForPoll = false waitingForPoll = false
waitingForPost = false waitingForPost = false
@ -484,8 +492,7 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
pongsMissed += 1 pongsMissed += 1
write("", withType: .ping, withData: []) write("", withType: .ping, withData: [])
let time = DispatchTime.now() + Double(Int64(pingInterval * Double(NSEC_PER_SEC))) / Double(NSEC_PER_SEC) engineQueue.asyncAfter(deadline: DispatchTime.now() + Double(pingInterval)) {[weak self] in self?.sendPing() }
DispatchQueue.main.asyncAfter(deadline: time) {[weak self] in self?.sendPing() }
} }
// Moves from long-polling to websockets // Moves from long-polling to websockets
@ -501,16 +508,16 @@ public final class SocketEngine : NSObject, URLSessionDelegate, SocketEnginePoll
/// Write a message, independent of transport. /// Write a message, independent of transport.
public func write(_ msg: String, withType type: SocketEnginePacketType, withData data: [Data]) { public func write(_ msg: String, withType type: SocketEnginePacketType, withData data: [Data]) {
emitQueue.async { engineQueue.async {
guard self.connected else { return } guard self.connected else { return }
if self.websocket { if self.websocket {
DefaultSocketLogger.Logger.log("Writing ws: %@ has data: %@", DefaultSocketLogger.Logger.log("Writing ws: %@ has data: %@",
type: self.logType, args: msg, data.count != 0) type: self.logType, args: msg, data.count != 0)
self.sendWebSocketMessage(msg, withType: type, withData: data) self.sendWebSocketMessage(msg, withType: type, withData: data)
} else if !self.probing { } else if !self.probing {
DefaultSocketLogger.Logger.log("Writing poll: %@ has data: %@", DefaultSocketLogger.Logger.log("Writing poll: %@ has data: %@",
type: self.logType, args: msg, data.count != 0) type: self.logType, args: msg, data.count != 0)
self.sendPollMessage(msg, withType: type, withData: data) self.sendPollMessage(msg, withType: type, withData: data)
} else { } else {
self.probeWait.append((msg, type, data)) self.probeWait.append((msg, type, data))

View File

@ -86,7 +86,7 @@ extension SocketEnginePollable {
req.httpBody = postData req.httpBody = postData
req.setValue(String(postData.count), forHTTPHeaderField: "Content-Length") req.setValue(String(postData.count), forHTTPHeaderField: "Content-Length")
return req as URLRequest return req
} }
public func doPoll() { public func doPoll() {
@ -94,11 +94,9 @@ extension SocketEnginePollable {
return return
} }
waitingForPoll = true
var req = URLRequest(url: urlPollingWithSid) var req = URLRequest(url: urlPollingWithSid)
req = addHeaders(for: req) req = addHeaders(for: req)
doLongPoll(for: req ) doLongPoll(for: req )
} }
@ -107,12 +105,15 @@ extension SocketEnginePollable {
return return
} }
DefaultSocketLogger.Logger.log("Doing polling request", type: "SocketEnginePolling") DefaultSocketLogger.Logger.log("Doing polling %@ %@", type: "SocketEnginePolling",
args: req.httpMethod ?? "", req)
session?.dataTask(with: req, completionHandler: callback).resume() session?.dataTask(with: req, completionHandler: callback).resume()
} }
func doLongPoll(for req: URLRequest) { func doLongPoll(for req: URLRequest) {
waitingForPoll = true
doRequest(for: req) {[weak self] data, res, err in doRequest(for: req) {[weak self] data, res, err in
guard let this = self, this.polling else { return } guard let this = self, this.polling else { return }
@ -129,9 +130,7 @@ extension SocketEnginePollable {
DefaultSocketLogger.Logger.log("Got polling response", type: "SocketEnginePolling") DefaultSocketLogger.Logger.log("Got polling response", type: "SocketEnginePolling")
if let str = String(data: data!, encoding: String.Encoding.utf8) { if let str = String(data: data!, encoding: String.Encoding.utf8) {
this.parseQueue.async { this.parsePollingMessage(str)
this.parsePollingMessage(str)
}
} }
this.waitingForPoll = false this.waitingForPoll = false
@ -173,11 +172,9 @@ extension SocketEnginePollable {
this.waitingForPost = false this.waitingForPost = false
this.emitQueue.async { if !this.fastUpgrade {
if !this.fastUpgrade { this.flushWaitingForPost()
this.flushWaitingForPost() this.doPoll()
this.doPoll()
}
} }
} }
} }
@ -189,11 +186,9 @@ extension SocketEnginePollable {
while reader.hasNext { while reader.hasNext {
if let n = Int(reader.readUntilOccurence(of: ":")) { if let n = Int(reader.readUntilOccurence(of: ":")) {
let str = reader.read(count: n) parseEngineMessage(reader.read(count: n), fromPolling: true)
handleQueue.async { self.parseEngineMessage(str, fromPolling: true) }
} else { } else {
handleQueue.async { self.parseEngineMessage(str, fromPolling: true) } parseEngineMessage(str, fromPolling: true)
break break
} }
} }

View File

@ -32,15 +32,13 @@ import Foundation
var connectParams: [String: Any]? { get set } var connectParams: [String: Any]? { get set }
var doubleEncodeUTF8: Bool { get } var doubleEncodeUTF8: Bool { get }
var cookies: [HTTPCookie]? { get } var cookies: [HTTPCookie]? { get }
var engineQueue: DispatchQueue { get }
var extraHeaders: [String: String]? { get } var extraHeaders: [String: String]? { get }
var fastUpgrade: Bool { get } var fastUpgrade: Bool { get }
var forcePolling: Bool { get } var forcePolling: Bool { get }
var forceWebsockets: Bool { get } var forceWebsockets: Bool { get }
var parseQueue: DispatchQueue { get }
var polling: Bool { get } var polling: Bool { get }
var probing: Bool { get } var probing: Bool { get }
var emitQueue: DispatchQueue { get }
var handleQueue: DispatchQueue { get }
var sid: String { get } var sid: String { get }
var socketPath: String { get } var socketPath: String { get }
var urlPolling: URL { get } var urlPolling: URL { get }