From 19508914cb0ef6a2fa1d89a9f54ec1a46f4db266 Mon Sep 17 00:00:00 2001 From: Josiah Glosson Date: Tue, 28 Jan 2025 16:18:52 -0600 Subject: [PATCH] Implement the socket loop Doesn't work though. Debugging time! --- apps/app-playground/src/main.rs | 2 +- packages/app-lib/src/state/friends.rs | 36 ++++++++++++++++++++++++--- packages/app-lib/src/state/tunnel.rs | 4 +-- 3 files changed, 36 insertions(+), 6 deletions(-) diff --git a/apps/app-playground/src/main.rs b/apps/app-playground/src/main.rs index 549d638cc..ff33dfb95 100644 --- a/apps/app-playground/src/main.rs +++ b/apps/app-playground/src/main.rs @@ -82,7 +82,7 @@ async fn main_client() -> theseus::Result<()> { .expect("Expected second CLI arg to be socket ID") .parse::()?; - tracing::info!("Listening on port 25565 to connect to {socket_id}"); + tracing::info!("Listening on port 25585 to connect to {socket_id}"); let tcp_stream = TcpListener::bind(SocketAddr::new("127.0.0.1".parse().unwrap(), 25585)) .await? diff --git a/packages/app-lib/src/state/friends.rs b/packages/app-lib/src/state/friends.rs index 698ec1ac2..c638849b1 100644 --- a/packages/app-lib/src/state/friends.rs +++ b/packages/app-lib/src/state/friends.rs @@ -24,7 +24,8 @@ use serde::{Deserialize, Serialize}; use std::net::SocketAddr; use std::ops::Deref; use std::sync::Arc; -use tokio::io::AsyncWriteExt; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::tcp::OwnedReadHalf; use tokio::net::TcpStream; use tokio::sync::{Mutex, RwLock}; use uuid::Uuid; @@ -180,7 +181,9 @@ impl FriendsSocket { if let Some(connected_to) = sockets.get(&to_socket) { if let InternalTunnelSocket::Listening(local_addr) = *connected_to.value().clone() { if let Ok(new_stream) = TcpStream::connect(local_addr).await { - sockets.insert(new_socket, Arc::new(InternalTunnelSocket::Connected(Mutex::new(new_stream)))); + let (read, write) = new_stream.into_split(); + sockets.insert(new_socket, Arc::new(InternalTunnelSocket::Connected(Mutex::new(write)))); + Self::socket_read_loop(write_handle.clone(), read, new_socket); continue; } } @@ -380,8 +383,9 @@ impl FriendsSocket { stream: TcpStream, ) -> crate::Result { let socket_id = Uuid::new_v4(); + let (read, write) = stream.into_split(); let socket = self.tunnel_sockets.entry(socket_id).insert(Arc::new( - InternalTunnelSocket::Connected(Mutex::new(stream)), + InternalTunnelSocket::Connected(Mutex::new(write)), )); Self::send_message( &self.write, @@ -391,6 +395,7 @@ impl FriendsSocket { }, ) .await?; + Self::socket_read_loop(self.write.clone(), read, socket_id); self.create_tunnel_socket(socket_id, socket) } @@ -407,6 +412,31 @@ impl FriendsSocket { }) } + fn socket_read_loop( + write: WriteSocket, + mut read_half: OwnedReadHalf, + socket_id: Uuid, + ) { + tokio::spawn(async move { + let mut read_buffer = [0u8; 8192]; + loop { + match read_half.read(&mut read_buffer).await { + Ok(0) | Err(_) => break, + Ok(n) => { + let _ = Self::send_message( + &write, + ClientToServerMessage::SocketSend { + socket: socket_id, + data: read_buffer[..n].to_vec(), + }, + ) + .await; + } + }; + } + }); + } + #[tracing::instrument(skip(write))] pub(super) async fn send_message( write: &WriteSocket, diff --git a/packages/app-lib/src/state/tunnel.rs b/packages/app-lib/src/state/tunnel.rs index 6f1dbee7b..696402563 100644 --- a/packages/app-lib/src/state/tunnel.rs +++ b/packages/app-lib/src/state/tunnel.rs @@ -4,13 +4,13 @@ use rust_common::networking::message::ClientToServerMessage; use std::net::SocketAddr; use std::sync::Arc; use tokio::io::AsyncWriteExt; -use tokio::net::TcpStream; +use tokio::net::tcp::OwnedWriteHalf; use tokio::sync::Mutex; use uuid::Uuid; pub(super) enum InternalTunnelSocket { Listening(SocketAddr), - Connected(Mutex), + Connected(Mutex), } pub struct TunnelSocket {