Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion Sources/SwiftNetwork/Protocols/ProtocolDatagramHandlers.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
}

Expand Down
11 changes: 10 additions & 1 deletion Sources/SwiftNetwork/Protocols/ProtocolStreamHandlers.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
}

Expand Down
83 changes: 44 additions & 39 deletions Sources/Tools/QUICTransfer/main.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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..<sendSize).map { _ in UInt8.random(in: 0...255) }
var index = 0
print("Running QUIC transfer, transferring \(iterations) packet\(iterations > 1 ? "s" : "")")
let timestart = DispatchTime.now().uptimeNanoseconds

Expand Down Expand Up @@ -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 {
Expand Down
Loading