diff --git a/.gitignore b/.gitignore index add4ead7..9bc6b9b8 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,7 @@ /target +# Out-of-workspace bench crate has its own target dir and lockfile. +/benches/rpc-grpc-rust/target +/benches/rpc-grpc-rust/Cargo.lock /conformance/bin/ # User-local Claude Code settings (project-level agents/commands ARE committed) diff --git a/Cargo.toml b/Cargo.toml index badf9c6d..9d563ed2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,8 @@ [workspace] members = ["connectrpc", "connectrpc-codegen", "connectrpc-build", "connectrpc-health", "connectrpc-reflection", "conformance", "examples/eliza", "examples/middleware", "examples/mtls-identity", "examples/multiservice", "examples/streaming-tour", "examples/wasm-client", "tests/streaming", "benches/rpc", "benches/rpc-tonic"] +# benches/rpc-grpc-rust needs the grpc-rust codegen toolchain (protoc 35.1 + +# a cmake-built C++ plugin); it is built on demand by the bench drivers, not CI. +exclude = ["benches/rpc-grpc-rust"] resolver = "2" [workspace.package] diff --git a/README.md b/README.md index 18afad9e..5dc4df6d 100644 --- a/README.md +++ b/README.md @@ -475,28 +475,67 @@ for `include_bytes!`) or by an existing `buffa_descriptor::DescriptorPool`. ## Performance -Comparison against [tonic](https://docs.rs/tonic/) 0.14 (the standard Rust gRPC -implementation, built on the same hyper/h2 stack). Measured on Intel Xeon -Platinum 8488C with [buffa](https://github.com/anthropics/buffa) as the proto -library. Higher is better unless noted. +Comparison against [tonic](https://docs.rs/tonic/) 0.14.6, the standard Rust +gRPC implementation built on the same hyper/h2 stack, in two configurations: +`tonic` is tonic with prost, and `tonic-protobuf` is tonic with +[grpc-rust](https://github.com/grpc/grpc-rust)'s codec over Google's +`protobuf` v4 runtime on the upb kernel. The `tonic` arm builds tonic from +crates.io and the `tonic-protobuf` arm from grpc-rust revision `7053afcd`, so +the two also differ by the handful of unreleased tonic commits at that +revision. connectrpc-rs uses [buffa](https://github.com/anthropics/buffa). +Unless a subsection says otherwise, the numbers were measured 2026-09 on a +bare-metal AWS c7i.metal-24xl (Intel Xeon Platinum 8488C, turbo disabled), +with client and servers on the same host over loopback. Higher is better +unless noted. + +grpc-rust's own `grpc` crate is a client channel with no server, so it does +not appear in the server tables; the [Client stacks](#client-stacks) table +compares it against the connectrpc-rs and tonic clients. The `tonic-protobuf` arm's first build compiles `protoc` and a +protoc plugin from C++ source, so it needs cmake and a C++17 compiler; see +[`benches/rpc-grpc-rust/README.md`](benches/rpc-grpc-rust/README.md). + +The short version: on small unary calls, echo throughput, and 10-message +client and server streams, the two tonic configurations are within 3% of +connectrpc-rs. The differences on log batches and large payloads come from the +proto library. At concurrency 1, a 50-record log-batch request takes 39% longer +with prost than with buffa's zero-copy views, and 16% longer with upb. Under +load, tonic serves 6–13% fewer log-batch requests per second than +connectrpc-rs (13% fewer at c=256), and tonic-protobuf is within 4% (−4% to ++3%), although its c=256 cell measured 23% lower on a second pass. The upb arm +also takes 11% longer on the 1 MB gzip'd payload. ### Single-request latency Criterion benchmarks at concurrency=1 (no h2 contention), measuring per-request -framework + proto work in isolation. Lower is better. +framework + proto work in isolation. Every arm is driven by the same +connectrpc-rs client, so the columns compare servers. Lower is better. ![Single-request latency](benches/charts/latency.svg)
Raw data (μs, lower is better) -| Benchmark | connectrpc-rs | tonic | ratio | +| Benchmark | connectrpc-rs | tonic | tonic-protobuf | |---|---:|---:|---:| -| unary_small (1 int32 + nested msg) | 87.6 | 170.8 | **1.95×** | -| unary_logs_50 (50 log records, ~15 KB) | 195.0 | 338.5 | **1.74×** | -| client_stream (10 messages) | 166.1 | 223.8 | **1.35×** | -| server_stream (10 messages) | 109.8 | 110.1 | 1.00× | - -Run with `task bench:cross:quick`. +| unary_small (1 int32 + nested msg) | 79.7 | 78.8 (−1%) | 78.4 (−2%) | +| unary_logs_50 (50 log records, ~22 KB) | 218.5 | 303.1 (+39%) | 253.5 (+16%) | +| unary_large (~1 MB payload, gzip request) | 4,486 | 4,521 (+1%) | 4,964 (+11%) | +| client_stream (10 messages) | 166.3 | 168.2 (+1%) | 162.2 (−2%) | +| server_stream (10 messages) | 108.0 | 107.0 (−1%) | 110.2 (+2%) | + +The same bench also runs [connect-go](https://github.com/connectrpc/connect-go) +over gRPC: 253 μs unary_small, 529 μs unary_logs_50, 406 μs client_stream, +1,080 μs server_stream. Over the Connect protocol, unary_small is 80.5 μs on +connectrpc-rs and 147 μs on connect-go. + +Compare the unary_large columns within a run, not across runs: the row +depends on what the bench process ran before it. Identical binaries on the +same instance type and kernel measured the connectrpc-rs arm at 4,486 μs in +the full suite and at 3,298 μs with the run filtered to unary_large, while +the tonic arm moved by under 1%. + +Run with `task bench:cross`. It builds the connect-go server with `go build`, +so it needs a Go toolchain unless `RPC_BENCH_BIN_DIR` points at prebuilt +server binaries (see [`benches/rpc/README.md`](benches/rpc/README.md)).
@@ -505,19 +544,56 @@ Run with `task bench:cross:quick`. 64-byte string echo, 8 h2 connections (to avoid single-connection mutex contention — see [h2 #531](https://github.com/hyperium/h2/issues/531)). Measures framework dispatch + envelope framing + proto encode/decode with -minimal handler work. +minimal handler work; the three stacks are within 1% of each other up to +c=64, and connectrpc-rs leads by 3% at c=256. ![Echo throughput](benches/charts/echo.svg)
Raw data (req/s) -| Concurrency | connectrpc-rs | tonic | -|---|---:|---:| -| c=16 | 170,292 | 168,811 (−1%) | -| c=64 | 238,498 | 234,304 (−2%) | -| c=256 | 252,000 | 247,167 (−2%) | +| Concurrency | connectrpc-rs | tonic | tonic-protobuf | +|---|---:|---:|---:| +| c=16 | 194,853 | 196,512 (+1%) | 195,235 | +| c=64 | 298,919 | 301,501 (+1%) | 298,403 | +| c=256 | 271,985 | 264,772 (−3%) | 262,920 (−3%) | + +A second pass in the same session reproduced every cell within 1.1%. + +Run with `task bench:echo -- --multi-conn=8`; the table shows the +`(8-conn)` rows of its output. + +
+ +### Client stacks -Run with `task bench:echo -- --multi-conn=8`. +The other tables in this section hold the client fixed (connectrpc-rs) and vary +the server; this one holds the server fixed (the connectrpc-rs echo server) and +varies the client: the generated connectrpc-rs client over `HttpClient` +(hyper-util's pooled client, one per connection) and over +`SharedHttp2Connection` (one raw h2 connection each, no pool), tonic's +generated client (tonic-prost, built from the same grpc-rust revision), and +grpc-rust's `grpc` channel with its `protobuf` codec. All four speak gRPC over +h2 to the same server; closed loop, 64-byte echo, requests in flight spread +round-robin over the connections, each cell the median-throughput run of three +10-second runs. + +
Raw data (req/s) + +| Connections | Requests in flight | connectrpc-rs `HttpClient` | connectrpc-rs `SharedHttp2Connection` | tonic | grpc-rust | +|---:|---:|---:|---:|---:|---:| +| 1 | 1 | 17,497 | 17,394 (−1%) | 17,281 (−1%) | 15,093 (−14%) | +| 1 | 16 | 37,111 | 36,923 (−1%) | 36,044 (−3%) | 36,780 (−1%) | +| 1 | 64 | 42,004 | 41,360 (−2%) | 36,414 (−13%) | 36,780 (−12%) | +| 8 | 16 | 195,213 | 197,060 (+1%) | 194,472 | 184,641 (−5%) | +| 8 | 64 | 307,749 | 304,224 (−1%) | 302,701 (−2%) | 294,227 (−4%) | + +At one request at a time, tonic and both connectrpc-rs clients take 56–58 μs +per call (p50), and grpc-rust's channel takes 66 μs. With 64 requests in flight on a single +connection, the two connectrpc-rs transports keep scaling to 41–42k req/s, +while tonic and grpc-rust level off at 36–37k. With the same 64 requests spread +over 8 connections, all four are within 5%. + +Run with `task bench:clients -- --repeat=3`.
@@ -525,29 +601,36 @@ Run with `task bench:echo -- --multi-conn=8`. 50 structured log records per request (~22 KB batch): varints, string fields, nested message, map entries. Handler iterates every field to force full decode. -This is where the proto library matters — buffa's zero-copy views avoid the -per-string allocations that prost's owned types require. +This is where the proto library matters: buffa's views borrow string data from +the request buffer, prost allocates a `String` per field and a `HashMap` per +map, and upb parses eagerly into a per-message arena, copying string bytes but +avoiding per-field heap allocations — which puts it much closer to buffa than +to prost. ![Log ingest throughput](benches/charts/log-ingest.svg)
Raw data (req/s) -| Concurrency | connectrpc-rs | tonic | -|---|---:|---:| -| c=16 | 32,257 | 28,110 (−13%) | -| c=64 | 73,313 | 68,690 (−6%) | -| c=256 | 112,027 | 84,171 (−25%) | +| Concurrency | connectrpc-rs | tonic | tonic-protobuf | +|---|---:|---:|---:| +| c=16 | 31,237 | 27,891 (−11%) | 30,561 (−2%) | +| c=64 | 76,678 | 71,910 (−6%) | 78,722 (+3%) | +| c=256 | 138,599 | 120,011 (−13%) | 133,558 (−4%) | -At c=256, connectrpc-rs decodes **5.6M records/sec** vs tonic's **4.2M**. +At c=256, connectrpc-rs decodes **6.9M records/sec**, tonic-protobuf 6.7M and +tonic 6.0M. A second pass in the same session reproduced every cell within +1.4% except tonic-protobuf at c=256, which measured 102,828 req/s (23% lower); +the table shows the first pass. **Raw mode (`strict_utf8_mapping`):** For trusted-source log ingestion where UTF-8 validation is unnecessary, buffa can emit `&[u8]` instead of `&str` for string fields (editions `utf8_validation = NONE` + the `strict_utf8_mapping` -codegen option). CPU profile shows this eliminates 11.8% of server CPU -(`str::from_utf8` drops to zero). End-to-end throughput gain in this benchmark -is smaller (~1%) because client encode becomes the bottleneck when both run on -one machine — in production with separate client/server, the server sees ~15% -more capacity. +codegen option). The 2026-03 CPU profile below attributes 11.2% of server CPU +to UTF-8 validation, which raw mode skips. The end-to-end gain in this +benchmark is within run-to-run noise (139.6k vs 138.6k req/s at c=256), +because client encode becomes the bottleneck when both run on one machine. In +production, with the client on another host, the server sees the CPU saving as +capacity. Run with `task bench:log`. @@ -559,7 +642,7 @@ Handler performs a network round-trip to a [valkey](https://valkey.io/) container (`HGETALL` of 12 fortune messages, ~800 bytes), adds an ephemeral record, sorts, and encodes a 13-message response. This is the shape of a typical read-mostly service: RPC framing + async I/O wait + moderate-size -response. All three servers use an 8-connection valkey pool; client uses +response. Every server uses an 8-connection valkey pool; client uses 8 h2 connections so protocol framing is the only variable.
Raw data (req/s, c=256) @@ -580,19 +663,22 @@ response. All three servers use an 8-connection valkey pool; client uses | gRPC | 69,706 | 157,481 | 199,574 | — | | gRPC-Web | 69,067 | 153,727 | 191,811 | — | -Connect's ~20% unary throughput advantage over gRPC at c=256 comes from +Connect's 23% unary throughput advantage over gRPC at c=256 comes from simpler framing: no envelope header, no trailing HEADERS frame. At 200k+ req/s, gRPC's trailer frame is ~200k extra h2 HEADERS encodes per second. The gap grows with throughput (5% @ c=16 → 23% @ c=256). Run with `task bench:fortunes:protocols:h2`. Requires `docker` for the -valkey sibling container (image pulled automatically on first run). +valkey sibling container (image pulled automatically on first run). These +fortunes figures are from the 2026-03 run and predate the `tonic-protobuf` +arm.
-### Where the advantage comes from +### Where the log-ingest difference comes from -CPU profile breakdown (log-ingest, c=64, 30s, `task profile:log`): +CPU profile breakdown (log-ingest, c=64, 30s, `task profile:log`, 2026-03 +run against tonic + prost): | Cost center | connectrpc-rs | tonic | |---|---:|---:| @@ -604,16 +690,20 @@ CPU profile breakdown (log-ingest, c=64, 30s, `task profile:log`): | **Total proto** | **27.1%** | **~24%** (+allocator) | | Allocator (malloc/free/realloc) | **3.6%** | **9.6%** | -connectrpc-rs spends a *larger fraction* of CPU in proto decode — because it -spends so much less everywhere else. buffa's view types borrow string data -directly from the request buffer (zero allocs per string field); `MapView` is -a flat `Vec<(K,V)>` scan with no hashing. tonic/prost must fully materialize -`String` + `HashMap` for every record before the handler runs. - -The framework itself contributes: codegen-emitted `FooServiceServer` with -compile-time `match` dispatch (no `Arc` vtable), a two-frame -`GrpcUnaryBody` for the common unary case, and stream-message batching into -fewer h2 DATA frames. +The difference is allocation: 3.6% of CPU in the allocator against 9.6%, and +nothing in `HashMap` operations against 8.5%. buffa's view types borrow +string data directly from the request buffer (zero allocs per string field); +`MapView` is a flat `Vec<(K,V)>` scan with no hashing. tonic/prost must fully +materialize `String` + `HashMap` for every record before the +handler runs. upb sits between the two: it copies string bytes into a +per-message arena but makes no per-field heap allocation, which is consistent +with it landing within a few percent of buffa in the throughput tables above. + +The framework layer itself — codegen-emitted `FooServiceServer` with +compile-time `match` dispatch, a two-frame `GrpcUnaryBody` for the common unary +case, and stream-message batching into fewer h2 DATA frames — measures level +with tonic's on the echo and small-unary benches above, so the decode path is +where the difference is made. ## Custom Compression diff --git a/Taskfile.yaml b/Taskfile.yaml index 6f5ec216..cf0ab83e 100644 --- a/Taskfile.yaml +++ b/Taskfile.yaml @@ -388,7 +388,7 @@ tasks: - cargo bench -p rpc-bench --bench rpc_bench -- --quick --warm-up-time 1 --measurement-time 3 bench:cross: - desc: Run cross-implementation benchmarks (connectrpc vs tonic vs connect-go) + desc: Run cross-implementation benchmarks (connectrpc vs tonic vs tonic-protobuf vs connect-go) cmds: - cargo bench -p rpc-bench --bench cross_impl_bench @@ -411,7 +411,7 @@ tasks: - cargo bench -p rpc-bench --bench cross_impl_bench -- --quick --warm-up-time 1 --measurement-time 3 bench:fortunes: - desc: Run fortunes benchmark (connectrpc vs tonic vs connect-go) + desc: Run fortunes benchmark (connectrpc vs tonic vs tonic-protobuf vs connect-go) cmds: - cargo run --release -p rpc-bench --bin fortune_bench @@ -431,25 +431,48 @@ tasks: - cargo run --release -p rpc-bench --bin fortune_bench -- --protocols --multi-conn=8 bench:echo: - desc: Run echo benchmark (connectrpc vs tonic, framework overhead only) + desc: Run echo benchmark (connectrpc vs tonic vs tonic-protobuf, framework overhead only) cmds: - - cargo run --release -p rpc-bench --bin echo_bench + - cargo run --release -p rpc-bench --bin echo_bench -- {{.CLI_ARGS}} bench:echo:quick: desc: Run echo benchmark with shorter duration cmds: - - cargo run --release -p rpc-bench --bin echo_bench -- --quick + - cargo run --release -p rpc-bench --bin echo_bench -- --quick {{.CLI_ARGS}} bench:log: - desc: Run log-ingest benchmark (decode-heavy, buffa-view vs prost-owned) + desc: Run log-ingest benchmark (decode-heavy, buffa-view vs prost-owned vs upb) cmds: - - cargo run --release -p rpc-bench --bin log_bench + - cargo run --release -p rpc-bench --bin log_bench -- {{.CLI_ARGS}} bench:log:quick: desc: Run log-ingest benchmark with shorter duration cmds: - cargo run --release -p rpc-bench --bin log_bench -- --quick + # GRPC_RUST_PROTOC_DIR (prebuilt protoc 35.1 + protoc-gen-rust-grpc) skips the + # cmake toolchain build; see benches/rpc-grpc-rust/README.md. + bench:grpc-rust:build: + desc: Build the grpc-rust bench crate (tonic-protobuf servers + client_bench); first run cmake-builds protoc + protoc-gen-rust-grpc + cmds: + - cargo build --release --manifest-path benches/rpc-grpc-rust/Cargo.toml --bins ${GRPC_RUST_PROTOC_DIR:+--no-default-features} + + bench:grpc-rust:lint: + desc: Clippy + fmt check for the out-of-workspace grpc-rust bench crate (CI does not cover it) + cmds: + - cargo fmt --manifest-path benches/rpc-grpc-rust/Cargo.toml --all -- --check + - cargo clippy --manifest-path benches/rpc-grpc-rust/Cargo.toml --all-targets ${GRPC_RUST_PROTOC_DIR:+--no-default-features} -- -D warnings + + bench:clients: + desc: Run client-stack benchmark (connectrpc-rs vs tonic vs grpc-rust `grpc`, same server) + cmds: + - cargo run --release --manifest-path benches/rpc-grpc-rust/Cargo.toml --bin client_bench ${GRPC_RUST_PROTOC_DIR:+--no-default-features} -- {{.CLI_ARGS}} + + bench:clients:quick: + desc: Run client-stack benchmark with shorter duration + cmds: + - cargo run --release --manifest-path benches/rpc-grpc-rust/Cargo.toml --bin client_bench ${GRPC_RUST_PROTOC_DIR:+--no-default-features} -- --quick {{.CLI_ARGS}} + bench:generate: desc: Regenerate benchmark Rust code from protos dir: "{{.ROOT_DIR}}/benches/rpc" @@ -494,14 +517,23 @@ tasks: cmds: - ./benches/profile_server.sh tonic {{.DURATION}} {{.CONCURRENCY}} + profile:tonic-protobuf: + desc: Profile tonic-protobuf (upb) fortune server (CPU + heap) + vars: + DURATION: '{{.DURATION | default "300"}}' + CONCURRENCY: '{{.CONCURRENCY | default "64"}}' + cmds: + - ./benches/profile_server.sh tonic-protobuf {{.DURATION}} {{.CONCURRENCY}} + profile:compare: - desc: Profile both servers and print comparative summary + desc: Profile all three fortune servers and print comparative summary vars: DURATION: '{{.DURATION | default "300"}}' CONCURRENCY: '{{.CONCURRENCY | default "64"}}' cmds: - ./benches/profile_server.sh connectrpc {{.DURATION}} {{.CONCURRENCY}} - ./benches/profile_server.sh tonic {{.DURATION}} {{.CONCURRENCY}} + - ./benches/profile_server.sh tonic-protobuf {{.DURATION}} {{.CONCURRENCY}} - echo "Results in /tmp/connectrpc-profile/" profile:echo: @@ -513,10 +545,11 @@ tasks: cmds: - ./benches/profile_server.sh echo-connectrpc {{.DURATION}} {{.CONCURRENCY}} {{.N_CONNS}} - ./benches/profile_server.sh echo-tonic {{.DURATION}} {{.CONCURRENCY}} {{.N_CONNS}} - - echo "Results in /tmp/connectrpc-profile/echo-{{'{'}}connectrpc,tonic{{'}'}}/" + - ./benches/profile_server.sh echo-tonic-protobuf {{.DURATION}} {{.CONCURRENCY}} {{.N_CONNS}} + - echo "Results in /tmp/connectrpc-profile/echo-{{'{'}}connectrpc,tonic,tonic-protobuf{{'}'}}/" profile:log: - desc: Profile log-ingest server (decode-heavy, buffa-view vs prost-owned) + desc: Profile log-ingest server (decode-heavy, buffa-view vs prost-owned vs upb) vars: DURATION: '{{.DURATION | default "30"}}' CONCURRENCY: '{{.CONCURRENCY | default "64"}}' @@ -525,7 +558,8 @@ tasks: cmds: - ./benches/profile_server.sh log-connectrpc {{.DURATION}} {{.CONCURRENCY}} {{.N_CONNS}} {{.RECORDS}} - ./benches/profile_server.sh log-tonic {{.DURATION}} {{.CONCURRENCY}} {{.N_CONNS}} {{.RECORDS}} - - echo "Results in /tmp/connectrpc-profile/log-{{'{'}}connectrpc,tonic{{'}'}}/" + - ./benches/profile_server.sh log-tonic-protobuf {{.DURATION}} {{.CONCURRENCY}} {{.N_CONNS}} {{.RECORDS}} + - echo "Results in /tmp/connectrpc-profile/log-{{'{'}}connectrpc,tonic,tonic-protobuf{{'}'}}/" # =========================================================================== # Cleanup @@ -535,4 +569,5 @@ tasks: desc: Remove all build artifacts cmds: - cargo clean + - cargo clean --manifest-path benches/rpc-grpc-rust/Cargo.toml - task: conformance:clean diff --git a/benches/charts/echo.svg b/benches/charts/echo.svg index 570592a7..713b4714 100644 --- a/benches/charts/echo.svg +++ b/benches/charts/echo.svg @@ -1,4 +1,4 @@ - + diff --git a/benches/charts/generate.py b/benches/charts/generate.py index 3a7c418a..5193657f 100644 --- a/benches/charts/generate.py +++ b/benches/charts/generate.py @@ -1,10 +1,10 @@ #!/usr/bin/env python3 -"""Generate SVG bar charts for connectrpc-rs vs tonic benchmarks. +"""Generate SVG bar charts for the connectrpc-rs vs tonic vs tonic-protobuf benchmarks. Reads benchmark data from this file's BENCHMARKS dict (update after -running `task bench:echo --multi-conn=8` and `task bench:log`) and -emits SVG charts to benches/charts/ plus a README-ready markdown -table block. +running `task bench:echo -- --multi-conn=8`, `task bench:log` and +`task bench:cross`) and emits SVG charts to benches/charts/ plus a +README-ready markdown table block. Usage: python3 benches/charts/generate.py @@ -16,9 +16,12 @@ from pathlib import Path # ── Benchmark data ────────────────────────────────────────────────────── -# Update these after running the benchmarks. Values are requests/sec. -# Source: task bench:echo -- --multi-conn=8 and task bench:log -# Machine: Intel Xeon Platinum 8488C, buffa @ 4edfba6 +# Update these after running the benchmarks. Values are requests/sec, except +# the latency block (microseconds). The README tables add descriptive row +# labels and a unary_large row that is kept out of the chart for scale. +# Source: task bench:echo -- --multi-conn=8, task bench:log, task bench:cross +# Machine: c7i.metal-24xl (Intel Xeon Platinum 8488C, bare metal, turbo off), +# 2026-09-25; tonic 0.14.6, tonic-protobuf from grpc-rust @ 7053afcd. BENCHMARKS = { # Echo: 64-byte string, pure framework overhead (8 h2 connections) @@ -27,8 +30,9 @@ "unit": "requests/sec", "groups": ["c=16", "c=64", "c=256"], "series": { - "connectrpc-rs": [170_292, 238_498, 252_000], - "tonic": [168_811, 234_304, 247_167], + "connectrpc-rs": [194_853, 298_919, 271_985], + "tonic": [196_512, 301_501, 264_772], + "tonic-protobuf": [195_235, 298_403, 262_920], }, }, # Log-ingest: 50 records × ~22KB, decode-heavy (8 h2 connections) @@ -37,8 +41,9 @@ "unit": "requests/sec", "groups": ["c=16", "c=64", "c=256"], "series": { - "connectrpc-rs": [32_257, 73_313, 112_027], - "tonic": [28_110, 68_690, 84_171], + "connectrpc-rs": [31_237, 76_678, 138_599], + "tonic": [27_891, 71_910, 120_011], + "tonic-protobuf": [30_561, 78_722, 133_558], }, }, # Single-request latency (criterion, no contention) @@ -47,8 +52,9 @@ "unit": "microseconds (lower is better)", "groups": ["unary_small", "unary_logs_50", "client_stream", "server_stream"], "series": { - "connectrpc-rs": [ 87.6, 195.0, 166.1, 109.8], - "tonic": [170.8, 338.5, 223.8, 110.1], + "connectrpc-rs": [ 79.7, 218.5, 166.3, 108.0], + "tonic": [ 78.8, 303.1, 168.2, 107.0], + "tonic-protobuf": [ 78.4, 253.5, 162.2, 110.2], }, }, } @@ -58,6 +64,7 @@ COLORS = { "connectrpc-rs": "#4C78A8", # blue (our primary) "tonic": "#F58518", # orange + "tonic-protobuf": "#54A24B", # green } # ── SVG generation (adapted from buffa/benchmarks/charts/generate.py) ── @@ -210,9 +217,9 @@ def _pct(val: float, baseline: float) -> str: v = f"{int(round(val)):,}" else: v = f"{val:.1f}" - if baseline == val: - return v diff = (val - baseline) / baseline * 100 + if abs(diff) < 0.5: + return v sign = "+" if diff > 0 else "\u2212" return f"{v} ({sign}{abs(diff):.0f}%)" diff --git a/benches/charts/latency.svg b/benches/charts/latency.svg index fdeb46a1..344fff5a 100644 --- a/benches/charts/latency.svg +++ b/benches/charts/latency.svg @@ -1,4 +1,4 @@ - + diff --git a/benches/charts/log-ingest.svg b/benches/charts/log-ingest.svg index 47b8e729..f7c3e3a7 100644 --- a/benches/charts/log-ingest.svg +++ b/benches/charts/log-ingest.svg @@ -1,4 +1,4 @@ - + diff --git a/benches/charts/tables.md b/benches/charts/tables.md index fc1280bd..1647fa9a 100644 --- a/benches/charts/tables.md +++ b/benches/charts/tables.md @@ -1,18 +1,18 @@ -| Concurrency | connectrpc-rs | tonic | -|---|---:|---:| -| c=16 | 170,292 | 168,811 (−1%) | -| c=64 | 238,498 | 234,304 (−2%) | -| c=256 | 252,000 | 247,167 (−2%) | +| Concurrency | connectrpc-rs | tonic | tonic-protobuf | +|---|---:|---:|---:| +| c=16 | 194,853 | 196,512 (+1%) | 195,235 | +| c=64 | 298,919 | 301,501 (+1%) | 298,403 | +| c=256 | 271,985 | 264,772 (−3%) | 262,920 (−3%) | -| Concurrency | connectrpc-rs | tonic | -|---|---:|---:| -| c=16 | 32,257 | 28,110 (−13%) | -| c=64 | 73,313 | 68,690 (−6%) | -| c=256 | 112,027 | 84,171 (−25%) | +| Concurrency | connectrpc-rs | tonic | tonic-protobuf | +|---|---:|---:|---:| +| c=16 | 31,237 | 27,891 (−11%) | 30,561 (−2%) | +| c=64 | 76,678 | 71,910 (−6%) | 78,722 (+3%) | +| c=256 | 138,599 | 120,011 (−13%) | 133,558 (−4%) | -| Benchmark | connectrpc-rs | tonic | -|---|---:|---:| -| unary_small | 87.6 | 170.8 (+95%) | -| unary_logs_50 | 195.0 | 338.5 (+74%) | -| client_stream | 166.1 | 223.8 (+35%) | -| server_stream | 109.8 | 110.1 (+0%) | +| Benchmark | connectrpc-rs | tonic | tonic-protobuf | +|---|---:|---:|---:| +| unary_small | 79.7 | 78.8 (−1%) | 78.4 (−2%) | +| unary_logs_50 | 218.5 | 303.1 (+39%) | 253.5 (+16%) | +| client_stream | 166.3 | 168.2 (+1%) | 162.2 (−2%) | +| server_stream | 108.0 | 107.0 (−1%) | 110.2 (+2%) | diff --git a/benches/profile_server.sh b/benches/profile_server.sh index 783f17ee..98a454ee 100755 --- a/benches/profile_server.sh +++ b/benches/profile_server.sh @@ -1,8 +1,12 @@ #!/usr/bin/env bash # -# CPU + allocator profiling harness for connectrpc-rs and tonic fortune servers. +# CPU + allocator profiling harness for the connectrpc-rs, tonic (prost) and +# tonic-protobuf (upb) bench servers. # -# Usage: profile_server.sh [connectrpc|tonic] [duration_secs] [concurrency] +# Usage: profile_server.sh [duration_secs] [concurrency] [n_conns] [records] +# target: connectrpc | tonic | tonic-protobuf (fortunes, needs docker) +# echo-{connectrpc,tonic,tonic-protobuf} +# log-{connectrpc,tonic,tonic-protobuf} | log-connectrpc-noutf8 # # Produces: # /tmp/connectrpc-profile//flamegraph.svg — CPU flamegraph @@ -28,6 +32,12 @@ ROOT_DIR="$(cd "$SCRIPT_DIR/.." && pwd)" # NEEDS_VALKEY marks targets whose server takes a valkey address as argv[1]. NEEDS_VALKEY=0 +# The tonic-protobuf servers live in the out-of-workspace benches/rpc-grpc-rust +# crate, which has its own manifest and target dir. +GRPC_RUST_DIR="$ROOT_DIR/benches/rpc-grpc-rust" +SERVER_PKG="" +SERVER_MANIFEST="" +SERVER_TARGET_DIR="$ROOT_DIR/target" case "$TARGET" in connectrpc) SERVER_PKG="rpc-bench" @@ -43,6 +53,14 @@ case "$TARGET" in LOAD_ARGS="$DURATION $CONCURRENCY grpc" NEEDS_VALKEY=1 ;; + tonic-protobuf) + SERVER_MANIFEST="$GRPC_RUST_DIR/Cargo.toml" + SERVER_TARGET_DIR="$GRPC_RUST_DIR/target" + SERVER_BIN_NAME="fortune-server-tonic-protobuf" + LOAD_BIN_NAME="fortune_load" + LOAD_ARGS="$DURATION $CONCURRENCY grpc" + NEEDS_VALKEY=1 + ;; echo-connectrpc) SERVER_PKG="rpc-bench" SERVER_BIN_NAME="echo_server" @@ -55,6 +73,13 @@ case "$TARGET" in LOAD_BIN_NAME="echo_load" LOAD_ARGS="$DURATION $CONCURRENCY $N_CONNS" ;; + echo-tonic-protobuf) + SERVER_MANIFEST="$GRPC_RUST_DIR/Cargo.toml" + SERVER_TARGET_DIR="$GRPC_RUST_DIR/target" + SERVER_BIN_NAME="echo-server-tonic-protobuf" + LOAD_BIN_NAME="echo_load" + LOAD_ARGS="$DURATION $CONCURRENCY $N_CONNS" + ;; log-connectrpc) SERVER_PKG="rpc-bench" SERVER_BIN_NAME="log_server" @@ -67,6 +92,13 @@ case "$TARGET" in LOAD_BIN_NAME="log_load" LOAD_ARGS="$DURATION $CONCURRENCY $N_CONNS $RECORDS" ;; + log-tonic-protobuf) + SERVER_MANIFEST="$GRPC_RUST_DIR/Cargo.toml" + SERVER_TARGET_DIR="$GRPC_RUST_DIR/target" + SERVER_BIN_NAME="log-server-tonic-protobuf" + LOAD_BIN_NAME="log_load" + LOAD_ARGS="$DURATION $CONCURRENCY $N_CONNS $RECORDS" + ;; log-connectrpc-noutf8) SERVER_PKG="rpc-bench" SERVER_BIN_NAME="log_server_noutf8" @@ -75,12 +107,12 @@ case "$TARGET" in ;; *) echo "Unknown target: $TARGET" >&2 - echo "Expected: connectrpc | tonic | echo-{connectrpc,tonic} | log-{connectrpc,tonic} | log-connectrpc-noutf8" >&2 + echo "Expected: connectrpc | tonic | tonic-protobuf | echo-{connectrpc,tonic,tonic-protobuf} | log-{connectrpc,tonic,tonic-protobuf} | log-connectrpc-noutf8" >&2 exit 1 ;; esac -SERVER_BIN="$ROOT_DIR/target/release/$SERVER_BIN_NAME" +SERVER_BIN="$SERVER_TARGET_DIR/release/$SERVER_BIN_NAME" LOAD_BIN="$ROOT_DIR/target/release/$LOAD_BIN_NAME" # ── Check prerequisites ───────────────────────────────────────────── @@ -123,7 +155,17 @@ echo "" # ── Build ──────────────────────────────────────────────────────────── echo "Building $TARGET server with debug info..." -CARGO_PROFILE_RELEASE_DEBUG=2 cargo build --release -p "$SERVER_PKG" --bin "$SERVER_BIN_NAME" +if [[ -n "$SERVER_MANIFEST" ]]; then + # Prebuilt protoc + plugin in GRPC_RUST_PROTOC_DIR skips the cmake toolchain + # build (see benches/rpc-grpc-rust/README.md). + GRPC_RUST_FEATURES=() + if [[ -n "${GRPC_RUST_PROTOC_DIR:-}" ]]; then + GRPC_RUST_FEATURES=(--no-default-features) + fi + CARGO_PROFILE_RELEASE_DEBUG=2 cargo build --release --manifest-path "$SERVER_MANIFEST" --target-dir "$SERVER_TARGET_DIR" --bin "$SERVER_BIN_NAME" "${GRPC_RUST_FEATURES[@]}" +else + CARGO_PROFILE_RELEASE_DEBUG=2 cargo build --release -p "$SERVER_PKG" --bin "$SERVER_BIN_NAME" +fi echo "Building load generator..." CARGO_PROFILE_RELEASE_DEBUG=2 cargo build --release -p rpc-bench --bin "$LOAD_BIN_NAME" diff --git a/benches/rpc-grpc-rust/Cargo.toml b/benches/rpc-grpc-rust/Cargo.toml new file mode 100644 index 00000000..5ff2d933 --- /dev/null +++ b/benches/rpc-grpc-rust/Cargo.toml @@ -0,0 +1,83 @@ +[package] +name = "rpc-bench-grpc-rust" +version = "0.1.0" +edition = "2024" +rust-version = "1.88" +license = "Apache-2.0" +publish = false + +# This crate is deliberately NOT a member of the connect-rust workspace (see +# `exclude` in the root Cargo.toml). Building it needs the grpc-rust codegen +# toolchain (protoc 35.1 + the C++ `protoc-gen-rust-grpc` plugin, built with +# cmake by default), which CI and ordinary contributors should not have to +# carry. The bench drivers in `benches/rpc` build it on demand with +# `cargo build --manifest-path`. +[workspace] + +[dependencies] +# grpc-rust is pinned to a git revision rather than the crates.io 0.9.0 +# preview because `tonic-protobuf` (the server-side codec) is unpublished, +# and mixing a crates.io `grpc` with a git `tonic-protobuf` would pull two +# `protobuf` runtimes. Bump every grpc-rust git dependency together. +grpc = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd" } +grpc-protobuf = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd" } +tonic = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd", features = ["gzip", "zstd"] } +tonic-protobuf = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd" } +# tonic + prost client for the client-stack bench; must be the same tonic as +# above or the generated stubs and `ProstCodec` would name two different +# `tonic` crates. +tonic-prost = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd" } +prost = "0.14" +# Generated message code hardcodes `::protobuf`, so the runtime must be a +# direct dependency at exactly the version grpc-protobuf-build generates for. +protobuf = "=4.35.1-release" +# The fortune server reuses rpc-bench's valkey pool, and the connectrpc-rs +# arms of the client-stack bench reuse its checked-in generated stubs. +rpc-bench = { path = "../rpc" } +connectrpc = { path = "../../connectrpc", features = ["client"] } +http = "1" +tokio = { version = "1", features = ["full"] } +tokio-stream = "0.1" +futures = "0.3" + +[build-dependencies] +grpc-protobuf-build = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd", default-features = false } +tonic-prost-build = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd" } +# Only to locate the cmake-built protoc for tonic-prost-build; see build.rs. +protoc-gen-rust-grpc = { git = "https://github.com/grpc/grpc-rust", rev = "7053afcd", optional = true } + +[features] +# Builds protoc 35.1 and protoc-gen-rust-grpc from C++ source with cmake +# (downloads protobuf + abseil). Disable with `--no-default-features` and set +# GRPC_RUST_PROTOC_DIR to a directory containing both binaries instead. +default = ["build-protoc-plugin"] +build-protoc-plugin = ["grpc-protobuf-build/build-plugin", "dep:protoc-gen-rust-grpc"] + +[[bin]] +name = "client_bench" + +[[bin]] +name = "bench-server-tonic-protobuf" +path = "src/bin/bench_server.rs" + +[[bin]] +name = "echo-server-tonic-protobuf" +path = "src/bin/echo_server.rs" + +[[bin]] +name = "log-server-tonic-protobuf" +path = "src/bin/log_server.rs" + +[[bin]] +name = "fortune-server-tonic-protobuf" +path = "src/bin/fortune_server.rs" + +[lints.rust] +missing_debug_implementations = "warn" +rust_2018_idioms = "warn" +unreachable_pub = "warn" + +[lints.clippy] +all = { level = "warn", priority = -1 } +dbg_macro = "warn" +uninlined_format_args = "warn" diff --git a/benches/rpc-grpc-rust/README.md b/benches/rpc-grpc-rust/README.md new file mode 100644 index 00000000..fe8253e2 --- /dev/null +++ b/benches/rpc-grpc-rust/README.md @@ -0,0 +1,34 @@ +# rpc-bench-grpc-rust + +Benchmark servers and a client-stack benchmark built on [grpc-rust](https://github.com/grpc/grpc-rust), Google's gRPC implementation for Rust, which builds on tonic. `publish = false`; nothing here ships. + +grpc-rust's `grpc` crate (0.9.0 preview on crates.io) is a client-side channel only, and its own benchmarks pair that client with a tonic server. The part of grpc-rust that changes what a server does is `tonic-protobuf`: a tonic codec over Google's official `protobuf` v4 Rust runtime, which wraps the upb C kernel and parses each message eagerly into a per-message arena. The server binaries here are the same echo, log-ingest, fortunes and `BenchService` handlers as `benches/rpc-tonic`, with prost swapped for that codec, so the bench drivers in `benches/rpc` report them as `tonic-protobuf` next to `tonic` (prost) and `connectrpc-rs` (buffa). One caveat when reading those two tonic rows against each other: `benches/rpc-tonic` builds tonic 0.14 from crates.io, while this crate builds tonic from the pinned grpc-rust revision, so the rows differ in tonic source as well as in codec. + +`client_bench` is the inverse experiment. It holds the server constant (the connectrpc-rs echo server unless `--server-bin=PATH` says otherwise) and drives it with each client stack in turn — connectrpc-rs over hyper-util's pooled client, connectrpc-rs over its own `Http2Connection`, tonic's `Channel` with prost, and grpc-rust's `grpc::client::Channel` with grpc-protobuf stubs — because a client benchmark is the only way to compare against the `grpc` crate at all. + +| Binary | Driven by | +|---|---| +| `bench-server-tonic-protobuf` | `cargo bench -p rpc-bench --bench cross_impl_bench` | +| `echo-server-tonic-protobuf` | `task bench:echo` | +| `log-server-tonic-protobuf` | `task bench:log` | +| `fortune-server-tonic-protobuf` | `task bench:fortunes` (needs docker for valkey) | +| `client_bench` | `task bench:clients -- [--quick] [--conns=1,8] [--repeat=3] [--server-bin=PATH]` | + +## Why this crate is outside the workspace + +The root `Cargo.toml` excludes this directory, and the bench drivers build it on demand with `cargo build --release --manifest-path benches/rpc-grpc-rust/Cargo.toml`, so it has its own `target/`. Code generation needs `protoc` 35.1 exactly (the `protobuf-codegen` crate refuses any other version) plus the C++ `protoc-gen-rust-grpc` plugin, and `tonic-protobuf` is unpublished, so every grpc-rust crate is a git dependency pinned to one revision. Keeping that toolchain out of the workspace means CI and ordinary contributors never need it. + +## Toolchain + +By default the first build compiles `protoc` and the plugin from C++ source through grpc-rust's `build-plugin` feature. That needs `cmake` ≥ 3.14, a C++17 compiler, and network access, because the cmake project downloads the protobuf and abseil sources; it takes about a minute on a 20-core machine and is cached in `target/` afterwards. The same `protoc` is handed to `tonic-prost-build` for the tonic client stubs, so no system `protoc` is needed. The `protobuf` runtime crate also compiles upb with `cc`, so a C compiler is always required. + +To use prebuilt binaries instead, put `protoc` (35.1) and `protoc-gen-rust-grpc` in one directory and export `GRPC_RUST_PROTOC_DIR` pointing at it. The bench drivers and `benches/profile_server.sh` then build this crate with `--no-default-features`, which skips the cmake build; a manual build needs the flag spelled out: + +```bash +GRPC_RUST_PROTOC_DIR=/path/to/dir cargo build --release --no-default-features \ + --manifest-path benches/rpc-grpc-rust/Cargo.toml +``` + +## Bumping grpc-rust + +Every grpc-rust git dependency in `Cargo.toml` must move to the same revision together, and the `protobuf = "=…"` pin must match the `protobuf-codegen` version that revision's `grpc-protobuf-build` uses, because the generated message code names `::protobuf` directly. When grpc-rust publishes `tonic-protobuf`, switch all of them to crates.io versions. diff --git a/benches/rpc-grpc-rust/build.rs b/benches/rpc-grpc-rust/build.rs new file mode 100644 index 00000000..94a78fa4 --- /dev/null +++ b/benches/rpc-grpc-rust/build.rs @@ -0,0 +1,34 @@ +fn main() { + println!("cargo:rerun-if-changed=build.rs"); + println!("cargo:rerun-if-changed=../rpc/proto/bench.proto"); + println!("cargo:rerun-if-changed=../rpc/proto/fortune.proto"); + println!("cargo:rerun-if-env-changed=GRPC_RUST_PROTOC_DIR"); + + // Google protobuf messages + grpc-rust client stubs + tonic-protobuf + // server stubs. + grpc_protobuf_build::CodeGen::new() + .include("../rpc/proto") + .inputs(["bench.proto", "fortune.proto"]) + .compile() + .expect("grpc-protobuf codegen failed"); + + // tonic + prost client stubs for the client-stack bench, compiled with the + // same protoc as the grpc-rust stubs: the cmake-built one, or the prebuilt + // one from GRPC_RUST_PROTOC_DIR when the plugin build is disabled. Only if + // neither applies does prost-build fall back to $PROTOC / PATH. + let mut prost_config = tonic_prost_build::Config::new(); + #[cfg(feature = "build-protoc-plugin")] + prost_config.protoc_executable(protoc_gen_rust_grpc::protoc()); + #[cfg(not(feature = "build-protoc-plugin"))] + if let Some(dir) = std::env::var_os("GRPC_RUST_PROTOC_DIR") { + prost_config.protoc_executable(std::path::Path::new(&dir).join("protoc")); + } + tonic_prost_build::configure() + .build_server(false) + .compile_with_config( + prost_config, + &["../rpc/proto/bench.proto"], + &["../rpc/proto"], + ) + .expect("tonic-prost codegen failed"); +} diff --git a/benches/rpc-grpc-rust/src/bin/bench_server.rs b/benches/rpc-grpc-rust/src/bin/bench_server.rs new file mode 100644 index 00000000..9164dc04 --- /dev/null +++ b/benches/rpc-grpc-rust/src/bin/bench_server.rs @@ -0,0 +1,122 @@ +//! `bench.v1.BenchService` server for the cross-implementation criterion +//! bench (tonic + tonic-protobuf on the upb kernel). Mirrors +//! `benches/rpc-tonic/src/main.rs` handler-for-handler. + +use rpc_bench_grpc_rust::pb::bench_service_server::{BenchService, BenchServiceServer}; +use rpc_bench_grpc_rust::pb::{BenchRequest, BenchResponse, LogRequest, LogResponse}; +use tokio::sync::mpsc; +use tokio_stream::StreamExt; +use tokio_stream::wrappers::ReceiverStream; +use tonic::codec::CompressionEncoding; +use tonic::codegen::BoxStream; +use tonic::transport::Server; +use tonic::{Request, Response, Status, Streaming}; + +struct BenchServiceImpl; + +/// Builds a response carrying a copy of the request's payload. upb messages +/// each own an arena, so the payload cannot be moved across as prost does; +/// setting from a view deep-copies it into the response arena. +fn echo_payload(req: &BenchRequest) -> BenchResponse { + let mut resp = BenchResponse::new(); + if let Some(payload) = req.payload_opt() { + resp.set_payload(payload); + } + resp +} + +fn log_count(req: &LogRequest) -> LogResponse { + let mut resp = LogResponse::new(); + resp.set_count(i32::try_from(req.records().len()).unwrap_or(i32::MAX)); + resp +} + +#[tonic::async_trait] +impl BenchService for BenchServiceImpl { + async fn unary( + &self, + request: Request, + ) -> Result, Status> { + Ok(Response::new(echo_payload(&request.into_inner()))) + } + + async fn server_stream( + &self, + request: Request, + ) -> Result>, Status> { + let req = request.into_inner(); + let count = req.response_count(); + let stream = futures::stream::unfold((req, 0), move |(req, i)| async move { + if i >= count { + return None; + } + Some((Ok(echo_payload(&req)), (req, i + 1))) + }); + Ok(Response::new(Box::pin(stream))) + } + + async fn client_stream( + &self, + request: Request>, + ) -> Result, Status> { + let mut stream = request.into_inner(); + let mut last = None; + while let Some(req) = stream.next().await { + last = Some(req?); + } + Ok(Response::new( + last.as_ref().map_or_else(BenchResponse::new, echo_payload), + )) + } + + async fn bidi_stream( + &self, + request: Request>, + ) -> Result>, Status> { + let mut stream = request.into_inner(); + let (tx, rx) = mpsc::channel(1); + tokio::spawn(async move { + while let Some(req) = stream.next().await { + let item = req.map(|req| echo_payload(&req)); + let failed = item.is_err(); + if tx.send(item).await.is_err() || failed { + break; + } + } + }); + Ok(Response::new(Box::pin(ReceiverStream::new(rx)))) + } + + async fn log_unary( + &self, + request: Request, + ) -> Result, Status> { + Ok(Response::new(log_count(&request.into_inner()))) + } + + async fn log_unary_owned( + &self, + request: Request, + ) -> Result, Status> { + Ok(Response::new(log_count(&request.into_inner()))) + } +} + +#[tokio::main] +async fn main() -> Result<(), Box> { + // As with the tonic/prost server: accept compressed requests, but do not + // send_compressed — tonic has no minimum-size threshold and would gzip + // every tiny streaming message, which connectrpc-rs does not. + let svc = BenchServiceServer::new(BenchServiceImpl) + .accept_compressed(CompressionEncoding::Gzip) + .accept_compressed(CompressionEncoding::Zstd); + + Server::builder() + .add_service(svc) + .serve_with_incoming_shutdown( + rpc_bench_grpc_rust::listen()?, + rpc_bench_grpc_rust::shutdown_signal(), + ) + .await?; + Ok(()) +} diff --git a/benches/rpc-grpc-rust/src/bin/client_bench.rs b/benches/rpc-grpc-rust/src/bin/client_bench.rs new file mode 100644 index 00000000..bc14df50 --- /dev/null +++ b/benches/rpc-grpc-rust/src/bin/client_bench.rs @@ -0,0 +1,555 @@ +//! Client-stack benchmark: connectrpc-rs vs tonic vs grpc-rust's `grpc` +//! channel, all driving the same server. +//! +//! The server-side benches in `benches/rpc` hold the client constant +//! (connectrpc-rs) and vary the server. This one holds the server constant +//! (the connectrpc-rs echo server by default) and varies the client, which is +//! the only way to compare against the `grpc` crate: it has no server. Every +//! arm sends the same 64-byte `EchoRequest` over gRPC on h2 with +//! `TCP_NODELAY`, closed-loop, one in-flight request per worker. +//! +//! Arms: +//! - `connectrpc-hyper`: `HttpClient::plaintext_http2_only()` (hyper-util +//! pooled client), one client per connection +//! - `connectrpc-h2`: `Http2Connection` shared with a 1024-deep buffer, +//! the transport the connectrpc-rs docs recommend for gRPC +//! - `tonic`: `tonic::transport::Channel` + prost stubs +//! - `grpc`: `grpc::client::Channel` (pick_first over a `dns:///` target, +//! `LocalChannelCredentials`) + grpc-protobuf stubs on upb messages +//! +//! `--conns=N[,N...]` (default `1,8`) runs the whole sweep once per value, +//! giving each arm N independent channels/connections with workers assigned +//! round-robin: N=1 exposes each stack's single-connection ceiling and N=8 +//! takes h2 connection-mutex contention out of the picture. `connectrpc-h2` +//! and `tonic` dial all N up front; `connectrpc-hyper` and `grpc` connect +//! lazily, so at concurrency below N they hold only as many connections as +//! there are workers. Concurrency 1 is the per-request latency floor. +//! `--repeat=R` runs the full sweep R times, interleaving the arms, and +//! reports the median-throughput run of each cell so that drift over the run +//! does not land on one arm. +//! +//! ```text +//! cargo run --release --manifest-path benches/rpc-grpc-rust/Cargo.toml \ +//! --bin client_bench -- [--quick] [--conns=1,8] [--repeat=3] [--server-bin=PATH] +//! ``` + +use std::future::Future; +use std::net::SocketAddr; +use std::process::{Command, Stdio}; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::time::{Duration, Instant}; + +use connectrpc::Protocol; +use connectrpc::client::{ClientConfig, Http2Connection, HttpClient, SharedHttp2Connection}; +use grpc::credentials::LocalChannelCredentials; +use rpc_bench::ServerProcess; +use rpc_bench_grpc_rust::{pb, tonic_pb}; + +// ── Configuration ──────────────────────────────────────────────────── + +const CONCURRENCY_LEVELS: &[usize] = &[1, 16, 64]; +const DEFAULT_WARMUP: Duration = Duration::from_secs(3); +const DEFAULT_MEASUREMENT: Duration = Duration::from_secs(10); +const QUICK_WARMUP: Duration = Duration::from_secs(1); +const QUICK_MEASUREMENT: Duration = Duration::from_secs(3); +/// Per-worker latency sample cap; every 10th request is sampled. +const MAX_LATENCY_SAMPLES_PER_WORKER: usize = 100_000; + +/// Same 64-byte payload as `benches/rpc/src/bin/echo_bench.rs`. +const PAYLOAD: &str = "lorem ipsum dolor sit amet, consectetur adipiscing elit sed do e"; + +/// Builds the connectrpc-rs echo server from the enclosing workspace, or +/// takes it from `RPC_BENCH_BIN_DIR` like the other drivers. +fn build_default_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("echo_server") { + return path; + } + eprintln!(" Building connectrpc-rs echo server..."); + let root = concat!(env!("CARGO_MANIFEST_DIR"), "/../.."); + // --target-dir pins the output so the returned path holds under + // CARGO_TARGET_DIR or a `[build] target-dir` config. + let status = Command::new("cargo") + .args([ + "build", + "--release", + "-p", + "rpc-bench", + "--bin", + "echo_server", + ]) + .args(["--message-format=short", "--manifest-path"]) + .arg(format!("{root}/Cargo.toml")) + .arg("--target-dir") + .arg(format!("{root}/target")) + .stderr(Stdio::inherit()) + .status() + .expect("failed to run cargo build for echo_server"); + assert!(status.success(), "failed to build echo_server"); + format!("{root}/target/release/echo_server") +} + +// ── Client arms ────────────────────────────────────────────────────── + +/// Outcome of one RPC; the error text is only built on failure. +type CallResult = Result<(), String>; + +/// One client stack under test. `connect` opens `n_conns` independent +/// channels up front (dialling eagerly where the stack allows it, so warmup +/// excludes handshakes); `worker` hands worker `idx` whatever per-task state +/// it needs, bound to channel `idx % n_conns`; `call` issues one echo RPC. +trait Arm { + const NAME: &'static str; + type Worker: Send + 'static; + + fn connect(addr: SocketAddr, n_conns: usize) -> impl Future + Send; + fn worker(&self, idx: usize) -> Self::Worker; + fn call(worker: &mut Self::Worker) -> impl Future + Send; +} + +fn buffa_echo_request() -> rpc_bench::EchoRequest { + rpc_bench::EchoRequest { + message: PAYLOAD.to_string(), + ..Default::default() + } +} + +fn server_uri(addr: SocketAddr) -> http::Uri { + format!("http://{addr}").parse().expect("valid server URL") +} + +fn grpc_config(uri: http::Uri) -> ClientConfig { + ClientConfig::new(uri).with_protocol(Protocol::Grpc) +} + +/// connectrpc-rs over hyper-util's pooled HTTP/2 client. +struct ConnectHyper { + clients: Vec>, +} + +impl Arm for ConnectHyper { + const NAME: &'static str = "connectrpc-hyper"; + type Worker = ( + rpc_bench::EchoServiceClient, + rpc_bench::EchoRequest, + ); + + async fn connect(addr: SocketAddr, n_conns: usize) -> Self { + // Each HttpClient has its own pool and so its own h2 connection; the + // pool connects lazily, during warmup. + let config = grpc_config(server_uri(addr)); + let clients = (0..n_conns) + .map(|_| { + rpc_bench::EchoServiceClient::new( + HttpClient::plaintext_http2_only(), + config.clone(), + ) + }) + .collect(); + Self { clients } + } + + fn worker(&self, idx: usize) -> Self::Worker { + ( + self.clients[idx % self.clients.len()].clone(), + buffa_echo_request(), + ) + } + + async fn call((client, req): &mut Self::Worker) -> CallResult { + // The generated stub takes the request by value, so a clone per call + // is inherent to the API (tonic is the same; grpc-rust takes a view). + client + .echo(req.clone()) + .await + .map(drop) + .map_err(|e| e.to_string()) + } +} + +/// connectrpc-rs over its own `Http2Connection` transport. +struct ConnectH2 { + clients: Vec>, +} + +impl Arm for ConnectH2 { + const NAME: &'static str = "connectrpc-h2"; + type Worker = ( + rpc_bench::EchoServiceClient, + rpc_bench::EchoRequest, + ); + + async fn connect(addr: SocketAddr, n_conns: usize) -> Self { + let uri = server_uri(addr); + let config = grpc_config(uri.clone()); + let mut clients = Vec::with_capacity(n_conns); + for _ in 0..n_conns { + let conn = Http2Connection::connect_plaintext(uri.clone()) + .await + .expect("connectrpc-h2 connect"); + // 1024 matches the buffer depth tonic's Channel uses. + clients.push(rpc_bench::EchoServiceClient::new( + conn.shared(1024), + config.clone(), + )); + } + Self { clients } + } + + fn worker(&self, idx: usize) -> Self::Worker { + ( + self.clients[idx % self.clients.len()].clone(), + buffa_echo_request(), + ) + } + + async fn call((client, req): &mut Self::Worker) -> CallResult { + client + .echo(req.clone()) + .await + .map(drop) + .map_err(|e| e.to_string()) + } +} + +/// tonic's `Channel` with prost messages. +struct Tonic { + clients: Vec>, +} + +impl Arm for Tonic { + const NAME: &'static str = "tonic"; + type Worker = ( + tonic_pb::echo_service_client::EchoServiceClient, + tonic_pb::EchoRequest, + ); + + async fn connect(addr: SocketAddr, n_conns: usize) -> Self { + let endpoint = tonic::transport::Endpoint::from_shared(format!("http://{addr}")) + .expect("valid server URL"); + let mut clients = Vec::with_capacity(n_conns); + for _ in 0..n_conns { + let channel = endpoint.connect().await.expect("tonic connect"); + clients.push(tonic_pb::echo_service_client::EchoServiceClient::new( + channel, + )); + } + Self { clients } + } + + fn worker(&self, idx: usize) -> Self::Worker { + let req = tonic_pb::EchoRequest { + message: PAYLOAD.to_string(), + }; + (self.clients[idx % self.clients.len()].clone(), req) + } + + async fn call((client, req): &mut Self::Worker) -> CallResult { + client + .echo(req.clone()) + .await + .map(drop) + .map_err(|e| e.to_string()) + } +} + +/// grpc-rust's `grpc::client::Channel` with grpc-protobuf stubs. +struct Grpc { + clients: Vec>, +} + +impl Arm for Grpc { + const NAME: &'static str = "grpc"; + type Worker = ( + pb::echo_service_client::EchoServiceClient, + pb::EchoRequest, + ); + + async fn connect(addr: SocketAddr, n_conns: usize) -> Self { + // A Channel is one pick_first subchannel, i.e. one h2 connection, and + // it connects lazily on the first RPC (during warmup). Local + // credentials are grpc-rust's plaintext option; they refuse + // non-loopback peers, which is all this bench uses. + let clients = (0..n_conns) + .map(|_| { + let channel = grpc::client::Channel::builder( + format!("dns:///{addr}"), + LocalChannelCredentials::new_arc(), + ) + .build(); + pb::echo_service_client::EchoServiceClient::new(channel) + }) + .collect(); + Self { clients } + } + + fn worker(&self, idx: usize) -> Self::Worker { + let mut req = pb::EchoRequest::new(); + req.set_message(PAYLOAD); + (self.clients[idx % self.clients.len()].clone(), req) + } + + async fn call((client, req): &mut Self::Worker) -> CallResult { + // The generated stub takes a message view, so unlike the other arms + // there is no per-call request clone to pay for. + // grpc-protobuf's StatusError has no Display impl; spell it out. + client + .echo(req.as_view()) + .await + .map(drop) + .map_err(|e| format!("{:?}: {}", e.code(), e.message())) + } +} + +// ── Benchmark runner ───────────────────────────────────────────────── + +#[derive(Clone)] +struct BenchResult { + arm: &'static str, + conns: usize, + concurrency: usize, + rps: f64, + p50_us: u64, + p99_us: u64, +} + +/// What each worker hands back when the run stops. +struct WorkerTally { + measured: u64, + failures: u64, + first_error: Option, + latencies_us: Vec, +} + +async fn run_arm( + addr: SocketAddr, + n_conns: usize, + concurrency: usize, + warmup: Duration, + measurement: Duration, +) -> BenchResult { + let arm = A::connect(addr, n_conns).await; + + let running = Arc::new(AtomicBool::new(true)); + let measuring = Arc::new(AtomicBool::new(false)); + // Shared only during warmup, to prove the arm works at all; the measured + // window counts per worker so no arm pays for a contended cache line in + // proportion to its own throughput. + let warmup_count = Arc::new(AtomicU64::new(0)); + + let handles: Vec<_> = (0..concurrency) + .map(|idx| { + let mut worker = arm.worker(idx); + let running = Arc::clone(&running); + let measuring = Arc::clone(&measuring); + let warmup_count = Arc::clone(&warmup_count); + tokio::spawn(async move { + let mut tally = WorkerTally { + measured: 0, + failures: 0, + first_error: None, + latencies_us: Vec::with_capacity(MAX_LATENCY_SAMPLES_PER_WORKER), + }; + while running.load(Ordering::Relaxed) { + let start = Instant::now(); + match A::call(&mut worker).await { + Ok(()) if measuring.load(Ordering::Relaxed) => { + tally.measured += 1; + if tally.measured.is_multiple_of(10) + && tally.latencies_us.len() < MAX_LATENCY_SAMPLES_PER_WORKER + { + tally.latencies_us.push(start.elapsed().as_micros() as u64); + } + } + Ok(()) => { + warmup_count.fetch_add(1, Ordering::Relaxed); + } + Err(e) => { + tally.failures += 1; + tally.first_error.get_or_insert(e); + } + } + } + tally + }) + }) + .collect(); + + tokio::time::sleep(warmup).await; + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + assert!( + warmup_count.load(Ordering::Relaxed) > 0, + "{}: no request succeeded during warmup", + A::NAME + ); + measuring.store(true, Ordering::Relaxed); + let measure_start = Instant::now(); + + tokio::time::sleep(measurement).await; + running.store(false, Ordering::Relaxed); + let elapsed = measure_start.elapsed(); + + let mut total = 0u64; + let mut failures = 0u64; + let mut first_error = None; + let mut latencies = Vec::new(); + for h in handles { + let tally = h.await.expect("worker panicked"); + total += tally.measured; + failures += tally.failures; + first_error = first_error.or(tally.first_error); + latencies.extend(tally.latencies_us); + } + // A cell with failed calls is not a throughput measurement; refuse to + // report it rather than let a degraded arm read as a slow one. + assert!( + failures == 0, + "{}: {failures} calls failed (first: {})", + A::NAME, + first_error.unwrap_or_default() + ); + latencies.sort_unstable(); + // Same percentile convention as echo_bench, so the two drivers compare. + let pct = |p: usize| { + latencies + .get(latencies.len() * p / 100) + .copied() + .unwrap_or(0) + }; + + BenchResult { + arm: A::NAME, + conns: n_conns, + concurrency, + rps: total as f64 / elapsed.as_secs_f64(), + p50_us: pct(50), + p99_us: pct(99), + } +} + +async fn sweep_arm( + server_bin: &str, + n_conns: usize, + warmup: Duration, + measurement: Duration, + results: &mut Vec, +) { + for &concurrency in CONCURRENCY_LEVELS { + // Fresh server per run, as in echo_bench, so no arm inherits another + // arm's connections or allocator state. + let server = ServerProcess::start(server_bin, &[]); + eprintln!( + " Benchmarking {} @ conns={n_conns} concurrency={concurrency}...", + A::NAME + ); + let r = run_arm::(server.addr, n_conns, concurrency, warmup, measurement).await; + eprintln!( + " => {:.0} req/s, p50={:.3}ms, p99={:.3}ms", + r.rps, + r.p50_us as f64 / 1000.0, + r.p99_us as f64 / 1000.0 + ); + results.push(r); + drop(server); + } +} + +/// Collapses repeated runs of the same (arm, conns, concurrency) cell to the +/// run with the median throughput, keeping first-seen cell order. +fn median_runs(results: &[BenchResult]) -> Vec { + let mut cells: Vec<(&str, usize, usize)> = Vec::new(); + for r in results { + let key = (r.arm, r.conns, r.concurrency); + if !cells.contains(&key) { + cells.push(key); + } + } + cells + .into_iter() + .map(|(arm, conns, concurrency)| { + let mut runs: Vec<&BenchResult> = results + .iter() + .filter(|r| r.arm == arm && r.conns == conns && r.concurrency == concurrency) + .collect(); + runs.sort_by(|a, b| a.rps.total_cmp(&b.rps)); + runs[runs.len() / 2].clone() + }) + .collect() +} + +fn parse_list(args: &[String], flag: &str, default: &[usize]) -> Vec { + let values: Vec = args + .iter() + .find_map(|a| a.strip_prefix(flag)) + .map(|s| { + s.split(',') + .map(|n| { + n.parse() + .unwrap_or_else(|_| panic!("{flag} takes comma-separated integers")) + }) + .collect() + }) + .unwrap_or_else(|| default.to_vec()); + assert!( + values.iter().all(|&n| n > 0), + "{flag} values must be positive" + ); + values +} + +#[tokio::main] +async fn main() { + let args: Vec = std::env::args().collect(); + let quick = args.iter().any(|a| a == "--quick"); + let conns = parse_list(&args, "--conns=", &[1, 8]); + let repeat = parse_list(&args, "--repeat=", &[1])[0]; + let server_bin = args + .iter() + .find_map(|a| a.strip_prefix("--server-bin=")) + .map_or_else(build_default_server, str::to_owned); + + let (warmup, measurement) = if quick { + eprintln!("Running in quick mode (1s warmup, 3s measurement)"); + (QUICK_WARMUP, QUICK_MEASUREMENT) + } else { + eprintln!("Running full benchmark (3s warmup, 10s measurement)"); + (DEFAULT_WARMUP, DEFAULT_MEASUREMENT) + }; + eprintln!("Server: {server_bin}, repeat={repeat}\n"); + + let mut results = Vec::new(); + for rep in 0..repeat { + if repeat > 1 { + eprintln!("Repetition {} of {repeat}", rep + 1); + } + for &n_conns in &conns { + sweep_arm::(&server_bin, n_conns, warmup, measurement, &mut results) + .await; + sweep_arm::(&server_bin, n_conns, warmup, measurement, &mut results).await; + sweep_arm::(&server_bin, n_conns, warmup, measurement, &mut results).await; + sweep_arm::(&server_bin, n_conns, warmup, measurement, &mut results).await; + } + } + + println!(); + if repeat > 1 { + println!("(median-throughput run of {repeat} per cell)"); + } + println!( + "{:<18} {:>6} {:>12} {:>14} {:>10} {:>10}", + "Client", "Conns", "Concurrency", "Requests/sec", "p50 (ms)", "p99 (ms)" + ); + println!("{}", "-".repeat(74)); + for r in median_runs(&results) { + println!( + "{:<18} {:>6} {:>12} {:>14.0} {:>10.3} {:>10.3}", + r.arm, + r.conns, + r.concurrency, + r.rps, + r.p50_us as f64 / 1000.0, + r.p99_us as f64 / 1000.0, + ); + } +} diff --git a/benches/rpc-grpc-rust/src/bin/echo_server.rs b/benches/rpc-grpc-rust/src/bin/echo_server.rs new file mode 100644 index 00000000..7e2f7bf9 --- /dev/null +++ b/benches/rpc-grpc-rust/src/bin/echo_server.rs @@ -0,0 +1,34 @@ +//! Minimal echo server for framework-overhead benchmarking (tonic + +//! tonic-protobuf on the upb kernel). + +use rpc_bench_grpc_rust::pb::echo_service_server::{EchoService, EchoServiceServer}; +use rpc_bench_grpc_rust::pb::{EchoRequest, EchoResponse}; +use tonic::transport::Server; +use tonic::{Request, Response, Status}; + +struct EchoImpl; + +#[tonic::async_trait] +impl EchoService for EchoImpl { + async fn echo(&self, req: Request) -> Result, Status> { + // upb messages own an arena each, so unlike prost the string cannot + // be moved from request to response; `set_message` copies it into + // the response's arena. + let req = req.into_inner(); + let mut resp = EchoResponse::new(); + resp.set_message(req.message()); + Ok(Response::new(resp)) + } +} + +#[tokio::main] +async fn main() -> Result<(), Box> { + Server::builder() + .add_service(EchoServiceServer::new(EchoImpl)) + .serve_with_incoming_shutdown( + rpc_bench_grpc_rust::listen()?, + rpc_bench_grpc_rust::shutdown_signal(), + ) + .await?; + Ok(()) +} diff --git a/benches/rpc-grpc-rust/src/bin/fortune_server.rs b/benches/rpc-grpc-rust/src/bin/fortune_server.rs new file mode 100644 index 00000000..2bf1f8b8 --- /dev/null +++ b/benches/rpc-grpc-rust/src/bin/fortune_server.rs @@ -0,0 +1,58 @@ +//! Fortunes server backed by valkey (tonic + tonic-protobuf on the upb +//! kernel). Mirrors `benches/rpc-tonic/src/bin/fortune_server.rs`, but takes +//! the valkey pool and query from `rpc_bench::fortune` rather than copying +//! them, since this crate already depends on `rpc-bench`. + +use rpc_bench::fortune::{ValkeyPool, query_fortunes}; +use rpc_bench_grpc_rust::pb::fortune_service_server::{FortuneService, FortuneServiceServer}; +use rpc_bench_grpc_rust::pb::{GetFortunesRequest, GetFortunesResponse}; +use tonic::transport::Server; +use tonic::{Request, Response, Status}; + +/// Same pool size as the connectrpc-rs and tonic/prost fortune servers. +const VALKEY_POOL_SIZE: usize = 8; + +struct FortuneServiceImpl { + pool: ValkeyPool, +} + +#[tonic::async_trait] +impl FortuneService for FortuneServiceImpl { + async fn get_fortunes( + &self, + _req: Request, + ) -> Result, Status> { + let mut conn = self.pool.get(); + let fortunes = query_fortunes(&mut conn) + .await + .map_err(|e| Status::internal(format!("valkey: {e}")))?; + + // Each message string is copied into the response arena; prost moves + // the `String`s into its `Fortune`s instead. + let mut resp = GetFortunesResponse::new(); + let mut out = resp.fortunes_mut(); + for (id, message) in fortunes { + let mut f = out.push_default(); + f.set_id(id); + f.set_message(message); + } + Ok(Response::new(resp)) + } +} + +#[tokio::main] +async fn main() -> Result<(), Box> { + let valkey_addr = std::env::args() + .nth(1) + .ok_or("usage: fortune-server-tonic-protobuf ")?; + let pool = ValkeyPool::connect(&valkey_addr, VALKEY_POOL_SIZE).await?; + + Server::builder() + .add_service(FortuneServiceServer::new(FortuneServiceImpl { pool })) + .serve_with_incoming_shutdown( + rpc_bench_grpc_rust::listen()?, + rpc_bench_grpc_rust::shutdown_signal(), + ) + .await?; + Ok(()) +} diff --git a/benches/rpc-grpc-rust/src/bin/log_server.rs b/benches/rpc-grpc-rust/src/bin/log_server.rs new file mode 100644 index 00000000..5c36b85a --- /dev/null +++ b/benches/rpc-grpc-rust/src/bin/log_server.rs @@ -0,0 +1,76 @@ +//! Log-ingest server for decode-heavy profiling (tonic + tonic-protobuf on +//! the upb kernel). +//! +//! Matches the connectrpc-rs `log_server` and tonic/prost `log-server-tonic` +//! handlers field-for-field so the measured difference is the proto library. +//! upb parses the whole `LogRequest` eagerly into an arena before the +//! handler runs (string data is copied into the arena, not borrowed from +//! the request buffer); the accessors below then read arena-backed views. + +use rpc_bench_grpc_rust::pb::log_ingest_service_server::{ + LogIngestService, LogIngestServiceServer, +}; +use rpc_bench_grpc_rust::pb::{LogIngestResponse, LogRequest}; +use tonic::transport::Server; +use tonic::{Request, Response, Status}; + +struct LogIngestImpl; + +#[tonic::async_trait] +impl LogIngestService for LogIngestImpl { + async fn ingest( + &self, + request: Request, + ) -> Result, Status> { + let req = request.into_inner(); + + let mut count = 0i32; + let mut total_message_bytes = 0i64; + let mut total_label_bytes = 0i64; + let mut max_severity = 0i32; + + for rec in req.records() { + count += 1; + + let sev = i32::from(rec.severity()); + if sev > max_severity { + max_severity = sev; + } + + total_message_bytes += rec.message().len() as i64; + total_message_bytes += rec.service_name().len() as i64; + total_message_bytes += rec.instance_id().len() as i64; + total_message_bytes += rec.trace_id().len() as i64; + total_message_bytes += rec.span_id().len() as i64; + + if let Some(src) = rec.source_opt() { + total_message_bytes += src.file().len() as i64; + total_message_bytes += src.function().len() as i64; + let _ = src.line(); + } + + for (k, v) in rec.labels() { + total_label_bytes += (k.len() + v.len()) as i64; + } + } + + let mut resp = LogIngestResponse::new(); + resp.set_count(count); + resp.set_total_message_bytes(total_message_bytes); + resp.set_total_label_bytes(total_label_bytes); + resp.set_max_severity(max_severity); + Ok(Response::new(resp)) + } +} + +#[tokio::main] +async fn main() -> Result<(), Box> { + Server::builder() + .add_service(LogIngestServiceServer::new(LogIngestImpl)) + .serve_with_incoming_shutdown( + rpc_bench_grpc_rust::listen()?, + rpc_bench_grpc_rust::shutdown_signal(), + ) + .await?; + Ok(()) +} diff --git a/benches/rpc-grpc-rust/src/lib.rs b/benches/rpc-grpc-rust/src/lib.rs new file mode 100644 index 00000000..d6d2e8ba --- /dev/null +++ b/benches/rpc-grpc-rust/src/lib.rs @@ -0,0 +1,58 @@ +//! Benchmark servers and clients built on grpc-rust — the `tonic-protobuf` +//! codec (Google's `protobuf` v4 runtime on the upb kernel) and the `grpc` +//! client channel — for comparison against connectrpc-rs + buffa and +//! tonic + prost. + +use std::net::{Ipv4Addr, SocketAddr}; + +use tonic::transport::server::TcpIncoming; + +/// Google-`protobuf` message types, grpc-rust client stubs and +/// tonic-protobuf server stubs for `bench.v1` and `fortune.v1`. +/// grpc-protobuf-build flattens every input proto's messages into one +/// `generated.rs`, so both packages share this namespace: `include_proto!` +/// pulls in that file plus `bench_grpc.pb.rs`, and only the per-proto service +/// file is left to include for `fortune`. +#[allow( + unreachable_pub, + missing_debug_implementations, + clippy::all, + clippy::pedantic, + rustdoc::all +)] +pub mod pb { + grpc::include_proto!("bench"); + include!(concat!(env!("OUT_DIR"), "/fortune_grpc.pb.rs")); +} + +/// prost message types and tonic client stubs for `bench.v1`, used by the +/// tonic arm of `client_bench`. +#[allow( + unreachable_pub, + missing_debug_implementations, + clippy::all, + clippy::pedantic +)] +pub mod tonic_pb { + tonic::include_proto!("bench.v1"); +} + +/// Binds an ephemeral loopback port with `TCP_NODELAY` and prints the +/// address on stdout, which is how the bench drivers in `benches/rpc` +/// discover every server they spawn. +/// +/// # Errors +/// +/// Returns the bind or `local_addr` failure. +pub fn listen() -> std::io::Result { + let incoming = + TcpIncoming::bind(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)))?.with_nodelay(Some(true)); + println!("{}", incoming.local_addr()?); + Ok(incoming) +} + +/// Resolves when the process receives Ctrl-C, for +/// `serve_with_incoming_shutdown`. +pub async fn shutdown_signal() { + let _ = tokio::signal::ctrl_c().await; +} diff --git a/benches/rpc/Cargo.toml b/benches/rpc/Cargo.toml index 0294c03d..d363f75e 100644 --- a/benches/rpc/Cargo.toml +++ b/benches/rpc/Cargo.toml @@ -57,6 +57,9 @@ name = "log_server_noutf8" [[bin]] name = "log_load_noutf8" +[[bin]] +name = "unary_load" + [[bench]] name = "rpc_bench" harness = false diff --git a/benches/rpc/README.md b/benches/rpc/README.md index f96ad7fa..cfb7de34 100644 --- a/benches/rpc/README.md +++ b/benches/rpc/README.md @@ -7,7 +7,7 @@ Benchmark crate for `connectrpc`. `publish = false`; nothing here ships. | Bench | What it measures | Run | |---|---|---| | `rpc_bench` | Full-stack unary/stream RPC over loopback HTTP, Connect/gRPC/gRPC-Web × proto/JSON. Needs a running server (`cargo run --bin echo_server` etc.). | `cargo bench --bench rpc_bench` | -| `cross_impl_bench` | Cross-implementation comparison vs tonic. | `cargo bench --bench cross_impl_bench` | +| `cross_impl_bench` | Cross-implementation comparison vs tonic (prost), tonic-protobuf (grpc-rust's upb codec, built from `../rpc-grpc-rust`) and connect-go. | `cargo bench --bench cross_impl_bench` | | `echo_bloat` | Codec-layer (no HTTP) `{owned,view}×{decode,encode}` sweep across five payload shapes + a 1→N fanout sweep. Motivates the future view-response handler API. | `cargo bench --bench echo_bloat` | | `reader_bench` | Request-body reader cost of client-streaming calls, without a network: envelope-framed messages go through `ConnectRpcService` from an in-memory body (ready, chunked, and pending before every frame). Run it on a dedicated machine; it has no client or server to contend with. | `cargo bench --bench reader_bench` | | `view_rope_encode` | Encode cost of a response view, contiguous vs a rope backed by the view's own buffer, swept either side of the segment threshold, and across field shapes the size alone cannot tell apart (many small fields, one large field among them). Shows what the segmented response path buys, what it costs below the threshold, and what its field probe costs. | `cargo bench --bench view_rope_encode` | @@ -18,8 +18,19 @@ Filter by criterion regex: `cargo bench --bench echo_bloat -- fanout` or `-- map ## Load-gen binaries `{echo,log,fortune}_server` / `_load` / `_bench` are standalone server + -client load generators for ad-hoc throughput testing. The `_noutf8` -variants exercise buffa's `respect_utf8_validation_feature(true)` mode. +client load generators for ad-hoc throughput testing. Each `_bench` driver +builds and runs the matching tonic (prost) and tonic-protobuf (upb) servers +from `../rpc-tonic` and `../rpc-grpc-rust` alongside ours; see +`../rpc-grpc-rust/README.md` for the toolchain the latter needs on first +build. The `_noutf8` variants exercise buffa's +`respect_utf8_validation_feature(true)` mode. + +To run the drivers on a host without the toolchains (a dedicated benchmark +box, say), build every server binary elsewhere, copy them into one +directory, and set `RPC_BENCH_BIN_DIR` to it: the drivers then take +`bench_server`, `echo-server-tonic`, `bench-connect-go`, +`log-server-tonic-protobuf` and the rest from there instead of running +`cargo build` / `go build`. ## echo_bloat: adding a payload shape diff --git a/benches/rpc/benches/cross_impl_bench.rs b/benches/rpc/benches/cross_impl_bench.rs index 14d562fa..fa2995a1 100644 --- a/benches/rpc/benches/cross_impl_bench.rs +++ b/benches/rpc/benches/cross_impl_bench.rs @@ -1,7 +1,5 @@ -use std::io::{BufRead, BufReader}; use std::net::SocketAddr; -use std::process::{Child, Command, Stdio}; -use std::time::Duration; +use std::process::{Command, Stdio}; use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; @@ -11,51 +9,6 @@ use rpc_bench::*; const STREAM_MSG_COUNT: i32 = 10; -// ── Server process management ───────────────────────────────────────── - -struct ServerProcess { - child: Child, - addr: SocketAddr, -} - -impl ServerProcess { - /// Start a server process and read the address from its first stdout line. - fn start(cmd: &str, args: &[&str]) -> Self { - let mut child = Command::new(cmd) - .args(args) - .stdout(Stdio::piped()) - .stderr(Stdio::null()) - .spawn() - .unwrap_or_else(|e| panic!("failed to start {cmd}: {e}")); - - let stdout = child.stdout.take().expect("no stdout"); - let mut reader = BufReader::new(stdout); - let mut line = String::new(); - reader - .read_line(&mut line) - .expect("failed to read server address"); - let addr: SocketAddr = line.trim().parse().unwrap_or_else(|e| { - panic!("failed to parse server address from {cmd} ({line:?}): {e}") - }); - - // Give the server a moment to be fully ready for connections. - std::thread::sleep(Duration::from_millis(50)); - - Self { child, addr } - } - - fn addr(&self) -> SocketAddr { - self.addr - } -} - -impl Drop for ServerProcess { - fn drop(&mut self) { - let _ = self.child.kill(); - let _ = self.child.wait(); - } -} - // ── Helpers ─────────────────────────────────────────────────────────── fn make_grpc_client(addr: SocketAddr) -> BenchServiceClient { @@ -73,6 +26,9 @@ fn make_connect_client(addr: SocketAddr) -> BenchServiceClient { // ── Server paths ────────────────────────────────────────────────────── fn connectrpc_server_path() -> String { + if let Some(path) = prebuilt_bin("bench_server") { + return path; + } // Build the server binary if needed and return its path. let output = Command::new("cargo") .args([ @@ -95,6 +51,9 @@ fn connectrpc_server_path() -> String { } fn tonic_server_path() -> String { + if let Some(path) = prebuilt_bin("rpc-bench-tonic") { + return path; + } let output = Command::new("cargo") .args([ "build", @@ -112,7 +71,14 @@ fn tonic_server_path() -> String { format!("{manifest_dir}/../../target/release/rpc-bench-tonic") } +fn tonic_protobuf_server_path() -> String { + build_grpc_rust_bin("bench-server-tonic-protobuf") +} + fn connect_go_server_path() -> String { + if let Some(path) = prebuilt_bin("bench-connect-go") { + return path; + } let manifest_dir = env!("CARGO_MANIFEST_DIR"); let go_dir = format!("{manifest_dir}/../rpc-go"); let bin_path = format!("{go_dir}/bench-connect-go"); @@ -128,29 +94,55 @@ fn connect_go_server_path() -> String { bin_path } -// ── Benchmarks ──────────────────────────────────────────────────────── - -/// Labels for each server implementation. -const IMPLS: [&str; 3] = ["connectrpc-rs", "tonic", "connect-go"]; +// ── Server sets ─────────────────────────────────────────────────────── + +/// One freshly started server per implementation that speaks gRPC, labelled +/// for the criterion report. `tonic` is tonic + prost; `tonic-protobuf` is +/// tonic + grpc-rust's codec on Google's upb-kernel `protobuf` runtime. +fn grpc_servers() -> Vec<(&'static str, ServerProcess)> { + vec![ + ( + "connectrpc-rs", + ServerProcess::start(&connectrpc_server_path(), &[]), + ), + ("tonic", ServerProcess::start(&tonic_server_path(), &[])), + ( + "tonic-protobuf", + ServerProcess::start(&tonic_protobuf_server_path(), &[]), + ), + ( + "connect-go", + ServerProcess::start(&connect_go_server_path(), &[]), + ), + ] +} -fn bench_unary_small_grpc(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let tonic_bin = tonic_server_path(); - let connect_go_bin = connect_go_server_path(); +/// One freshly started server per implementation that speaks the Connect +/// protocol. tonic serves gRPC only. +fn connect_servers() -> Vec<(&'static str, ServerProcess)> { + vec![ + ( + "connectrpc-rs", + ServerProcess::start(&connectrpc_server_path(), &[]), + ), + ( + "connect-go", + ServerProcess::start(&connect_go_server_path(), &[]), + ), + ] +} - let servers = [ - ServerProcess::start(&connectrpc_bin, &[]), - ServerProcess::start(&tonic_bin, &[]), - ServerProcess::start(&connect_go_bin, &[]), - ]; +// ── Benchmarks ──────────────────────────────────────────────────────── +fn bench_unary_small_grpc(c: &mut Criterion) { + let servers = grpc_servers(); let req = small_request(); let rt = tokio::runtime::Runtime::new().unwrap(); let mut group = c.benchmark_group("cross/unary_small_grpc"); - for (impl_name, server) in IMPLS.iter().zip(servers.iter()) { - let client = make_grpc_client(server.addr()); + for (impl_name, server) in &servers { + let client = make_grpc_client(server.addr); group.bench_function(BenchmarkId::from_parameter(impl_name), |b| { b.to_async(&rt) .iter(|| async { client.unary(req.clone()).await.expect("unary RPC failed") }); @@ -158,26 +150,17 @@ fn bench_unary_small_grpc(c: &mut Criterion) { } group.finish(); - drop(servers); } fn bench_unary_small_connect(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let connect_go_bin = connect_go_server_path(); - - // tonic doesn't support Connect protocol, only gRPC - let servers = [ - ("connectrpc-rs", ServerProcess::start(&connectrpc_bin, &[])), - ("connect-go", ServerProcess::start(&connect_go_bin, &[])), - ]; - + let servers = connect_servers(); let req = small_request(); let rt = tokio::runtime::Runtime::new().unwrap(); let mut group = c.benchmark_group("cross/unary_small_connect"); for (impl_name, server) in &servers { - let client = make_connect_client(server.addr()); + let client = make_connect_client(server.addr); group.bench_function(BenchmarkId::from_parameter(impl_name), |b| { b.to_async(&rt) .iter(|| async { client.unary(req.clone()).await.expect("unary RPC failed") }); @@ -188,16 +171,7 @@ fn bench_unary_small_connect(c: &mut Criterion) { } fn bench_unary_large_grpc(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let tonic_bin = tonic_server_path(); - let connect_go_bin = connect_go_server_path(); - - let servers = [ - ServerProcess::start(&connectrpc_bin, &[]), - ServerProcess::start(&tonic_bin, &[]), - ServerProcess::start(&connect_go_bin, &[]), - ]; - + let servers = grpc_servers(); let req = large_request(); let payload_size = { use buffa::Message; @@ -208,8 +182,8 @@ fn bench_unary_large_grpc(c: &mut Criterion) { let mut group = c.benchmark_group("cross/unary_large_grpc"); group.throughput(Throughput::Bytes(payload_size)); - for (impl_name, server) in IMPLS.iter().zip(servers.iter()) { - let config = ClientConfig::new(format!("http://{}", server.addr()).parse().unwrap()) + for (impl_name, server) in &servers { + let config = ClientConfig::new(format!("http://{}", server.addr).parse().unwrap()) .with_protocol(Protocol::Grpc) .compress_requests("gzip"); let client = BenchServiceClient::new(HttpClient::plaintext_http2_only(), config); @@ -220,20 +194,10 @@ fn bench_unary_large_grpc(c: &mut Criterion) { } group.finish(); - drop(servers); } fn bench_server_stream_grpc(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let tonic_bin = tonic_server_path(); - let connect_go_bin = connect_go_server_path(); - - let servers = [ - ServerProcess::start(&connectrpc_bin, &[]), - ServerProcess::start(&tonic_bin, &[]), - ServerProcess::start(&connect_go_bin, &[]), - ]; - + let servers = grpc_servers(); let base_req = BenchRequest { response_count: STREAM_MSG_COUNT, payload: small_payload().into(), @@ -244,8 +208,8 @@ fn bench_server_stream_grpc(c: &mut Criterion) { let mut group = c.benchmark_group("cross/server_stream_grpc"); group.throughput(Throughput::Elements(STREAM_MSG_COUNT as u64)); - for (impl_name, server) in IMPLS.iter().zip(servers.iter()) { - let client = make_grpc_client(server.addr()); + for (impl_name, server) in &servers { + let client = make_grpc_client(server.addr); group.bench_function(BenchmarkId::from_parameter(impl_name), |b| { b.to_async(&rt).iter(|| async { let mut stream = client @@ -267,28 +231,18 @@ fn bench_server_stream_grpc(c: &mut Criterion) { } group.finish(); - drop(servers); } fn bench_client_stream_grpc(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let tonic_bin = tonic_server_path(); - let connect_go_bin = connect_go_server_path(); - - let servers = [ - ServerProcess::start(&connectrpc_bin, &[]), - ServerProcess::start(&tonic_bin, &[]), - ServerProcess::start(&connect_go_bin, &[]), - ]; - + let servers = grpc_servers(); let messages: Vec = (0..STREAM_MSG_COUNT).map(|_| small_request()).collect(); let rt = tokio::runtime::Runtime::new().unwrap(); let mut group = c.benchmark_group("cross/client_stream_grpc"); group.throughput(Throughput::Elements(STREAM_MSG_COUNT as u64)); - for (impl_name, server) in IMPLS.iter().zip(servers.iter()) { - let client = make_grpc_client(server.addr()); + for (impl_name, server) in &servers { + let client = make_grpc_client(server.addr); group.bench_function(BenchmarkId::from_parameter(impl_name), |b| { b.to_async(&rt).iter(|| async { client @@ -300,20 +254,10 @@ fn bench_client_stream_grpc(c: &mut Criterion) { } group.finish(); - drop(servers); } fn bench_unary_logs_grpc(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let tonic_bin = tonic_server_path(); - let connect_go_bin = connect_go_server_path(); - - let servers = [ - ServerProcess::start(&connectrpc_bin, &[]), - ServerProcess::start(&tonic_bin, &[]), - ServerProcess::start(&connect_go_bin, &[]), - ]; - + let servers = grpc_servers(); let req = log_request(50); let payload_size = { use buffa::Message; @@ -324,8 +268,8 @@ fn bench_unary_logs_grpc(c: &mut Criterion) { let mut group = c.benchmark_group("cross/unary_logs_50_grpc"); group.throughput(Throughput::Bytes(payload_size)); - for (impl_name, server) in IMPLS.iter().zip(servers.iter()) { - let client = make_grpc_client(server.addr()); + for (impl_name, server) in &servers { + let client = make_grpc_client(server.addr); group.bench_function(BenchmarkId::from_parameter(impl_name), |b| { b.to_async(&rt) .iter(|| async { client.log_unary(req.clone()).await.expect("log RPC failed") }); @@ -333,18 +277,10 @@ fn bench_unary_logs_grpc(c: &mut Criterion) { } group.finish(); - drop(servers); } fn bench_unary_logs_connect(c: &mut Criterion) { - let connectrpc_bin = connectrpc_server_path(); - let connect_go_bin = connect_go_server_path(); - - let servers = [ - ("connectrpc-rs", ServerProcess::start(&connectrpc_bin, &[])), - ("connect-go", ServerProcess::start(&connect_go_bin, &[])), - ]; - + let servers = connect_servers(); let req = log_request(50); let payload_size = { use buffa::Message; @@ -356,7 +292,7 @@ fn bench_unary_logs_connect(c: &mut Criterion) { group.throughput(Throughput::Bytes(payload_size)); for (impl_name, server) in &servers { - let client = make_connect_client(server.addr()); + let client = make_connect_client(server.addr); group.bench_function(BenchmarkId::from_parameter(impl_name), |b| { b.to_async(&rt) .iter(|| async { client.log_unary(req.clone()).await.expect("log RPC failed") }); diff --git a/benches/rpc/src/bin/bench_server.rs b/benches/rpc/src/bin/bench_server.rs index b1075b29..2d8826c5 100644 --- a/benches/rpc/src/bin/bench_server.rs +++ b/benches/rpc/src/bin/bench_server.rs @@ -1,18 +1,28 @@ +//! `bench.v1.BenchService` server for the cross-implementation criterion +//! bench, on the codegen `BenchServiceServer` dispatcher — the same +//! monomorphic path `echo_server` and `log_server` use. Pass `--router` to +//! serve through the dynamic `Router` instead, for comparing the two. + use std::sync::Arc; -use connectrpc::Router; -use rpc_bench::{BenchServiceExt, BenchServiceImpl}; +use connectrpc::{ConnectRpcService, Router}; +use rpc_bench::{BenchServiceExt, BenchServiceImpl, BenchServiceServer}; #[tokio::main] async fn main() -> Result<(), Box> { - let router = Router::new(); - let router = Arc::new(BenchServiceImpl).register(router); + let use_router = std::env::args().any(|a| a == "--router"); let bound = connectrpc::server::Server::bind("127.0.0.1:0").await?; let addr = bound.local_addr()?; // Print the address to stdout for the benchmark harness. println!("{addr}"); - bound.serve(router).await?; + if use_router { + let router = Arc::new(BenchServiceImpl).register(Router::new()); + bound.serve(router).await?; + } else { + let service = ConnectRpcService::new(BenchServiceServer::new(BenchServiceImpl)); + bound.serve_with_service(service).await?; + } Ok(()) } diff --git a/benches/rpc/src/bin/echo_bench.rs b/benches/rpc/src/bin/echo_bench.rs index 1fd7a9ce..2396320d 100644 --- a/benches/rpc/src/bin/echo_bench.rs +++ b/benches/rpc/src/bin/echo_bench.rs @@ -1,10 +1,11 @@ -//! Echo benchmark: connectrpc-rs vs tonic, framework overhead only. +//! Echo benchmark: connectrpc-rs vs tonic (prost and upb-protobuf codecs), +//! framework overhead only. //! //! Unlike fortune_bench, this eliminates the database, spawn_blocking, and //! complex message encoding from the critical path. The handler just copies //! a short string from request to response. What remains is: //! -//! - HTTP/2 framing (h2 crate — shared between both frameworks) +//! - HTTP/2 framing (h2 crate — shared by all three servers) //! - Protocol detection + envelope framing //! - Dispatch (monomorphic `FooServiceServer` vs tonic's match) //! - Proto encode/decode of one string field @@ -13,15 +14,15 @@ //! This makes framework-specific overhead a much larger fraction of //! per-request CPU, so small improvements become visible. -use std::io::{BufRead, BufReader}; use std::net::SocketAddr; -use std::process::{Child, Command, Stdio}; +use std::process::{Command, Stdio}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::{Duration, Instant}; use connectrpc::Protocol; use connectrpc::client::{ClientConfig, Http2Connection, HttpClient, SharedHttp2Connection}; +use rpc_bench::ServerProcess; use rpc_bench::connect::bench::v1::*; use rpc_bench::proto::bench::v1::*; @@ -39,51 +40,12 @@ const MAX_LATENCY_SAMPLES: usize = 500_000; /// short enough that memcpy doesn't dominate. const PAYLOAD: &str = "lorem ipsum dolor sit amet, consectetur adipiscing elit sed do e"; -// ── Server process management ──────────────────────────────────────── - -struct ServerProcess { - child: Child, - addr: SocketAddr, -} - -impl ServerProcess { - fn start(cmd: &str, args: &[&str]) -> Self { - let mut child = Command::new(cmd) - .args(args) - .stdout(Stdio::piped()) - .stderr(Stdio::null()) - .spawn() - .unwrap_or_else(|e| panic!("failed to start {cmd}: {e}")); - - let stdout = child.stdout.take().expect("no stdout"); - let mut reader = BufReader::new(stdout); - let mut line = String::new(); - let bytes_read = reader - .read_line(&mut line) - .expect("failed to read server address"); - if bytes_read == 0 { - panic!("server {cmd} exited before printing its address"); - } - let addr: SocketAddr = line.trim().parse().unwrap_or_else(|e| { - panic!("failed to parse server address from {cmd} ({line:?}): {e}") - }); - - std::thread::sleep(Duration::from_millis(50)); - - Self { child, addr } - } -} - -impl Drop for ServerProcess { - fn drop(&mut self) { - let _ = self.child.kill(); - let _ = self.child.wait(); - } -} - // ── Build helpers ──────────────────────────────────────────────────── fn build_connectrpc_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("echo_server") { + return path; + } eprintln!(" Building connectrpc-rs echo server..."); let output = Command::new("cargo") .args([ @@ -105,6 +67,9 @@ fn build_connectrpc_server() -> String { } fn build_tonic_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("echo-server-tonic") { + return path; + } eprintln!(" Building tonic echo server..."); let output = Command::new("cargo") .args([ @@ -125,6 +90,10 @@ fn build_tonic_server() -> String { format!("{manifest_dir}/../../target/release/echo-server-tonic") } +fn build_tonic_protobuf_server() -> String { + rpc_bench::build_grpc_rust_bin("echo-server-tonic-protobuf") +} + // ── Benchmark result ───────────────────────────────────────────────── struct BenchResult { @@ -193,7 +162,13 @@ async fn bench_server( // Warmup phase. tokio::time::sleep(warmup).await; - count.store(0, Ordering::Relaxed); + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + let warmed = count.swap(0, Ordering::Relaxed); + assert!( + warmed > 0, + "{impl_name}: no request succeeded during warmup" + ); latencies.lock().await.clear(); let measure_start = Instant::now(); @@ -299,7 +274,13 @@ async fn bench_server_multiconn( } tokio::time::sleep(warmup).await; - count.store(0, Ordering::Relaxed); + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + let warmed = count.swap(0, Ordering::Relaxed); + assert!( + warmed > 0, + "{impl_name}: no request succeeded during warmup" + ); latencies.lock().await.clear(); let measure_start = Instant::now(); @@ -361,10 +342,12 @@ async fn main() { eprintln!("\nBuilding servers..."); let connectrpc_bin = build_connectrpc_server(); let tonic_bin = build_tonic_server(); + let tonic_protobuf_bin = build_tonic_protobuf_server(); let servers = [ ("connectrpc-rs", &connectrpc_bin, Protocol::Grpc), ("tonic", &tonic_bin, Protocol::Grpc), + ("tonic-protobuf", &tonic_protobuf_bin, Protocol::Grpc), ]; let mut results = Vec::new(); diff --git a/benches/rpc/src/bin/echo_load.rs b/benches/rpc/src/bin/echo_load.rs index 28b87f08..313741c8 100644 --- a/benches/rpc/src/bin/echo_load.rs +++ b/benches/rpc/src/bin/echo_load.rs @@ -1,6 +1,6 @@ //! Standalone load generator for profiling echo servers. //! -//! Usage: `echo_load [duration_secs] [concurrency] [n_conns]` +//! Usage: `echo_load [duration_secs] [concurrency] [n_conns] [payload_bytes]` //! //! Uses gRPC (HTTP/2) always. `n_conns` spreads load across N //! SharedHttp2Connection instances to reduce h2 mutex contention, @@ -24,6 +24,12 @@ async fn main() { let duration = args.get(2).and_then(|s| s.parse().ok()).unwrap_or(30u64); let concurrency: usize = args.get(3).and_then(|s| s.parse().ok()).unwrap_or(64); let n_conns: usize = args.get(4).and_then(|s| s.parse().ok()).unwrap_or(8); + // Default is the 64-byte PAYLOAD; a larger value repeats it to that length, + // e.g. to straddle h2's 256-byte DATA-frame chain threshold. + let payload_bytes: usize = args + .get(5) + .and_then(|s| s.parse().ok()) + .unwrap_or(PAYLOAD.len()); let uri: http::Uri = format!("http://{addr}").parse().unwrap(); let config = ClientConfig::new(uri.clone()).with_protocol(Protocol::Grpc); @@ -42,7 +48,7 @@ async fn main() { } let request = EchoRequest { - message: PAYLOAD.to_string(), + message: PAYLOAD.repeat(payload_bytes.div_ceil(PAYLOAD.len()))[..payload_bytes].to_string(), ..Default::default() }; diff --git a/benches/rpc/src/bin/fortune_bench.rs b/benches/rpc/src/bin/fortune_bench.rs index e952a085..08f39d19 100644 --- a/benches/rpc/src/bin/fortune_bench.rs +++ b/benches/rpc/src/bin/fortune_bench.rs @@ -1,19 +1,20 @@ -//! Fortunes benchmark: connectrpc-rs vs tonic vs connect-go. +//! Fortunes benchmark: connectrpc-rs vs tonic (prost and upb-protobuf codecs) +//! vs connect-go. //! //! Adapted from the TechEmpower Web Framework Benchmarks. Measures a realistic //! workload: network round-trip to a valkey backing store, string processing, //! sorting, and response encoding via a single unary RPC. A valkey container //! is spawned as a sibling process and shared across all server runs. -use std::io::{BufRead, BufReader}; use std::net::SocketAddr; -use std::process::{Child, Command, Stdio}; +use std::process::{Command, Stdio}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::{Duration, Instant}; use connectrpc::Protocol; use connectrpc::client::{ClientConfig, Http2Connection, HttpClient, SharedHttp2Connection}; +use rpc_bench::ServerProcess; use rpc_bench::connect::fortune::v1::*; use rpc_bench::fortune; use rpc_bench::proto::fortune::v1::*; @@ -31,48 +32,6 @@ const QUICK_WARMUP: Duration = Duration::from_secs(1); const QUICK_MEASUREMENT: Duration = Duration::from_secs(3); const MAX_LATENCY_SAMPLES: usize = 500_000; -// ── Server process management ──────────────────────────────────────── - -struct ServerProcess { - child: Child, - addr: SocketAddr, -} - -impl ServerProcess { - fn start(cmd: &str, args: &[&str]) -> Self { - let mut child = Command::new(cmd) - .args(args) - .stdout(Stdio::piped()) - .stderr(Stdio::null()) - .spawn() - .unwrap_or_else(|e| panic!("failed to start {cmd}: {e}")); - - let stdout = child.stdout.take().expect("no stdout"); - let mut reader = BufReader::new(stdout); - let mut line = String::new(); - let bytes_read = reader - .read_line(&mut line) - .expect("failed to read server address"); - if bytes_read == 0 { - panic!("server {cmd} exited before printing its address"); - } - let addr: SocketAddr = line.trim().parse().unwrap_or_else(|e| { - panic!("failed to parse server address from {cmd} ({line:?}): {e}") - }); - - std::thread::sleep(Duration::from_millis(50)); - - Self { child, addr } - } -} - -impl Drop for ServerProcess { - fn drop(&mut self) { - let _ = self.child.kill(); - let _ = self.child.wait(); - } -} - // ── Valkey container ──────────────────────────────────────────────── /// Managed valkey container: spawns via `docker run` on an ephemeral host @@ -152,6 +111,9 @@ impl Drop for ValkeyContainer { // ── Build helpers ──────────────────────────────────────────────────── fn build_connectrpc_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("fortune_server") { + return path; + } eprintln!(" Building connectrpc-rs fortune server..."); let output = Command::new("cargo") .args([ @@ -173,6 +135,9 @@ fn build_connectrpc_server() -> String { } fn build_tonic_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("fortune-server-tonic") { + return path; + } eprintln!(" Building tonic fortune server..."); let output = Command::new("cargo") .args([ @@ -196,7 +161,14 @@ fn build_tonic_server() -> String { format!("{manifest_dir}/../../target/release/fortune-server-tonic") } +fn build_tonic_protobuf_server() -> String { + rpc_bench::build_grpc_rust_bin("fortune-server-tonic-protobuf") +} + fn build_connect_go_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("fortune-connect-go") { + return path; + } eprintln!(" Building connect-go fortune server..."); let manifest_dir = env!("CARGO_MANIFEST_DIR"); let go_dir = format!("{manifest_dir}/../rpc-go"); @@ -301,7 +273,13 @@ async fn bench_server( // Warmup phase. tokio::time::sleep(warmup).await; - count.store(0, Ordering::Relaxed); + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + let warmed = count.swap(0, Ordering::Relaxed); + assert!( + warmed > 0, + "{impl_name}: no request succeeded during warmup" + ); latencies.lock().await.clear(); let measure_start = Instant::now(); @@ -409,7 +387,13 @@ async fn bench_server_multiconn( } tokio::time::sleep(warmup).await; - count.store(0, Ordering::Relaxed); + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + let warmed = count.swap(0, Ordering::Relaxed); + assert!( + warmed > 0, + "{impl_name}: no request succeeded during warmup" + ); latencies.lock().await.clear(); let measure_start = Instant::now(); @@ -497,6 +481,7 @@ async fn main() { eprintln!("\nBuilding servers..."); let connectrpc_bin = build_connectrpc_server(); let tonic_bin = build_tonic_server(); + let tonic_protobuf_bin = build_tonic_protobuf_server(); let connect_go_bin = build_connect_go_server(); // Each server + the set of protocols to hit it with. The connectrpc-rs @@ -513,12 +498,14 @@ async fn main() { vec![ ("connectrpc-rs", &connectrpc_bin, all_three), ("tonic", &tonic_bin, grpc_only), + ("tonic-protobuf", &tonic_protobuf_bin, grpc_only), ("connect-go", &connect_go_bin, connect_and_grpc), ] } else { vec![ ("connectrpc-rs", &connectrpc_bin, grpc_only), ("tonic", &tonic_bin, grpc_only), + ("tonic-protobuf", &tonic_protobuf_bin, grpc_only), ("connect-go", &connect_go_bin, grpc_only), ] }; diff --git a/benches/rpc/src/bin/log_bench.rs b/benches/rpc/src/bin/log_bench.rs index 782ba051..18be1262 100644 --- a/benches/rpc/src/bin/log_bench.rs +++ b/benches/rpc/src/bin/log_bench.rs @@ -1,4 +1,5 @@ -//! Log-ingest benchmark: connectrpc-rs (buffa views) vs tonic (prost owned). +//! Log-ingest benchmark: connectrpc-rs (buffa views) vs tonic (prost owned) +//! vs tonic-protobuf (Google `protobuf` v4 on the upb arena). //! //! Decode-heavy workload: each request carries a batch of structured log //! records (50 by default, ~15 KB encoded). The handler iterates every @@ -7,18 +8,20 @@ //! This is where the proto library cost becomes visible: //! - buffa/connectrpc-rs: zero-copy view decode, no string allocs //! - prost/tonic: fully-materialized owned types, ~10 string allocs/record +//! - upb/tonic-protobuf: eager parse into a per-message arena; strings are +//! copied into the arena rather than heap-allocated one by one //! //! Per 50-record batch that's ~450 varints + ~400 string fields decoded. -use std::io::{BufRead, BufReader}; use std::net::SocketAddr; -use std::process::{Child, Command, Stdio}; +use std::process::{Command, Stdio}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::{Duration, Instant}; use connectrpc::Protocol; use connectrpc::client::{ClientConfig, Http2Connection, SharedHttp2Connection}; +use rpc_bench::ServerProcess; use rpc_bench::connect::bench::v1::*; use rpc_bench::log_request; @@ -30,40 +33,12 @@ const QUICK_MEASUREMENT: Duration = Duration::from_secs(3); const DEFAULT_RECORDS: usize = 50; const MAX_LATENCY_SAMPLES: usize = 500_000; -// ── Server process management ──────────────────────────────────────── - -struct ServerProcess { - child: Child, - addr: SocketAddr, -} - -impl ServerProcess { - fn start(cmd: &str) -> Self { - let mut child = Command::new(cmd) - .stdout(Stdio::piped()) - .stderr(Stdio::null()) - .spawn() - .unwrap_or_else(|e| panic!("failed to start {cmd}: {e}")); - let stdout = child.stdout.take().expect("no stdout"); - let mut reader = BufReader::new(stdout); - let mut line = String::new(); - reader.read_line(&mut line).expect("read addr"); - let addr: SocketAddr = line.trim().parse().expect("parse addr"); - std::thread::sleep(Duration::from_millis(50)); - Self { child, addr } - } -} - -impl Drop for ServerProcess { - fn drop(&mut self) { - let _ = self.child.kill(); - let _ = self.child.wait(); - } -} - // ── Build helpers ──────────────────────────────────────────────────── fn build_connectrpc_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("log_server") { + return path; + } eprintln!(" Building connectrpc-rs log server..."); let output = Command::new("cargo") .args([ @@ -84,6 +59,9 @@ fn build_connectrpc_server() -> String { } fn build_tonic_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("log-server-tonic") { + return path; + } eprintln!(" Building tonic log server..."); let output = Command::new("cargo") .args([ @@ -103,7 +81,14 @@ fn build_tonic_server() -> String { format!("{manifest_dir}/../../target/release/log-server-tonic") } +fn build_tonic_protobuf_server() -> String { + rpc_bench::build_grpc_rust_bin("log-server-tonic-protobuf") +} + fn build_noutf8_server() -> String { + if let Some(path) = rpc_bench::prebuilt_bin("log_server_noutf8") { + return path; + } eprintln!(" Building connectrpc-rs log server (no-utf8)..."); let output = Command::new("cargo") .args([ @@ -190,7 +175,13 @@ async fn bench_server( } tokio::time::sleep(warmup).await; - count.store(0, Ordering::Relaxed); + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + let warmed = count.swap(0, Ordering::Relaxed); + assert!( + warmed > 0, + "{impl_name}: no request succeeded during warmup" + ); latencies.lock().await.clear(); let measure_start = Instant::now(); @@ -338,7 +329,13 @@ async fn bench_server_noutf8( } tokio::time::sleep(warmup).await; - count.store(0, Ordering::Relaxed); + // A server that rejects every call (e.g. an unimplemented method after a + // stub regen) would otherwise report 0 req/s instead of failing. + let warmed = count.swap(0, Ordering::Relaxed); + assert!( + warmed > 0, + "{impl_name}: no request succeeded during warmup" + ); latencies.lock().await.clear(); let measure_start = Instant::now(); @@ -412,13 +409,18 @@ async fn main() { let connectrpc_bin = build_connectrpc_server(); let noutf8_bin = build_noutf8_server(); let tonic_bin = build_tonic_server(); + let tonic_protobuf_bin = build_tonic_protobuf_server(); let mut results = Vec::new(); - // Benchmark utf8 servers (connectrpc-rs + tonic). - for (impl_name, bin) in [("connectrpc-rs", &connectrpc_bin), ("tonic", &tonic_bin)] { + // Benchmark utf8 servers (connectrpc-rs + tonic + tonic-protobuf). + for (impl_name, bin) in [ + ("connectrpc-rs", &connectrpc_bin), + ("tonic", &tonic_bin), + ("tonic-protobuf", &tonic_protobuf_bin), + ] { for &concurrency in CONCURRENCY_LEVELS { - let server = ServerProcess::start(bin); + let server = ServerProcess::start(bin, &[]); eprintln!( " Benchmarking {impl_name} @ concurrency={concurrency} ({n_conns} conns, {records} records)..." ); @@ -446,7 +448,7 @@ async fn main() { // Benchmark noutf8 connectrpc-rs variant (separate bench fn, different types). for &concurrency in CONCURRENCY_LEVELS { - let server = ServerProcess::start(&noutf8_bin); + let server = ServerProcess::start(&noutf8_bin, &[]); let impl_name = "connectrpc-noutf8"; eprintln!( " Benchmarking {impl_name} @ concurrency={concurrency} ({n_conns} conns, {records} records)..." diff --git a/benches/rpc/src/bin/unary_load.rs b/benches/rpc/src/bin/unary_load.rs new file mode 100644 index 00000000..d046cf3e --- /dev/null +++ b/benches/rpc/src/bin/unary_load.rs @@ -0,0 +1,88 @@ +//! Standalone closed-loop load generator for `BenchService.Unary` with the +//! `small_request()` payload — the same call `cross_impl_bench`'s +//! `unary_small_grpc` measures — for profiling a server (or this client) under +//! `perf`/`strace` without criterion in the picture. +//! +//! Usage: `unary_load [duration_secs] [concurrency] [hyper|h2]` +//! +//! `hyper` (default) uses `HttpClient::plaintext_http2_only()`, exactly as the +//! criterion bench does; `h2` uses one `Http2Connection` shared by all tasks. + +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::time::{Duration, Instant}; + +use connectrpc::Protocol; +use connectrpc::client::{ClientConfig, Http2Connection, HttpClient}; +use rpc_bench::{BenchServiceClient, small_request}; + +#[tokio::main] +async fn main() { + let args: Vec = std::env::args().collect(); + let addr = args + .get(1) + .expect("usage: unary_load [duration] [concurrency] [hyper|h2]"); + let duration = args.get(2).and_then(|s| s.parse().ok()).unwrap_or(30u64); + let concurrency: usize = args.get(3).and_then(|s| s.parse().ok()).unwrap_or(1); + let transport = args.get(4).map_or("hyper", String::as_str); + + let uri: http::Uri = format!("http://{addr}").parse().expect("valid addr"); + let config = ClientConfig::new(uri.clone()).with_protocol(Protocol::Grpc); + let request = small_request(); + + eprintln!( + "unary_load: {concurrency} task(s), {transport} transport, {duration}s against {addr}" + ); + + let running = Arc::new(AtomicBool::new(true)); + let count = Arc::new(AtomicU64::new(0)); + let mut handles = Vec::new(); + + // One call site per transport keeps the client types monomorphic. + macro_rules! spawn_workers { + ($client:expr) => { + for _ in 0..concurrency { + let client = $client.clone(); + let running = Arc::clone(&running); + let count = Arc::clone(&count); + let request = request.clone(); + handles.push(tokio::spawn(async move { + while running.load(Ordering::Relaxed) { + if client.unary(request.clone()).await.is_ok() { + count.fetch_add(1, Ordering::Relaxed); + } + } + })); + } + }; + } + + match transport { + "h2" => { + let conn = Http2Connection::connect_plaintext(uri) + .await + .expect("connect") + .shared(1024); + let client = BenchServiceClient::new(conn, config.clone()); + spawn_workers!(client); + } + _ => { + let client = BenchServiceClient::new(HttpClient::plaintext_http2_only(), config); + spawn_workers!(client); + } + } + + let start = Instant::now(); + tokio::time::sleep(Duration::from_secs(duration)).await; + running.store(false, Ordering::Relaxed); + for h in handles { + let _ = h.await; + } + let elapsed = start.elapsed().as_secs_f64(); + let total = count.load(Ordering::Relaxed); + eprintln!( + "{total} requests in {elapsed:.1}s = {:.0} req/s, {:.1} us/req at c={concurrency}", + total as f64 / elapsed, + elapsed * 1e6 / total.max(1) as f64 * concurrency as f64 + ); +} diff --git a/benches/rpc/src/lib.rs b/benches/rpc/src/lib.rs index 2c495238..0f80fbd1 100644 --- a/benches/rpc/src/lib.rs +++ b/benches/rpc/src/lib.rs @@ -160,6 +160,122 @@ pub async fn start_server() -> (SocketAddr, tokio::task::JoinHandle<()>) { (addr, handle) } +/// A benchmark server running as a child process, killed on drop. +/// +/// Every server binary in this suite binds an ephemeral loopback port and +/// prints the bound address as its first stdout line; `start` reads that +/// line to learn where to connect. +pub struct ServerProcess { + child: std::process::Child, + /// The address the server reported on startup. + pub addr: SocketAddr, +} + +impl ServerProcess { + /// Spawns `cmd args...` and blocks until it has printed its address. + /// + /// # Panics + /// + /// Panics if the process cannot be spawned, exits before printing an + /// address, or prints something that is not a socket address. + pub fn start(cmd: &str, args: &[&str]) -> Self { + use std::io::BufRead; + + let mut child = std::process::Command::new(cmd) + .args(args) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::null()) + .spawn() + .unwrap_or_else(|e| panic!("failed to start {cmd}: {e}")); + + let stdout = child.stdout.take().expect("no stdout"); + let mut line = String::new(); + let bytes_read = std::io::BufReader::new(stdout) + .read_line(&mut line) + .expect("failed to read server address"); + assert!( + bytes_read != 0, + "server {cmd} exited before printing its address" + ); + let addr: SocketAddr = line.trim().parse().unwrap_or_else(|e| { + panic!("failed to parse server address from {cmd} ({line:?}): {e}") + }); + + // Give the server a moment to be fully ready for connections. + std::thread::sleep(std::time::Duration::from_millis(50)); + + Self { child, addr } + } +} + +impl Drop for ServerProcess { + fn drop(&mut self) { + let _ = self.child.kill(); + let _ = self.child.wait(); + } +} + +/// Returns `$RPC_BENCH_BIN_DIR/` when that variable is set, so the +/// bench drivers can run from a directory of prebuilt server binaries (for +/// example on a dedicated benchmarking host with no Rust, Go or C++ +/// toolchain) instead of invoking `cargo build` / `go build` themselves. +/// +/// # Panics +/// +/// Panics if the variable is set but the binary is missing, since silently +/// falling back to a local build would defeat the point. +pub fn prebuilt_bin(bin: &str) -> Option { + let dir = std::env::var_os("RPC_BENCH_BIN_DIR").filter(|d| !d.is_empty())?; + let path = std::path::Path::new(&dir).join(bin); + assert!( + path.is_file(), + "RPC_BENCH_BIN_DIR is set but {} does not exist", + path.display() + ); + Some(path.to_string_lossy().into_owned()) +} + +/// Builds `bin` from the out-of-workspace `benches/rpc-grpc-rust` crate +/// (tonic + tonic-protobuf servers on Google's upb-kernel `protobuf` +/// runtime) in release mode and returns the binary's path. +/// +/// That crate is excluded from the workspace because its codegen needs the +/// grpc-rust toolchain (protoc 35.1 plus the C++ `protoc-gen-rust-grpc` +/// plugin, cmake-built on first use), so it has its own `target/`. When +/// `GRPC_RUST_PROTOC_DIR` names a directory of prebuilt binaries the cmake +/// build is skipped (`--no-default-features`), matching that crate's README. +/// +/// # Panics +/// +/// Panics if the build fails; the bench cannot proceed without the server. +pub fn build_grpc_rust_bin(bin: &str) -> String { + if let Some(path) = prebuilt_bin(bin) { + return path; + } + eprintln!(" Building {bin} (benches/rpc-grpc-rust)..."); + let crate_dir = concat!(env!("CARGO_MANIFEST_DIR"), "/../rpc-grpc-rust"); + let mut cmd = std::process::Command::new("cargo"); + // --target-dir pins the output location even under CARGO_TARGET_DIR or a + // `[build] target-dir` config, so the returned path is always right. + cmd.args(["build", "--release", "--bin", bin, "--message-format=short"]) + .arg("--manifest-path") + .arg(format!("{crate_dir}/Cargo.toml")) + .arg("--target-dir") + .arg(format!("{crate_dir}/target")); + if std::env::var_os("GRPC_RUST_PROTOC_DIR").is_some() { + cmd.arg("--no-default-features"); + } + let status = cmd + .stderr(std::process::Stdio::inherit()) + .status() + .unwrap_or_else(|e| panic!("failed to run cargo build for {bin}: {e}")); + assert!( + status.success(), + "failed to build {bin} (see benches/rpc-grpc-rust/README.md for toolchain setup)" + ); + format!("{crate_dir}/target/release/{bin}") +} + /// Create a client for the given address, protocol, and codec format. /// /// gRPC requires HTTP/2 (uses `HttpClient::plaintext_http2_only()`), while Connect and