Skip to content

Commit 330454d

Browse files
committed
feat: perf tests and fix: pin connections and streams
1 parent 875334a commit 330454d

7 files changed

Lines changed: 143 additions & 15 deletions

File tree

lsquic.nimble

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ task format, "Format nim code using nph":
2121

2222
task test, "Run tests":
2323
when defined(windows):
24-
exec "nim c --mm:refc -d:nimDebugDlOpen -r --threads:on tests/test_connection.nim"
24+
exec "nim c --mm:refc -d:nimDebugDlOpen --threads:on tests/test_connection.nim"
2525
else:
26-
exec "nim c --mm:refc -r --threads:on tests/test_connection.nim"
26+
exec "nim c --mm:refc --threads:on tests/test_connection.nim"
27+
exec "./tests/test_connection --output-level=VERBOSE"

lsquic/connection.nim

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -91,6 +91,9 @@ method incomingStream*(
9191
method incomingStream*(
9292
connection: IncomingConnection
9393
): Future[Stream] {.async: (raises: [CancelledError, ConnectionError]).} =
94+
if connection.isClosed:
95+
raise newException(ConnectionError, "connection is closed")
96+
9497
let closedFut = connection.closed.wait()
9598
let incomingFut = connection.quicConn.incomingStream()
9699
let raceFut = await race(closedFut, incomingFut)
@@ -108,6 +111,8 @@ method openStream*(
108111
method openStream*(
109112
connection: OutgoingConnection
110113
): Future[Stream] {.async: (raises: [CancelledError, ConnectionError]).} =
114+
if connection.isClosed:
115+
raise newException(ConnectionError, "connection is closed")
111116
let s = Stream.new()
112117
let created = connection.quicConn.addPendingStream(s)
113118
connection.quicContext.makeStream(connection.quicConn)

lsquic/context/client.nim

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ proc onConnClosed(conn: ptr lsquic_conn_t) {.cdecl.} =
5252
)
5353
quicClientConn.cancelPending()
5454
quicClientConn.onClose()
55+
GC_unref(quicClientConn)
5556
lsquic_conn_set_ctx(conn, nil)
5657

5758
proc onNewStream(
@@ -90,7 +91,7 @@ method dial*(
9091

9192
let quicClientConn =
9293
QuicClientConn(connectedFut: connectedFut, local: local, remote: remote)
93-
94+
GC_ref(quicClientConn) # Keep it pinned until on_conn_closed is called
9495
let conn = lsquic_engine_connect(
9596
ctx.engine,
9697
N_LSQVER,

lsquic/context/server.nim

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ proc onNewConn(
1919
remote: remote.toTransportAddress(),
2020
lsquicConn: conn,
2121
)
22+
GC_ref(quicServerConn) # Keep it pinned until on_conn_closed is called
2223
let serverCtx = cast[ServerContext](stream_if_ctx)
2324
serverCtx.incoming.putNoWait(quicServerConn)
2425
cast[ptr lsquic_conn_ctx_t](quicServerConn)

lsquic/context/stream.nim

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,10 @@ proc onClose*(stream: ptr lsquic_stream_t, ctx: ptr lsquic_stream_ctx_t) {.cdecl
1212
let streamCtx = cast[Stream](ctx)
1313
if not streamCtx.closeWrite:
1414
streamCtx.isEof = true
15+
echo "FIRING CLOSE 3"
1516
streamCtx.closed.fire()
16-
streamCtx.abortPendingWrites("stream closed")
17+
streamCtx.abortPendingWrites("stream closed 4")
18+
GC_unref(streamCtx)
1719

1820
type StreamReadContext = object
1921
stream: ptr lsquic_stream_t
@@ -40,12 +42,11 @@ proc onRead*(stream: ptr lsquic_stream_t, ctx: ptr lsquic_stream_ctx_t) {.cdecl.
4042
if nread < 0:
4143
error "could not read from stream", nread, streamId = lsquic_stream_id(stream)
4244
streamCtx.abort()
43-
45+
4446
if lsquic_stream_wantread(stream, 0) == -1:
4547
error "could not set stream wantread", streamId = lsquic_stream_id(stream)
4648
streamCtx.abort()
4749

48-
4950
proc onWrite*(stream: ptr lsquic_stream_t, ctx: ptr lsquic_stream_ctx_t) {.cdecl.} =
5051
trace "onWrite"
5152

lsquic/stream.nim

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -18,11 +18,13 @@ type Stream* = ref object
1818
toWrite*: seq[WriteTask]
1919

2020
proc new*(T: typedesc[Stream], quicStream: ptr lsquic_stream_t = nil): T =
21-
Stream(
21+
let s = Stream(
2222
quicStream: quicStream,
2323
incoming: newAsyncQueue[seq[byte]](),
2424
closed: newAsyncEvent(),
2525
)
26+
GC_ref(s) # Keep it pinned until stream_if.on_close is executed
27+
s
2628

2729
proc abortPendingWrites*(stream: Stream, reason: string = "") =
2830
for pendingWrite in stream.toWrite.mitems:
@@ -33,6 +35,7 @@ proc abortPendingWrites*(stream: Stream, reason: string = "") =
3335
proc abort*(stream: Stream) =
3436
if stream.closeWrite and stream.isEof:
3537
if not stream.closed.isSet():
38+
echo "FIRING CLOSE 1"
3639
stream.closed.fire()
3740
stream.abortPendingWrites("stream aborted")
3841
return
@@ -44,16 +47,19 @@ proc abort*(stream: Stream) =
4447
stream.closeWrite = true
4548
stream.isEof = true
4649
stream.abortPendingWrites("stream aborted")
50+
echo "FIRING CLOSE 2"
4751
stream.closed.fire()
4852

4953
proc close*(stream: Stream) =
5054
if stream.closeWrite:
5155
return
5256

5357
# Closing only the write side
58+
echo "SHUTDOWN WRITE"
5459
let ret = lsquic_stream_shutdown(stream.quicStream, 1)
5560
if ret == 0:
5661
if stream.isEof:
62+
echo "CLOSE DUE TO EOF"
5763
if lsquic_stream_close(stream.quicStream) != 0:
5864
stream.abort()
5965
raise newException(StreamError, "could not close the stream")
@@ -78,7 +84,7 @@ proc read*(
7884
await incomingFut.cancelAndWait()
7985
stream.isEof = true
8086
stream.closeWrite = true
81-
raise newException(StreamError, "stream closed")
87+
raise newException(StreamError, "stream closed 1")
8288

8389
let incoming = await incomingFut
8490
if incoming.len == 0:
@@ -95,15 +101,16 @@ proc write*(
95101
stream: Stream, data: seq[byte]
96102
) {.async: (raises: [CancelledError, StreamError]).} =
97103
if stream.closeWrite:
98-
raise newException(StreamError, "stream is closed")
104+
raise newException(StreamError, "stream closed 3")
99105

100106
let closedFut = stream.closed.wait()
101107
let doneFut = Future[void].Raising([CancelledError, StreamError]).init()
102108
stream.toWrite.add(WriteTask(data: data, doneFut: doneFut))
103109
discard lsquic_stream_wantwrite(stream.quicStream, 1)
104110
let raceFut = await race(closedFut, doneFut)
105111
if raceFut == closedFut:
106-
doneFut.fail(newException(StreamError, "stream closed"))
112+
if not doneFut.finished:
113+
doneFut.fail(newException(StreamError, "stream closed 2"))
107114
stream.closeWrite = true
108115

109116
await doneFut

tests/test_connection.nim

Lines changed: 117 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,8 @@ import ./helpers/certificate
1111
import lsquic/certificateverifier
1212
import lsquic/stream
1313
import lsquic/lsquic_ffi
14+
import stew/endians2
15+
import sequtils
1416

1517
proc logging(ctx: pointer, buf: cstring, len: csize_t): cint {.cdecl.} =
1618
echo $buf
@@ -21,10 +23,118 @@ proc certificateCb(
2123
): bool {.gcsafe.} =
2224
return derCertificates.len > 0
2325

24-
suite "connections":
25-
setup:
26-
let address = initTAddress("127.0.0.1:12345")
26+
let address = initTAddress("127.0.0.1:12345")
27+
28+
const
29+
runs = 1
30+
uploadSize = 100000 # 100KB
31+
downloadSize = 100000000 # 100MB
32+
chunkSize = 65536 # 64KB chunks like perf
33+
34+
proc runPerf(): Future[Duration] {.async.} =
35+
let customCertVerif: CertificateVerifier =
36+
CustomCertificateVerifier.init(certificateCb)
37+
let clientTLSConfig = TLSConfig.new(
38+
testCertificate(),
39+
testPrivateKey(),
40+
@["test"].toHashSet(),
41+
Opt.some(customCertVerif),
42+
)
43+
let serverTLSConfig = TLSConfig.new(
44+
testCertificate(),
45+
testPrivateKey(),
46+
@["test"].toHashSet(),
47+
Opt.some(customCertVerif),
48+
)
49+
let client = QuicClient.new(clientTLSConfig)
50+
let server = QuicServer.new(serverTLSConfig)
51+
let listener = server.listen(address)
52+
let accepting = listener.accept()
53+
let dialing = client.dial(address)
54+
55+
let outgoingConn = await dialing
56+
let incomingConn = await accepting
57+
58+
let serverHandler = proc() {.async.} =
59+
let stream = await incomingConn.incomingStream()
60+
61+
# Step 1: Read download size (8 bytes)
62+
let clientDownloadSize = await stream.read()
63+
64+
# Step 2: Read upload data until EOF
65+
var totalBytesRead = 0
66+
while true:
67+
let chunk = await stream.read()
68+
if chunk.len == 0:
69+
break
70+
totalBytesRead += chunk.len
71+
72+
# Step 3: Send download data back
73+
var remainingToSend = uint64.fromBytesBE(clientDownloadSize)
74+
while remainingToSend > 0:
75+
let toSend = min(remainingToSend, chunkSize)
76+
try:
77+
await stream.write(newSeq[byte](toSend))
78+
except StreamError:
79+
echo "STREAM ERROR 1"
80+
quit(1)
81+
remainingToSend -= toSend
82+
83+
stream.close()
84+
85+
# Start server handler
86+
asyncSpawn serverHandler()
87+
88+
let startTime = Moment.now()
89+
90+
# Step 1: Send download size, activate stream first
91+
let clientStream = await outgoingConn.openStream()
92+
try:
93+
await clientStream.write(toSeq(downloadSize.uint64.toBytesBE()))
94+
except StreamError:
95+
echo "STREAM ERROR 2"
96+
quit(1)
97+
# Step 2: Send upload data in chunks
98+
var remainingToSend = uploadSize
99+
while remainingToSend > 0:
100+
let toSend = min(remainingToSend, chunkSize)
101+
try:
102+
await clientStream.write(newSeq[byte](toSend))
103+
except StreamError:
104+
echo "STRAM ERROR 4"
105+
quit(1)
106+
remainingToSend -= toSend
107+
108+
# Step 3: Close write side
109+
clientStream.close()
110+
111+
# Step 4: Start reading download data
112+
var totalDownloaded = 0
113+
while totalDownloaded < downloadSize:
114+
let chunk = await clientStream.read()
115+
totalDownloaded += chunk.len
116+
117+
let duration = Moment.now() - startTime
118+
119+
# TODO: use waitgroup
120+
await sleepAsync(5.seconds)
121+
122+
await listener.stop()
123+
await client.stop()
124+
125+
return duration
126+
127+
suite "perf protocol simulation":
128+
asyncTest "test":
129+
var total: Duration
130+
for i in 0 ..< runs:
131+
let duration = await runPerf()
132+
total += duration
133+
echo "\trun #" & $(i + 1) & " duration: " & $duration
27134

135+
echo "\tavrg duration: " & $(total div runs)
136+
137+
suite "tests":
28138
asyncTest "test":
29139
let logger = struct_lsquic_logger_if(log_buf: logging)
30140
discard lsquic_set_log_level("debug")
@@ -76,7 +186,6 @@ suite "connections":
76186

77187
let incomingBehaviour = proc() {.async.} =
78188
try:
79-
echo "HERE"
80189
let stream = await incomingConn.incomingStream()
81190
echo "Received stream in server"
82191

@@ -107,12 +216,15 @@ suite "connections":
107216
outgoingConn.close()
108217
incomingConn.close()
109218

219+
# Cannot create a stream once closed
220+
expect ConnectionError:
221+
discard await outgoingConn.openStream()
222+
110223
await sleepAsync(1.seconds)
111224

112225
await client.stop()
113226
await listener.stop()
114227

115-
# TODO: perf example
116228
# TODO: destructors: (nice to have:)
117229
# - lsquic_global_cleanup() to free global resources.
118230
# - lsquic_engine_destroy(engine)

0 commit comments

Comments
 (0)