diff --git a/candid_server/src/can_reader.rs b/candid_server/src/can_reader.rs index 14d763a..a084b46 100644 --- a/candid_server/src/can_reader.rs +++ b/candid_server/src/can_reader.rs @@ -1,4 +1,3 @@ -use std::ops::Deref; use std::sync::{mpsc, Arc, Mutex}; use std::thread; @@ -25,8 +24,17 @@ impl CANReader { let frame = can_socket.lock().unwrap().read_frame().unwrap(); // Write it to each client - for client in clients.lock().unwrap().deref() { - client.send(frame.clone()); + let mut dropped_clients = Vec::new(); + for (i, client) in clients.lock().unwrap().iter().enumerate() { + client.send(frame.clone()).unwrap_or_else(|_| { + // Indicate that this client's connection has been dropped + dropped_clients.push(i); + }); + } + + // Clean up dropped clients + for i in dropped_clients { + clients.lock().unwrap().remove(i); } } }); diff --git a/candid_server/src/client_handler.rs b/candid_server/src/client_handler.rs index c73efd5..7056b00 100644 --- a/candid_server/src/client_handler.rs +++ b/candid_server/src/client_handler.rs @@ -24,15 +24,20 @@ impl ClientHandler { }; let handle = thread::spawn(move || { + let ip = handler.connection.peer_addr().unwrap(); + loop { // Get an incoming frame from the reader let frame = handler.can_reader.recv().unwrap(); // Relay to the client - handler - .connection - .write_fmt(format_args!("{:X}\n", frame)) - .unwrap(); + match handler.connection.write_fmt(format_args!("{:X}\n", frame)) { + Ok(_) => {} + Err(_) => { + println!("Connection to client {:?} dropped.", ip); + break; + } + } } });