From aeb14be0b61e585d263f7a38386fa684637caff3 Mon Sep 17 00:00:00 2001 From: Lucas Schumacher Date: Thu, 28 Dec 2023 16:25:52 -0500 Subject: [PATCH] Switch to async implementation --- Cargo.lock | 351 ++++++++++++++++++++++++++++++++++++++++++++++++++++ Cargo.toml | 1 + src/main.rs | 153 +++++------------------ 3 files changed, 380 insertions(+), 125 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e537e9d..257c705 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,357 @@ # It is not intended for manual editing. version = 3 +[[package]] +name = "addr2line" +version = "0.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8a30b2e23b9e17a9f90641c7ab1549cd9b44f296d3ccbf309d2863cfe398a0cb" +dependencies = [ + "gimli", +] + +[[package]] +name = "adler" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f26201604c87b1e01bd3d98f8d5d9a8fcbb815e8cedb41ffccbeb4bf593a35fe" + +[[package]] +name = "autocfg" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d468802bab17cbc0cc575e9b053f41e72aa36bfa6b7f55e3529ffa43161b97fa" + +[[package]] +name = "backtrace" +version = "0.3.69" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2089b7e3f35b9dd2d0ed921ead4f6d318c27680d4a5bd167b3ee120edb105837" +dependencies = [ + "addr2line", + "cc", + "cfg-if", + "libc", + "miniz_oxide", + "object", + "rustc-demangle", +] + +[[package]] +name = "bitflags" +version = "1.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" + +[[package]] +name = "bytes" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a2bd12c1caf447e69cd4528f47f94d203fd2582878ecb9e9465484c4148a8223" + +[[package]] +name = "cc" +version = "1.0.83" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1174fb0b6ec23863f8b971027804a42614e347eafb0a95bf0b12cdae21fc4d0" +dependencies = [ + "libc", +] + +[[package]] +name = "cfg-if" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "baf1de4339761588bc0619e3cbc0120ee582ebb74b53b4efbf79117bd2da40fd" + +[[package]] +name = "gimli" +version = "0.28.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4271d37baee1b8c7e4b708028c57d816cf9d2434acb33a549475f78c181f6253" + +[[package]] +name = "hermit-abi" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d77f7ec81a6d05a3abb01ab6eb7590f6083d08449fe5a1c8b1e620283546ccb7" + [[package]] name = "kissdummy" version = "0.1.0" +dependencies = [ + "tokio", +] + +[[package]] +name = "libc" +version = "0.2.151" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "302d7ab3130588088d277783b1e2d2e10c9e9e4a16dd9050e6ec93fb3e7048f4" + +[[package]] +name = "lock_api" +version = "0.4.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c168f8615b12bc01f9c17e2eb0cc07dcae1940121185446edc3744920e8ef45" +dependencies = [ + "autocfg", + "scopeguard", +] + +[[package]] +name = "memchr" +version = "2.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f665ee40bc4a3c5590afb1e9677db74a508659dfd71e126420da8274909a0167" + +[[package]] +name = "miniz_oxide" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e7810e0be55b428ada41041c41f32c9f1a42817901b4ccf45fa3d4b6561e74c7" +dependencies = [ + "adler", +] + +[[package]] +name = "mio" +version = "0.8.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f3d0b296e374a4e6f3c7b0a1f5a51d748a0d34c85e7dc48fc3fa9a87657fe09" +dependencies = [ + "libc", + "wasi", + "windows-sys", +] + +[[package]] +name = "num_cpus" +version = "1.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4161fcb6d602d4d2081af7c3a45852d875a03dd337a6bfdd6e06407b61342a43" +dependencies = [ + "hermit-abi", + "libc", +] + +[[package]] +name = "object" +version = "0.32.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6a622008b6e321afc04970976f62ee297fdbaa6f95318ca343e3eebb9648441" +dependencies = [ + "memchr", +] + +[[package]] +name = "parking_lot" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3742b2c103b9f06bc9fff0a37ff4912935851bee6d36f3c02bcc755bcfec228f" +dependencies = [ + "lock_api", + "parking_lot_core", +] + +[[package]] +name = "parking_lot_core" +version = "0.9.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4c42a9226546d68acdd9c0a280d17ce19bfe27a46bf68784e4066115788d008e" +dependencies = [ + "cfg-if", + "libc", + "redox_syscall", + "smallvec", + "windows-targets", +] + +[[package]] +name = "pin-project-lite" +version = "0.2.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8afb450f006bf6385ca15ef45d71d2288452bc3683ce2e2cacc0d18e4be60b58" + +[[package]] +name = "proc-macro2" +version = "1.0.71" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75cb1540fadbd5b8fbccc4dddad2734eba435053f725621c070711a14bb5f4b8" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "quote" +version = "1.0.33" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5267fca4496028628a95160fc423a33e8b2e6af8a5302579e322e4b520293cae" +dependencies = [ + "proc-macro2", +] + +[[package]] +name = "redox_syscall" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4722d768eff46b75989dd134e5c353f0d6296e5aaa3132e776cbdb56be7731aa" +dependencies = [ + "bitflags", +] + +[[package]] +name = "rustc-demangle" +version = "0.1.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d626bb9dae77e28219937af045c257c28bfd3f69333c512553507f5f9798cb76" + +[[package]] +name = "scopeguard" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" + +[[package]] +name = "signal-hook-registry" +version = "1.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d8229b473baa5980ac72ef434c4415e70c4b5e71b423043adb4ba059f89c99a1" +dependencies = [ + "libc", +] + +[[package]] +name = "smallvec" +version = "1.11.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4dccd0940a2dcdf68d092b8cbab7dc0ad8fa938bf95787e1b916b0e3d0e8e970" + +[[package]] +name = "socket2" +version = "0.5.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b5fac59a5cb5dd637972e5fca70daf0523c9067fcdc4842f053dae04a18f8e9" +dependencies = [ + "libc", + "windows-sys", +] + +[[package]] +name = "syn" +version = "2.0.43" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee659fb5f3d355364e1f3e5bc10fb82068efbf824a1e9d1c9504244a6469ad53" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + +[[package]] +name = "tokio" +version = "1.35.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c89b4efa943be685f629b149f53829423f8f5531ea21249408e8e2f8671ec104" +dependencies = [ + "backtrace", + "bytes", + "libc", + "mio", + "num_cpus", + "parking_lot", + "pin-project-lite", + "signal-hook-registry", + "socket2", + "tokio-macros", + "windows-sys", +] + +[[package]] +name = "tokio-macros" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b8a1e28f2deaa14e508979454cb3a223b10b938b45af148bc0986de36f1923b" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "unicode-ident" +version = "1.0.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3354b9ac3fae1ff6755cb6db53683adb661634f67557942dea4facebec0fee4b" + +[[package]] +name = "wasi" +version = "0.11.0+wasi-snapshot-preview1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" + +[[package]] +name = "windows-sys" +version = "0.48.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "677d2418bec65e3338edb076e806bc1ec15693c5d0104683f2efe857f61056a9" +dependencies = [ + "windows-targets", +] + +[[package]] +name = "windows-targets" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a2fa6e2155d7247be68c096456083145c183cbbbc2764150dda45a87197940c" +dependencies = [ + "windows_aarch64_gnullvm", + "windows_aarch64_msvc", + "windows_i686_gnu", + "windows_i686_msvc", + "windows_x86_64_gnu", + "windows_x86_64_gnullvm", + "windows_x86_64_msvc", +] + +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b38e32f0abccf9987a4e3079dfb67dcd799fb61361e53e2882c3cbaf0d905d8" + +[[package]] +name = "windows_aarch64_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc35310971f3b2dbbf3f0690a219f40e2d9afcf64f9ab7cc1be722937c26b4bc" + +[[package]] +name = "windows_i686_gnu" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a75915e7def60c94dcef72200b9a8e58e5091744960da64ec734a6c6e9b3743e" + +[[package]] +name = "windows_i686_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f55c233f70c4b27f66c523580f78f1004e8b5a8b659e05a4eb49d4166cca406" + +[[package]] +name = "windows_x86_64_gnu" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53d40abd2583d23e4718fddf1ebec84dbff8381c07cae67ff7768bbf19c6718e" + +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b7b52767868a23d5bab768e390dc5f5c55825b6d30b86c844ff2dc7414044cc" + +[[package]] +name = "windows_x86_64_msvc" +version = "0.48.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed94fce61571a4006852b7389a063ab983c02eb1bb37b47f8272ce92d06d9538" diff --git a/Cargo.toml b/Cargo.toml index 09098fe..8a041f1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,3 +6,4 @@ edition = "2021" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html [dependencies] +tokio = { version = "1.35.1", features = ["full"] } diff --git a/src/main.rs b/src/main.rs index 522dd90..cdf1774 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,134 +1,37 @@ -use std::io::Read; -use std::io::Write; -use std::net::TcpListener; -use std::net::TcpStream; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpListener; +use tokio::sync::broadcast; -mod message_bus { - use std::sync::mpsc::{channel, Receiver, SendError, Sender}; +#[tokio::main] +async fn main() -> std::io::Result<()> { + let (tx, _rx) = broadcast::channel(128); + let listener = TcpListener::bind("127.0.0.1:8001").await?; + loop { + let (stream, addr) = listener.accept().await?; + println!("{addr} Connected"); + let (mut rx_tcp, mut tx_tcp) = stream.into_split(); + let tx_bus = tx.clone(); + let mut rx_bus = tx_bus.subscribe(); - pub struct MsgSender(Sender); - impl MsgSender { - pub fn send(&self, data: Vec) -> Result<(), SendError> { - self.0.send(Msg::Data(data)) - } - } - //pub struct MsgReceiver(Receiver>); - pub type MsgReceiver = Receiver>; - - pub enum Msg { - Data(Vec), - Subscribe(Sender>), - } - - pub struct MessageBus { - input: Sender, - //outputs: Vec>, - } - - fn output_loop(main_rx: Receiver) { - let mut outputs: Vec>> = vec![]; - loop { - match main_rx.recv() { - Ok(Msg::Data(data)) => { - let mut i = 0; - while i < outputs.len() { - if outputs[i].send(data.clone()).is_err() { - outputs.swap_remove(i); - } else { - i += 1; - } - } - } - Ok(Msg::Subscribe(sender)) => { - outputs.push(sender); - } - Err(_e) => { - println!("Error bus died") + tokio::task::spawn(async move { + let mut buffer = vec![0_u8; 2048]; + while let Ok(len) = rx_tcp.read(&mut buffer).await { + let data = buffer[..len].to_vec(); + if tx_bus.send(data).is_err() { + break; } } - } - } - impl MessageBus { - pub fn new() -> Self { - let (input, input_rx) = channel(); + }); - std::thread::spawn(move || output_loop(input_rx)); - Self { input, /*outputs*/ } - } - - pub fn subscribe(&self) -> (MsgSender, Receiver>) { - let (tx, rx) = channel(); - let input = self.input.clone(); - - input.send(Msg::Subscribe(tx)).unwrap(); - (MsgSender(input), rx) - } - } -} - -fn main() -> std::io::Result<()> { - let bus = message_bus::MessageBus::new(); - - /* - let (tx1, rx1) = bus.subscribe(); - let (tx2, rx2) = bus.subscribe(); - - let payload = "Hello world!".to_string().as_bytes().to_vec(); - - tx1.send(payload.clone()).unwrap(); - assert_eq!(rx2.recv().unwrap(), rx1.recv().unwrap()); - - println!("sleeping 5"); - std::thread::sleep(std::time::Duration::from_secs(5)); - println!("dropping rx2 and sleeping 5"); - drop(rx2); - drop(tx2); - std::thread::sleep(std::time::Duration::from_secs(5)); - println!("sending payload and sleeping 5"); - tx1.send(payload.clone()).unwrap(); - std::thread::sleep(std::time::Duration::from_secs(5)); - println!("creating new rx2 and sleeping 5"); - let (tx2, rx2) = bus.subscribe(); - std::thread::sleep(std::time::Duration::from_secs(5)); - println!("sending payload and sleeping 5"); - tx2.send(payload.clone()).unwrap(); - std::thread::sleep(std::time::Duration::from_secs(5)); - */ - - let listener = TcpListener::bind("127.0.0.1:8001")?; - for stream in listener.incoming() { - match stream { - Ok(tx_tcp) => { - let (tx_bus, rx_bus) = bus.subscribe(); - let rx_tcp = tx_tcp.try_clone().unwrap(); - - std::thread::spawn(move || publish_thread(rx_tcp, tx_bus)); - - std::thread::spawn(move || consume_thread(tx_tcp, rx_bus)); + tokio::task::spawn(async move { + while let Ok(data) = rx_bus.recv().await { + let owned_data = data.to_vec(); + if tx_tcp.write_all(&owned_data).await.is_err() { + break; + } } - Err(e) => { - println!("Error accepting connection: {}", e); - } - } + }); } - Ok(()) -} - -fn publish_thread(mut rx_tcp: TcpStream, tx_bus: message_bus::MsgSender) { - let mut buffer = vec![0_u8; 2048]; - while let Ok(len) = rx_tcp.read(&mut buffer) { - let data = buffer[..len].to_vec(); - if tx_bus.send(data).is_err() { - break; - } - } -} - -fn consume_thread(mut tx_tcp: TcpStream, rx_bus: message_bus::MsgReceiver) { - while let Ok(data) = rx_bus.recv() { - if tx_tcp.write_all(&data).is_err() { - break; - } - } + //Ok(()) }