diff --git a/Sources/SwiftNetwork/Protocols/ProtocolDatagramHandlers.swift b/Sources/SwiftNetwork/Protocols/ProtocolDatagramHandlers.swift index 5d334d6..443369c 100644 --- a/Sources/SwiftNetwork/Protocols/ProtocolDatagramHandlers.swift +++ b/Sources/SwiftNetwork/Protocols/ProtocolDatagramHandlers.swift @@ -121,7 +121,16 @@ extension AutomaticUpperDatagramProcessing where Self: ~Copyable { /// Notifies the upper protocol that frames are available in `upperReceiveQueue`. public func serviceUpperReceiveQueue() { guard !upperReceiveQueue.isEmpty else { return } - upper.deliverInboundDataAvailableEvent(reference) + if upper.isDetached { + // Enqueue pending event until the daragram is completely attached. + // This will ensure that the inboundDataAvailable event is sent for all new datagrams carrying data + let selfReference = self.reference + selfReference.enqueuePendingEventForUpperProtocol( + event: .inboundDataAvailable(selfReference, upper.reference) + ) + } else { + upper.deliverInboundDataAvailableEvent(reference) + } } } diff --git a/Sources/SwiftNetwork/Protocols/ProtocolStreamHandlers.swift b/Sources/SwiftNetwork/Protocols/ProtocolStreamHandlers.swift index 2fde116..2ab7016 100644 --- a/Sources/SwiftNetwork/Protocols/ProtocolStreamHandlers.swift +++ b/Sources/SwiftNetwork/Protocols/ProtocolStreamHandlers.swift @@ -133,7 +133,16 @@ extension AutomaticUpperStreamProcessing where Self: ~Copyable { /// Notifies the upper protocol that frames are available in `upperReceiveQueue`. public func serviceUpperReceiveQueue() { guard !upperReceiveQueue.isEmpty else { return } - upper.deliverInboundDataAvailableEvent(reference) + if upper.isDetached { + // Enqueue pending event until the stream is completely attached. + // This will ensure that the inboundDataAvailable event is sent for all new streams carrying data + let selfReference = self.reference + selfReference.enqueuePendingEventForUpperProtocol( + event: .inboundDataAvailable(selfReference, upper.reference) + ) + } else { + upper.deliverInboundDataAvailableEvent(reference) + } } } diff --git a/Sources/Tools/QUICTransfer/main.swift b/Sources/Tools/QUICTransfer/main.swift index 6e7c47a..9ae73cf 100644 --- a/Sources/Tools/QUICTransfer/main.swift +++ b/Sources/Tools/QUICTransfer/main.swift @@ -62,7 +62,6 @@ final class QUICTransfer { // Create a random payload to send back and forth var payload = [UInt8](repeating: 0, count: sendSize) payload = (0.. 1 ? "s" : "")") let timestart = DispatchTime.now().uptimeNanoseconds @@ -251,53 +250,59 @@ final class QUICTransfer { } var serverStream: StreamUpperHarness? + var writeIndex = 0 + var writeSucceeded = true var totalReadSize = 0 - while index < iterations { - group.enter() + let totalExpectedSize = iterations * payload.count + let doneSemaphore = DispatchSemaphore(value: 0) + + // Client write loop to perform all writes until finished + func writeLoop() { + guard writeIndex < iterations else { return } + guard clientStream.write(payload) else { + loggingHandle.log("Client failed to write at iteration: \(writeIndex)") + writeSucceeded = false + return + } + writeIndex += 1 context.async { - guard clientStream.write(payload) else { - loggingHandle.log("Client failed to write at iteration: \(index)") - group.leave() - return - } - if serverStream == nil { - group.enter() - serverInput.waitForNewFlow { - loggingHandle.log("Server got new inbound flow") - serverStream = serverInput.upperHarnesses.last - group.leave() - } - } - group.leave() + writeLoop() } - group.wait() + } - guard let serverStream else { - return 0 + // Server read loop: keeps draining inbound data as it arrives. + // Readloop used for multiple iterations + func readLoop(stream: StreamUpperHarness) { + stream.waitForInboundDataAvailable { available in + guard available else { return } + totalReadSize += stream.readAndDrop() + if totalReadSize >= totalExpectedSize { + doneSemaphore.signal() + } else { + readLoop(stream: stream) + } } + } - group.enter() - context.async { - var serverReadDataSizeForIteration = 0 - var serverReadCompletion: ((Bool) -> Void)? = nil - serverReadCompletion = { _ in - let readBytes = serverStream.readAndDrop() - if readBytes > 0 { - serverReadDataSizeForIteration += readBytes - totalReadSize += readBytes - } - if serverReadDataSizeForIteration >= payload.count { - serverReadCompletion = nil - index += 1 - group.leave() - } else { - serverStream.waitForInboundDataAvailable(completion: serverReadCompletion!) - } + context.async { + // Setup inbound flow observer and then start the client write loop + serverInput.waitForNewFlow { + serverStream = serverInput.upperHarnesses.last + if let serverStream { + readLoop(stream: serverStream) + } else { + doneSemaphore.signal() } - serverStream.waitForInboundDataAvailable(completion: serverReadCompletion!) } - group.wait() + writeLoop() } + doneSemaphore.wait() + + guard serverStream != nil, writeSucceeded else { + return 0 + } + + let index = min(totalReadSize / payload.count, iterations) group.enter() context.async {