use crate::ipc::{self, ClientMessage, ClientRole, ServerMessage}; use interprocess::local_socket::traits::{ListenerExt as _, Stream as _}; use interprocess::local_socket::{ListenerOptions, Stream}; use std::io::{BufReader, BufWriter}; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::mpsc; use std::thread; static NEXT_CLIENT_ID: AtomicUsize = AtomicUsize::new(1); pub enum IpcEvent { ClientConnected { client_id: usize, role: ClientRole, width: u16, height: u16, tx: mpsc::SyncSender, }, ClientDisconnected { client_id: usize, }, Message { client_id: usize, msg: ClientMessage, }, } pub fn start(socket_path: &str, event_tx: mpsc::Sender) { // På unix: städa bort en ev. kvarlämnad socket-fil från en tidigare körning. #[cfg(unix)] let _ = std::fs::remove_file(socket_path); let name = match ipc::socket_name(socket_path) { Ok(n) => n, Err(e) => { eprintln!("Ogiltigt socket-namn {}: {}", socket_path, e); return; } }; let listener = match ListenerOptions::new().name(name).create_sync() { Ok(l) => l, Err(e) => { eprintln!("Kunde inte binda lokal socket {}: {}", socket_path, e); return; } }; let socket_path = socket_path.to_string(); thread::spawn(move || { for stream in listener.incoming() { match stream { Ok(stream) => { let client_id = NEXT_CLIENT_ID.fetch_add(1, Ordering::SeqCst); let event_tx = event_tx.clone(); let sp = socket_path.clone(); thread::spawn(move || { handle_client(client_id, stream, event_tx, sp); }); } // Transienta accept-fel (kan hända för named pipes) — fortsätt lyssna. Err(_) => continue, } } }); } fn handle_client( client_id: usize, stream: Stream, event_tx: mpsc::Sender, socket_path: String, ) { let (recv_half, send_half) = stream.split(); let mut reader = BufReader::new(recv_half); let (tx, rx) = mpsc::sync_channel::(64); // Writer thread thread::spawn(move || { let mut writer = BufWriter::new(send_half); for msg in rx { if ipc::write_message(&mut writer, &msg).is_err() { break; } } }); // Read Hello let hello: ClientMessage = match ipc::read_message(&mut reader) { Ok(m) => m, Err(_) => return, }; let (role, width, height) = match &hello { ClientMessage::Hello { role, width, height, .. } => (role.clone(), *width, *height), _ => return, }; let _ = tx.send(ServerMessage::HelloOk { version: ipc::VERSION.to_string(), socket_path: socket_path.clone(), }); let _ = event_tx.send(IpcEvent::ClientConnected { client_id, role, width, height, tx: tx.clone(), }); loop { match ipc::read_message::<_, ClientMessage>(&mut reader) { Ok(msg) => { if event_tx.send(IpcEvent::Message { client_id, msg }).is_err() { break; } } Err(_) => break, } } let _ = event_tx.send(IpcEvent::ClientDisconnected { client_id }); }