优雅退出 TcpListener.incoming()

WJM*_*WJM 13 tcp rust

来自 rust std net 库:

let listener = TcpListener::bind(("127.0.0.1", port)).unwrap();

info!("Opened socket on localhost port {}", port);

// accept connections and process them serially
for stream in listener.incoming() {
    break;
}

info!("closed socket");
Run Code Online (Sandbox Code Playgroud)

怎样才能让听者不再听呢?API 中表示,当侦听器被删除时,它会停止。incoming()但如果是阻塞调用,我们该如何删除它呢?最好没有像 tokio/mio 这样的外部板条箱。

euc*_*lio 13

标准库没有为此提供 API,但您可以使用一些策略来解决它:

关闭套接字上的读取

您可以使用特定于平台的 API 来关闭套接字上的读取,这将导致incoming迭代器返回错误。然后,您可以在收到错误时中断处理连接。例如,在 Unix 系统上:

use std::net::TcpListener;
use std::os::unix::io::AsRawFd;
use std::thread;

let listener = TcpListener::bind("localhost:0")?;

let fd = listener.as_raw_fd();

let handle = thread::spawn(move || {
  for connection in listener.incoming() {
    match connection {
      Ok(connection) => { /* handle connection */ }
      Err(_) => break,
  }
});

libc::shutdown(fd, libc::SHUT_RD);

handle.join();
Run Code Online (Sandbox Code Playgroud)

强迫听者醒来

另一个(跨平台)技巧是设置一个变量指示您要停止监听,然后自己连接到套接字以强制唤醒监听线程。当侦听线程唤醒时,它会检查“停止侦听”变量,如果已设置,则干净地退出。

use std::net::{TcpListener, TcpStream};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;

let listener = TcpListener::bind("localhost:0")?;
let local_addr = listener.local_addr()?;

let shutdown = Arc::new(AtomicBool::new(false));
let server_shutdown = shutdown.clone();
let handle = thread::spawn(move || {
    for connection in listener.incoming() {
        if server_shutdown.load(Ordering::Relaxed) {
            return;
        }

        match connection {
            Ok(connection) => { /* handle connection */ }
            Err(_) => break,
        }
    }
});

shutdown.store(true, Ordering::Relaxed);
let _ = TcpStream::connect(local_addr);

handle.join().unwrap();
Run Code Online (Sandbox Code Playgroud)


eff*_*ect 9

您需要使用 set_nonblocking() 方法将 TcpListener 置于非阻塞模式,如下所示:

use std::io;
use std::net::TcpListener;

let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
listener.set_nonblocking(true).expect("Cannot set non-blocking");

for stream in listener.incoming() {
    match stream {
        Ok(s) => {
            // do something with the TcpStream
            handle_connection(s);
        }
        Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
            // Decide if we should exit
            break;
            // Decide if we should try to accept a connection again
            continue;
        }
        Err(e) => panic!("encountered IO error: {}", e),
    }
}
Run Code Online (Sandbox Code Playgroud)

incoming() 调用将立即返回 Result<> 类型,而不是等待连接。如果 Result 为 Ok(),则表示已建立连接,您可以对其进行处理。如果结果是 Err(WouldBlock),这实际上并不是一个错误,只是在传入()检查套接字的那一刻没有挂起的连接。

请注意,在 WouldBlock 情况下,您可能需要在继续之前放置 sleep() 或其他内容,否则您的程序将快速轮询 incoming() 函数以检查连接,从而导致 CPU 使用率较高。

代码示例改编自这里