git.delta.rocks / fleet / refs/commits / 78196e185cf1

difftreelog

source

remowt/crates/remowt-client/src/editor.rs5.5 KiBsourcehistory
1use std::collections::HashMap;2use std::net::SocketAddr;3use std::sync::atomic::{AtomicU64, Ordering};4use std::sync::{Arc, Mutex};56use remowt_endpoints::forward::ForwardClient;7use remowt_link_shared::editor::{EditorBackend, Error};8use remowt_link_shared::BifConfig;9use tokio::net::{TcpListener, UdpSocket, UnixListener};10use tracing::error;1112use crate::Remowt;1314pub struct SshEditor {15	pub conn: Remowt,16}17impl EditorBackend for SshEditor {18	async fn open_editor(&self, socket_path: String) -> Result<(), Error> {19		let local = std::env::temp_dir().join(format!("remowt-nvim-{}.sock", uuid::Uuid::new_v4()));20		let _ = std::fs::remove_file(&local);21		let listener = UnixListener::bind(&local).map_err(|e| Error::Failed(e.to_string()))?;2223		let conn = self.conn.clone();24		let forward = tokio::spawn(async move {25			loop {26				let Ok((mut stream, _)) = listener.accept().await else {27					break;28				};29				let conn = conn.clone();30				let remote = socket_path.clone();31				tokio::spawn(async move {32					// Rides the iroh fast tunnel when established, else an33					// ssh-forwarded unix socket.34					let (forwarded, tunnel) = match conn.bind_fast_tunnel("editor", false).await {35						Ok(v) => v,36						Err(e) => {37							error!("editor: bind tunnel failed: {e}");38							return;39						}40					};41					let fclient: ForwardClient<BifConfig> = conn.endpoints();42					match fclient.connect_unix(tunnel, remote).await {43						Ok(Ok(())) => {}44						Ok(Err(e)) => {45							error!("editor: agent connect_unix failed: {e}");46							return;47						}48						Err(e) => {49							error!("editor: connect_unix rpc failed: {e}");50							return;51						}52					}53					match forwarded.accept().await {54						Ok(mut remote) => {55							let _ = tokio::io::copy_bidirectional(&mut stream, &mut remote).await;56						}57						Err(e) => error!("editor: accept tunnel failed: {e}"),58					}59				});60			}61		});6263		let status = tokio::process::Command::new("neovide")64			.arg("--no-fork")65			.arg("--server")66			.arg(&local)67			.status()68			.await69			.map_err(|e| Error::Failed(format!("spawning neovide: {e}")));7071		forward.abort();72		let _ = std::fs::remove_file(&local);7374		match status? {75			s if s.success() => Ok(()),76			s => Err(Error::Failed(format!("neovide exited with {s}"))),77		}78	}7980	async fn expose_tcp(&self, addr: String) -> Result<u16, Error> {81		let listener = TcpListener::bind(("127.0.0.1", 0))82			.await83			.map_err(|e| Error::Failed(e.to_string()))?;84		let local = listener85			.local_addr()86			.map_err(|e| Error::Failed(e.to_string()))?87			.port();8889		let conn = self.conn.clone();90		tokio::spawn(async move {91			loop {92				let Ok((mut tcp, _)) = listener.accept().await else {93					break;94				};95				let conn = conn.clone();96				let addr = addr.clone();97				tokio::spawn(async move {98					let (forwarded, tunnel) = match conn.bind_fast_tunnel("forward", false).await {99						Ok(v) => v,100						Err(e) => {101							error!("forward: bind tunnel failed: {e}");102							return;103						}104					};105					let fclient: ForwardClient<BifConfig> = conn.endpoints();106					match fclient.connect_tcp(tunnel, addr).await {107						Ok(Ok(())) => {}108						Ok(Err(e)) => {109							error!("forward: agent connect_tcp failed: {e}");110							return;111						}112						Err(e) => {113							error!("forward: connect_tcp rpc failed: {e}");114							return;115						}116					}117					match forwarded.accept().await {118						Ok(mut stream) => {119							let _ = tokio::io::copy_bidirectional(&mut tcp, &mut stream).await;120						}121						Err(e) => error!("forward: accept tunnel failed: {e}"),122					}123				});124			}125		});126127		Ok(local)128	}129130	async fn expose_udp(&self, addr: String) -> Result<u16, Error> {131		let router = self.conn.datagram_router().ok_or_else(|| {132			Error::Failed(133				"udp forward requires the iroh fast path, which is not established".into(),134			)135		})?;136137		let fclient: ForwardClient<BifConfig> = self.conn.endpoints();138		let session = fclient139			.open_udp(addr)140			.await141			.map_err(|e| Error::Failed(format!("open_udp rpc: {e}")))?142			.map_err(|e| Error::Failed(format!("agent open_udp: {e}")))?;143144		let sock = Arc::new(145			UdpSocket::bind(("127.0.0.1", 0))146				.await147				.map_err(|e| Error::Failed(e.to_string()))?,148		);149		let local = sock150			.local_addr()151			.map_err(|e| Error::Failed(e.to_string()))?152			.port();153154		let sub_for_source: Arc<Mutex<HashMap<SocketAddr, u64>>> =155			Arc::new(Mutex::new(HashMap::new()));156		let source_for_sub: Arc<Mutex<HashMap<u64, SocketAddr>>> =157			Arc::new(Mutex::new(HashMap::new()));158		let next_sub = Arc::new(AtomicU64::new(0));159		let mut rx = router.register(session);160161		let up_sock = sock.clone();162		let up_router = router.clone();163		let down_source_for_sub = source_for_sub.clone();164		tokio::spawn(async move {165			let mut buf = vec![0u8; 65535];166			loop {167				let (n, src) = match up_sock.recv_from(&mut buf).await {168					Ok(v) => v,169					Err(_) => break,170				};171				let sub = {172					let mut by_src = sub_for_source.lock().expect("lock");173					if let Some(&sub) = by_src.get(&src) {174						sub175					} else {176						let sub = next_sub.fetch_add(1, Ordering::Relaxed);177						by_src.insert(src, sub);178						source_for_sub.lock().expect("lock").insert(sub, src);179						sub180					}181				};182				if up_router.send(session, sub, &buf[..n]).is_err() {183					break;184				}185			}186			up_router.unregister(session);187		});188189		let down_sock = sock.clone();190		tokio::spawn(async move {191			while let Some((sub, payload)) = rx.recv().await {192				let dst = down_source_for_sub.lock().expect("lock").get(&sub).copied();193				if let Some(dst) = dst {194					let _ = down_sock.send_to(&payload, dst).await;195				}196			}197		});198199		Ok(local)200	}201}