diff --git a/Cargo.lock b/Cargo.lock index bb3f440..3a9398f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -109,7 +109,6 @@ dependencies = [ "clap", "itertools", "serde", - "tempfile", "tokio", "tracing", "tracing-log", @@ -122,15 +121,6 @@ version = "1.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f107b87b6afc2a64fd13cac55fe06d6c8859f12d4b14cbcdd2c67d0976781be" -[[package]] -name = "fastrand" -version = "1.8.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a7a407cfaa3385c4ae6b23e84623d48c2798d06e3e6a1878f7f59f17b3f86499" -dependencies = [ - "instant", -] - [[package]] name = "hashbrown" version = "0.12.3" @@ -162,15 +152,6 @@ dependencies = [ "hashbrown", ] -[[package]] -name = "instant" -version = "0.1.12" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7a5bbe824c507c5da5956355e86a746d82e0e1464f65d862cc5e71da70e94b2c" -dependencies = [ - "cfg-if", -] - [[package]] name = "itertools" version = "0.10.3" @@ -364,15 +345,6 @@ version = "0.6.27" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a3f87b73ce11b1619a3c6332f45341e0047173771e8b8b73f87bfeefb7b56244" -[[package]] -name = "remove_dir_all" -version = "0.5.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3acd125665422973a33ac9d3dd2df85edad0f4ae9b00dafb1a05e43a9f5ef8e7" -dependencies = [ - "winapi", -] - [[package]] name = "scopeguard" version = "1.1.0" @@ -450,20 +422,6 @@ dependencies = [ "unicode-ident", ] -[[package]] -name = "tempfile" -version = "3.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5cdb1ef4eaeeaddc8fbd371e5017057064af0911902ef36b39801f67cc6d79e4" -dependencies = [ - "cfg-if", - "fastrand", - "libc", - "redox_syscall", - "remove_dir_all", - "winapi", -] - [[package]] name = "termcolor" version = "1.1.3" diff --git a/Cargo.toml b/Cargo.toml index 3a48b5e..19055ee 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,5 +16,4 @@ anyhow = "1" itertools = "0.10.3" serde = { version = "1.0", features = ["derive"] } bincode = "1.3.3" -tempfile = "3" tokio = { version = "1.20.1", features = ["full"] } \ No newline at end of file diff --git a/src/client.rs b/src/client.rs index 6f109c3..19dbdae 100644 --- a/src/client.rs +++ b/src/client.rs @@ -1,8 +1,10 @@ use anyhow::Result; -use std::fs; use tokio::{ io::BufReader, - net::{UnixListener, UnixStream}, + net::{ + unix::{OwnedReadHalf, OwnedWriteHalf}, + UnixStream, + }, signal, sync::mpsc, }; @@ -24,15 +26,13 @@ impl Client { pub async fn run(&mut self, server_socket_path: &str) -> Result<()> { log::warn!("🚀 - Starting echo client"); - let mut server_stream = UnixStream::connect(server_socket_path).await?; - - let (response_socket_file, response_socket) = Client::create_response_socket()?; - let response_file_path = response_socket_file.as_ref().to_str().unwrap().to_string(); - EchoCommand::Hello(response_file_path).send(&mut server_stream).await?; + let server_stream = UnixStream::connect(server_socket_path).await?; + let (reader, mut writer) = server_stream.into_split(); + EchoCommand::Hello().send(&mut writer).await?; tokio::select! { - _ = Client::handle_server_response(response_socket) => {}, - _ = Client::handle_input(server_stream) => {}, + _ = Client::handle_server_response(reader) => {}, + _ = Client::handle_input(writer) => {}, _ = signal::ctrl_c() => {} }; log::warn!("👋 - Quitting echo client"); @@ -40,19 +40,7 @@ impl Client { Ok(()) } - fn create_response_socket() -> Result<(tempfile::NamedTempFile, UnixListener)> { - let response_socket = tempfile::NamedTempFile::new()?; - let response_socket_path = response_socket.as_ref(); - - log::info!("Response Socket: {}", response_socket_path.display()); - - let _ = fs::remove_file(&response_socket_path); - let listener = UnixListener::bind(response_socket_path)?; - // Bind and listen on the response socket - Ok((response_socket, listener)) - } - - async fn handle_input(mut stream: UnixStream) { + async fn handle_input(mut writer: OwnedWriteHalf) { log::info!("Starting Input Task"); let (shutdown_complete_tx, mut shutdown_complete_rx) = mpsc::channel::(1); @@ -78,7 +66,7 @@ impl Client { loop { match shutdown_complete_rx.recv().await { Some(line) => { - EchoCommand::Message(line).send(&mut stream).await.unwrap(); + EchoCommand::Message(line).send(&mut writer).await.unwrap(); } None => { // Empty string sent means EOF, closing time. @@ -92,11 +80,9 @@ impl Client { log::info!("⌨️ - Ending Input Task"); } - async fn handle_server_response(listener: UnixListener) -> Result<()> { + async fn handle_server_response(reader: OwnedReadHalf) -> Result<()> { log::info!("⌨️ - Starting Server Response Task"); - let stream = listener.accept().await; - let (stream, _) = stream?; - let mut reader = BufReader::new(stream); + let mut reader = BufReader::new(reader); loop { match EchoResponse::read(&mut reader).await? { diff --git a/src/message.rs b/src/message.rs index 24fbdf4..1e8b633 100644 --- a/src/message.rs +++ b/src/message.rs @@ -4,7 +4,7 @@ use tokio::io::{AsyncReadExt, AsyncWriteExt}; #[derive(Serialize, Deserialize, Debug, PartialEq, Eq)] pub enum EchoCommand { - Hello(String), + Hello(), Message(String), Goodbye(), } @@ -79,11 +79,7 @@ mod tests { #[tokio::test] async fn round_trip_message() { - for message in [ - EchoCommand::Hello("foo.socket".to_owned()), - EchoCommand::Message("Hello World".to_owned()), - EchoCommand::Goodbye(), - ] { + for message in [EchoCommand::Hello(), EchoCommand::Message("Hello World".to_owned()), EchoCommand::Goodbye()] { let buff: Cursor> = Cursor::new(vec![]); tokio::pin!(buff); message.send(&mut buff).await.unwrap(); diff --git a/src/server.rs b/src/server.rs index 5ef8752..4fcd9f7 100644 --- a/src/server.rs +++ b/src/server.rs @@ -75,12 +75,11 @@ impl Server { struct Handler { stream: UnixStream, - response_socket: Option, } impl Handler { pub fn new(stream: UnixStream) -> Self { - Handler { stream, response_socket: None } + Handler { stream } } pub async fn run(&mut self, mut shutdown: Shutdown) -> Result<()> { @@ -94,7 +93,7 @@ impl Handler { } // They may have already shut down the socket, so ignore any errors log::info!("👋 - Sending Goodbye"); - let _ = EchoResponse::Goodbye().send(self.response_socket.as_mut().unwrap()).await; + let _ = EchoResponse::Goodbye().send(&mut self.stream).await; Ok(()) } @@ -106,19 +105,16 @@ impl Handler { let mut reader = BufReader::new(read); match EchoCommand::read(&mut reader).await { - Ok(EchoCommand::Hello(response_path)) => { + Ok(EchoCommand::Hello()) => { log::info!("👋 - Connection Started"); - let response_socket = UnixStream::connect(response_path.clone()).await.unwrap(); - self.response_socket = Some(response_socket); } Ok(EchoCommand::Message(msg)) => { - EchoResponse::EchoResponse(msg).send(self.response_socket.as_mut().unwrap()).await.unwrap(); + EchoResponse::EchoResponse(msg).send(&mut self.stream).await.unwrap(); } Ok(EchoCommand::Goodbye()) => { // Close // They may have already shut down the socket, so ignore any errors - let _ = EchoResponse::Goodbye().send(self.response_socket.as_mut().unwrap()); - self.response_socket = None; + let _ = EchoResponse::Goodbye().send(&mut self.stream); log::info!("🚪 - Connection Closed"); return Ok(()); }