diff options
Diffstat (limited to 'core/src/service')
| -rw-r--r-- | core/src/service/cli_addition/mod.rs | 2 | ||||
| -rw-r--r-- | core/src/service/cli_addition/runtime_consoles.rs | 60 | ||||
| -rw-r--r-- | core/src/service/cli_addition/utils.rs | 58 | ||||
| -rw-r--r-- | core/src/service/mod.rs | 4 | ||||
| -rw-r--r-- | core/src/service/service_runner.rs | 31 | ||||
| -rw-r--r-- | core/src/service/service_types.rs | 13 | ||||
| -rw-r--r-- | core/src/service/tcp_network/long_connection.rs | 272 | ||||
| -rw-r--r-- | core/src/service/tcp_network/mod.rs | 6 | ||||
| -rw-r--r-- | core/src/service/tcp_network/pad_client/implements.rs | 156 | ||||
| -rw-r--r-- | core/src/service/tcp_network/pad_client/mod.rs | 2 | ||||
| -rw-r--r-- | core/src/service/tcp_network/pad_client/structs.rs | 8 | ||||
| -rw-r--r-- | core/src/service/tcp_network/pad_server/implements.rs | 186 | ||||
| -rw-r--r-- | core/src/service/tcp_network/pad_server/mod.rs | 2 | ||||
| -rw-r--r-- | core/src/service/tcp_network/pad_server/structs.rs | 12 | ||||
| -rw-r--r-- | core/src/service/tcp_network/utils/mod.rs | 2 | ||||
| -rw-r--r-- | core/src/service/tcp_network/utils/stream_utils.rs | 44 | ||||
| -rw-r--r-- | core/src/service/tcp_network/utils/tokio_utils.rs | 10 |
17 files changed, 868 insertions, 0 deletions
diff --git a/core/src/service/cli_addition/mod.rs b/core/src/service/cli_addition/mod.rs new file mode 100644 index 0000000..ced68da --- /dev/null +++ b/core/src/service/cli_addition/mod.rs @@ -0,0 +1,2 @@ +pub mod runtime_consoles; +mod utils; diff --git a/core/src/service/cli_addition/runtime_consoles.rs b/core/src/service/cli_addition/runtime_consoles.rs new file mode 100644 index 0000000..76ff20c --- /dev/null +++ b/core/src/service/cli_addition/runtime_consoles.rs @@ -0,0 +1,60 @@ +use crate::service::cli_addition::utils::read_cli; +use clap::{Command, FromArgMatches}; +use std::sync::{Arc, Mutex}; +use tokio::{join, spawn}; +use crate::service::service_runner::NoGamepadsService; + +pub struct RuntimeConsole<PadService, Cmd> +where Cmd: FromArgMatches { + command: Command, + prefix: String, + service: Arc<Mutex<PadService>>, + process_command: fn(Arc<Mutex<PadService>>, Cmd), +} + +impl<PadService: Send + 'static, Cmd: FromArgMatches + 'static> RuntimeConsole<PadService, Cmd> { + pub fn build(command: Command, + prefix: String, + service: Arc<Mutex<PadService>>, + process_command: fn(Arc<Mutex<PadService>>, Cmd), + ) -> RuntimeConsole<PadService, Cmd> { + RuntimeConsole { command, prefix, service, process_command } + } + + pub fn build_entry(self) -> NoGamepadsService { + let arc = Arc::new(self); + + let entry = async move { + let console_main = spawn({ + let console = Arc::clone(&arc); + async move { + Self::console_main(console).await + } + }); + + // Join + let _ = join!(console_main); + }; + + Box::pin(entry) + } + + async fn console_main(self: Arc<RuntimeConsole<PadService, Cmd>>) { + let prefix_uppercase = self.prefix.to_uppercase(); + let prefix_lowercase = self.prefix.to_lowercase(); + + loop { + let option : Option<Cmd> = read_cli( + format!("{}> ", prefix_uppercase), + &prefix_lowercase, + self.command.clone() + ).await; + match option { + None => {} + Some(cmd) => { + (self.process_command)(Arc::clone(&self.service), cmd); + } + } + } + } +}
\ No newline at end of file diff --git a/core/src/service/cli_addition/utils.rs b/core/src/service/cli_addition/utils.rs new file mode 100644 index 0000000..05969a5 --- /dev/null +++ b/core/src/service/cli_addition/utils.rs @@ -0,0 +1,58 @@ +use clap::{Command, FromArgMatches}; +use clap::ColorChoice::Auto; +use tokio::io::{AsyncBufReadExt, AsyncWriteExt}; + +pub async fn read_cli<Cmd>(prefix: String, entry: &String, cmd: Command) -> Option<Cmd> +where Cmd: FromArgMatches { + let input: String = { + let mut buffer = String::new(); + let mut stdin = tokio::io::BufReader::new(tokio::io::stdin()); + let mut stdout = tokio::io::stdout(); + + stdout.write_all(prefix.as_bytes()).await.ok().unwrap(); + stdout.flush().await.ok().unwrap(); + + stdin.read_line(&mut buffer).await.ok().unwrap(); + buffer.trim().to_string() + }; + + process_debug_cli(entry, input, cmd).await +} + +async fn process_debug_cli<Cmd>(entry: &String, input: String, cmd: Command) -> Option<Cmd> +where Cmd: FromArgMatches { + if input.trim().is_empty() { + return None; + } + + let cmd = cmd + .color(Auto) + .help_template( + "{subcommands}{options}" + ) + .disable_help_flag(true) + .disable_version_flag(true); + + let args = shell_words::split(input.as_str()).unwrap_or_else(|_e| { + ["".to_string()].to_vec() + }); + + let full_args = std::iter::once(entry.into()).chain(args); + + match cmd.try_get_matches_from(full_args) { + Ok(matches) => { + match Cmd::from_arg_matches(&matches) { + Ok(cmd) => { + Some(cmd) + } + Err(_err) => { + None + } + } + } + Err(err) => { + println!("{}", err); + None + } + } +}
\ No newline at end of file diff --git a/core/src/service/mod.rs b/core/src/service/mod.rs new file mode 100644 index 0000000..2631af8 --- /dev/null +++ b/core/src/service/mod.rs @@ -0,0 +1,4 @@ +pub mod cli_addition; +pub mod tcp_network; +pub mod service_types; +pub mod service_runner; diff --git a/core/src/service/service_runner.rs b/core/src/service/service_runner.rs new file mode 100644 index 0000000..bd2d1d9 --- /dev/null +++ b/core/src/service/service_runner.rs @@ -0,0 +1,31 @@ +use std::pin::Pin; +use crate::service::tcp_network::utils::tokio_utils::build_tokio_runtime; +use tokio::spawn; + +#[macro_export] +#[allow(unused_macros)] +macro_rules! run_services { + ($($service:expr),+ $(,)?) => { + nogamepads_core::service::service_runner::ServiceRunner::run(Vec::from([$($service),+])); + }; +} + +pub type NoGamepadsService = Pin<Box<dyn Future<Output = ()> + Send>>; + +pub struct ServiceRunner; + +impl ServiceRunner { + pub fn run(futures: Vec<NoGamepadsService>) { + let runtime = build_tokio_runtime("nogamepads".to_string()); + runtime.block_on(async { + let mut handles = Vec::new(); + for fut in futures { + handles.push(spawn(fut)); + } + + for handle in handles { + handle.await.expect("Task panicked"); + } + }); + } +}
\ No newline at end of file diff --git a/core/src/service/service_types.rs b/core/src/service/service_types.rs new file mode 100644 index 0000000..dda18bd --- /dev/null +++ b/core/src/service/service_types.rs @@ -0,0 +1,13 @@ +use bincode::{Decode, Encode}; +use crate::data::message::traits::MessageEncoder; +use crate::encoder; + +#[derive(Default, Encode, Decode, PartialEq, Debug, Clone, Eq, Hash)] +pub enum ServiceType { + #[default] + TCPConnection, + BlueTooth, + USB, +} + +encoder!(ServiceType);
\ No newline at end of file diff --git a/core/src/service/tcp_network/long_connection.rs b/core/src/service/tcp_network/long_connection.rs new file mode 100644 index 0000000..ada9751 --- /dev/null +++ b/core/src/service/tcp_network/long_connection.rs @@ -0,0 +1,272 @@ +use crate::data::player::structs::Player; +use crate::service::tcp_network::pad_client::structs::PadClientNetwork; +use crate::service::tcp_network::pad_server::structs::PadServerNetwork; +use std::sync::Arc; +use std::sync::atomic::Ordering::SeqCst; +use log::{error, info, trace, warn}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::tcp::{OwnedReadHalf, OwnedWriteHalf}; +use tokio::net::TcpStream; +use tokio::spawn; +use nogamepads::entry_mutex; +use crate::data::message::enums::{ControlMessage, GameMessage}; +use crate::data::message::enums::ExitReason::GameOver; +use crate::data::message::enums::GameMessage::{End, LetExit}; +use crate::data::message::traits::{MessageEncoder, MessageManager}; +use crate::service::service_types::ServiceType; +use crate::service::service_types::ServiceType::TCPConnection; + +impl PadServerNetwork { + + pub async fn start_long_connection(self: Arc<Self>, player: Player, stream: TcpStream) { + let (reader, writer) = stream.into_split(); + spawn(Self::read_task(Arc::clone(&self), player.clone(), reader)); + spawn(Self::write_task(Arc::clone(&self), player.clone(), writer)); + } + + async fn read_task(self: Arc<Self>, player: Player, mut reader: OwnedReadHalf) { + info!("[TCP Server] [Runtime] Reader started."); + entry_mutex!(self.runtime, |guard| { + guard.reader_count += 1; + }); + + let mut buffer = [0u8; 1024]; + let mut err_message_counter = 0; + + loop { + let read = reader.read(&mut buffer).await; + match read { + Ok(size) => { + let message : ControlMessage = ControlMessage::de(buffer[0..size].to_vec()); + + // Preprocess messages: handle exit messages. + match message { + ControlMessage::Exit => { + info!("[TCP Server] [Runtime] Player {} exited.", player.account.id); + entry_mutex!(self.runtime, |guard| { + guard.send((player.account.clone(), End), player.account.clone(), ServiceType::TCPConnection); + }); + break; + } + + ControlMessage::Err => { + info!("[TCP Server] [Runtime] Received error message from {}.", player.account.id); + if err_message_counter < 16 { + err_message_counter += 1; + } else { + warn!("[TCP Server] [Runtime] Too many error messages! Connection closed."); + entry_mutex!(self.runtime, |guard| { + guard.send((player.account.clone(), End), player.account.clone(), ServiceType::TCPConnection); + }); + break; + } + } + + _ => {} + } + + // Process messages + entry_mutex!(self.runtime, |guard| { + trace!("[TCP Server] [Runtime] Received: {:?}", &message); + guard.put_into_receive_list((player.account.clone(), message), player.account.clone(), ServiceType::TCPConnection); + }); + } + + Err(error) => { + warn!("[TCP Server] [Runtime] Error reading from socket: {:?}", error); + break; + } + } + + // Check close + entry_mutex!(self.runtime, |guard| { + if guard.data.close.load(SeqCst) { + break; + } + }); + } + + info!("[TCP Server] [Runtime] Reader between {} closed.", player.account.id); + entry_mutex!(self.runtime, |guard| { + guard.send((player.account.clone(), End), player.account.clone(), TCPConnection); + guard.reader_count -= 1; + }) + } + + async fn write_task(self: Arc<Self>, player: Player, mut writer: OwnedWriteHalf) { + info!("[TCP Server] [Runtime] Writer started."); + entry_mutex!(self.runtime, |guard| { + guard.writer_count += 1; + }); + + let mut closed = false; + + loop { + // Check close + entry_mutex!(self.runtime, |guard| { + if guard.data.close.load(SeqCst) && !closed { + guard.send((player.account.clone(), LetExit(GameOver)), player.account.clone(), ServiceType::TCPConnection); + closed = true; + } + }); + + let mut message = None; + entry_mutex!(self.runtime, |guard| { + message = guard.pop_from_send_list(player.account.clone(), ServiceType::TCPConnection); + }); + + if let Some(message) = message { + + // Preprocess messages: handle end messages. + match message.1 { + End => { + break; + } + GameMessage::Err => { + continue; + } + _ => {} + } + + // Process messages + match writer.write_all(GameMessage::en(&message.1).as_slice()).await { + Ok(_) => { + trace!("[TCP Server] [Runtime] Sent {:?} to {}", &message.1, player.account.id); + writer.flush().await.expect("[TCP Client] [Runtime] Writer encountered an error"); + } + Err(error) => { + warn!("[TCP Server] [Runtime] Sent {:?} to {} failed: {}", &message.1, player.account.id, error); + break; + } + } + } + } + + info!("[TCP Server] [Runtime] Writer between {} closed.", player.account.id); + entry_mutex!(self.runtime, |guard| { + guard.data.sign_player_online_status(&player, TCPConnection, false); + guard.writer_count -= 1; + }) + } +} + +impl PadClientNetwork { + + pub async fn start_long_connection(self: Arc<Self>, stream: TcpStream) { + let (reader, writer) = stream.into_split(); + spawn(Self::read_task(Arc::clone(&self), reader)); + spawn(Self::write_task(Arc::clone(&self), writer)); + } + + async fn read_task(self: Arc<Self>, mut reader: OwnedReadHalf) { + info!("[TCP Client] [Runtime] Reader started."); + + let mut buffer = [0u8; 1024]; + let mut err_message_counter = 0; + loop { + // Check close + entry_mutex!(self.runtime, |guard| { + if guard.close.load(SeqCst) { + break; + } + }); + + let read = reader.read(&mut buffer).await; + match read { + Ok(size) => { + let message : GameMessage = GameMessage::de(buffer[0..size].to_vec()); + + // Preprocess messages: handle exit messages. + match message { + LetExit(reason) => { + info!("[TCP Client] [Runtime] Server let you exit: {:?}", reason); + entry_mutex!(self.runtime, |guard| { + guard.send(ControlMessage::End, 0, ServiceType::TCPConnection); + }); + break; + } + + GameMessage::Err => { + info!("[TCP Client] [Runtime] Received error message from server."); + if err_message_counter < 16 { + err_message_counter += 1; + } else { + warn!("[TCP Client] [Runtime] Too many error messages! Connection closed."); + entry_mutex!(self.runtime, |guard| { + guard.send(ControlMessage::End, 0, ServiceType::TCPConnection); + }); + break; + } + } + + _ => {} + } + + // Process messages + entry_mutex!(self.runtime, |guard| { + trace!("[TCP Client] [Runtime] Received: {:?}", &message); + guard.put_into_receive_list(message, 0, ServiceType::TCPConnection); + }); + } + Err(err) => { + error!("[TCP Client] [Runtime] Reader encountered an error: {}", err); + break; + } + } + } + + info!("[TCP Client] [Runtime] Reader closed."); + } + + async fn write_task(self: Arc<Self>, mut writer: OwnedWriteHalf) { + info!("[TCP Client] [Runtime] Writer started."); + + let mut closed = false; + + loop { + // Check close + if !closed { + entry_mutex!(self.runtime, |guard| { + if guard.close.load(SeqCst) { + guard.send(ControlMessage::Exit, 0, ServiceType::TCPConnection); + closed = true; + } + }); + } + + let mut message = None; + entry_mutex!(self.runtime, |guard| { + message = guard.pop_from_send_list(0, ServiceType::TCPConnection); + }); + + if let Some(message) = message { + + // Preprocess messages: handle exit messages. + match message { + ControlMessage::End => { break; } + ControlMessage::Err => { + continue; + } + _ => {} + } + + // Process messages + match writer.write_all(ControlMessage::en(&message).as_slice()).await { + Ok(_) => { + trace!("[TCP Client] [Runtime] Sent {:?}.", &message); + writer.flush().await.expect("[TCP Client] [Runtime] Writer encountered an error"); + } + Err(error) => { + warn!("[TCP Client] [Runtime] Sent {:?} failed: {}", &message, error); + break; + } + } + } + } + + info!("[TCP Client] [Runtime] Writer closed."); + entry_mutex!(self.runtime, |guard| { + guard.close(); + }) + } +}
\ No newline at end of file diff --git a/core/src/service/tcp_network/mod.rs b/core/src/service/tcp_network/mod.rs new file mode 100644 index 0000000..4954c1c --- /dev/null +++ b/core/src/service/tcp_network/mod.rs @@ -0,0 +1,6 @@ +pub mod utils; +pub mod pad_client; +pub mod pad_server; +pub mod long_connection; + +pub const DEFAULT_PORT : u16 = 5989;
\ No newline at end of file diff --git a/core/src/service/tcp_network/pad_client/implements.rs b/core/src/service/tcp_network/pad_client/implements.rs new file mode 100644 index 0000000..584bc83 --- /dev/null +++ b/core/src/service/tcp_network/pad_client/implements.rs @@ -0,0 +1,156 @@ +use std::net::{IpAddr, SocketAddr}; +use std::sync::{Arc, Mutex}; +use std::sync::atomic::Ordering::SeqCst; +use std::time::Duration; +use log::{error, info, warn}; +use tokio::{join, spawn}; +use tokio::time::sleep; +use nogamepads::entry_mutex; +use crate::data::controller::runtime::structs::ControllerRuntime; +use crate::data::message::enums::ConnectionMessage::{Join, RequestGameInfos}; +use crate::data::message::enums::ConnectionResponseMessage; +use crate::service::service_runner::NoGamepadsService; +use crate::service::tcp_network::pad_client::structs::PadClientNetwork; +use crate::service::tcp_network::DEFAULT_PORT; +use crate::service::tcp_network::utils::stream_utils::{read_msg, send_msg}; +use crate::service::tcp_network::utils::tokio_utils::build_tokio_runtime; + +macro_rules! connect_once { + ($addr:expr, |$conn:ident| $code:block) => {{ + use tokio::net::TcpStream; + match TcpStream::connect($addr).await { + Ok(mut $conn) => { + $code + true + }, + Err(e) => { + error!("[TCP Client] [Main] Connection failed {:?}", e); + false + } + } + }} +} + +impl PadClientNetwork { + + pub fn build(runtime: Arc<Mutex<ControllerRuntime>>) -> PadClientNetwork { + PadClientNetwork { + addr: SocketAddr::from(([127, 0, 0, 1], DEFAULT_PORT)), + runtime + } + } + + pub fn bind_addr(&mut self, addr: SocketAddr) -> &mut PadClientNetwork { + self.addr = addr; + self + } + + pub fn bind_ip(&mut self, addr: IpAddr) -> &mut PadClientNetwork { + self.addr.set_ip(addr); + self + } + + pub fn bind_port(&mut self, port: u16) -> &mut PadClientNetwork { + self.addr.set_port(port); + self + } + + pub fn build_entry(self) -> NoGamepadsService { + let arc = Arc::new(self); + + let entry = async move { + // Connection thread: Download the relevant resources, verify connection eligibility, and attempt to join the game. + let connection_thread = spawn({ + let client = Arc::clone(&arc); + async move { + Self::connection_thread(client).await + } + }); + + // Join + let _ = join!(connection_thread); + }; + + Box::pin(entry) + } + + pub fn connect(self) { + let runtime = build_tokio_runtime("padclient_tcp".to_string()); + + info!("[TCP Client] Connecting to {}:{}", self.addr.ip().to_string(), self.addr.port()); + runtime.block_on(self.build_entry()); + } +} + +impl PadClientNetwork { + + async fn connection_thread(self: Arc<PadClientNetwork>) { + let mut buffer = [0; 1024]; + + // Requests game infos + if !connect_once!(self.addr, |stream| { + info!("[TCP Client] [Main] Requesting game infos."); + send_msg(&mut stream, RequestGameInfos).await; + let response : ConnectionResponseMessage = read_msg(&mut buffer, &mut stream).await; + match response { + ConnectionResponseMessage::GameInfos(infos) => { + entry_mutex!(self.runtime, |guard| { + guard.game_info = infos; + }); + info!("[TCP Client] [Main] Download game infos successfully."); + } + ConnectionResponseMessage::Err => { + warn!("[TCP Client] [Main] Download game infos failed."); + } + _ => { + warn!("[TCP Client] [Main] Not found game infos."); + } + } + }) { + return; + } + + // TODO :: Download game layouts + + // TODO :: Download skin assets + + // Try to join game + let _ = connect_once!(self.addr, |connection| { + let mut player = None; + entry_mutex!(self.runtime, |guard| { + player = Some(guard.player.clone()); + }); + if player.is_some() { + info!("[TCP Client] [Main] Trying to join game."); + send_msg(&mut connection, Join(player.unwrap())).await; + let response : ConnectionResponseMessage = read_msg(&mut buffer, &mut connection).await; + match response { + ConnectionResponseMessage::Welcome => { + + // Long Connection + info!("[TCP Client] [Main] Welcome"); + spawn(Self::start_long_connection(Arc::clone(&self), connection)); + } + ConnectionResponseMessage::Deny(why) => { + error!("[TCP Client] [Main] Connection denied: {:?}", why); + } + _ => { } + } + } else { + error!("[TCP Client] [Main] No player found."); + return; + } + }); + + loop { + sleep(Duration::from_millis(1000)).await; + entry_mutex!(self.runtime, |guard| { + if guard.close.load(SeqCst) { + break; + } + }) + } + + info!("[TCP Client] [Main] Main thread closed."); + } +}
\ No newline at end of file diff --git a/core/src/service/tcp_network/pad_client/mod.rs b/core/src/service/tcp_network/pad_client/mod.rs new file mode 100644 index 0000000..0ff870f --- /dev/null +++ b/core/src/service/tcp_network/pad_client/mod.rs @@ -0,0 +1,2 @@ +pub mod implements; +pub mod structs;
\ No newline at end of file diff --git a/core/src/service/tcp_network/pad_client/structs.rs b/core/src/service/tcp_network/pad_client/structs.rs new file mode 100644 index 0000000..10c93c4 --- /dev/null +++ b/core/src/service/tcp_network/pad_client/structs.rs @@ -0,0 +1,8 @@ +use std::net::SocketAddr; +use std::sync::{Arc, Mutex}; +use crate::data::controller::runtime::structs::ControllerRuntime; + +pub struct PadClientNetwork { + pub(crate) addr: SocketAddr, + pub(crate) runtime: Arc<Mutex<ControllerRuntime>> +}
\ No newline at end of file diff --git a/core/src/service/tcp_network/pad_server/implements.rs b/core/src/service/tcp_network/pad_server/implements.rs new file mode 100644 index 0000000..c4d8cfc --- /dev/null +++ b/core/src/service/tcp_network/pad_server/implements.rs @@ -0,0 +1,186 @@ +use std::net::{IpAddr, SocketAddr}; +use std::sync::{Arc, Mutex}; +use std::sync::atomic::Ordering::SeqCst; +use std::time::Duration; +use log::{error, info, trace, warn}; +use tokio::{join, select, spawn}; +use tokio::net::{TcpListener, TcpStream}; +use tokio::sync::watch::{channel, Receiver, Sender}; +use tokio::time::sleep; +use nogamepads::entry_mutex; +use crate::data::game::runtime::structs::GameRuntime; +use crate::data::message::enums::ConnectionMessage; +use crate::data::message::enums::ConnectionMessage::{Join, RequestGameInfos, RequestLayoutConfigure, RequestSkinPackage, Ready}; +use crate::data::message::enums::ConnectionResponseMessage::{Deny, GameInfos, Welcome}; +use crate::service::service_runner::NoGamepadsService; +use crate::service::tcp_network::DEFAULT_PORT; +use crate::service::tcp_network::pad_server::structs::PadServerNetwork; +use crate::service::tcp_network::utils::stream_utils::{get_target_address, read_msg, send_msg}; +use crate::service::tcp_network::utils::tokio_utils::build_tokio_runtime; + +impl PadServerNetwork { + + pub fn build(runtime: Arc<Mutex<GameRuntime>>) -> PadServerNetwork { + let (close_tx, close_rx) = channel(false); + PadServerNetwork { + addr: SocketAddr::from(([127, 0, 0, 1], DEFAULT_PORT)), + runtime, + close_tx, + close_rx + } + } + + pub fn bind_ip(&mut self, ip: IpAddr) -> &mut PadServerNetwork { + self.addr.set_ip(ip); + self + } + + pub fn bind_port(&mut self, port: u16) -> &mut PadServerNetwork { + self.addr.set_port(port); + self + } + + pub fn build_entry(self) -> NoGamepadsService { + let arc = Arc::new(self); + + let entry = async move { + // Main thread: Used to handle connection requests, data requests, and transfer skin assets + let main_thread = spawn({ + let server = Arc::clone(&arc); + async move { + Self::main_thread(server).await + } + }); + + let close_checker = { + let server = Arc::clone(&arc); + async move { + Self::close_checker(server).await + } + }; + + // Join + let _ = join!(close_checker, main_thread); + }; + + Box::pin(entry) + } + + pub fn listening_block_on(self) { + let runtime = build_tokio_runtime("padserver_tcp".to_string()); + + info!("[TCP Server] Server start."); + runtime.block_on(self.build_entry()); + info!("[TCP Server] Finished."); + } +} + +impl PadServerNetwork { + + async fn main_thread(self: Arc<PadServerNetwork>) { + + info!("[TCP Server] [Main] Server listening at {}", self.addr.to_string()); + + let listener = TcpListener::bind(self.addr).await; + if listener.is_err() { + error!("[TCP Server] [Main] Failed to bind to {}", self.addr.to_string()); + return; + } + + let listener = listener.unwrap(); + info!("[TCP Server] [Main] Listener created, start listening."); + + let mut local_close_rx = self.close_rx.clone(); + + loop { + select! { + _ = local_close_rx.changed() => { + if *local_close_rx.borrow() { + break; + } + } + + accept = listener.accept() => { + match accept { + Ok((stream, _)) => { + spawn(Self::process_connection(Arc::clone(&self), stream)); + } + Err(error) => { + warn!("[TCP Server] [Main] Failed to accept TCP connections: {}", error); + } + } + } + } + } + + info!("[TCP Server] [Main] Main thread closed."); + } + + async fn process_connection(self: Arc<Self>, mut stream: TcpStream) { + let mut buffer = [0; 1024]; + let message: ConnectionMessage = read_msg(&mut buffer, &mut stream).await; + let from_address = get_target_address(&stream); + + match message { + + Join(player) => { + trace!("[TCP Server] [Main] Trying to join Player \"{}\"", &player.account.id); + let mut result = Ok(()); + entry_mutex!(self.runtime, |guard| { + match guard.try_join_player(player.clone()) { + Ok(_) => { result = Ok(()); } + Err(why) => { result = Err(why); } + } + }); + if result.is_err() { + let fail_message = result.unwrap_err(); + error!("[TCP Server] [Main] Player join failed: {:?}", &fail_message); + send_msg(&mut stream, Deny(fail_message)).await; + } else { + + // Long Connection + info!("[TCP Server] [Main] Player joined, begin long connection."); + send_msg(&mut stream, Welcome).await; + spawn(Self::start_long_connection(Arc::clone(&self), player, stream)); + } + } + + RequestGameInfos => { + info!("[TCP Server] [Main] Client({}) requests game infos.", from_address); + let mut info = Default::default(); + entry_mutex!(self.runtime, |guard| { + info = guard.info.clone(); + }); + send_msg(&mut stream, GameInfos(info)).await; + info!("[TCP Server] [Main] Game infos sent."); + } + + RequestLayoutConfigure => { + info!("[TCP Server] [Main] Client({}) requests layout configures.", from_address); + } + + RequestSkinPackage => { + info!("[TCP Server] [Main] Client({}) requests to download skin package.", from_address); + } + + Ready => { + info!("[TCP Server] [Main] Client({}) is ready!", from_address); + warn!("[TCP Server] [Main] But I don't know who he is.....") + } + + _ => { } + } + } + + async fn close_checker(self: Arc<Self>) { + loop { + sleep(Duration::from_millis(1000)).await; + entry_mutex!(self.runtime, |guard| { + if guard.data.close.load(SeqCst) { + let _ = self.close_tx.send(true); + break; + } + }) + } + } +}
\ No newline at end of file diff --git a/core/src/service/tcp_network/pad_server/mod.rs b/core/src/service/tcp_network/pad_server/mod.rs new file mode 100644 index 0000000..b37b320 --- /dev/null +++ b/core/src/service/tcp_network/pad_server/mod.rs @@ -0,0 +1,2 @@ +pub mod implements; +pub mod structs; diff --git a/core/src/service/tcp_network/pad_server/structs.rs b/core/src/service/tcp_network/pad_server/structs.rs new file mode 100644 index 0000000..f46608a --- /dev/null +++ b/core/src/service/tcp_network/pad_server/structs.rs @@ -0,0 +1,12 @@ +use crate::data::game::runtime::structs::GameRuntime; +use std::net::SocketAddr; +use std::sync::{Arc, Mutex}; +use tokio::sync::watch::{Receiver, Sender}; + +pub struct PadServerNetwork { + pub(crate) addr: SocketAddr, + pub(crate) runtime: Arc<Mutex<GameRuntime>>, + + pub(crate) close_tx: Sender<bool>, + pub(crate) close_rx: Receiver<bool>, +}
\ No newline at end of file diff --git a/core/src/service/tcp_network/utils/mod.rs b/core/src/service/tcp_network/utils/mod.rs new file mode 100644 index 0000000..ad06dac --- /dev/null +++ b/core/src/service/tcp_network/utils/mod.rs @@ -0,0 +1,2 @@ +pub mod stream_utils; +pub mod tokio_utils;
\ No newline at end of file diff --git a/core/src/service/tcp_network/utils/stream_utils.rs b/core/src/service/tcp_network/utils/stream_utils.rs new file mode 100644 index 0000000..9ccff3b --- /dev/null +++ b/core/src/service/tcp_network/utils/stream_utils.rs @@ -0,0 +1,44 @@ +use std::fmt::Debug; +use bincode::{Decode, Encode}; +use log::{error, trace}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpStream; +use crate::data::message::traits::MessageEncoder; + +pub async fn send_msg<Message>( + stream: &mut TcpStream, + msg: impl MessageEncoder<Message> + Encode + Decode<()> + Default + Debug +) +where Message: MessageEncoder<Message> + Encode + Decode<()> + Default + Debug { + match stream.write_all(MessageEncoder::en(&msg).as_slice()).await { + Ok(_) => { trace!("[Message Sender] Sent {:?} to {}", msg, get_target_address(stream)); } + Err(err) => { error!("[Message Sender] Failed to send message: {}", err); } + } +} + +pub async fn read_msg<Message>( + buffer: &mut [u8], + stream: &mut TcpStream +) -> Message +where Message: MessageEncoder<Message> + Encode + Decode<()> + Default + Debug { + match stream.read(buffer).await { + Ok(read) => { + let received = Message::de(Vec::from(&buffer[..read])); + trace!("[Message Reader] Received {:?} from {}", received, get_target_address(stream)); + received + } + Err(err) => { + error!("[Message Reader] Error reading from stream: {}", err); + Message::err_result_decode() + } + } +} + +pub fn get_target_address(stream: &TcpStream) -> String { + let p = stream.peer_addr(); + if p.is_ok() { + p.unwrap().to_string() + } else { + "Unknown".to_string() + } +}
\ No newline at end of file diff --git a/core/src/service/tcp_network/utils/tokio_utils.rs b/core/src/service/tcp_network/utils/tokio_utils.rs new file mode 100644 index 0000000..efa68c3 --- /dev/null +++ b/core/src/service/tcp_network/utils/tokio_utils.rs @@ -0,0 +1,10 @@ +use tokio::runtime::{Builder, Runtime}; + +pub fn build_tokio_runtime(name: String) -> Runtime { + Builder::new_multi_thread() + .thread_name(name) + .thread_stack_size(32 * 1024 * 1024) + .enable_all() + .build() + .unwrap() +}
\ No newline at end of file |
