Skip to content

Commit 8939ae1

Browse files
authored
Merge pull request #16 from apple/agnosticDev/FixQUICTransfer
Fix QUICTransfer
2 parents 64c7d51 + 7fbd2cf commit 8939ae1

1 file changed

Lines changed: 118 additions & 65 deletions

File tree

Sources/Tools/QUICTransfer/main.swift

Lines changed: 118 additions & 65 deletions
Original file line numberDiff line numberDiff line change
@@ -43,9 +43,18 @@ final class QUICTransfer {
4343
let NSEC_PER_MSEC = UInt64(Duration.milliseconds(1) / Duration.nanoseconds(1))
4444
var serverSigningKey = P256.Signing.PrivateKey()
4545

46-
func run(iterations: Int, loggingHandle: LoggingHandle, group: DispatchGroup, sendSize: Int) -> Double {
47-
let ipv4Client = Endpoint(address: IPv4Address(localIPv4Address)!, port: 0)
48-
let ipv4Server = Endpoint(address: IPv4Address(remoteIPv4Address)!, port: 0)
46+
func run(
47+
iterations: Int,
48+
loggingHandle: LoggingHandle,
49+
group: DispatchGroup,
50+
sendSize: Int,
51+
linkDelay: NetworkDuration = .zero
52+
) -> Double {
53+
let ipv4Client = Endpoint(address: IPv4Address(localIPv4Address)!, port: 1234)
54+
let ipv4Server = Endpoint(address: IPv4Address(remoteIPv4Address)!, port: 2345)
55+
var clientStream: StreamUpperHarness? = nil
56+
var clientInput: NewStreamFlowHarness? = nil
57+
var serverInput: NewStreamFlowHarness? = nil
4958
// Create a random payload to send back and forth
5059
var payload = [UInt8](repeating: 0, count: sendSize)
5160
payload = (0..<sendSize).map { _ in UInt8.random(in: 0...255) }
@@ -57,10 +66,15 @@ final class QUICTransfer {
5766
var clientParameters = Parameters()
5867
let context = NetworkContext(identifier: "QUICTransfer")
5968
clientParameters.context = context
69+
let path = PathProperties(parameters: clientParameters)
70+
71+
var serverParameters = Parameters()
72+
serverParameters.isServer = true
73+
serverParameters.context = context
74+
6075
context.activate()
6176
context.async {
6277
// Client
63-
let path = PathProperties(parameters: clientParameters)
6478
let clientIP = IPProtocol.instance(context: clientParameters.context)
6579
let clientIPOptions = IPProtocol.options()
6680
clientIPOptions.setLogID(prefix: "C", parent: "1", protocolLogIDNumber: 3)
@@ -87,8 +101,14 @@ final class QUICTransfer {
87101
clientQUICOptions.setProtocolInstance(clientQUIC)
88102
clientParameters.defaultStack.prepend(applicationProtocol: .quic(clientQUICOptions))
89103

104+
let clientOutput = BridgeDatagramProtocol.instance(context: clientParameters.context)
105+
let bridgeOptions = BridgeDatagramProtocol.options()
106+
bridgeOptions.linkDelay = linkDelay
107+
bridgeOptions.setProtocolInstance(clientOutput)
108+
clientParameters.defaultStack.link = .custom(bridgeOptions)
109+
90110
let clientListenerLinkage = StreamListenerLinkage(reference: clientQUIC)
91-
let clientInput = StreamUpperHarness(
111+
clientInput = NewStreamFlowHarness(
92112
identifier: "Client",
93113
local: ipv4Client,
94114
remote: ipv4Server,
@@ -97,14 +117,20 @@ final class QUICTransfer {
97117
context: context,
98118
listenerProtocol: clientListenerLinkage
99119
)
120+
121+
clientStream = StreamUpperHarness(
122+
identifier: "C1",
123+
local: ipv4Client,
124+
remote: ipv4Server,
125+
parameters: clientParameters,
126+
path: path,
127+
context: context,
128+
listenerProtocol: clientListenerLinkage
129+
)
100130
guard let clientInput else {
131+
group.leave()
101132
return
102133
}
103-
104-
let clientOutput = DatagramLowerHarness(
105-
identifier: "Client",
106-
context: clientParameters.context
107-
)
108134
do {
109135
try clientQUIC.attachLowerDatagramProtocolForNewPath(
110136
clientUDP,
@@ -121,22 +147,20 @@ final class QUICTransfer {
121147
path: path
122148
)
123149
try clientIP.attachLowerDatagramProtocol(
124-
clientOutput.reference,
150+
clientOutput,
125151
remote: ipv4Server,
126152
local: ipv4Client,
127153
parameters: clientParameters,
128154
path: path
129155
)
130156
} catch {
131157
loggingHandle.log("Failed to attach client IP to lower protocol")
158+
group.leave()
132159
return
133160
}
134161
// Server
135-
var serverParameters = Parameters()
136-
serverParameters.isServer = true
137-
serverParameters.context = context
138162
let serverPath = PathProperties(parameters: serverParameters)
139-
let serverIP = IPProtocol.instance(context: clientParameters.context)
163+
let serverIP = IPProtocol.instance(context: context)
140164
let serverIPOptions = IPProtocol.options()
141165
serverIPOptions.setLogID(prefix: "L", parent: "1", protocolLogIDNumber: 3)
142166
clientIPOptions.setProtocolInstance(serverIP)
@@ -161,8 +185,14 @@ final class QUICTransfer {
161185
serverQUICOptions.setProtocolInstance(serverQUIC)
162186
serverParameters.defaultStack.prepend(applicationProtocol: .quic(serverQUICOptions))
163187

188+
let serverOutput = BridgeDatagramProtocol.instance(context: context)
189+
let serverBridgeOptions = BridgeDatagramProtocol.options()
190+
serverBridgeOptions.linkDelay = linkDelay
191+
serverBridgeOptions.setProtocolInstance(serverOutput)
192+
serverParameters.defaultStack.link = .custom(serverBridgeOptions)
193+
164194
let serverListenerLinkage = StreamListenerLinkage(reference: serverQUIC)
165-
let serverInput = NewStreamFlowHarness(
195+
serverInput = NewStreamFlowHarness(
166196
identifier: "Server",
167197
local: ipv4Server,
168198
remote: ipv4Client,
@@ -171,14 +201,11 @@ final class QUICTransfer {
171201
context: serverParameters.context,
172202
listenerProtocol: serverListenerLinkage
173203
)
204+
174205
guard let serverInput else {
206+
group.leave()
175207
return
176208
}
177-
178-
let serverOutput = DatagramLowerHarness(
179-
identifier: "Server",
180-
context: clientParameters.context
181-
)
182209
do {
183210
try serverQUIC.attachLowerDatagramProtocolForNewPath(
184211
serverUDP,
@@ -196,74 +223,89 @@ final class QUICTransfer {
196223
path: path
197224
)
198225
try serverIP.attachLowerDatagramProtocol(
199-
serverOutput.reference,
226+
serverOutput,
200227
remote: ipv4Client,
201228
local: ipv4Server,
202229
parameters: serverParameters,
203230
path: serverPath
204231
)
205232
} catch {
206233
loggingHandle.log("Failed to attach server IP to lower protocol")
234+
group.leave()
207235
return
208236
}
209-
var serverConnected: Bool = false
210237
serverInput.start { connected in
211-
serverConnected = connected
238+
// Server connected event
239+
group.leave()
212240
}
213241
clientInput.start()
214-
// The client will complete TLS before the server
215-
while !serverConnected {
216-
let _ = self.dataBenchmarkUtility.loopOutputHandlerPackets(
217-
sender: clientOutput,
218-
receiver: serverOutput,
219-
maximumBurst: 10
220-
)
221-
let _ = self.dataBenchmarkUtility.loopOutputHandlerPackets(
222-
sender: serverOutput,
223-
receiver: clientOutput,
224-
maximumBurst: 10
225-
)
226-
}
227-
// Transfer all of the data
228-
for _ in 0..<iterations {
229-
group.enter()
230-
let writeSuccess = clientInput.write(payload)
231-
guard writeSuccess else {
242+
clientStream?.start()
243+
}
244+
group.wait()
245+
guard let serverInput, let clientInput, let clientStream else {
246+
return 0
247+
}
248+
var serverStream: StreamUpperHarness?
249+
250+
var totalReadSize = 0
251+
while index < iterations {
252+
group.enter()
253+
context.async {
254+
guard clientStream.write(payload) else {
255+
loggingHandle.log("Client failed to write at iteration: \(index)")
232256
group.leave()
233-
print("Issue took place writing to the client")
234-
break
257+
return
235258
}
236-
var payloadReceived = false
237-
var readDataSize = 0
238-
while !payloadReceived {
239-
let _ = self.dataBenchmarkUtility.loopOutputHandlerPackets(
240-
sender: clientOutput,
241-
receiver: serverOutput,
242-
maximumBurst: 50
243-
)
244-
let _ = self.dataBenchmarkUtility.loopOutputHandlerPackets(
245-
sender: serverOutput,
246-
receiver: clientOutput,
247-
maximumBurst: 50
248-
)
249-
250-
if let readBytes = serverInput.upperHarnesses.first?.readAndDrop() {
251-
readDataSize += readBytes
252-
if readDataSize == sendSize {
253-
payloadReceived = true
254-
}
259+
if serverStream == nil {
260+
group.enter()
261+
serverInput.waitForNewFlow {
262+
loggingHandle.log("Server got new inbound flow")
263+
serverStream = serverInput.upperHarnesses.last
264+
group.leave()
255265
}
256266
}
257-
index += 1
258267
group.leave()
259268
}
269+
group.wait()
270+
271+
guard let serverStream else {
272+
return 0
273+
}
274+
275+
group.enter()
276+
context.async {
277+
var serverReadDataSizeForIteration = 0
278+
var serverReadCompletion: ((Bool) -> Void)? = nil
279+
serverReadCompletion = { _ in
280+
let readBytes = serverStream.readAndDrop()
281+
if readBytes > 0 {
282+
serverReadDataSizeForIteration += readBytes
283+
totalReadSize += readBytes
284+
}
285+
if serverReadDataSizeForIteration >= payload.count {
286+
serverReadCompletion = nil
287+
index += 1
288+
group.leave()
289+
} else {
290+
serverStream.waitForInboundDataAvailable(completion: serverReadCompletion!)
291+
}
292+
}
293+
serverStream.waitForInboundDataAvailable(completion: serverReadCompletion!)
294+
}
295+
group.wait()
296+
}
297+
298+
group.enter()
299+
context.async {
300+
clientStream.stop()
260301
clientInput.stop()
261302
clientInput.teardown()
262303
serverInput.stop()
263304
serverInput.teardown()
264305
group.leave()
265306
}
266307
group.wait()
308+
267309
print("Completed \(index) / \(iterations) transfers")
268310
// Short circuit and return 0 if index does not match iterations
269311
if index != iterations {
@@ -281,6 +323,7 @@ final class QUICTransfer {
281323
var iterations = 10000 // 5gb total (if 500000 sendSize)
282324
var loggingHandler: LoggingHandle = LoggingHandle(loggingType: .none)
283325
var sendSize = 500000 // 500kb
326+
var linkDelay = NetworkDuration.zero
284327
var arguments = CommandLine.arguments.dropFirst(0)
285328
if arguments.contains("-iterations"),
286329
let index = arguments.firstIndex(of: "-iterations")
@@ -310,6 +353,15 @@ if arguments.contains("-size"),
310353
}
311354
}
312355
}
356+
if arguments.contains("-link-delay-ms"),
357+
let index = arguments.firstIndex(of: "-link-delay-ms")
358+
{
359+
if arguments.count >= (index + 2) {
360+
if let linkDelayOption = Int(arguments[index + 1]) {
361+
linkDelay = .milliseconds(linkDelayOption)
362+
}
363+
}
364+
}
313365

314366
// Create and run the transfers
315367
let quicTransfer = QUICTransfer()
@@ -319,7 +371,8 @@ let totalTime = quicTransfer.run(
319371
iterations: iterations,
320372
loggingHandle: loggingHandler,
321373
group: group,
322-
sendSize: sendSize
374+
sendSize: sendSize,
375+
linkDelay: linkDelay
323376
)
324377
if totalTime > 0 {
325378
print("Finished all (\(iterations)) transfers in \(totalTime) seconds")

0 commit comments

Comments
 (0)