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 33 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}