diff options
| author | 1992414357@qq.com <1992414357@qq.com> | 2025-06-09 01:44:36 +0800 |
|---|---|---|
| committer | 1992414357@qq.com <1992414357@qq.com> | 2025-06-09 01:44:36 +0800 |
| commit | 49192dbb98e0ab1f2a66b4786fac86d63f69be64 (patch) | |
| tree | 7cdf27c7cab5fb6b1aabc2ac7c44b7b492558945 /core/src/service/tcp_network | |
| parent | c9dbba0d288becb7f05cebe526be25c76e5a850a (diff) | |
重构所有部分
Diffstat (limited to 'core/src/service/tcp_network')
| -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 |
11 files changed, 700 insertions, 0 deletions
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 |
