Rust基于TCP流实现Pub/Sub多读写端的问题与解决方案
TCP流默认代表本地套接字与单个远程套接字之间的一对一连接,本次需要开发一套API实现标准Pub/Sub发布订阅效果:
- 用户向单个本地写套接字写入的数据,可被多个远程读套接字同时读取
- 支持创建多个绑定到同一个远程写套接字的本地读套接字,所有读套接字都能收到写入端发送的全量数据
- 读写端在独立进程中创建时,上述逻辑必须保持一致
测试用例
以下测试用例定义了API的预期使用方式与行为:
#[test] fn single_writer_single_reader() { //reader 1 should print "hello world" create_reader(|s| {println!("reader 1 {}", s)}); let mut writer = create_writer(); writer.write_all(b"hello world").unwrap(); thread::sleep(Duration::from_millis(1000)); } #[test] fn multi_writer_single_reader() { //reader 1 should print "hello" twice create_reader(|s| {println!("reader 1 {}", s)}); let mut writer1 = create_writer(); let mut writer2 = create_writer(); writer1.write_all(b"hello").unwrap(); writer2.write_all(b"hello").unwrap(); thread::sleep(Duration::from_millis(1000)); } #[test] fn single_writer_multi_reader() { //each reader should print "hello world" create_reader(|s| {println!("reader 1 {}", s)}); create_reader(|s| {println!("reader 2 {}", s)}); let mut writer = create_writer(); writer.write_all(b"hello world").unwrap(); thread::sleep(Duration::from_millis(1000)); }
现有实现逻辑
读端实现
create_reader接口会启动新线程,运行绑定在WRITER_LISTENER_ADDR地址的服务(即writer监听器),持续接收所有传入的写端连接,为每个连接创建独立的读线程。读线程会先等待reader监听器接收当前reader的注册连接,之后持续从对应写端连接读取数据,每收到一条消息就调用传入的回调函数处理,对应实现代码如下:
pub fn create_reader(callback: fn(String)) { std::thread::spawn(move || {listen_for_writers_and_make_readers(callback)}); } fn listen_for_writers_and_make_readers(callback: fn(String)) { //listen for remote sockets that are connecting to WRITER_LISTENER_ADDR. let writer_listener = TcpListener::bind(WRITER_LISTENER_ADDR).expect("writer_listener: failed to bind"); println!("writer_listener: started on {}", WRITER_LISTENER_ADDR); loop { let (stream, _) = writer_listener.accept().unwrap(); println!("writer_listener: New writer connected"); let _reader_thread = std::thread::spawn( move || {read_loop(stream, callback)}); } } fn read_loop(stream: TcpStream, callback: fn(String)) { let mut incoming_stream = stream; //stream from the writer //wait for the reader listener to accept this reader (this is to signal the reader listener that a new reader is created) TcpStream::connect(READER_LISTENER_ADDR).expect("reader: failed to connect to reader listener"); println!("reader: connected to reader listener"); println!("reader: starting read loop"); loop { let mut buf = [0; 1024]; match incoming_stream.read(&mut buf) { Ok(0) => { println!("reader: writer disconnected"); break; }, Ok(_) => { let msg = String::from_utf8(buf.to_vec()).unwrap(); println!("reader: calling callback"); callback(msg); } Err(e) => { println!("{}", e); break; } } } }
写端实现
create_writer接口会创建一个TcpStream实例(即writer对象),等待writer监听器接收该连接后,启动新线程运行绑定在READER_LISTENER_ADDR地址的服务(即reader监听器),持续接收所有传入的读端注册连接,每接收一个连接就克隆一次writer实例。原设计假设对原始writer调用write()方法时,数据会同步到所有克隆的writer实例,进而发送给所有关联的reader端,对应实现代码如下:
fn create_writer() -> TcpStream { //wait for the writer listener to accept this writer let writer = TcpStream::connect(WRITER_LISTENER_ADDR).expect("writer: failed to connect to writer listener"); println!("writer: connected to the writer listener"); //create a thread that listens for readers and clones the writer let writer_clone = writer.try_clone().unwrap(); std::thread::spawn( move || {listen_for_readers_and_clone_writer(&writer_clone)}); return writer } fn listen_for_readers_and_clone_writer(writer: &TcpStream) { //listen for remote sockets that are connecting to READER_LISTENER_ADDR let reader_listener = TcpListener::bind(READER_LISTENER_ADDR).expect("reader_listener: failed to bind"); println!("reader_listener: started on {}", READER_LISTENER_ADDR); let mut writers = vec![]; loop { let (_stream, _) = reader_listener.accept().unwrap(); println!("reader_listener: New reader connected"); writers.push(writer.try_clone().unwrap()); } }
遇到的问题
当前single_writer_single_reader单写单读测试可正常通过,但其余两个测试均抛出地址占用错误:
- 多写单读
multi_writer_single_reader测试报错:reader_listener: failed to bind: Os { code: 10048, kind: AddrInUse, message: "Only one usage of each socket address (protocol/network address/port) is normally permitted." } - 单写多读
single_writer_multi_reader测试报错:writer_listener: failed to bind: Os { code: 10048, kind: AddrInUse, message: "Only one usage of each socket address (protocol/network address/port) is normally permitted." }
核心疑问:
- 是否支持多个
TcpListener实例绑定到同一个套接字地址? - 如果不支持,该发布订阅场景的正确实现方案是什么?
补充调整尝试
曾尝试通过给监听器分配可选地址段的方式修复地址占用问题,代码如下:
#[derive(Clone, Copy)] struct AddrRange([&'static str; 10]); static WRITER_LISTENER_ADDR_RANGE : AddrRange = AddrRange([ "127.0.0.1:80", "127.0.0.1:81", "127.0.0.1:82", "127.0.0.1:83", "127.0.0.1:84", "127.0.0.1:85", "127.0.0.1:86", "127.0.0.1:87", "127.0.0.1:88", "127.0.0.1:89" ]); static READER_LISTENER_ADDR_RANGE : AddrRange = AddrRange([ "127.0.0.1:90", "127.0.0.1:91", "127.0.0.1:92", "127.0.0.1:93", "127.0.0.1:94", "127.0.0.1:95", "127.0.0.1:96", "127.0.0.1:97", "127.0.0.1:98", "127.0.0.1:99" ]);
修改后多写单读测试可正常通过,但单写多读测试中,克隆writer实例时并不会如预期创建新的网络连接,新的疑问:如何让writer_listener感知到writer被克隆的事件,进而启动新的read_loop线程处理对应连接?
内容的提问来源于stack exchange,提问作者Blue7

