use std::{str::FromStr, time::Duration}; use anyhow::Result; use hashbrown::HashMap; use parking_lot::RwLock; use tokio::{io::{AsyncBufReadExt, AsyncWriteExt, BufReader}, select, sync::mpsc::{self, UnboundedReceiver}, task::JoinHandle, time}; use yazi_macro::try_format; use yazi_shared::Id; use yazi_shim::cell::RoCell; use crate::{Client, ClientWriter, Payload, Peer, STATE, Stream, ember::{Ember, EmberBye, EmberHey}}; pub(super) static CLIENTS: RoCell>> = RoCell::new(); pub(super) struct Server; impl Server { pub(super) async fn make() -> Result> { CLIENTS.write().clear(); let listener = Stream::bind().await?; Ok(tokio::spawn(async move { while let Ok((stream, _)) = listener.accept().await { let (tx, mut rx) = mpsc::unbounded_channel::(); let (reader, mut writer) = tokio::io::split(stream); tokio::spawn(async move { let mut id = None; let mut lines = BufReader::new(reader).lines(); loop { select! { Some(payload) = rx.recv() => { if writer.write_all(payload.as_bytes()).await.is_err() { break; } } _ = time::sleep(Duration::from_secs(5)) => { if writer.write_u8(b'\n').await.is_err() { break; } } Ok(Some(mut line)) = lines.next_line() => { if line.starts_with("hi,") { Self::handle_hi(line, &mut id, tx.clone()); continue; } let Some(id) = id else { continue }; if line.starts_with("bye,") { Self::handle_bye(id, rx, writer).await; break; } let mut parts = line.splitn(4, ','); let Some(kind) = parts.next() else { continue }; let Some(receiver) = parts.next().and_then(|s| s.parse::().ok()) else { continue }; let Some(sender) = parts.next().and_then(|s| s.parse::().ok()) else { continue }; let clients = CLIENTS.read(); let clients: Vec<_> = if receiver == 0 { clients.values().filter(|c| c.able(kind)).collect() } else if let Some(c) = clients.get(&receiver).filter(|c| c.able(kind)) { vec![c] } else { vec![] }; if clients.is_empty() { continue; } if receiver == 0 && kind.starts_with('@') { let Some(body) = parts.next() else { continue }; if !STATE.set(kind, sender, body) { continue } } line.push('\n'); clients.into_iter().for_each(|c| _ = c.tx.send(line.clone())); } else => break } } let mut clients = CLIENTS.write(); if id.and_then(|id| clients.remove(&id)).is_some() { Self::handle_hey(&clients); } }); } })) } fn handle_hi(s: String, id: &mut Option, tx: mpsc::UnboundedSender) { let Ok(payload) = Payload::from_str(&s) else { return }; let Ember::Hi(hi) = payload.body else { return }; let mut clients = CLIENTS.write(); if let Some(old_id) = id.replace(payload.sender) { clients.remove(&old_id); } else if let Some(state) = &*STATE.read() { state.values().for_each(|s| _ = tx.send(s.clone())); } clients.insert(payload.sender, Client { id: payload.sender, tx, abilities: hi.abilities.into_iter().map(|s| s.into_owned()).collect(), }); Self::handle_hey(&clients); } fn handle_hey(clients: &HashMap) { let payload = Payload::new(EmberHey::owned( clients.values().map(|c| (c.id, Peer::new(&c.abilities))).collect(), )); if let Ok(s) = try_format!("{payload}\n") { clients.values().for_each(|c| _ = c.tx.send(s.clone())); } } async fn handle_bye(id: Id, mut rx: UnboundedReceiver, mut writer: ClientWriter) { while let Ok(payload) = rx.try_recv() { if writer.write_all(payload.as_bytes()).await.is_err() { break; } } let bye = EmberBye::borrowed().with_receiver(id).with_sender(Id::ZERO); if let Ok(s) = try_format!("{bye}") { writer.write_all(s.as_bytes()).await.ok(); writer.flush().await.ok(); } } }