Skip to content

perf(http2): significantly improve http2 multi-core performance - #892

Merged
0x676e67 merged 1 commit into
mainfrom
build
Aug 12, 2025
Merged

perf(http2): significantly improve http2 multi-core performance#892
0x676e67 merged 1 commit into
mainfrom
build

Conversation

@0x676e67

@0x676e67 0x676e67 commented Aug 12, 2025

Copy link
Copy Markdown
Owner

from: 0x676e67/http2#62

Benchmark

//! This example runs a server that responds to any request with "Hello, world!"

use std::{convert::Infallible, error::Error, time::Duration};

use bytes::Bytes;
use http::{Request, Response, header::CONTENT_TYPE};
use http_body_util::{BodyExt, Full, combinators::BoxBody};
use hyper::{body::Incoming, service::service_fn};
use hyper_util::{
    rt::{TokioExecutor, TokioIo},
    server::conn::auto::Builder,
};
use tokio::{
    net::{TcpListener, TcpStream},
    time::Instant,
};
use url::Url;

const NUM_REQUESTS_TO_SEND: usize = 100_000;

// The actual server.
async fn server(addr: &str) -> Result<(), Box<dyn Error + Send + Sync>> {
    let listener = TcpListener::bind(addr).await?;

    loop {
        if let Ok((socket, _peer_addr)) = listener.accept().await {
            tokio::spawn(async move {
                if let Err(e) = serve(socket).await {
                    println!("  -> err={:?}", e);
                }
            });
        }
    }
}

async fn serve(stream: TcpStream) -> Result<(), Box<dyn Error + Send + Sync>> {
    let result = Builder::new(TokioExecutor::new())
        .serve_connection(TokioIo::new(stream), service_fn(handle_request))
        .await;

    if let Err(e) = result {
        eprintln!("error serving: {e}");
    }

    Ok(())
}

async fn handle_request(
    _request: Request<Incoming>,
) -> Result<Response<BoxBody<Bytes, Infallible>>, Infallible> {
    let response = Response::builder()
        .header(CONTENT_TYPE, "text/plain")
        .body(Full::new(Bytes::from("Hello, world!\n")).boxed())
        .expect("values provided to the builder should be valid");

    Ok(response)
}

fn spawn_single_thread_server(addr: &'static str) {
    println!("\n\n===============================");
    println!("Starting single-threaded server at {addr}");
    println!("===============================");
    std::thread::spawn(move || {
        let rt = tokio::runtime::Builder::new_current_thread()
            .enable_all()
            .build()
            .unwrap();
        rt.block_on(server(addr)).unwrap();
    });
    std::thread::sleep(Duration::from_millis(500));
}

fn spawn_multi_thread_server(addr: &'static str) {
    println!("\n\n===============================");
    println!("Starting multi-threaded server at {addr}");
    println!("===============================");
    std::thread::spawn(move || {
        let rt = tokio::runtime::Builder::new_multi_thread()
            .worker_threads(4)
            .enable_all()
            .build()
            .unwrap();
        rt.block_on(server(addr)).unwrap();
    });
    std::thread::sleep(Duration::from_millis(500));
}

fn run_single_thread_client<F: Future>(desc: &str, addr: &str, future: F) {
    println!("-------------------------------");
    println!("Single-threaded client: {desc} at {addr}");
    let rt = tokio::runtime::Builder::new_current_thread()
        .enable_all()
        .build()
        .unwrap();
    rt.block_on(future);
}

fn run_multi_thread_client<F: Future>(desc: &str, addr: &str, future: F) {
    println!("-------------------------------");
    println!("Multi-threaded client: {desc} at {addr}");
    let rt = tokio::runtime::Builder::new_multi_thread()
        .worker_threads(4)
        .enable_all()
        .build()
        .unwrap();
    rt.block_on(future);
}

// The benchmark
async fn wreq_send_requests(addr: &str) -> Result<(), Box<dyn Error>> {
    let client = wreq::Client::builder()
        .no_proxy()
        .http2_only()
        .build()
        .unwrap();
    let url = Url::parse(&format!("http://{addr}")).unwrap();
    let mut handles = Vec::with_capacity(NUM_REQUESTS_TO_SEND);
    for _i in 0..NUM_REQUESTS_TO_SEND {
        let url = url.clone();
        let client = client.clone();
        let task = tokio::spawn(async move {
            let instant = Instant::now();
            let mut response = client.get(url).send().await.unwrap();
            while let Ok(Some(_chunk)) = response.chunk().await {}
            instant.elapsed()
        });
        handles.push(task);
    }

    let instant = Instant::now();
    let mut result = Vec::with_capacity(NUM_REQUESTS_TO_SEND);
    for handle in handles {
        result.push(handle.await.unwrap());
    }
    let mut sum = Duration::new(0, 0);
    for r in result.iter() {
        sum = sum.checked_add(*r).unwrap();
    }

    println!("Overall: {}ms.", instant.elapsed().as_millis());
    println!("Fastest: {}ms", result.iter().min().unwrap().as_millis());
    println!("Slowest: {}ms", result.iter().max().unwrap().as_millis());
    println!(
        "Avg    : {}ms",
        sum.div_f64(NUM_REQUESTS_TO_SEND as f64).as_millis()
    );
    Ok(())
}

async fn reqwest_send_requests(addr: &str) -> Result<(), Box<dyn Error>> {
    let client = reqwest::Client::builder()
        .no_proxy()
        .http2_prior_knowledge()
        .build()
        .unwrap();
    let url = Url::parse(&format!("http://{addr}")).unwrap();
    let mut handles = Vec::with_capacity(NUM_REQUESTS_TO_SEND);
    for _i in 0..NUM_REQUESTS_TO_SEND {
        let url = url.clone();
        let client = client.clone();
        let task = tokio::spawn(async move {
            let instant = Instant::now();
            let mut response = client.get(url).send().await.unwrap();
            while let Ok(Some(_chunk)) = response.chunk().await {}
            instant.elapsed()
        });
        handles.push(task);
    }

    let instant = Instant::now();
    let mut result = Vec::with_capacity(NUM_REQUESTS_TO_SEND);
    for handle in handles {
        result.push(handle.await.unwrap());
    }
    let mut sum = Duration::new(0, 0);
    for r in result.iter() {
        sum = sum.checked_add(*r).unwrap();
    }

    println!("Overall: {}ms.", instant.elapsed().as_millis());
    println!("Fastest: {}ms", result.iter().min().unwrap().as_millis());
    println!("Slowest: {}ms", result.iter().max().unwrap().as_millis());
    println!(
        "Avg    : {}ms",
        sum.div_f64(NUM_REQUESTS_TO_SEND as f64).as_millis()
    );
    Ok(())
}

fn main() {
    println!("===============================");
    println!("Benchmarking with concurrency = {}", NUM_REQUESTS_TO_SEND);
    println!("===============================");

    let addr = "127.0.0.1:5928";
    spawn_single_thread_server(addr);
    run_single_thread_client("wreq_send_requests", addr, wreq_send_requests(addr));
    run_single_thread_client("reqwest_send_requests", addr, reqwest_send_requests(addr));

    // Single-threaded server, multi-threaded client
    run_multi_thread_client("wreq_send_requests", addr, wreq_send_requests(addr));
    run_multi_thread_client("reqwest_send_requests", addr, reqwest_send_requests(addr));

    // Multi-threaded server, single-threaded client
    let addr = "127.0.0.1:5929";
    spawn_multi_thread_server(addr);
    run_single_thread_client("wreq_send_requests", addr, wreq_send_requests(addr));
    run_single_thread_client("reqwest_send_requests", addr, reqwest_send_requests(addr));

    // Multi-threaded server, multi-threaded client
    run_multi_thread_client("wreq_send_requests", addr, wreq_send_requests(addr));
    run_multi_thread_client("reqwest_send_requests", addr, reqwest_send_requests(addr));
}
===============================
Benchmarking with concurrency = 100000
===============================


===============================
Starting single-threaded server at 127.0.0.1:5928
===============================
-------------------------------
Single-threaded client: wreq_send_requests at 127.0.0.1:5928
Overall: 843ms.
Fastest: 355ms
Slowest: 586ms
Avg    : 463ms
-------------------------------
Single-threaded client: reqwest_send_requests at 127.0.0.1:5928
Overall: 922ms.
Fastest: 341ms
Slowest: 678ms
Avg    : 507ms
-------------------------------
Multi-threaded client: wreq_send_requests at 127.0.0.1:5928
Overall: 511ms.
Fastest: 31ms
Slowest: 353ms
Avg    : 222ms
-------------------------------
Multi-threaded client: reqwest_send_requests at 127.0.0.1:5928
Overall: 674ms.
Fastest: 33ms
Slowest: 508ms
Avg    : 280ms


===============================
Starting multi-threaded server at 127.0.0.1:5929
===============================
-------------------------------
Single-threaded client: wreq_send_requests at 127.0.0.1:5929
Overall: 878ms.
Fastest: 348ms
Slowest: 631ms
Avg    : 484ms
-------------------------------
Single-threaded client: reqwest_send_requests at 127.0.0.1:5929
Overall: 984ms.
Fastest: 339ms
Slowest: 742ms
Avg    : 539ms
-------------------------------
Multi-threaded client: wreq_send_requests at 127.0.0.1:5929
Overall: 845ms.
Fastest: 57ms
Slowest: 606ms
Avg    : 354ms
-------------------------------
Multi-threaded client: reqwest_send_requests at 127.0.0.1:5929
Overall: 1015ms.
Fastest: 25ms
Slowest: 810ms
Avg    : 417ms

@0x676e67
0x676e67 marked this pull request as ready for review August 12, 2025 09:28
@0x676e67
0x676e67 merged commit 2c3f873 into main Aug 12, 2025
7 checks passed
@0x676e67
0x676e67 deleted the build branch August 12, 2025 09:36
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant