a lot
This commit is contained in:
@@ -10,6 +10,12 @@ final class PeerSession {
|
||||
var onStop: ((PeerSession) -> Void)?
|
||||
private let queue: DispatchQueue
|
||||
private var stopped = false
|
||||
private struct Pending { let data: Data; let bulk: Bool }
|
||||
private var priorityQueue: [Pending] = []
|
||||
private var bulkQueue: [Pending] = []
|
||||
private var bulkBytes = 0
|
||||
private var sending = false
|
||||
private static let maximumPendingBulkBytes = 512 * 1024
|
||||
|
||||
init(connection: NWConnection, outbound: Bool, queue: DispatchQueue) {
|
||||
self.connection = connection; self.outbound = outbound; self.queue = queue
|
||||
@@ -30,9 +36,13 @@ final class PeerSession {
|
||||
|
||||
func send(_ message: WireMessage) {
|
||||
guard !stopped, let data = try? FrameCodec.encode(WireEnvelope(message: message)) else { return }
|
||||
connection.send(content: data, completion: .contentProcessed { [weak self] error in
|
||||
if error != nil { self?.stop() }
|
||||
})
|
||||
let bulk: Bool
|
||||
if case .packet(let packet) = message { bulk = packet.trafficClass == .nativeIPv6 } else { bulk = false }
|
||||
if bulk {
|
||||
guard bulkBytes + data.count <= Self.maximumPendingBulkBytes else { return }
|
||||
bulkBytes += data.count; bulkQueue.append(.init(data: data, bulk: true))
|
||||
} else { priorityQueue.append(.init(data: data, bulk: false)) }
|
||||
pumpWrites()
|
||||
}
|
||||
|
||||
func stop() {
|
||||
@@ -54,6 +64,22 @@ final class PeerSession {
|
||||
}
|
||||
}
|
||||
|
||||
private func pumpWrites() {
|
||||
guard !stopped, !sending else { return }
|
||||
let pending: Pending
|
||||
if !priorityQueue.isEmpty { pending = priorityQueue.removeFirst() }
|
||||
else if !bulkQueue.isEmpty { pending = bulkQueue.removeFirst(); bulkBytes -= pending.data.count }
|
||||
else { return }
|
||||
sending = true
|
||||
connection.send(content: pending.data, completion: .contentProcessed { [weak self] error in
|
||||
guard let self else { return }
|
||||
self.queue.async {
|
||||
self.sending = false
|
||||
if error != nil { self.stop() } else { self.pumpWrites() }
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
private func receiveExactly(_ count: Int, accumulated: Data, completion: @escaping (Data?) -> Void) {
|
||||
if accumulated.count == count { completion(accumulated); return }
|
||||
connection.receive(minimumIncompleteLength: 1, maximumLength: count - accumulated.count) { [weak self] data, _, complete, error in
|
||||
|
||||
Reference in New Issue
Block a user