You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Rust TCP服务器并发请求时出现Connection Reset by Peer错误排查

问题分析与解决方案

核心问题根源

你的服务器仅读取到HTTP头结束就停止读取,但客户端(reqwest)会完整发送请求头+请求体。当服务器在客户端尚未完成请求体发送时就关闭连接,客户端的后续写操作会触发connection reset by peer错误。同时,代码中大量忽略IO错误的处理方式,也会掩盖潜在问题、增加调试难度。

具体修复步骤

1. 完整读取HTTP请求(含请求体)

当前代码找到\r\n\r\n后立即停止读取,仅获取了请求头,完全忽略了请求体。需要根据请求头中的Content-Length字段,读取完整的请求内容:

// 在读取完headers后,补充解析请求体的逻辑
let request_str = String::from_utf8_lossy(&buffer);
let mut parts = request_str.split("\r\n\r\n");
let headers_part = parts.next().unwrap_or_default();
let mut body_part = parts.next().unwrap_or_default().to_string();

// 从headers中提取Content-Length,读取剩余的请求体
if let Some(content_length) = extract_content_length(headers_part) {
    let remaining = content_length.saturating_sub(body_part.len());
    if remaining > 0 {
        let mut remaining_buffer = vec![0; remaining];
        match stream.read_exact(&mut remaining_buffer).await {
            Ok(_) => {
                body_part.push_str(&String::from_utf8_lossy(&remaining_buffer));
            }
            Err(e) => {
                eprintln!("Failed to read request body: {}", e);
                return;
            }
        }
    }
}

// 辅助函数:解析headers中的Content-Length
fn extract_content_length(headers: &str) -> Option<usize> {
    headers.lines()
        .find(|line| line.starts_with("Content-Length:"))
        .and_then(|line| line.split(':').nth(1))
        .and_then(|val| val.trim().parse().ok())
}

2. 修正连接关闭逻辑

不要在未确认客户端完成数据发送前关闭连接,且只需关闭连接的写方向即可(告知客户端服务器不再发送数据),让客户端自行处理读方向的关闭:

use std::net::Shutdown;

// 替换原有的响应发送与关闭逻辑
if let Err(e) = stream.write(response.to_string().as_bytes()).await {
    eprintln!("Failed to send response: {}", e);
}
if let Err(e) = stream.flush().await {
    eprintln!("Failed to flush stream: {}", e);
}
// 仅关闭写方向,避免提前终止客户端的写操作
if let Err(e) = stream.shutdown(Shutdown::Write).await {
    eprintln!("Failed to shutdown stream write: {}", e);
}

3. 完善错误处理

原代码用let _ = ...忽略所有IO错误,这会导致问题被掩盖。修改为显式处理错误,便于调试:

// 替换原有的读取循环
loop {
    match stream.read(&mut read_buffer).await {
        Ok(n) => {
            if n == 0 { break; }
            buffer.extend_from_slice(&read_buffer[..n]);
            if buffer.windows(4).any(|window| window == b"\r\n\r\n") {
                headers_ended = true;
                break;
            }
        }
        Err(e) => {
            eprintln!("Stream read failed: {}", e);
            return;
        }
    }
}

4. 限制并发连接数(可选)

当并发请求过高时,可能耗尽系统文件描述符,导致连接异常。可以通过信号量限制最大并发连接数:

use tokio::sync::Semaphore;

pub async fn run(self) {
    let listener = TcpListener::bind(&self.address).await.expect("Failed to bind");
    println!("Server listening on {}", self.address);
    // 限制最大并发连接为100,可根据系统配置调整
    let semaphore = Arc::new(Semaphore::new(100));

    loop {
        let (mut stream, _) = match listener.accept().await {
            Ok(conn) => conn,
            Err(e) => {
                eprintln!("Accept connection failed: {}", e);
                continue;
            }
        };

        let router = Arc::clone(&self.router);
        let permit = semaphore.clone().acquire_owned().await.unwrap();

        tokio::spawn(async move {
            let _permit = permit; // 持有许可直到任务结束,自动释放
            // ... 原有请求处理逻辑
        });
    }
}

验证

修改完成后重新启动服务器,执行测试代码。此时服务器会完整读取请求内容,等待客户端完成数据发送后再关闭连接,即可避免connection reset by peer错误。

内容的提问来源于stack exchange,提问作者DevExodus

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 10:55:38