diff --git a/yazi-cli/src/args.rs b/yazi-cli/src/args.rs index cb8f39bd..16e084be 100644 --- a/yazi-cli/src/args.rs +++ b/yazi-cli/src/args.rs @@ -22,6 +22,8 @@ pub(super) enum Command { PubStatic(CommandPubStatic), /// Manage packages. Pack(CommandPack), + /// Subscribe to messages from all remote instances. + Sub(CommandSub), } #[derive(clap::Args)] @@ -109,3 +111,10 @@ pub(super) struct CommandPack { #[arg(short = 'u', long)] pub(super) upgrade: bool, } + +#[derive(clap::Args)] +pub(super) struct CommandSub { + /// The kind of messages to subscribe to, separated by commas if multiple. + #[arg(index = 1)] + pub(super) kinds: String, +} diff --git a/yazi-cli/src/main.rs b/yazi-cli/src/main.rs index e7aacaf0..9fa6b84c 100644 --- a/yazi-cli/src/main.rs +++ b/yazi-cli/src/main.rs @@ -46,6 +46,13 @@ async fn main() -> anyhow::Result<()> { package::Package::add_to_config(repo).await?; } } + + Command::Sub(cmd) => { + yazi_dds::init(); + yazi_dds::Client::draw(cmd.kinds.split(',').collect()).await?; + + tokio::signal::ctrl_c().await?; + } } Ok(()) diff --git a/yazi-config/preset/yazi.toml b/yazi-config/preset/yazi.toml index 47b37cf2..db6f72df 100644 --- a/yazi-config/preset/yazi.toml +++ b/yazi-config/preset/yazi.toml @@ -29,8 +29,8 @@ ueberzug_offset = [ 0, 0, 0, 0 ] [opener] edit = [ { run = '${EDITOR:=vi} "$@"', desc = "$EDITOR", block = true, for = "unix" }, - { run = 'code "%*"', orphan = true, desc = "code", for = "windows" }, - { run = 'code -w "%*"', block = true, desc = "code (block)", for = "windows" }, + { run = 'code %*', orphan = true, desc = "code", for = "windows" }, + { run = 'code -w %*', block = true, desc = "code (block)", for = "windows" }, ] open = [ { run = 'xdg-open "$1"', desc = "Open", for = "linux" }, diff --git a/yazi-dds/src/body/hi.rs b/yazi-dds/src/body/hi.rs index c60a2370..c80e6f11 100644 --- a/yazi-dds/src/body/hi.rs +++ b/yazi-dds/src/body/hi.rs @@ -5,15 +5,17 @@ use serde::{Deserialize, Serialize}; use super::Body; +/// The client handshake #[derive(Debug, Serialize, Deserialize)] pub struct BodyHi<'a> { - pub abilities: HashSet>, + /// Specifies the kinds of events that the client can handle + pub abilities: HashSet>, pub version: String, } impl<'a> BodyHi<'a> { #[inline] - pub fn borrowed(abilities: HashSet<&'a String>) -> Body<'a> { + pub fn borrowed(abilities: HashSet<&'a str>) -> Body<'a> { Self { abilities: abilities.into_iter().map(Cow::Borrowed).collect(), version: Self::version(), diff --git a/yazi-dds/src/client.rs b/yazi-dds/src/client.rs index b9bf8d18..ea79d6d1 100644 --- a/yazi-dds/src/client.rs +++ b/yazi-dds/src/client.rs @@ -1,6 +1,6 @@ use std::{collections::{HashMap, HashSet}, mem, str::FromStr}; -use anyhow::{bail, Result}; +use anyhow::{bail, Context, Result}; use parking_lot::RwLock; use serde::{Deserialize, Serialize}; use tokio::{io::AsyncWriteExt, select, sync::mpsc, task::JoinHandle, time}; @@ -27,6 +27,7 @@ pub struct Peer { } impl Client { + /// Connect to an existing server or start a new one. pub(super) fn serve() { let mut rx = QUEUE_RX.drop(); while rx.try_recv().is_ok() {} @@ -52,7 +53,7 @@ impl Client { if line.is_empty() { continue; } else if line.starts_with("hey,") { - Self::handle_hey(line); + Self::handle_hey(&line); } else { Payload::from_str(&line).map(|p| p.emit()).ok(); } @@ -62,6 +63,7 @@ impl Client { }); } + /// Connect to an existing server to send a single message. pub async fn shot(kind: &str, receiver: u64, severity: Option, body: &str) -> Result<()> { Body::validate(kind)?; @@ -100,6 +102,40 @@ impl Client { Ok(()) } + /// Connect to an existing server and listen in on the messages that are being + /// sent by other yazi instances: + /// - If no server is running, fail right away; + /// - If a server is closed, attempt to reconnect forever. + pub async fn draw(kinds: HashSet<&str>) -> Result<()> { + async fn make(kinds: &HashSet<&str>) -> Result { + let (lines, mut writer) = Stream::connect().await?; + let hi = Payload::new(BodyHi::borrowed(kinds.clone())); + writer.write_all(format!("{hi}\n").as_bytes()).await?; + writer.flush().await?; + Ok(lines) + } + + let mut lines = make(&kinds).await.context("No running Yazi instance found")?; + loop { + match lines.next_line().await? { + Some(s) => { + let kind = s.split(',').next(); + if matches!(kind, Some(kind) if kinds.contains(kind)) { + println!("{s}"); + } + } + None => loop { + if let Ok(new) = make(&kinds).await { + lines = new; + break; + } else { + time::sleep(time::Duration::from_secs(1)).await; + } + }, + } + } + } + #[inline] pub(super) fn push<'a>(payload: impl Into>) { QUEUE_TX.send(format!("{}\n", payload.into())).ok(); @@ -136,8 +172,8 @@ impl Client { Self::connect(server).await } - fn handle_hey(s: String) { - if let Ok(Body::Hey(mut hey)) = Payload::from_str(&s).map(|p| p.body) { + fn handle_hey(s: &str) { + if let Ok(Body::Hey(mut hey)) = Payload::from_str(s).map(|p| p.body) { hey.peers.retain(|&id, _| id != *ID); *PEERS.write() = hey.peers; } diff --git a/yazi-dds/src/pubsub.rs b/yazi-dds/src/pubsub.rs index 61f7230c..4021a02b 100644 --- a/yazi-dds/src/pubsub.rs +++ b/yazi-dds/src/pubsub.rs @@ -88,7 +88,7 @@ impl Pubsub { pub fn pub_from_hi() -> bool { let abilities = REMOTE.read().keys().cloned().collect(); - let abilities = BOOT.remote_events.union(&abilities).collect(); + let abilities = BOOT.remote_events.union(&abilities).map(|s| s.as_str()).collect(); Client::push(BodyHi::borrowed(abilities)); true