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

Rust使用interprocess的Unix Domain Socket间歇性出现缓冲区填充失败问题

问题描述

我有一个Rust程序,启动服务进程监听唯一命名的Unix Domain Socket,客户端连接后发送单请求、等待响应后断开,使用interprocess crate管理套接字。

服务进程资源差异大:部分启动耗时60秒、占用10-15GB内存,部分启动耗时<1秒、仅占数百MB。外部Deno进程根据消息队列需求启停这些Rust进程,空闲或需内存时发送SIGTERM终止。

架构在开发机(Arch)和生产服务器(Ubuntu 24.04.1 LTS主机+Debian系Docker容器denoland/deno:2.0.0)可正常运行,但运行一段时间后会间歇性出现failed to fill whole buffer错误,后续连接返回Connection refused (os error 111),疑似服务端套接字已关闭。

仅生产环境复现,多次创建关闭套接字后触发,无明确规律,本地无法复现。目前唯一线索是ss -a输出:异常套接字处于ESTAB状态,正常套接字为LISTENING状态。


Rust服务端代码

use interprocess::local_socket::{prelude::*, GenericNamespaced, ListenerOptions, Name, Stream};

pub struct IpcHandler<'a> {
    socket_print_name: String,
    socket_name: Name<'a>,
}

impl<'a> IpcHandler<'a> {
    pub fn init(socket_name: Option<String>) -> Result<Self, IpcHandlerError> {
        let socket_name = socket_name.map_or("1".to_string(), |v| {
            v.chars()
                .map(|c| if c.is_alphanumeric() { c } else { '-' })
                .collect::<String>()
        });
        let socket_print_name = if GenericNamespaced::is_supported() {
            format!("ridi-router-{}.socket", socket_name)
        } else {
            format!("/tmp/ridi-router-{}.socket", socket_name)
        };

        let socket_name = socket_print_name
            .clone()
            .to_ns_name::<GenericNamespaced>()
            .map_err(|error| IpcHandlerError::NamespaceName { error })?;

        Ok(Self {
            socket_print_name,
            socket_name,
        })
    }

    pub fn listen<T>(&self, message_handler: T) -> Result<(), IpcHandlerError>
    where
        T: Fn(RequestMessage) -> ResponseMessage + Sync + Send + Copy + 'static,
    {
        let opts = ListenerOptions::new().name(self.socket_name.clone());

        let listener = match opts.create_sync() {
            Err(e) if e.kind() == io::ErrorKind::AddrInUse => {
                return Err(IpcHandlerError::SocketAddressInUse { error: e });
            }
            x => x.map_err(|error| IpcHandlerError::CreateListener { error })?,
        };

        println!(";RIDI_ROUTER SERVER READY;"); // 通知调用进程服务已就绪

        for conn in listener.incoming() {
            rayon::spawn(move || match conn {
                Err(e) => {
                    warn!("Incoming connection failed {}", e);
                }
                Ok(conn) => {
                    trace!("received connection");
                    let req = match IpcHandler::process_request(&conn) {
                        Err(err) => {
                            warn!("error from connection {:?}", err);
                            return;
                        }
                        Ok(req) => req,
                    };
                    let resp = message_handler(req);
                    if let Err(error) = IpcHandler::process_response(&conn, &resp) {
                        warn!("error from connection {:?}", error);
                    }
                }
            });
        }

        Ok(())
    }

    fn process_request(conn: &Stream) -> Result<RequestMessage, IpcHandlerError> {
        let start = SystemTime::now();
        let req_timestamp = start
            .duration_since(UNIX_EPOCH)
            .expect("Time went backwards")
            .as_millis();

        let mut conn = BufReader::new(conn);

        let mut mes_len_buf = [0u8; 8];
        conn.read_exact(&mut mes_len_buf)
            .map_err(|error| IpcHandlerError::ReadLine { error })?;

        let mut buffer = vec![0; u64::from_ne_bytes(mes_len_buf) as usize];
        conn.read_exact(&mut buffer[..])
            .map_err(|error| IpcHandlerError::ReadLine { error })?;

        let string_message =
            std::str::from_utf8(&buffer).map_err(|error| IpcHandlerError::Utf8Message { error })?;

        let request_message: RequestMessage = serde_json::from_str(&string_message)
            .map_err(|error| IpcHandlerError::DeserializeMessage { error })?;

        Ok(request_message)
    }
    fn process_response(
        conn: &Stream,
        response_message: &ResponseMessage,
    ) -> Result<(), IpcHandlerError> {
        let mut conn = BufReader::new(conn);

        let string_message = serde_json::to_string(response_message)
            .map_err(|error| IpcHandlerError::SerializeMessage { error })?;

        let buffer = string_message.as_bytes();

        let mes_len_bytes: u64 = buffer.len() as u64;
        conn.get_mut()
            .write_all(&mes_len_bytes.to_ne_bytes()[..])
            .map_err(|error| IpcHandlerError::WriteAll { error })?;

        conn.get_mut()
            .write_all(buffer)
            .map_err(|error| IpcHandlerError::WriteLine { error })?;

        Ok(())
    }
}

Rust客户端代码

pub fn connect(
    &self,
    routing_mode: &RoutingMode,
    rules: RouterRules,
    route_req_id: Option<String>,
) -> Result<ResponseMessage, IpcHandlerError> {
    let conn = Stream::connect(self.socket_name.clone())
        .map_err(|error| IpcHandlerError::Connect { error })?;

    let mut conn = BufReader::new(conn);

    let req_msg = RequestMessage {
        id: route_req_id.map_or(String::from("default-request-id"), |v| v.to_string()),
        routing_mode: routing_mode.clone(),
        rules,
    };
    let string_req = serde_json::to_string(&req_msg)
        .map_err(|error| IpcHandlerError::SerializeMessage { error })?;

    let req_buf = string_req.as_bytes();

    let mes_len_bytes: u64 = req_buf.len() as u64;
    conn.get_mut()
        .write_all(&mes_len_bytes.to_ne_bytes()[..])
        .map_err(|error| IpcHandlerError::WriteAll { error })?;

    conn.get_mut()
        .write_all(req_buf)
        .map_err(|error| IpcHandlerError::WriteAll { error })?;

    let mut mes_len_buf = [0u8; 8];
    conn.read_exact(&mut mes_len_buf)
        .map_err(|error| IpcHandlerError::ReadLine { error })?;

    let mut resp_buf = vec![0; u64::from_ne_bytes(mes_len_buf) as usize];
    conn.read_exact(&mut resp_buf[..])
        .map_err(|error| IpcHandlerError::ReadLine { error })?;

    let string_resp = std::str::from_utf8(&resp_buf)
        .map_err(|error| IpcHandlerError::Utf8Message { error })?;

    let resp_msg: ResponseMessage = serde_json::from_str(string_resp)
        .map_err(|error| IpcHandlerError::DeserializeMessage { error })?;

    Ok(resp_msg)
}

可能原因

  1. 套接字资源泄漏:Rust进程被SIGTERM终止时未正确关闭监听套接字,导致套接字处于ESTAB状态残留,后续无法重新绑定。
  2. 进程终止资源清理不及时:未捕获SIGTERM信号,进程直接退出,监听套接字未主动关闭,系统回收延迟。
  3. 并发连接竞态条件:rayon::spawn处理连接时,Stream的生命周期管理不当,导致监听套接字意外关闭。
  4. Docker容器套接字管理问题:容器内套接字文件未被正确清理,或命名空间隔离导致残留套接字无法被新进程识别。
  5. interprocess crate潜在bug:高频创建/销毁套接字场景下,库的资源管理存在疏漏,导致套接字状态异常。

调试步骤

  1. 添加信号处理逻辑:在Rust服务中捕获SIGTERM,主动关闭监听套接字并记录终止时的套接字状态日志。
  2. 监控套接字生命周期:Deno进程发送SIGTERM后,等待Rust进程退出,再检查并删除对应套接字文件。
  3. 增强日志输出:在Rust服务的listen、连接处理及终止逻辑中添加详细日志,包括套接字名称、连接状态、错误栈信息。
  4. 生产环境临时调试:异常时执行lsof -U | grep ridi-router查看套接字关联进程,确认是否有僵尸进程;用strace跟踪Rust进程的套接字系统调用。
  5. 模拟高频启停场景:编写脚本批量启停Rust进程,模拟生产环境队列波动,尝试复现问题。

解决方案建议

  1. 强制清理套接字文件:
    • Rust服务启动前检查对应套接字文件是否存在,存在则删除;
    • Deno进程在Rust进程退出后,强制删除对应套接字文件。
  2. 正确处理SIGTERM信号:
    use std::sync::Arc;
    use signal_hook::{consts::SIGTERM, iterator::Signals};
    
    pub fn listen<T>(&self, message_handler: T) -> Result<(), IpcHandlerError>
    where
        T: Fn(RequestMessage) -> ResponseMessage + Sync + Send + Copy + 'static,
    {
        let opts = ListenerOptions::new().name(self.socket_name.clone());
        let listener = Arc::new(opts.create_sync()?);
        let listener_clone = listener.clone();
    
        // 启动信号处理线程
        std::thread::spawn(move || {
            let mut signals = Signals::new(&[SIGTERM]).unwrap();
            for _ in signals.forever() {
                info!("Received SIGTERM, shutting down listener");
                drop(listener_clone); // 关闭监听套接字
                break;
            }
        });
    
        println!(";RIDI_ROUTER SERVER READY;");
        for conn in listener.incoming() {
            // 原连接处理逻辑
        }
        Ok(())
    }
    
  3. 优化连接生命周期管理:
    • 将Stream所有权转移到rayon::spawn的子任务中,避免引用悬空;
    • 处理完请求/响应后主动关闭Stream,确保连接资源释放。
  4. 统一套接字命名方式:
    • 直接使用/tmp下的绝对路径套接字,便于管理和清理;
    • 为每个进程分配唯一的套接字名称,避免冲突。
  5. 升级interprocess版本:检查是否有最新版本修复了套接字管理相关的bug。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 10:10:56