Skip to content

Commit 6080fef

Browse files
authored
Speed up http client (#633)
The http client currently uses an inefficient byte-by-byte copy in its `BoundedStreamReader` - the real solution here is to get rid of the copy completely but until then, we can use bulk-copy the data at least. Ditto `read`, `readN` and similar helpers - this ~doubles throughput for bulk reading. Pre: ``` | Small/small | 0.070s | 1000 | 14226.009 | 0.031 MB | 0.016 MB | 0.434 MB/s | 0.231 MB/s | | Medium/small | 1.834s | 1000 | 545.267 | 1000.000 MB | 0.021 MB | 545.267 MB/s | 0.011 MB/s | | Small/Medium | 2.238s | 1000 | 446.735 | 0.031 MB | 1000.000 MB | 0.014 MB/s | 446.735 MB/s | | Medium/Medium | 2.583s | 1000 | 387.196 | 1000.000 MB | 1000.000 MB | 387.196 MB/s | 387.196 MB/s | ``` Post: ``` | Small/small | 0.066s | 1000 | 15038.890 | 0.031 MB | 0.016 MB | 0.459 MB/s | 0.244 MB/s | | Medium/small | 0.954s | 1000 | 1048.475 | 1000.000 MB | 0.021 MB | 1048.475 MB/s | 0.022 MB/s | | Small/Medium | 1.264s | 1000 | 791.318 | 0.031 MB | 1000.000 MB | 0.024 MB/s | 791.318 MB/s | | Medium/Medium | 1.615s | 1000 | 619.083 | 1000.000 MB | 1000.000 MB | 619.083 MB/s | 619.083 MB/s | ```
1 parent 539767c commit 6080fef

7 files changed

Lines changed: 249 additions & 78 deletions

File tree

benchmarks/bench_http_fetch.nim

Lines changed: 179 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,179 @@
1+
# Chronos Test Suite
2+
# (c) Copyright 2021-Present
3+
# Status Research & Development GmbH
4+
#
5+
# Licensed under either of
6+
# Apache License, version 2.0, (LICENSE-APACHEv2)
7+
# MIT license (LICENSE-MIT)
8+
import
9+
std/[strformat, times],
10+
11+
../chronos,
12+
../chronos/threadsync,
13+
../chronos/apps/http/[httpserver, httpclient, httpcommon]
14+
15+
{.used.}
16+
17+
# Benchmark configuration
18+
const
19+
SmallRequestSize = 32
20+
MediumRequestSize = 1024 * 1024 # 1MB
21+
MediumResponseSize = 1024 * 1024 # 1MB
22+
NumClients = 10 # Number of concurrent clients
23+
NumRequestsPerClient = 100 # Number of requests per client
24+
25+
type BenchmarkResult = object
26+
testName: string
27+
totalTime: times.Duration
28+
totalBytesSent: int64
29+
totalBytesReceived: int64
30+
numRequests: int
31+
32+
# Create a message of specified size
33+
proc createMessage(size: int): seq[byte] =
34+
let message = "Hello, World! This is a benchmark message. "
35+
result = newSeq[byte](size)
36+
for i in 0 ..< size:
37+
result[i] = byte(message[i mod len(message)])
38+
39+
40+
const
41+
smallRequest = createMessage(SmallRequestSize)
42+
mediumMessage = createMessage(MediumRequestSize)
43+
44+
proc inSecondsFloat(d: times.Duration): float =
45+
d.inNanoseconds() / 1000000000
46+
47+
# Create a simple HTTP server for benchmarking
48+
proc createBenchmarkServer(address: TransportAddress): HttpServerRef =
49+
proc process(
50+
r: RequestFence
51+
): Future[HttpResponseRef] {.async: (raises: [CancelledError]).} =
52+
if r.isOk():
53+
let request = r.get()
54+
try:
55+
case request.uri.path
56+
of "/small":
57+
# Small request, small response
58+
let data = await request.getBody()
59+
await request.respond(Http200, "Received " & $len(data) & " bytes")
60+
of "/medium":
61+
# Small request, medium response
62+
discard await request.getBody()
63+
await request.respond(Http200, mediumMessage)
64+
else:
65+
await request.respond(Http404, "Not found")
66+
except HttpError as exc:
67+
defaultResponse(exc)
68+
else:
69+
defaultResponse()
70+
71+
let socketFlags = {ServerFlags.TcpNoDelay, ServerFlags.ReuseAddr}
72+
let res = HttpServerRef.new(address, process, socketFlags = socketFlags)
73+
res.get()
74+
75+
var
76+
serverAddress: TransportAddress
77+
tsp: ThreadSignalPtr
78+
79+
proc runServer() {.thread, nimcall.} =
80+
let server = createBenchmarkServer(initTAddress("127.0.0.1:0"))
81+
server.start()
82+
serverAddress = server.instance.localAddress()
83+
discard tsp.fireSync().expect("ok")
84+
runForever()
85+
86+
# Create an HTTP client session
87+
proc createBenchmarkClient(): HttpSessionRef =
88+
HttpSessionRef.new({HttpClientFlag.Http11Pipeline})
89+
90+
# Client worker that sends requests and measures performance
91+
proc clientWorker(
92+
session: HttpSessionRef,
93+
address: TransportAddress,
94+
testPath: string,
95+
requestData: seq[byte],
96+
results: ref BenchmarkResult,
97+
) {.async.} =
98+
let ha = getAddress(address, HttpClientScheme.NonSecure, testPath)
99+
100+
for i in 0 ..< NumRequestsPerClient:
101+
var req = HttpClientRequestRef.new(session, ha, MethodPost, body = requestData)
102+
103+
let response = await fetch(req)
104+
assert response.status == 200, "No failing requests in this benchmark"
105+
106+
inc(results.numRequests)
107+
results.totalBytesReceived += int64(len(response.data))
108+
109+
await req.closeWait()
110+
111+
# Run a benchmark test
112+
proc runBenchmark(
113+
testName: string, testPath: string, requestData: seq[byte]
114+
): Future[BenchmarkResult] {.async.} =
115+
var session = createBenchmarkClient()
116+
var futures = newSeq[Future[void]](NumClients)
117+
var results = (ref BenchmarkResult)(
118+
testName: testName,
119+
totalBytesSent: int64(len(requestData)) * int64(NumClients * NumRequestsPerClient),
120+
numRequests: 0,
121+
)
122+
123+
let startTime = getTime()
124+
125+
# Start client workers
126+
for i in 0 ..< NumClients:
127+
futures[i] = clientWorker(session, serverAddress, testPath, requestData, results)
128+
129+
# Wait for all workers to complete
130+
await allFutures(futures)
131+
132+
await session.closeWait()
133+
let endTime = getTime()
134+
135+
results.totalTime = endTime - startTime
136+
137+
results[]
138+
139+
# Print benchmark results as a table
140+
proc print(results: BenchmarkResult) =
141+
let v = (
142+
name: results.testName,
143+
totalTime: results.totalTime.inSecondsFloat,
144+
numRequests: results.numRequests,
145+
reqsps: ((results.numRequests.float / results.totalTime.inSecondsFloat)),
146+
sent: results.totalBytesSent / (1024 * 1024),
147+
received: results.totalBytesReceived / (1024 * 1024),
148+
sendSpeed:
149+
(results.totalBytesSent.float / results.totalTime.inSecondsFloat / (1024 * 1024)),
150+
recvSpeed: (
151+
results.totalBytesReceived.float / results.totalTime.inSecondsFloat / (
152+
1024 * 1024
153+
)
154+
),
155+
)
156+
157+
echo &"| {v.name:14} | {v.totalTime:8.3f}s | {v.numRequests:6} | {v.reqsps:8.3f} | {v.sent:8.3f} MB | {v.received:8.3f} MB | {v.sendSpeed:8.3f} MB/s | {v.recvSpeed:8.3f} MB/s |"
158+
159+
# Main benchmark function
160+
proc runBenchmarks() {.async.} =
161+
echo " Clients: " & $NumClients
162+
echo " Requests per client: " & $NumRequestsPerClient
163+
echo " Total requests: " & $(NumClients * NumRequestsPerClient)
164+
echo ""
165+
echo "| Benchmark | Time (s) | Reqs | Req/s | Bytes Sent | Bytes Recv | Send Speed | Recv Speed |"
166+
echo "|----------------|-----------|--------|----------|-------------|-------------|---------------|---------------|"
167+
168+
print(await runBenchmark("Small/small", "/small", smallRequest))
169+
print(await runBenchmark("Medium/small", "/small", mediumMessage))
170+
print(await runBenchmark("Small/Medium", "/medium", smallRequest))
171+
print(await runBenchmark("Medium/Medium", "/medium", mediumMessage))
172+
173+
when isMainModule:
174+
tsp = ThreadSignalPtr.new()[]
175+
var server: Thread[void]
176+
createThread(server, runServer)
177+
discard tsp.waitSync().expect("ok")
178+
179+
waitFor runBenchmarks()

chronos.nimble

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ skipDirs = @["tests"]
99

1010
requires "nim >= 1.6.16",
1111
"results",
12-
"stew",
12+
"stew >= 0.5.0",
1313
"bearssl >= 0.2.7",
1414
"httputils",
1515
"unittest2"

chronos/bipbuffer.nim

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@
2525

2626
{.push raises: [].}
2727

28+
import stew/[arrayops, ptrops]
29+
2830
type
2931
BipPos = object
3032
start: Natural
@@ -138,3 +140,12 @@ iterator regions*(bp: var BipBuffer): tuple[data: ptr byte, size: Natural] =
138140
yield (addr bp.data[bp.a.start], len(bp.a))
139141
if len(bp.b) > 0:
140142
yield (addr bp.data[bp.b.start], len(bp.b))
143+
144+
proc copyInto*(bp: var BipBuffer, tgt: var openArray[byte]): Natural =
145+
## Copy `min(tgt.len, bp.len)` bytes into `tgt`, consuming them in the process
146+
##
147+
## Returns the number of copied bytes
148+
var n = 0
149+
for (region, rsize) in bp.regions():
150+
n += tgt.toOpenArray(n, tgt.high()).copyFrom(region.makeOpenArray(rsize))
151+
n

chronos/streams/asyncstream.nim

Lines changed: 21 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
{.push raises: [].}
1111

1212
import std/strutils
13+
import stew/[ptrops, shims/sequninit]
1314
import ../[config, asyncloop, asyncsync, bipbuffer]
1415
import ../transports/[common, stream]
1516
export asyncloop, asyncsync, stream, common
@@ -340,22 +341,17 @@ proc readExactly*(rstream: AsyncStreamReader, pbytes: pointer,
340341
await readExactly(rstream.rsource, pbytes, nbytes)
341342
else:
342343
var
343-
index = 0
344-
pbuffer = pbytes.toUnchecked()
344+
total = 0
345+
pbuffer = cast[ptr byte](pbytes)
345346
readLoop():
346347
if len(rstream.buffer.backend) == 0:
347348
if rstream.atEof():
348349
raise newAsyncStreamIncompleteError()
349-
var bytesRead = 0
350-
for (region, rsize) in rstream.buffer.backend.regions():
351-
let count = min(nbytes - index, rsize)
352-
bytesRead += count
353-
if count > 0:
354-
copyMem(addr pbuffer[index], region, count)
355-
index += count
356-
if index == nbytes:
357-
break
358-
(consumed: bytesRead, done: index == nbytes)
350+
let consumed =
351+
rstream.buffer.backend.copyInto(pbuffer.makeOpenArray(nbytes - total))
352+
pbuffer = pbuffer.offset(consumed)
353+
total += consumed
354+
(consumed: consumed, done: total == nbytes)
359355

360356
proc readOnce*(rstream: AsyncStreamReader, pbytes: pointer,
361357
nbytes: int): Future[int] {.
@@ -378,20 +374,15 @@ proc readOnce*(rstream: AsyncStreamReader, pbytes: pointer,
378374
return await readOnce(rstream.rsource, pbytes, nbytes)
379375
else:
380376
var
381-
pbuffer = pbytes.toUnchecked()
382-
index = 0
377+
total = 0
378+
pbuffer = cast[ptr byte](pbytes)
383379
readLoop():
384380
if len(rstream.buffer.backend) == 0:
385381
(0, rstream.atEof())
386382
else:
387-
for (region, rsize) in rstream.buffer.backend.regions():
388-
let size = min(rsize, nbytes - index)
389-
copyMem(addr pbuffer[index], region, size)
390-
index += size
391-
if index >= nbytes:
392-
break
393-
(index, true)
394-
index
383+
total = rstream.buffer.backend.copyInto(pbuffer.makeOpenArray(nbytes))
384+
(total, true)
385+
total
395386

396387
proc readUntil*(rstream: AsyncStreamReader, pbytes: pointer, nbytes: int,
397388
sep: seq[byte]): Future[int] {.
@@ -530,10 +521,10 @@ proc read*(rstream: AsyncStreamReader): Future[seq[byte]] {.
530521
if rstream.atEof():
531522
(0, true)
532523
else:
533-
var bytesRead = 0
534-
for (region, rsize) in rstream.buffer.backend.regions():
535-
bytesRead += rsize
536-
res.add(region.toUnchecked().toOpenArray(0, rsize - 1))
524+
var pos = res.len
525+
res.setLenUninit(pos + rstream.buffer.backend.len())
526+
let bytesRead =
527+
rstream.buffer.backend.copyInto(res.toOpenArray(pos, res.high()))
537528
(bytesRead, false)
538529
res
539530

@@ -562,11 +553,10 @@ proc read*(rstream: AsyncStreamReader, n: int): Future[seq[byte]] {.
562553
if rstream.atEof():
563554
(0, true)
564555
else:
565-
var bytesRead = 0
566-
for (region, rsize) in rstream.buffer.backend.regions():
567-
let count = min(rsize, n - len(res))
568-
bytesRead += count
569-
res.add(region.toUnchecked().toOpenArray(0, count - 1))
556+
var pos = res.len
557+
res.setLenUninit(pos + min(rstream.buffer.backend.len(), n - res.len))
558+
let bytesRead =
559+
rstream.buffer.backend.copyInto(res.toOpenArray(pos, res.high()))
570560
(bytesRead, len(res) == n)
571561
res
572562

chronos/streams/boundstream.nim

Lines changed: 18 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
{.push raises: [].}
1919

2020
import results
21+
import stew/[arrayops, ptrops]
2122
import ../[asyncloop, timer, bipbuffer, config]
2223
import asyncstream, ../transports/[stream, common]
2324
export asyncloop, asyncstream, stream, timer, common
@@ -68,7 +69,7 @@ proc readUntilBoundary(rstream: AsyncStreamReader, pbytes: pointer,
6869
var state = 0
6970
var pbuffer = cast[ptr UncheckedArray[byte]](pbytes)
7071

71-
proc predicate(data: openArray[byte]): tuple[consumed: int, done: bool] =
72+
proc predicateSep(data: openArray[byte]): tuple[consumed: int, done: bool] =
7273
if len(data) == 0:
7374
(0, true)
7475
else:
@@ -80,16 +81,24 @@ proc readUntilBoundary(rstream: AsyncStreamReader, pbytes: pointer,
8081
inc(index)
8182
pbuffer[k] = ch
8283
inc(k)
83-
if len(sep) > 0:
84-
if sep[state] == ch:
85-
inc(state)
86-
if state == len(sep):
87-
break
88-
else:
89-
state = 0
84+
if sep[state] == ch:
85+
inc(state)
86+
if state == len(sep):
87+
break
88+
else:
89+
state = 0
90+
(index, (state == len(sep)) or (k == nbytes))
91+
92+
proc predicate(data: openArray[byte]): tuple[consumed: int, done: bool] =
93+
if len(data) == 0:
94+
(0, true)
95+
else:
96+
let index = pbytes.offset(k).makeOpenArray(byte, nbytes - k).copyFrom(data)
97+
k += index
9098
(index, (state == len(sep)) or (k == nbytes))
9199

92-
await rstream.readMessage(predicate)
100+
await rstream.readMessage(if sep.len == 0: predicate else: predicateSep)
101+
93102
return k
94103

95104
func endsWith(s, suffix: openArray[byte]): bool =

0 commit comments

Comments
 (0)