diff --git a/yazi-dds/src/client.rs b/yazi-dds/src/client.rs index 3f0ceca3..4b5db8dc 100644 --- a/yazi-dds/src/client.rs +++ b/yazi-dds/src/client.rs @@ -81,10 +81,13 @@ impl Client { server.take().map(|h| h.abort()); *server = Server::make().await.ok(); + if server.is_some() { + super::STATE.load().await.ok(); + } + if mem::replace(&mut first, false) && server.is_some() { continue; } - time::sleep(time::Duration::from_secs(1)).await; } } diff --git a/yazi-dds/src/server.rs b/yazi-dds/src/server.rs index 2cbd035c..679f8634 100644 --- a/yazi-dds/src/server.rs +++ b/yazi-dds/src/server.rs @@ -13,6 +13,7 @@ pub(super) struct Server; impl Server { pub(super) async fn make() -> Result> { + CLIENTS.write().clear(); let listener = Self::bind().await?; Ok(tokio::spawn(async move { @@ -62,7 +63,7 @@ impl Server { if receiver == 0 && severity > 0 { let Some(body) = parts.next() else { continue }; - STATE.lock().add(format!("{}_{severity}_{kind}", Body::tab(kind, body)), &line); + STATE.add(format!("{}_{severity}_{kind}", Body::tab(kind, body)), &line); } line.push('\n'); @@ -96,6 +97,10 @@ impl Server { let mut clients = CLIENTS.write(); id.replace(hi.id).and_then(|id| clients.remove(&id)); + if let Some(ref state) = *STATE.read() { + state.values().for_each(|s| _ = tx.send(format!("{s}\n"))); + } + clients.insert(hi.id, Client { id: hi.id, tx, abilities: hi.abilities }); Self::handle_hey(&clients); } diff --git a/yazi-dds/src/state.rs b/yazi-dds/src/state.rs index ed2b3739..7d15af8e 100644 --- a/yazi-dds/src/state.rs +++ b/yazi-dds/src/state.rs @@ -1,51 +1,87 @@ -use std::{collections::HashMap, io::{BufRead, BufReader, BufWriter, Write}, mem}; +use std::{collections::HashMap, mem, ops::Deref, sync::atomic::{AtomicU64, Ordering}, time::UNIX_EPOCH}; use anyhow::Result; -use parking_lot::Mutex; +use parking_lot::RwLock; +use tokio::{fs::{self, File, OpenOptions}, io::{AsyncBufReadExt, AsyncWriteExt, BufReader, BufWriter}}; use yazi_boot::BOOT; -use yazi_shared::RoCell; +use yazi_shared::{timestamp_us, RoCell}; -use crate::{body::Body, QUEUE}; +use crate::{body::Body, CLIENTS}; -pub static STATE: RoCell> = RoCell::new(); +pub static STATE: RoCell = RoCell::new(); #[derive(Default)] pub struct State { - inner: HashMap, + inner: RwLock>>, + last: AtomicU64, +} + +impl Deref for State { + type Target = RwLock>>; + + fn deref(&self) -> &Self::Target { &self.inner } } impl State { - pub fn add(&mut self, key: String, value: &str) { self.inner.insert(key, value.to_owned()); } + pub fn add(&self, key: String, value: &str) { + if let Some(ref mut inner) = *self.inner.write() { + inner.insert(key, value.to_owned()); + self.last.store(timestamp_us(), Ordering::Relaxed); + } + } - pub fn load(&mut self) -> Result<()> { - let mut buf = BufReader::new(std::fs::File::open(BOOT.state_dir.join("state"))?); + pub async fn load(&self) -> Result<()> { + let mut buf = BufReader::new(File::open(BOOT.state_dir.join(".dds")).await?); let mut line = String::new(); - while buf.read_line(&mut line)? > 0 { + let mut inner = HashMap::new(); + while buf.read_line(&mut line).await? > 0 { let mut parts = line.splitn(4, ','); let Some(kind) = parts.next() else { continue }; let Some(_) = parts.next() else { continue }; let Some(severity) = parts.next().and_then(|s| s.parse::().ok()) else { continue }; let Some(body) = parts.next() else { continue }; - - self.inner.insert(format!("{}_{severity}_{kind}", Body::tab(kind, body)), line.clone()); - QUEUE.send(mem::take(&mut line)).ok(); + inner.insert(format!("{}_{severity}_{kind}", Body::tab(kind, body)), mem::take(&mut line)); } + + let clients = CLIENTS.read(); + for payload in inner.values() { + clients.values().for_each(|c| _ = c.tx.send(format!("{payload}\n"))); + } + + self.inner.write().replace(inner); + self.last.store(timestamp_us(), Ordering::Relaxed); Ok(()) } - pub fn drain(&mut self) -> Result<()> { + pub async fn drain(&self) -> Result<()> { + let Some(inner) = self.inner.write().take() else { return Ok(()) }; + if self.skip().await.unwrap_or(false) { + return Ok(()); + } + let mut buf = BufWriter::new( - std::fs::OpenOptions::new() + OpenOptions::new() .write(true) .create(true) .truncate(true) - .open(BOOT.state_dir.join("state"))?, + .open(BOOT.state_dir.join(".dds")) + .await?, ); - let mut state = mem::take(&mut self.inner).into_iter().collect::>(); + let mut state = inner.into_iter().collect::>(); state.sort_unstable_by(|(a, _), (b, _)| a.cmp(b)); - state.into_iter().for_each(|(_, v)| _ = writeln!(buf, "{v}")); + for (_, v) in state { + buf.write_all(v.as_bytes()).await?; + buf.write_u8(b'\n').await?; + } + Ok(()) } + + async fn skip(&self) -> Result { + let meta = fs::symlink_metadata(BOOT.state_dir.join(".dds")).await?; + let modified = meta.modified()?.duration_since(UNIX_EPOCH)?.as_micros(); + Ok(modified >= self.last.load(Ordering::Relaxed) as u128) + } } diff --git a/yazi-fm/src/app/commands/quit.rs b/yazi-fm/src/app/commands/quit.rs index 5004ea3b..655e1dd4 100644 --- a/yazi-fm/src/app/commands/quit.rs +++ b/yazi-fm/src/app/commands/quit.rs @@ -9,8 +9,8 @@ impl App { pub(crate) fn quit(&mut self, opt: EventQuit) -> ! { self.cx.tasks.shutdown(); self.cx.manager.shutdown(); + futures::executor::block_on(yazi_dds::STATE.drain()).ok(); - yazi_dds::STATE.lock().drain().ok(); if !opt.no_cwd_file { self.cwd_to_file(); }