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) }
可能原因
- 套接字资源泄漏:Rust进程被
SIGTERM终止时未正确关闭监听套接字,导致套接字处于ESTAB状态残留,后续无法重新绑定。 - 进程终止资源清理不及时:未捕获
SIGTERM信号,进程直接退出,监听套接字未主动关闭,系统回收延迟。 - 并发连接竞态条件:
rayon::spawn处理连接时,Stream的生命周期管理不当,导致监听套接字意外关闭。 - Docker容器套接字管理问题:容器内套接字文件未被正确清理,或命名空间隔离导致残留套接字无法被新进程识别。
interprocesscrate潜在bug:高频创建/销毁套接字场景下,库的资源管理存在疏漏,导致套接字状态异常。
调试步骤
- 添加信号处理逻辑:在Rust服务中捕获
SIGTERM,主动关闭监听套接字并记录终止时的套接字状态日志。 - 监控套接字生命周期:Deno进程发送
SIGTERM后,等待Rust进程退出,再检查并删除对应套接字文件。 - 增强日志输出:在Rust服务的
listen、连接处理及终止逻辑中添加详细日志,包括套接字名称、连接状态、错误栈信息。 - 生产环境临时调试:异常时执行
lsof -U | grep ridi-router查看套接字关联进程,确认是否有僵尸进程;用strace跟踪Rust进程的套接字系统调用。 - 模拟高频启停场景:编写脚本批量启停Rust进程,模拟生产环境队列波动,尝试复现问题。
解决方案建议
- 强制清理套接字文件:
- Rust服务启动前检查对应套接字文件是否存在,存在则删除;
- Deno进程在Rust进程退出后,强制删除对应套接字文件。
- 正确处理
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(()) } - 优化连接生命周期管理:
- 将
Stream所有权转移到rayon::spawn的子任务中,避免引用悬空; - 处理完请求/响应后主动关闭
Stream,确保连接资源释放。
- 将
- 统一套接字命名方式:
- 直接使用
/tmp下的绝对路径套接字,便于管理和清理; - 为每个进程分配唯一的套接字名称,避免冲突。
- 直接使用
- 升级
interprocess版本:检查是否有最新版本修复了套接字管理相关的bug。
内容的提问来源于stack exchange,提问作者Toms Jansons
相关产品推荐
相关产品推荐

