From c470dff02c2c4def065efa59dd04026ca50dc090 Mon Sep 17 00:00:00 2001 From: "1992414357@qq.com" <1992414357@qq.com> Date: Tue, 27 May 2025 10:48:06 +0800 Subject: 将pad_io更名为pad_service MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- core/examples/start_client_console.rs | 2 +- core/examples/start_server_console.rs | 2 +- core/src/lib.rs | 2 +- core/src/pad_io/client.rs | 447 --------------------- core/src/pad_io/client_debug_cli.rs | 92 ----- core/src/pad_io/mod.rs | 5 - core/src/pad_io/server.rs | 643 ------------------------------- core/src/pad_io/server_debug_cli.rs | 204 ---------- core/src/pad_service/client.rs | 447 +++++++++++++++++++++ core/src/pad_service/client_debug_cli.rs | 92 +++++ core/src/pad_service/mod.rs | 5 + core/src/pad_service/server.rs | 643 +++++++++++++++++++++++++++++++ core/src/pad_service/server_debug_cli.rs | 204 ++++++++++ 13 files changed, 1394 insertions(+), 1394 deletions(-) delete mode 100644 core/src/pad_io/client.rs delete mode 100644 core/src/pad_io/client_debug_cli.rs delete mode 100644 core/src/pad_io/mod.rs delete mode 100644 core/src/pad_io/server.rs delete mode 100644 core/src/pad_io/server_debug_cli.rs create mode 100644 core/src/pad_service/client.rs create mode 100644 core/src/pad_service/client_debug_cli.rs create mode 100644 core/src/pad_service/mod.rs create mode 100644 core/src/pad_service/server.rs create mode 100644 core/src/pad_service/server_debug_cli.rs (limited to 'core') diff --git a/core/examples/start_client_console.rs b/core/examples/start_client_console.rs index 3d6653b..d138047 100644 --- a/core/examples/start_client_console.rs +++ b/core/examples/start_client_console.rs @@ -1,6 +1,6 @@ use std::net::{IpAddr, Ipv4Addr}; use nogamepads_core::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; -use nogamepads_core::pad_io::client::nogamepads_client::PadClient; +use nogamepads_core::pad_service::client::nogamepads_client::PadClient; const PASSWORD : &str = "password"; diff --git a/core/examples/start_server_console.rs b/core/examples/start_server_console.rs index 4b34cdd..4b42849 100644 --- a/core/examples/start_server_console.rs +++ b/core/examples/start_server_console.rs @@ -1,5 +1,5 @@ use nogamepads_core::pad_data::game_profile::game_profile::GameProfile; -use nogamepads_core::pad_io::server::nogamepads_server::PadServer; +use nogamepads_core::pad_service::server::nogamepads_server::PadServer; use nogamepads_core::DEFAULT_PORT; use std::net::{IpAddr, Ipv4Addr}; diff --git a/core/src/lib.rs b/core/src/lib.rs index f70fc4e..1abb191 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -1,7 +1,7 @@ use bincode::config; use bincode::config::Configuration; -pub mod pad_io; +pub mod pad_service; pub mod pad_data; pub const DEFAULT_PORT : u16 = 5989; diff --git a/core/src/pad_io/client.rs b/core/src/pad_io/client.rs deleted file mode 100644 index bd8bc2c..0000000 --- a/core/src/pad_io/client.rs +++ /dev/null @@ -1,447 +0,0 @@ -pub mod nogamepads_client { - use std::collections::VecDeque; - use crate::pad_data::pad_messages::nogamepads_message_transfer::{read_msg, send_msg}; - use crate::pad_data::pad_messages::nogamepads_messages::{ConnectionCallbackMessage, ConnectionErrorType, ConnectionMessage, ControlMessage, GameMessage, LeaveReason}; - use crate::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; - use log::{error, info}; - use std::net::{IpAddr, Ipv4Addr}; - use std::process::exit; - use std::sync::atomic::AtomicBool; - use std::sync::atomic::Ordering::SeqCst; - use std::sync::{Arc, Mutex}; - use std::time::Duration; - use clap::CommandFactory; - use tokio::io::{AsyncReadExt, AsyncWriteExt, ReadHalf, WriteHalf}; - use tokio::net::TcpStream; - use tokio::{io, spawn}; - use nogamepads::console_utils::debug_console::read_cli; - use nogamepads::convert_utils::convert_deque_to_vec; - use nogamepads::logger_utils::logger_build; - use crate::pad_io::client_debug_cli::{process_debug_cmd, Pcc}; - use crate::DEFAULT_PORT; - use crate::pad_data::game_profile::game_profile::GameProfile; - use crate::pad_data::pad_messages::nogamepads_message_encoder::NgpdMessageEncoder; - - type WriteList = Arc>>; - type ReadList = Arc>>; - - pub struct PadClient { - - // --- 主要参数 --- - - // 目标地址 - target_address: IpAddr, - - // 目标端口 - #[allow(dead_code)] - target_port: u16, - - // 绑定的玩家 - bind_player: PlayerInfo, - - // 调试模式 - enable_console: bool, - - // 保持安静,不初始化 env_logger - quiet: bool, - - // --- 运行时参数 --- - - // 发送信息列表 - write_list: WriteList, - - // 读取信息列表 - read_list: ReadList, - - // 是否退出 - exit: AtomicBool, - } - - impl Default for PadClient { - fn default() -> Self { - PadClient { - enable_console: false, - target_address: IpAddr::from(Ipv4Addr::new(127, 0, 0, 1)), - target_port: DEFAULT_PORT, - bind_player: PlayerInfo::new(), - quiet: false, - - write_list: WriteList::default(), - read_list: ReadList::default(), - exit: AtomicBool::new(false) - } - } - } - - // 客户端构建部分 - impl PadClient { - - pub fn bind_addr(address: IpAddr) -> PadClient { - PadClient { - target_address: address, - ..PadClient::default() - } - } - - pub fn bind_addr_with_port(address: IpAddr, port: u16) -> PadClient { - PadClient { - target_address: address, - target_port: port, - ..PadClient::default() - } - } - - pub fn enable_console(&mut self) { - self.enable_console = true; - } - - pub fn quiet(&mut self) -> &mut PadClient { - self.quiet = true; - self - } - - pub fn bind_player(&mut self, player: PlayerInfo) { - self.bind_player = player; - } - } - - // 客户端消息管理 - impl PadClient { - - - pub fn put_msg(&self, msg: ControlMessage) { - let mut guard = self.write_list.lock().unwrap(); - guard.push_back(msg); - } - - pub fn pop_a_msg(&self) -> Option { - let mut guard = self.read_list.lock().unwrap(); - if !guard.is_empty() { - guard.pop_front() - } else { - None - } - } - - pub fn pop_msg_or(&self, or: GameMessage) -> GameMessage { - self.pop_a_msg().unwrap_or(or) - } - - pub fn list_received(&self) -> Vec { - match self.read_list.lock() { - Ok(guard) => { - convert_deque_to_vec(&guard.to_owned()) - } - Err(_) => { Vec::new() } - } - } - } - - // 客户端状态控制 - impl PadClient { - - pub fn connect(self) { - - self.exit.store(false, SeqCst); - - // 构建 Logger - if !self.quiet { - logger_build(); - } - - // 调试模式 - let debug = self.enable_console; - - // 客户端对象的 Arc - let arc_client = Arc::new(self); - - // 部署环境 - let runtime = tokio::runtime::Builder::new_multi_thread() - .thread_name("nogpad-pad_io") - .thread_stack_size(32 * 1024 * 1024) - .enable_time() - .enable_io() - .build() - .unwrap(); - - info!("Starting \"NoGamepads Client\"."); - - // 入口 - let entry = async move { - let main_thread = spawn({ - let client = Arc::clone(&arc_client); - async move { - Self::main_client_thread(client).await - } - }); - - let background_thread = spawn({ - let client = Arc::clone(&arc_client); - async move { - Self::background_thread(client).await - } - }); - - if debug { - let debug_cli = spawn({ - let client = Arc::clone(&arc_client); - async move { - Self::process_debug_cli(client).await - } - }); - let _ = tokio::join!(debug_cli, main_thread, background_thread); - } else { - let _ = tokio::join!(main_thread, background_thread); - } - }; - - // 阻塞运行 - runtime.block_on(entry); - } - - pub fn exit_server(&self) { - self.exit.store(true, SeqCst); - } - - async fn main_client_thread(self: Arc) { - let mut buffer : [u8; 1024] = [0; 1024]; - let addr_str = format!("{}:{}", self.target_address.to_string(), DEFAULT_PORT); - - info!("Connected to {}", &addr_str); - - // 下载服务端配置文件 - { - info!("Check: Downloaded game profile."); - let profile = self.check_server_profile(&mut buffer, addr_str.clone()).await; - if profile.is_some() { - info!("Success: Downloaded."); - let profile = profile.unwrap_or(GameProfile::default()); - for line in profile.to_string().split('\n') { - info!("{}", line); - } - } - else { - error!("Failed: Can't download profile!"); - self.exit_server(); - } - } - - // 尝试加入服务端,并建立长连接 - { - if !self.try_join_game(&mut buffer, addr_str.clone()).await { - error!("Failed: Can't join the game!"); - self.exit_server(); - return; - } - } - } - - async fn check_server_profile(self: &Arc, buffer: &mut [u8], addr_str: String) -> Option { - match TcpStream::connect(&addr_str).await { - Ok(mut stream) => { - send_msg(&mut stream, ConnectionMessage::RequestProfile).await; - let callback : ConnectionCallbackMessage = read_msg(buffer, &mut stream).await; - match callback { - ConnectionCallbackMessage::Profile(profile) => { - Some(profile) - } - ConnectionCallbackMessage::Deny(err_type) => { - error!("Request failed: Server denied your request! ({:?})", err_type); - None - } - ConnectionCallbackMessage::Err => { - error!("Connection failed: Can't connect to server!"); - None - } - _ => { None } - } - } - Err(_err) => { - None - } - } - } - - async fn try_join_game(self: &Arc, buffer: &mut [u8], addr_str: String) -> bool { - - match TcpStream::connect(&addr_str).await { - Ok(mut stream) => { - - // 发送连接请求 - let info = self.bind_player.clone(); - send_msg(&mut stream, ConnectionMessage::Connection(info)).await; - - // 读取回调 - let callback : ConnectionCallbackMessage = read_msg(buffer, &mut stream).await; - match callback { - ConnectionCallbackMessage::Deny(error) => { - match error { - ConnectionErrorType::ContainSamePlayer => { - error!("Connection failed: Contains same player!"); - false - } - ConnectionErrorType::PlayerBanned => { - error!("Connection failed: You are banned!"); - false - } - ConnectionErrorType::Timeout => { - error!("Connection failed: Timeout!"); - false - } - ConnectionErrorType::GameLocked => { - error!("Connection failed: Game was locked!"); - false - } - _ => { false } - } - } - ConnectionCallbackMessage::Ok => { - - // 服务端检查完毕,发送 Ready 以示加入游戏 - send_msg(&mut stream, ConnectionMessage::Ready).await; - let callback : ConnectionCallbackMessage = read_msg(buffer, &mut stream).await; - match callback { - ConnectionCallbackMessage::Welcome => { - info!("Welcome!"); - Self::long_connection(Arc::clone(&self), stream).await; - } - ConnectionCallbackMessage::Deny(_error) => { - error!("Request failed: Server denied your request",); - } - _ => {} - } - true - } - _ => { false } - } - } - Err(err) => { - error!("Failed to connect to server: {}", err); - false - } - } - } - - async fn long_connection(self: Arc, stream: TcpStream) { - let (reader, writer) = io::split(stream); - spawn(Self::read_task(Arc::clone(&self), reader)); - spawn(Self::write_task(Arc::clone(&self), writer)); - } - - async fn read_task(self: Arc, mut reader: ReadHalf) { - let mut buf = [0u8; 1024]; - loop { - match reader.read(&mut buf).await { - Ok(0) => break, - Ok(n) => { - let msg = GameMessage::de(buf[0..n].to_vec()); - { - match self.read_list.lock() { - Ok(mut guard) => { - match &msg { - GameMessage::Leave(reason) => { - match reason { - LeaveReason::GameOver => { - info!("Leave Game: Game Over!"); - self.exit_server(); - } - LeaveReason::ServerClosed => { - info!("Leave Game: Server closed!"); - self.exit_server(); - } - LeaveReason::YouAreKicked => { - error!("Kick Game: You are kicked!"); - self.exit_server(); - } - LeaveReason::YouAreBanned => { - error!("Kick Game: You are banned!"); - self.exit_server(); - } - } - } - _ => { - info!("{:?}", &msg); - guard.push_back(msg); - } - } - } - Err(_) => {} - } - } - } - Err(e) => { - error!("Error reading from stream: {}", e); - self.exit_server(); - break; - } - } - } - } - - async fn write_task(self: Arc, mut writer: WriteHalf) { - loop { - let msg : Option; - { - let lock = self.write_list.lock(); - match lock { - Ok(mut guard) => { - if ! guard.is_empty() { - msg = guard.pop_front(); - } else { - msg = None; - } - } - Err(_) => { - msg = None; - } - } - } - if msg.is_some() { - let msg = msg.unwrap(); - match &writer.write_all(NgpdMessageEncoder::en(&msg).as_slice()).await { - Ok(_) => { - info!("Sent {:?}", msg); - } - Err(_error) => { - error!("Sent {:?} failed!", msg); - } - } - } - } - } - - async fn background_thread(self: Arc) { - loop { - // 退出程序的监听 - if self.exit.load(SeqCst) { - tokio::time::sleep(Duration::from_secs(1)).await; - info!("Main thread exited."); - exit(0); - } - } - } - - async fn process_debug_cli(self: Arc) { - loop { - if self.exit.load(SeqCst) { - info!("Debug console exited"); - break - } - tokio::time::sleep(Duration::from_secs_f64(0.2)).await; - let option: Option = read_cli( - format!("CLIENT {}/{}> ", - self.target_address.to_string(), - self.bind_player.account.id).as_str(), - "pcc".to_string(), - Pcc::command() - ).await; - match option { - None => {} - Some(cmd) => { - process_debug_cmd(cmd, Arc::clone(&self)); - } - } - } - } - } -} \ No newline at end of file diff --git a/core/src/pad_io/client_debug_cli.rs b/core/src/pad_io/client_debug_cli.rs deleted file mode 100644 index b34bef5..0000000 --- a/core/src/pad_io/client_debug_cli.rs +++ /dev/null @@ -1,92 +0,0 @@ -use crate::pad_io::client::nogamepads_client::PadClient; -use crate::pad_data::pad_messages::nogamepads_messages::{ControlMessage, GameMessage}; -use clap::{Args, Parser, Subcommand}; -use std::sync::Arc; -use log::info; - -/// NoGamePads Client - Cli -#[derive(Parser, Debug)] -#[command(author, version, about, long_about = None)] -pub struct Pcc { - #[command(subcommand)] - command: Commands, -} - -/// 主要命令 -#[derive(Subcommand, Debug)] -enum Commands { - - // 清屏 - #[command(about = "Clean the screen")] - Clear, - - // 断开当前连接 - #[command(about = "Exit from server")] - Exit, - - // 检查收到的消息 - #[command(about = "Check received")] - Received(ReceivedArgs), - - // 取出一条消息 - #[command(about = "Pop a message")] - Pop(PopArgs), - - // 发送消息 - #[command(about = "Send Message")] - Msg(MsgArgs), -} - -#[derive(Args, Debug)] -struct ReceivedArgs { - - #[arg(long)] - list: bool -} - -/// 发送消息 参数 -#[derive(Args, Debug)] -struct MsgArgs { - - // 消息内容 - #[arg(value_name = "CONTENT")] - message: String, -} - -#[derive(Args, Debug)] -struct PopArgs { } - -pub fn process_debug_cmd (cmd: Pcc, client: Arc) { - match cmd.command { - Commands::Clear => { - clearscreen::clear().expect("Failed to clear screen"); - } - - Commands::Exit => { - client.exit_server(); - } - - Commands::Received(args) => { - if args.list { - for msg in client.list_received() { - info!("{:?}", msg); - } - } else { - info!("Total {} messsage(s)!", client.list_received().iter().count()); - } - } - - Commands::Pop(_args) => { - info!("{:?}", client.pop_msg_or(GameMessage::Err)); - } - - Commands::Msg(args) => { - client.put_msg(ControlMessage::Msg(args.message)); - } - } -} - -#[allow(dead_code)] -fn put_to_list(client: Arc, message: ControlMessage) { - client.put_msg(message); -} \ No newline at end of file diff --git a/core/src/pad_io/mod.rs b/core/src/pad_io/mod.rs deleted file mode 100644 index ae10f2d..0000000 --- a/core/src/pad_io/mod.rs +++ /dev/null @@ -1,5 +0,0 @@ -pub mod client; -pub mod client_debug_cli; - -pub mod server; -pub mod server_debug_cli; \ No newline at end of file diff --git a/core/src/pad_io/server.rs b/core/src/pad_io/server.rs deleted file mode 100644 index 59241f1..0000000 --- a/core/src/pad_io/server.rs +++ /dev/null @@ -1,643 +0,0 @@ -pub mod nogamepads_server { - use std::collections::{HashMap, VecDeque}; - use crate::pad_data::pad_messages::nogamepads_message_encoder::NgpdMessageEncoder; - use crate::pad_data::pad_messages::nogamepads_message_transfer::{read_msg, send_msg}; - use crate::pad_data::pad_messages::nogamepads_messages::{ConnectionCallbackMessage, ConnectionMessage, ControlMessage, GameMessage, LeaveReason}; - use log::{error, info, warn}; - use std::net::{IpAddr, Ipv4Addr, SocketAddr}; - use std::process::exit; - use std::sync::atomic::AtomicBool; - use std::sync::atomic::Ordering::SeqCst; - use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; - use std::time::Duration; - use clap::CommandFactory; - use tokio::io::{AsyncReadExt, AsyncWriteExt, ReadHalf, WriteHalf}; - use tokio::net::{TcpListener, TcpStream}; - use tokio::{io, spawn}; - use tokio::runtime::Runtime; - use nogamepads::console_utils::debug_console::read_cli; - use nogamepads::convert_utils::convert_deque_to_vec; - use nogamepads::logger_utils::logger_build; - use crate::DEFAULT_PORT; - use crate::pad_data::game_profile::game_profile::GameProfile; - use crate::pad_data::pad_messages::nogamepads_messages::ConnectionErrorType::{ContainSamePlayer, GameLocked, PlayerBanned, WhatTheHell}; - use crate::pad_data::pad_messages::nogamepads_messages::GameMessage::Leave; - use crate::pad_data::pad_messages::nogamepads_messages::LeaveReason::ServerClosed; - use crate::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; - use crate::pad_io::server_debug_cli::{process_debug_cmd, Psc}; - - type PlayerMap = Arc>>; - type WriteList = Arc>>>; - type ReadList = Arc>>>; - - pub struct PadServer { - - // --- 主要参数 --- - - // 本地监听地址 - address: IpAddr, - - // 游戏信息 - game_profile: GameProfile, - - // 绑定端口 - port: u16, - - // 调试模式 - enable_console: bool, - - // 保持安静,不初始化 env_logger - quiet: bool, - - // --- 运行时参数 --- - - // 发送信息列表 - write_list: WriteList, - - // 读取信息列表 - read_list: ReadList, - - // 在线玩家 - online_players: PlayerMap, - - // 被封禁的玩家 - banned_players: PlayerMap, - - // 是否锁定该游戏:禁止后续玩家加入 - game_locked: AtomicBool, - - // 是否停止服务器 - stop: AtomicBool, - } - - impl Clone for PadServer { - fn clone(&self) -> Self { - PadServer { - address: self.address.clone(), - game_profile: self.game_profile.clone(), - port: self.port.clone(), - enable_console: self.enable_console, - quiet: self.quiet, - - write_list: self.write_list.clone(), - read_list: self.read_list.clone(), - online_players: self.online_players.clone(), - banned_players: self.banned_players.clone(), - game_locked: AtomicBool::new((&self.game_locked.load(SeqCst)).clone()), - stop: AtomicBool::new((&self.stop.load(SeqCst)).clone()), - } - } - } - - impl Default for PadServer { - fn default() -> Self { - PadServer { - address: IpAddr::from(Ipv4Addr::new(127, 0, 0, 1)), - game_profile: GameProfile::default(), - port: DEFAULT_PORT, - enable_console: false, - quiet: false, - - write_list: WriteList::default(), - read_list: ReadList::default(), - online_players: PlayerMap::default(), - banned_players: PlayerMap::default(), - game_locked: AtomicBool::new(false), - stop: AtomicBool::new(false), - } - } - } - - // 服务端构建部分 - impl PadServer { - - pub fn build_simple() -> Arc { - Arc::new(Self::default() - .addr(IpAddr::from(Ipv4Addr::new(127, 0, 0, 1)), DEFAULT_PORT) - .put_profile(GameProfile::default()).to_owned()) - } - - pub fn addr(&mut self, ip_addr: IpAddr, port: u16) -> &mut PadServer { - self.ip_addr(ip_addr).port(port) - } - - pub fn socket_addr(&mut self, socket_addr: SocketAddr) -> &mut PadServer { - self.ip_addr(socket_addr.ip()).port(socket_addr.port()) - } - - pub fn port(&mut self, port: u16) -> &mut PadServer { - self.port = port; - self - } - - pub fn ip_addr(&mut self, ip_addr: IpAddr) -> &mut PadServer { - self.address = ip_addr; - self - } - - pub fn put_profile(&mut self, profile: GameProfile) -> &mut PadServer { - self.game_profile = profile; - self - } - - pub fn enable_console(&mut self) -> &mut PadServer { - self.enable_console = true; - self - } - - pub fn quiet(&mut self) -> &mut PadServer { - self.quiet = true; - self - } - - pub fn build(&self) -> Arc { - Arc::new(self.clone()) - } - - } - - // 服务端消息管理 - impl PadServer { - - pub fn put_msg_to(&self, msg: GameMessage, player: &PlayerInfo) { - match self.write_list.lock() { - Ok(mut guard) => { - let hash = &player.account.player_hash.clone(); - if ! guard.contains_key(hash.as_str()) { - guard.entry(player.account.player_hash.clone()) - .or_insert_with(VecDeque::new) - .push_back(msg); - } - } - Err(_) => { - error!("Cannot lock \"{:?}\" in write_list", player.account.player_hash); - } - } - } - - pub fn put_msg_to_all(&self, msg: &GameMessage) { - match self.list_players() { - Ok(list) => { - for player in list { - self.put_msg_to(msg.clone(), &player); - } - } - Err(_) => { - error!("Cannot put GameMessage with no players."); - } - } - } - - pub fn pop_a_msg(&self, player: &PlayerInfo) -> Option { - match self.read_list.lock() { - Ok(mut guard) => { - match guard.get_mut(&player.account.player_hash) { - None => { None } - Some(queue) => { - if ! queue.is_empty() { - queue.pop_front() - } else { - guard.remove(&player.account.player_hash); - None - } - } - } - } - Err(_) => { - error!("Cannot lock \"{:?}\" in read_list", player.account.player_hash); - None - } - } - } - - pub fn pop_msg_or(&self, player: &PlayerInfo, or: ControlMessage) -> ControlMessage { - self.pop_a_msg(player).unwrap_or(or) - } - - pub fn list_received(&self, player: &PlayerInfo) -> Vec { - match self.read_list.lock() { - Ok(guard) => { - match guard.get_key_value(player.account.player_hash.as_str()) { - None => { Vec::new() } - Some(result) => { - convert_deque_to_vec(result.1) - } - } - } - Err(_) => { Vec::new() } - } - } - } - - // 服务端玩家管理 - impl PadServer { - - pub fn is_player_online (&self, player: &PlayerInfo) -> bool { - let guard = self.online_players.lock().unwrap(); - guard.contains_key(&player.account.player_hash) - } - - pub fn set_player_online (&self, player: &PlayerInfo, online: bool) { - let online_current = self.is_player_online(player); - if online_current && !online { - let mut guard = self.online_players.lock().unwrap(); - guard.remove(&player.account.player_hash); - info!("{} is OFFLINE!", player.account.id); - } else if !online_current && online { - let mut guard = self.online_players.lock().unwrap(); - guard.insert(player.account.player_hash.clone(), player.clone()); - info!("{} is ONLINE!", player.account.id); - } - } - - pub fn is_player_banned (&self, player: &PlayerInfo) -> bool { - let guard = self.banned_players.lock().unwrap(); - guard.contains_key(&player.account.player_hash) - } - - pub fn kick_player(&self, player: &PlayerInfo) { - if self.is_player_online(player) { - self.put_msg_to(Leave(LeaveReason::YouAreKicked), player); - } - } - - pub fn ban_player(&self, player: &PlayerInfo) { - self.set_player_banned(player, true); - if self.is_player_online(player) { - self.put_msg_to(Leave(LeaveReason::YouAreBanned), player); - } - } - - pub fn pardon_player(&self, player: &PlayerInfo) { - self.set_player_banned(player, false); - } - - fn set_player_banned (&self, player: &PlayerInfo, banned: bool) { - let banned_current = self.is_player_banned(player); - if banned_current && !banned { - let mut guard = self.banned_players.lock().unwrap(); - guard.remove(&player.account.player_hash); - info!("Pardoned player {}", player.account.id); - } else if !banned_current && banned { - let mut guard = self.banned_players.lock().unwrap(); - guard.insert(player.account.player_hash.clone(), player.clone()); - info!("Banned player {}!", player.account.id); - } - } - - pub fn list_players(&self) -> Result, PoisonError>>> { - match self.online_players.lock() { - Ok(guard) => { - Ok(guard.values().cloned().collect()) - } - Err(err) => Err(err) - } - } - - pub fn list_players_banned(&self) -> Result, PoisonError>>> { - match self.banned_players.lock() { - Ok(guard) => { - Ok(guard.values().cloned().collect()) - } - Err(err) => Err(err) - } - } - } - - // 服务端状态控制 - #[allow(dead_code)] - impl PadServer { - - pub fn stop_server(&self) { - self.put_msg_to_all(&Leave(ServerClosed)); - self.stop.store(true, SeqCst); - } - - pub fn start_server(self: Arc) { - - // 构建 Logger - if ! self.quiet { - logger_build(); - } - - // 运行时 - let runtime = Self::get_runtime(); - - info!("Starting \"NoGamepads Server\"."); - - // 入口 - let console = self.enable_console; - let entry = self.get_entry(console); - - // 阻塞运行 - runtime.block_on(entry); - } - - fn get_runtime() -> Runtime { - tokio::runtime::Builder::new_multi_thread() - .thread_name("nogpad-server") - .thread_stack_size(32 * 1024 * 1024) - .enable_time() - .enable_io() - .build() - .unwrap() - } - - fn get_entry(self: Arc, debug: bool) -> impl Future + Send + 'static { - async move { - let main_thread = spawn({ - let client = Arc::clone(&self); - async move { - Self::main_request_thread(client).await - } - }); - - let background_thread = spawn({ - let client = Arc::clone(&self); - async move { - Self::background_thread(client).await - } - }); - - if debug { - let debug_cli = spawn({ - let client = Arc::clone(&self); - async move { - Self::process_debug_cli(client).await - } - }); - - let _ = tokio::join!(debug_cli, main_thread, background_thread); - } else { - let _ = tokio::join!(main_thread, background_thread); - } - } - } - - fn lock_game(&self) { - self.game_locked.store(true, SeqCst); - } - - fn unlock_game(&self) { - self.game_locked.store(false, SeqCst); - } - - fn is_game_locked(&self) -> bool { - self.game_locked.load(SeqCst) - } - - async fn main_request_thread(self: Arc) { - - let addr_str = format!("{}:{}", self.address.to_string(), self.port); - info!("Server listening at {}", addr_str); - - // Tcp 监听器 - let listener : TcpListener; - match TcpListener::bind(&addr_str).await { - Ok(result) => { - info!("Listener created."); - listener = result; - } - Err(_) => { - error!("Server listening at {} failed!", addr_str); - exit(1); - } - } - - // 请求信息循环 - loop { - match listener.accept().await { - Ok((stream, _)) => { - spawn(Self::process_request(Arc::clone(&self), stream)); - } - Err(error) => { - error!("Error: {}", error); - } - } - } - } - - async fn process_request(self: Arc, mut stream: TcpStream) { - let mut buffer = [0; 1024]; - let connection_msg : ConnectionMessage = read_msg(&mut buffer, &mut stream).await; - match connection_msg { - - // 客户端请求加入游戏,并建立长连接 - ConnectionMessage::Connection(info) => { - - // 加入游戏资格检测 - info!("Account {} trying to connect.", info.account.player_hash); - - // 0. 当前游戏是否已经锁定? - if self.is_game_locked() { - // 当前游戏已经锁定,禁止加入玩家,发送失败信息,并断开连接 - send_msg(&mut stream, ConnectionCallbackMessage::Deny(GameLocked)).await; - return; - } - - // 1. 是否存在重复玩家? - let online = self.is_player_online(&info); - if online { - // 当前玩家已在线,发送失败信息,并断开连接 - send_msg(&mut stream, ConnectionCallbackMessage::Deny(ContainSamePlayer)).await; - return; - } - - // 2. 该玩家是否被封禁? - let banned = self.is_player_banned(&info); - if banned { - // 当前玩家已被封禁,发送失败信息,并断开连接 - send_msg(&mut stream, ConnectionCallbackMessage::Deny(PlayerBanned)).await; - return; - } - - // OK!若执行到此处,说明该玩家具有加入资格,Welcome! - - send_msg(&mut stream, ConnectionCallbackMessage::Ok).await; - let callback : ConnectionMessage = read_msg(&mut buffer, &mut stream).await; - - match callback { - // 玩家已就绪,发送 Welcome 信息以邀请该玩家加入游戏 - ConnectionMessage::Ready => { - info!("Player \"{}\" is ready!", info.account.id); - - // 发送 Welcome - send_msg(&mut stream, ConnectionCallbackMessage::Welcome).await; - - // 注册该玩家到在线列表 - self.set_player_online(&info, true); - - // 启动控制循环 - - spawn(Self::long_connection(Arc::clone(&self), stream, info)); - }, - _ => { - send_msg(&mut stream, ConnectionCallbackMessage::Deny(WhatTheHell)).await; // WTH ? - } - } - } - - // 客户端请求获得游戏信息 - ConnectionMessage::RequestProfile => { - // 发送游戏信息到客户端 - send_msg(&mut stream, ConnectionCallbackMessage::Profile(self.game_profile.clone())).await; - } - - // 客户端发来了错误信息 - ConnectionMessage::Err => { - match stream.peer_addr() { - Ok(addr) => { - warn!("Received an error message from {}.", addr.to_string()); - } - Err(_) => { - warn!("Received an error message from unknown pad_io."); - } - } - } - - // 客户端发来了不相干的信息 - _ => { - match stream.peer_addr() { - Ok(addr) => { - warn!("Received unknown connection message from {}.", addr.to_string()); - } - Err(_) => { - warn!("Received unknown connection message from unknown pad_io."); - } - } - } - } - } - - async fn long_connection(self: Arc, stream: TcpStream, player_info: PlayerInfo) { - let player_info_arc = Arc::new(player_info); - let (reader, writer) = io::split(stream); - spawn(Self::read_task(Arc::clone(&self), reader, Arc::clone(&player_info_arc))); - spawn(Self::write_task(Arc::clone(&self), writer, Arc::clone(&player_info_arc))); - } - - async fn read_task(self: Arc, - mut reader: ReadHalf, - player_info: Arc) { - let player_hash = player_info.account.player_hash.clone(); - let mut buf = [0u8; 1024]; - loop { - match reader.read(&mut buf).await { - Ok(0) => break, - Ok(n) => { - let msg = ControlMessage::de(buf[0..n].to_vec()); - { - match self.read_list.lock() { - Ok(mut guard) => { - info!("{:?} from {}({})", &msg, player_info.customize.nickname, player_info.account.id); - guard - .entry(player_hash.clone()) - .or_insert_with(VecDeque::new) - .push_back(msg); - - } - Err(_) => { - } - } - } - } - Err(e) => { - warn!("Error reading from stream: {}", e); - - self.set_player_online(&player_info, false); - - // 放入一条错误信息到队列,使 write_task 及时发现该玩家离开 - self.put_msg_to(GameMessage::Err, &player_info); - - break; - } - } - } - } - - async fn write_task(self: Arc, - mut writer: WriteHalf, - player_info: Arc) { - let player_hash = player_info.account.player_hash.clone(); - let mut exit = false; - loop { - let msg : Option; - match self.write_list.lock() { - Ok(mut hash_map) => { - if ! hash_map.is_empty() { - match hash_map.get_mut(&player_hash) { - None => { - msg = None; - } - Some(queue) => { - if ! queue.is_empty() { - msg = queue.pop_front(); - } else { - msg = None; - hash_map.remove(&player_hash); - } - } - } - } - else { msg = None; } - } - Err(_) => { - msg = None; - } - } - if msg.is_some() { - let msg = msg.unwrap(); - match &writer.write_all(NgpdMessageEncoder::en(&msg).as_slice()).await { - Ok(_) => { - - info!("Sent {:?} to {}", msg, &player_info.account.id); - } - Err(error) => { - warn!("Sent {:?} to {} failed!", msg, &player_info.account.id); - warn!("{:?}", error); - - exit = true; - } - } - } - if exit { - warn!("Long connection between \"{}\" closed.", &player_info.account.id); - break - } - } - } - - async fn background_thread(self: Arc) { - loop { - // 退出程序的监听 - if self.stop.load(SeqCst) { - tokio::time::sleep(Duration::from_secs(1)).await; - - info!("Main thread exited."); - exit(0); - } - } - } - - async fn process_debug_cli(self: Arc) { - loop { - if self.stop.load(SeqCst) { - info!("Debug console exited"); - return; - } - tokio::time::sleep(Duration::from_secs_f64(0.2)).await; - let option: Option = read_cli( - format!("SERVER {}> ", self.address.to_string()).as_str(), - "psc".to_string(), - Psc::command() - ).await; - match option { - None => {} - Some(cmd) => { - process_debug_cmd(cmd, Arc::clone(&self)); - } - } - } - } - } -} \ No newline at end of file diff --git a/core/src/pad_io/server_debug_cli.rs b/core/src/pad_io/server_debug_cli.rs deleted file mode 100644 index 39e1868..0000000 --- a/core/src/pad_io/server_debug_cli.rs +++ /dev/null @@ -1,204 +0,0 @@ -use std::collections::HashMap; -use crate::pad_data::pad_messages::nogamepads_messages::{ControlMessage, GameMessage}; -use crate::pad_io::server::nogamepads_server::PadServer; -use clap::{Args, Parser, Subcommand}; -use std::ops::{Index}; -use std::sync::{Arc, MutexGuard, PoisonError}; -use log::{error, info}; -use crate::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; - -/// NoGamePads Server - Cli -#[derive(Parser, Debug)] -#[command(author, version, about, long_about = None)] -pub struct Psc { - #[command(subcommand)] - command: Commands, -} - -#[derive(Subcommand, Debug)] -enum Commands { - - // 清屏 - #[command(about = "Clean the screen")] - Clear, - - // 关闭服务器 - #[command(about = "Close the server")] - Stop, - - // 展示所有玩家 - #[command(about = "List all online players")] - List, - - // 展示所有封禁的玩家 - #[command(about = "List all banned players")] - Banned, - - // 检查收到的消息 - #[command(about = "Check received")] - Received(ReceivedArgs), - - // 取出一条消息 - #[command(about = "Pop a message")] - Pop(PlayerArgs), - - // 踢出玩家 - #[command(about = "Kick a player")] - Kick(PlayerArgs), - - // 封禁玩家 - #[command(about = "Ban a player")] - Ban(PlayerArgs), - - // 解封(赦免)玩家 - #[command(about = "Pardon a player")] - Pardon(PlayerArgs), - - // 激活事件触发器 - #[command(about = "Send SkinEventTrigger")] - Event(EventArgs) -} - -// 检查收到的消息 -#[derive(Args, Debug)] -struct ReceivedArgs { - - #[arg(default_value = "0")] - player: usize, - - #[arg(long)] - list: bool, -} - -/// 激活事件触发器 参数 -#[derive(Args, Debug)] -struct EventArgs { - - // 玩家序号 - #[arg(value_name = "PLAYER_INDEX")] - index: usize, - - // 事件编号 - #[arg(value_name = "CONTENT")] - message: u8, -} - -#[derive(Args, Debug)] -struct PlayerArgs { - - // 玩家序号 - #[arg(value_name = "PLAYER_INDEX")] - index: usize -} - -pub fn process_debug_cmd (cmd: Psc, server: Arc) { - match cmd.command { - - Commands::Clear => { - clearscreen::clear().expect("Failed to clear screen"); - } - - Commands::Stop => { - server.stop_server(); - } - - Commands::List => { - print_player_list(server.list_players()); - } - - Commands::Banned => { - print_player_list(server.list_players_banned()); - } - - Commands::Received(args) => { - let players = server.list_players().unwrap_or(Vec::new()); - let player = players.index(args.player.clamp(0, players.iter().count() -1)); - if args.list { - for msg in server.list_received(player) { - info!("{:?}", msg); - } - } else { - info!("Total {} messsage(s)!", server.list_received(player).iter().count()); - } - } - - Commands::Pop(args) => { - match get_player_by_index(&server, args.index) { - None => { - error!("Pup message failed : Player index \"{}\" not found!", args.index); - } - Some(player) => { - let message = server.pop_msg_or(&player, ControlMessage::Err); - info!("{:?}", message); - } - } - } - - Commands::Kick(args) => { - let player = get_player_by_index(&server, args.index); - if player.is_some() { - let player = player.unwrap(); - server.kick_player(&player); - } - } - - Commands::Ban(args) => { - let player = get_player_by_index(&server, args.index); - if player.is_some() { - let player = player.unwrap(); - server.ban_player(&player); - } - } - - Commands::Pardon(args) => { - let player = get_player_by_ban_index(&server, args.index); - if player.is_some() { - let player = player.unwrap(); - server.pardon_player(&player); - } - } - - Commands::Event(args) => { - put_to_list(server, args.index, GameMessage::SkinEventTrigger(args.message)); - } - } -} - -fn put_to_list(server: Arc, player_index: usize, message: GameMessage) { - match get_player_by_index(&server, player_index) { - None => { - error!("Put message failed : Player index \"{}\" not found!", player_index); - } - Some(player) => { - server.put_msg_to(message, &player); - } - } -} - -fn get_player_by_index(server: &Arc, index: usize) -> Option { - let list = server.list_players().unwrap_or(Vec::new()); - let max = list.iter().count(); - let index = if max > 0 { index.clamp(0, max - 1) } else { 0 }; - - let result = list.get(index).cloned(); - result -} - -fn get_player_by_ban_index(server: &Arc, index: usize) -> Option { - let list = server.list_players_banned().unwrap_or(Vec::new()); - let max = list.iter().count(); - let index = if max > 0 { index.clamp(0, max - 1) } else { 0 }; - - let result = list.get(index).cloned(); - result -} - -fn print_player_list(list: Result, PoisonError>>>) { - let list = list.unwrap_or(Vec::new()); - let mut i = 0; - for player in list { - let n = player.customize.nickname; - info!("({}){} ", i, n); - i += 1; - } -} \ No newline at end of file diff --git a/core/src/pad_service/client.rs b/core/src/pad_service/client.rs new file mode 100644 index 0000000..3114330 --- /dev/null +++ b/core/src/pad_service/client.rs @@ -0,0 +1,447 @@ +pub mod nogamepads_client { + use std::collections::VecDeque; + use crate::pad_data::pad_messages::nogamepads_message_transfer::{read_msg, send_msg}; + use crate::pad_data::pad_messages::nogamepads_messages::{ConnectionCallbackMessage, ConnectionErrorType, ConnectionMessage, ControlMessage, GameMessage, LeaveReason}; + use crate::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; + use log::{error, info}; + use std::net::{IpAddr, Ipv4Addr}; + use std::process::exit; + use std::sync::atomic::AtomicBool; + use std::sync::atomic::Ordering::SeqCst; + use std::sync::{Arc, Mutex}; + use std::time::Duration; + use clap::CommandFactory; + use tokio::io::{AsyncReadExt, AsyncWriteExt, ReadHalf, WriteHalf}; + use tokio::net::TcpStream; + use tokio::{io, spawn}; + use nogamepads::console_utils::debug_console::read_cli; + use nogamepads::convert_utils::convert_deque_to_vec; + use nogamepads::logger_utils::logger_build; + use crate::pad_service::client_debug_cli::{process_debug_cmd, Pcc}; + use crate::DEFAULT_PORT; + use crate::pad_data::game_profile::game_profile::GameProfile; + use crate::pad_data::pad_messages::nogamepads_message_encoder::NgpdMessageEncoder; + + type WriteList = Arc>>; + type ReadList = Arc>>; + + pub struct PadClient { + + // --- 主要参数 --- + + // 目标地址 + target_address: IpAddr, + + // 目标端口 + #[allow(dead_code)] + target_port: u16, + + // 绑定的玩家 + bind_player: PlayerInfo, + + // 调试模式 + enable_console: bool, + + // 保持安静,不初始化 env_logger + quiet: bool, + + // --- 运行时参数 --- + + // 发送信息列表 + write_list: WriteList, + + // 读取信息列表 + read_list: ReadList, + + // 是否退出 + exit: AtomicBool, + } + + impl Default for PadClient { + fn default() -> Self { + PadClient { + enable_console: false, + target_address: IpAddr::from(Ipv4Addr::new(127, 0, 0, 1)), + target_port: DEFAULT_PORT, + bind_player: PlayerInfo::new(), + quiet: false, + + write_list: WriteList::default(), + read_list: ReadList::default(), + exit: AtomicBool::new(false) + } + } + } + + // 客户端构建部分 + impl PadClient { + + pub fn bind_addr(address: IpAddr) -> PadClient { + PadClient { + target_address: address, + ..PadClient::default() + } + } + + pub fn bind_addr_with_port(address: IpAddr, port: u16) -> PadClient { + PadClient { + target_address: address, + target_port: port, + ..PadClient::default() + } + } + + pub fn enable_console(&mut self) { + self.enable_console = true; + } + + pub fn quiet(&mut self) -> &mut PadClient { + self.quiet = true; + self + } + + pub fn bind_player(&mut self, player: PlayerInfo) { + self.bind_player = player; + } + } + + // 客户端消息管理 + impl PadClient { + + + pub fn put_msg(&self, msg: ControlMessage) { + let mut guard = self.write_list.lock().unwrap(); + guard.push_back(msg); + } + + pub fn pop_a_msg(&self) -> Option { + let mut guard = self.read_list.lock().unwrap(); + if !guard.is_empty() { + guard.pop_front() + } else { + None + } + } + + pub fn pop_msg_or(&self, or: GameMessage) -> GameMessage { + self.pop_a_msg().unwrap_or(or) + } + + pub fn list_received(&self) -> Vec { + match self.read_list.lock() { + Ok(guard) => { + convert_deque_to_vec(&guard.to_owned()) + } + Err(_) => { Vec::new() } + } + } + } + + // 客户端状态控制 + impl PadClient { + + pub fn connect(self) { + + self.exit.store(false, SeqCst); + + // 构建 Logger + if !self.quiet { + logger_build(); + } + + // 调试模式 + let debug = self.enable_console; + + // 客户端对象的 Arc + let arc_client = Arc::new(self); + + // 部署环境 + let runtime = tokio::runtime::Builder::new_multi_thread() + .thread_name("nogpad-pad_service") + .thread_stack_size(32 * 1024 * 1024) + .enable_time() + .enable_io() + .build() + .unwrap(); + + info!("Starting \"NoGamepads Client\"."); + + // 入口 + let entry = async move { + let main_thread = spawn({ + let client = Arc::clone(&arc_client); + async move { + Self::main_client_thread(client).await + } + }); + + let background_thread = spawn({ + let client = Arc::clone(&arc_client); + async move { + Self::background_thread(client).await + } + }); + + if debug { + let debug_cli = spawn({ + let client = Arc::clone(&arc_client); + async move { + Self::process_debug_cli(client).await + } + }); + let _ = tokio::join!(debug_cli, main_thread, background_thread); + } else { + let _ = tokio::join!(main_thread, background_thread); + } + }; + + // 阻塞运行 + runtime.block_on(entry); + } + + pub fn exit_server(&self) { + self.exit.store(true, SeqCst); + } + + async fn main_client_thread(self: Arc) { + let mut buffer : [u8; 1024] = [0; 1024]; + let addr_str = format!("{}:{}", self.target_address.to_string(), DEFAULT_PORT); + + info!("Connected to {}", &addr_str); + + // 下载服务端配置文件 + { + info!("Check: Downloaded game profile."); + let profile = self.check_server_profile(&mut buffer, addr_str.clone()).await; + if profile.is_some() { + info!("Success: Downloaded."); + let profile = profile.unwrap_or(GameProfile::default()); + for line in profile.to_string().split('\n') { + info!("{}", line); + } + } + else { + error!("Failed: Can't download profile!"); + self.exit_server(); + } + } + + // 尝试加入服务端,并建立长连接 + { + if !self.try_join_game(&mut buffer, addr_str.clone()).await { + error!("Failed: Can't join the game!"); + self.exit_server(); + return; + } + } + } + + async fn check_server_profile(self: &Arc, buffer: &mut [u8], addr_str: String) -> Option { + match TcpStream::connect(&addr_str).await { + Ok(mut stream) => { + send_msg(&mut stream, ConnectionMessage::RequestProfile).await; + let callback : ConnectionCallbackMessage = read_msg(buffer, &mut stream).await; + match callback { + ConnectionCallbackMessage::Profile(profile) => { + Some(profile) + } + ConnectionCallbackMessage::Deny(err_type) => { + error!("Request failed: Server denied your request! ({:?})", err_type); + None + } + ConnectionCallbackMessage::Err => { + error!("Connection failed: Can't connect to server!"); + None + } + _ => { None } + } + } + Err(_err) => { + None + } + } + } + + async fn try_join_game(self: &Arc, buffer: &mut [u8], addr_str: String) -> bool { + + match TcpStream::connect(&addr_str).await { + Ok(mut stream) => { + + // 发送连接请求 + let info = self.bind_player.clone(); + send_msg(&mut stream, ConnectionMessage::Connection(info)).await; + + // 读取回调 + let callback : ConnectionCallbackMessage = read_msg(buffer, &mut stream).await; + match callback { + ConnectionCallbackMessage::Deny(error) => { + match error { + ConnectionErrorType::ContainSamePlayer => { + error!("Connection failed: Contains same player!"); + false + } + ConnectionErrorType::PlayerBanned => { + error!("Connection failed: You are banned!"); + false + } + ConnectionErrorType::Timeout => { + error!("Connection failed: Timeout!"); + false + } + ConnectionErrorType::GameLocked => { + error!("Connection failed: Game was locked!"); + false + } + _ => { false } + } + } + ConnectionCallbackMessage::Ok => { + + // 服务端检查完毕,发送 Ready 以示加入游戏 + send_msg(&mut stream, ConnectionMessage::Ready).await; + let callback : ConnectionCallbackMessage = read_msg(buffer, &mut stream).await; + match callback { + ConnectionCallbackMessage::Welcome => { + info!("Welcome!"); + Self::long_connection(Arc::clone(&self), stream).await; + } + ConnectionCallbackMessage::Deny(_error) => { + error!("Request failed: Server denied your request",); + } + _ => {} + } + true + } + _ => { false } + } + } + Err(err) => { + error!("Failed to connect to server: {}", err); + false + } + } + } + + async fn long_connection(self: Arc, stream: TcpStream) { + let (reader, writer) = io::split(stream); + spawn(Self::read_task(Arc::clone(&self), reader)); + spawn(Self::write_task(Arc::clone(&self), writer)); + } + + async fn read_task(self: Arc, mut reader: ReadHalf) { + let mut buf = [0u8; 1024]; + loop { + match reader.read(&mut buf).await { + Ok(0) => break, + Ok(n) => { + let msg = GameMessage::de(buf[0..n].to_vec()); + { + match self.read_list.lock() { + Ok(mut guard) => { + match &msg { + GameMessage::Leave(reason) => { + match reason { + LeaveReason::GameOver => { + info!("Leave Game: Game Over!"); + self.exit_server(); + } + LeaveReason::ServerClosed => { + info!("Leave Game: Server closed!"); + self.exit_server(); + } + LeaveReason::YouAreKicked => { + error!("Kick Game: You are kicked!"); + self.exit_server(); + } + LeaveReason::YouAreBanned => { + error!("Kick Game: You are banned!"); + self.exit_server(); + } + } + } + _ => { + info!("{:?}", &msg); + guard.push_back(msg); + } + } + } + Err(_) => {} + } + } + } + Err(e) => { + error!("Error reading from stream: {}", e); + self.exit_server(); + break; + } + } + } + } + + async fn write_task(self: Arc, mut writer: WriteHalf) { + loop { + let msg : Option; + { + let lock = self.write_list.lock(); + match lock { + Ok(mut guard) => { + if ! guard.is_empty() { + msg = guard.pop_front(); + } else { + msg = None; + } + } + Err(_) => { + msg = None; + } + } + } + if msg.is_some() { + let msg = msg.unwrap(); + match &writer.write_all(NgpdMessageEncoder::en(&msg).as_slice()).await { + Ok(_) => { + info!("Sent {:?}", msg); + } + Err(_error) => { + error!("Sent {:?} failed!", msg); + } + } + } + } + } + + async fn background_thread(self: Arc) { + loop { + // 退出程序的监听 + if self.exit.load(SeqCst) { + tokio::time::sleep(Duration::from_secs(1)).await; + info!("Main thread exited."); + exit(0); + } + } + } + + async fn process_debug_cli(self: Arc) { + loop { + if self.exit.load(SeqCst) { + info!("Debug console exited"); + break + } + tokio::time::sleep(Duration::from_secs_f64(0.2)).await; + let option: Option = read_cli( + format!("CLIENT {}/{}> ", + self.target_address.to_string(), + self.bind_player.account.id).as_str(), + "pcc".to_string(), + Pcc::command() + ).await; + match option { + None => {} + Some(cmd) => { + process_debug_cmd(cmd, Arc::clone(&self)); + } + } + } + } + } +} \ No newline at end of file diff --git a/core/src/pad_service/client_debug_cli.rs b/core/src/pad_service/client_debug_cli.rs new file mode 100644 index 0000000..988bf39 --- /dev/null +++ b/core/src/pad_service/client_debug_cli.rs @@ -0,0 +1,92 @@ +use crate::pad_service::client::nogamepads_client::PadClient; +use crate::pad_data::pad_messages::nogamepads_messages::{ControlMessage, GameMessage}; +use clap::{Args, Parser, Subcommand}; +use std::sync::Arc; +use log::info; + +/// NoGamePads Client - Cli +#[derive(Parser, Debug)] +#[command(author, version, about, long_about = None)] +pub struct Pcc { + #[command(subcommand)] + command: Commands, +} + +/// 主要命令 +#[derive(Subcommand, Debug)] +enum Commands { + + // 清屏 + #[command(about = "Clean the screen")] + Clear, + + // 断开当前连接 + #[command(about = "Exit from server")] + Exit, + + // 检查收到的消息 + #[command(about = "Check received")] + Received(ReceivedArgs), + + // 取出一条消息 + #[command(about = "Pop a message")] + Pop(PopArgs), + + // 发送消息 + #[command(about = "Send Message")] + Msg(MsgArgs), +} + +#[derive(Args, Debug)] +struct ReceivedArgs { + + #[arg(long)] + list: bool +} + +/// 发送消息 参数 +#[derive(Args, Debug)] +struct MsgArgs { + + // 消息内容 + #[arg(value_name = "CONTENT")] + message: String, +} + +#[derive(Args, Debug)] +struct PopArgs { } + +pub fn process_debug_cmd (cmd: Pcc, client: Arc) { + match cmd.command { + Commands::Clear => { + clearscreen::clear().expect("Failed to clear screen"); + } + + Commands::Exit => { + client.exit_server(); + } + + Commands::Received(args) => { + if args.list { + for msg in client.list_received() { + info!("{:?}", msg); + } + } else { + info!("Total {} messsage(s)!", client.list_received().iter().count()); + } + } + + Commands::Pop(_args) => { + info!("{:?}", client.pop_msg_or(GameMessage::Err)); + } + + Commands::Msg(args) => { + client.put_msg(ControlMessage::Msg(args.message)); + } + } +} + +#[allow(dead_code)] +fn put_to_list(client: Arc, message: ControlMessage) { + client.put_msg(message); +} \ No newline at end of file diff --git a/core/src/pad_service/mod.rs b/core/src/pad_service/mod.rs new file mode 100644 index 0000000..ae10f2d --- /dev/null +++ b/core/src/pad_service/mod.rs @@ -0,0 +1,5 @@ +pub mod client; +pub mod client_debug_cli; + +pub mod server; +pub mod server_debug_cli; \ No newline at end of file diff --git a/core/src/pad_service/server.rs b/core/src/pad_service/server.rs new file mode 100644 index 0000000..3c60864 --- /dev/null +++ b/core/src/pad_service/server.rs @@ -0,0 +1,643 @@ +pub mod nogamepads_server { + use std::collections::{HashMap, VecDeque}; + use crate::pad_data::pad_messages::nogamepads_message_encoder::NgpdMessageEncoder; + use crate::pad_data::pad_messages::nogamepads_message_transfer::{read_msg, send_msg}; + use crate::pad_data::pad_messages::nogamepads_messages::{ConnectionCallbackMessage, ConnectionMessage, ControlMessage, GameMessage, LeaveReason}; + use log::{error, info, warn}; + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + use std::process::exit; + use std::sync::atomic::AtomicBool; + use std::sync::atomic::Ordering::SeqCst; + use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; + use std::time::Duration; + use clap::CommandFactory; + use tokio::io::{AsyncReadExt, AsyncWriteExt, ReadHalf, WriteHalf}; + use tokio::net::{TcpListener, TcpStream}; + use tokio::{io, spawn}; + use tokio::runtime::Runtime; + use nogamepads::console_utils::debug_console::read_cli; + use nogamepads::convert_utils::convert_deque_to_vec; + use nogamepads::logger_utils::logger_build; + use crate::DEFAULT_PORT; + use crate::pad_data::game_profile::game_profile::GameProfile; + use crate::pad_data::pad_messages::nogamepads_messages::ConnectionErrorType::{ContainSamePlayer, GameLocked, PlayerBanned, WhatTheHell}; + use crate::pad_data::pad_messages::nogamepads_messages::GameMessage::Leave; + use crate::pad_data::pad_messages::nogamepads_messages::LeaveReason::ServerClosed; + use crate::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; + use crate::pad_service::server_debug_cli::{process_debug_cmd, Psc}; + + type PlayerMap = Arc>>; + type WriteList = Arc>>>; + type ReadList = Arc>>>; + + pub struct PadServer { + + // --- 主要参数 --- + + // 本地监听地址 + address: IpAddr, + + // 游戏信息 + game_profile: GameProfile, + + // 绑定端口 + port: u16, + + // 调试模式 + enable_console: bool, + + // 保持安静,不初始化 env_logger + quiet: bool, + + // --- 运行时参数 --- + + // 发送信息列表 + write_list: WriteList, + + // 读取信息列表 + read_list: ReadList, + + // 在线玩家 + online_players: PlayerMap, + + // 被封禁的玩家 + banned_players: PlayerMap, + + // 是否锁定该游戏:禁止后续玩家加入 + game_locked: AtomicBool, + + // 是否停止服务器 + stop: AtomicBool, + } + + impl Clone for PadServer { + fn clone(&self) -> Self { + PadServer { + address: self.address.clone(), + game_profile: self.game_profile.clone(), + port: self.port.clone(), + enable_console: self.enable_console, + quiet: self.quiet, + + write_list: self.write_list.clone(), + read_list: self.read_list.clone(), + online_players: self.online_players.clone(), + banned_players: self.banned_players.clone(), + game_locked: AtomicBool::new((&self.game_locked.load(SeqCst)).clone()), + stop: AtomicBool::new((&self.stop.load(SeqCst)).clone()), + } + } + } + + impl Default for PadServer { + fn default() -> Self { + PadServer { + address: IpAddr::from(Ipv4Addr::new(127, 0, 0, 1)), + game_profile: GameProfile::default(), + port: DEFAULT_PORT, + enable_console: false, + quiet: false, + + write_list: WriteList::default(), + read_list: ReadList::default(), + online_players: PlayerMap::default(), + banned_players: PlayerMap::default(), + game_locked: AtomicBool::new(false), + stop: AtomicBool::new(false), + } + } + } + + // 服务端构建部分 + impl PadServer { + + pub fn build_simple() -> Arc { + Arc::new(Self::default() + .addr(IpAddr::from(Ipv4Addr::new(127, 0, 0, 1)), DEFAULT_PORT) + .put_profile(GameProfile::default()).to_owned()) + } + + pub fn addr(&mut self, ip_addr: IpAddr, port: u16) -> &mut PadServer { + self.ip_addr(ip_addr).port(port) + } + + pub fn socket_addr(&mut self, socket_addr: SocketAddr) -> &mut PadServer { + self.ip_addr(socket_addr.ip()).port(socket_addr.port()) + } + + pub fn port(&mut self, port: u16) -> &mut PadServer { + self.port = port; + self + } + + pub fn ip_addr(&mut self, ip_addr: IpAddr) -> &mut PadServer { + self.address = ip_addr; + self + } + + pub fn put_profile(&mut self, profile: GameProfile) -> &mut PadServer { + self.game_profile = profile; + self + } + + pub fn enable_console(&mut self) -> &mut PadServer { + self.enable_console = true; + self + } + + pub fn quiet(&mut self) -> &mut PadServer { + self.quiet = true; + self + } + + pub fn build(&self) -> Arc { + Arc::new(self.clone()) + } + + } + + // 服务端消息管理 + impl PadServer { + + pub fn put_msg_to(&self, msg: GameMessage, player: &PlayerInfo) { + match self.write_list.lock() { + Ok(mut guard) => { + let hash = &player.account.player_hash.clone(); + if ! guard.contains_key(hash.as_str()) { + guard.entry(player.account.player_hash.clone()) + .or_insert_with(VecDeque::new) + .push_back(msg); + } + } + Err(_) => { + error!("Cannot lock \"{:?}\" in write_list", player.account.player_hash); + } + } + } + + pub fn put_msg_to_all(&self, msg: &GameMessage) { + match self.list_players() { + Ok(list) => { + for player in list { + self.put_msg_to(msg.clone(), &player); + } + } + Err(_) => { + error!("Cannot put GameMessage with no players."); + } + } + } + + pub fn pop_a_msg(&self, player: &PlayerInfo) -> Option { + match self.read_list.lock() { + Ok(mut guard) => { + match guard.get_mut(&player.account.player_hash) { + None => { None } + Some(queue) => { + if ! queue.is_empty() { + queue.pop_front() + } else { + guard.remove(&player.account.player_hash); + None + } + } + } + } + Err(_) => { + error!("Cannot lock \"{:?}\" in read_list", player.account.player_hash); + None + } + } + } + + pub fn pop_msg_or(&self, player: &PlayerInfo, or: ControlMessage) -> ControlMessage { + self.pop_a_msg(player).unwrap_or(or) + } + + pub fn list_received(&self, player: &PlayerInfo) -> Vec { + match self.read_list.lock() { + Ok(guard) => { + match guard.get_key_value(player.account.player_hash.as_str()) { + None => { Vec::new() } + Some(result) => { + convert_deque_to_vec(result.1) + } + } + } + Err(_) => { Vec::new() } + } + } + } + + // 服务端玩家管理 + impl PadServer { + + pub fn is_player_online (&self, player: &PlayerInfo) -> bool { + let guard = self.online_players.lock().unwrap(); + guard.contains_key(&player.account.player_hash) + } + + pub fn set_player_online (&self, player: &PlayerInfo, online: bool) { + let online_current = self.is_player_online(player); + if online_current && !online { + let mut guard = self.online_players.lock().unwrap(); + guard.remove(&player.account.player_hash); + info!("{} is OFFLINE!", player.account.id); + } else if !online_current && online { + let mut guard = self.online_players.lock().unwrap(); + guard.insert(player.account.player_hash.clone(), player.clone()); + info!("{} is ONLINE!", player.account.id); + } + } + + pub fn is_player_banned (&self, player: &PlayerInfo) -> bool { + let guard = self.banned_players.lock().unwrap(); + guard.contains_key(&player.account.player_hash) + } + + pub fn kick_player(&self, player: &PlayerInfo) { + if self.is_player_online(player) { + self.put_msg_to(Leave(LeaveReason::YouAreKicked), player); + } + } + + pub fn ban_player(&self, player: &PlayerInfo) { + self.set_player_banned(player, true); + if self.is_player_online(player) { + self.put_msg_to(Leave(LeaveReason::YouAreBanned), player); + } + } + + pub fn pardon_player(&self, player: &PlayerInfo) { + self.set_player_banned(player, false); + } + + fn set_player_banned (&self, player: &PlayerInfo, banned: bool) { + let banned_current = self.is_player_banned(player); + if banned_current && !banned { + let mut guard = self.banned_players.lock().unwrap(); + guard.remove(&player.account.player_hash); + info!("Pardoned player {}", player.account.id); + } else if !banned_current && banned { + let mut guard = self.banned_players.lock().unwrap(); + guard.insert(player.account.player_hash.clone(), player.clone()); + info!("Banned player {}!", player.account.id); + } + } + + pub fn list_players(&self) -> Result, PoisonError>>> { + match self.online_players.lock() { + Ok(guard) => { + Ok(guard.values().cloned().collect()) + } + Err(err) => Err(err) + } + } + + pub fn list_players_banned(&self) -> Result, PoisonError>>> { + match self.banned_players.lock() { + Ok(guard) => { + Ok(guard.values().cloned().collect()) + } + Err(err) => Err(err) + } + } + } + + // 服务端状态控制 + #[allow(dead_code)] + impl PadServer { + + pub fn stop_server(&self) { + self.put_msg_to_all(&Leave(ServerClosed)); + self.stop.store(true, SeqCst); + } + + pub fn start_server(self: Arc) { + + // 构建 Logger + if ! self.quiet { + logger_build(); + } + + // 运行时 + let runtime = Self::get_runtime(); + + info!("Starting \"NoGamepads Server\"."); + + // 入口 + let console = self.enable_console; + let entry = self.get_entry(console); + + // 阻塞运行 + runtime.block_on(entry); + } + + fn get_runtime() -> Runtime { + tokio::runtime::Builder::new_multi_thread() + .thread_name("nogpad-server") + .thread_stack_size(32 * 1024 * 1024) + .enable_time() + .enable_io() + .build() + .unwrap() + } + + fn get_entry(self: Arc, debug: bool) -> impl Future + Send + 'static { + async move { + let main_thread = spawn({ + let client = Arc::clone(&self); + async move { + Self::main_request_thread(client).await + } + }); + + let background_thread = spawn({ + let client = Arc::clone(&self); + async move { + Self::background_thread(client).await + } + }); + + if debug { + let debug_cli = spawn({ + let client = Arc::clone(&self); + async move { + Self::process_debug_cli(client).await + } + }); + + let _ = tokio::join!(debug_cli, main_thread, background_thread); + } else { + let _ = tokio::join!(main_thread, background_thread); + } + } + } + + fn lock_game(&self) { + self.game_locked.store(true, SeqCst); + } + + fn unlock_game(&self) { + self.game_locked.store(false, SeqCst); + } + + fn is_game_locked(&self) -> bool { + self.game_locked.load(SeqCst) + } + + async fn main_request_thread(self: Arc) { + + let addr_str = format!("{}:{}", self.address.to_string(), self.port); + info!("Server listening at {}", addr_str); + + // Tcp 监听器 + let listener : TcpListener; + match TcpListener::bind(&addr_str).await { + Ok(result) => { + info!("Listener created."); + listener = result; + } + Err(_) => { + error!("Server listening at {} failed!", addr_str); + exit(1); + } + } + + // 请求信息循环 + loop { + match listener.accept().await { + Ok((stream, _)) => { + spawn(Self::process_request(Arc::clone(&self), stream)); + } + Err(error) => { + error!("Error: {}", error); + } + } + } + } + + async fn process_request(self: Arc, mut stream: TcpStream) { + let mut buffer = [0; 1024]; + let connection_msg : ConnectionMessage = read_msg(&mut buffer, &mut stream).await; + match connection_msg { + + // 客户端请求加入游戏,并建立长连接 + ConnectionMessage::Connection(info) => { + + // 加入游戏资格检测 + info!("Account {} trying to connect.", info.account.player_hash); + + // 0. 当前游戏是否已经锁定? + if self.is_game_locked() { + // 当前游戏已经锁定,禁止加入玩家,发送失败信息,并断开连接 + send_msg(&mut stream, ConnectionCallbackMessage::Deny(GameLocked)).await; + return; + } + + // 1. 是否存在重复玩家? + let online = self.is_player_online(&info); + if online { + // 当前玩家已在线,发送失败信息,并断开连接 + send_msg(&mut stream, ConnectionCallbackMessage::Deny(ContainSamePlayer)).await; + return; + } + + // 2. 该玩家是否被封禁? + let banned = self.is_player_banned(&info); + if banned { + // 当前玩家已被封禁,发送失败信息,并断开连接 + send_msg(&mut stream, ConnectionCallbackMessage::Deny(PlayerBanned)).await; + return; + } + + // OK!若执行到此处,说明该玩家具有加入资格,Welcome! + + send_msg(&mut stream, ConnectionCallbackMessage::Ok).await; + let callback : ConnectionMessage = read_msg(&mut buffer, &mut stream).await; + + match callback { + // 玩家已就绪,发送 Welcome 信息以邀请该玩家加入游戏 + ConnectionMessage::Ready => { + info!("Player \"{}\" is ready!", info.account.id); + + // 发送 Welcome + send_msg(&mut stream, ConnectionCallbackMessage::Welcome).await; + + // 注册该玩家到在线列表 + self.set_player_online(&info, true); + + // 启动控制循环 + + spawn(Self::long_connection(Arc::clone(&self), stream, info)); + }, + _ => { + send_msg(&mut stream, ConnectionCallbackMessage::Deny(WhatTheHell)).await; // WTH ? + } + } + } + + // 客户端请求获得游戏信息 + ConnectionMessage::RequestProfile => { + // 发送游戏信息到客户端 + send_msg(&mut stream, ConnectionCallbackMessage::Profile(self.game_profile.clone())).await; + } + + // 客户端发来了错误信息 + ConnectionMessage::Err => { + match stream.peer_addr() { + Ok(addr) => { + warn!("Received an error message from {}.", addr.to_string()); + } + Err(_) => { + warn!("Received an error message from unknown pad_service."); + } + } + } + + // 客户端发来了不相干的信息 + _ => { + match stream.peer_addr() { + Ok(addr) => { + warn!("Received unknown connection message from {}.", addr.to_string()); + } + Err(_) => { + warn!("Received unknown connection message from unknown pad_service."); + } + } + } + } + } + + async fn long_connection(self: Arc, stream: TcpStream, player_info: PlayerInfo) { + let player_info_arc = Arc::new(player_info); + let (reader, writer) = io::split(stream); + spawn(Self::read_task(Arc::clone(&self), reader, Arc::clone(&player_info_arc))); + spawn(Self::write_task(Arc::clone(&self), writer, Arc::clone(&player_info_arc))); + } + + async fn read_task(self: Arc, + mut reader: ReadHalf, + player_info: Arc) { + let player_hash = player_info.account.player_hash.clone(); + let mut buf = [0u8; 1024]; + loop { + match reader.read(&mut buf).await { + Ok(0) => break, + Ok(n) => { + let msg = ControlMessage::de(buf[0..n].to_vec()); + { + match self.read_list.lock() { + Ok(mut guard) => { + info!("{:?} from {}({})", &msg, player_info.customize.nickname, player_info.account.id); + guard + .entry(player_hash.clone()) + .or_insert_with(VecDeque::new) + .push_back(msg); + + } + Err(_) => { + } + } + } + } + Err(e) => { + warn!("Error reading from stream: {}", e); + + self.set_player_online(&player_info, false); + + // 放入一条错误信息到队列,使 write_task 及时发现该玩家离开 + self.put_msg_to(GameMessage::Err, &player_info); + + break; + } + } + } + } + + async fn write_task(self: Arc, + mut writer: WriteHalf, + player_info: Arc) { + let player_hash = player_info.account.player_hash.clone(); + let mut exit = false; + loop { + let msg : Option; + match self.write_list.lock() { + Ok(mut hash_map) => { + if ! hash_map.is_empty() { + match hash_map.get_mut(&player_hash) { + None => { + msg = None; + } + Some(queue) => { + if ! queue.is_empty() { + msg = queue.pop_front(); + } else { + msg = None; + hash_map.remove(&player_hash); + } + } + } + } + else { msg = None; } + } + Err(_) => { + msg = None; + } + } + if msg.is_some() { + let msg = msg.unwrap(); + match &writer.write_all(NgpdMessageEncoder::en(&msg).as_slice()).await { + Ok(_) => { + + info!("Sent {:?} to {}", msg, &player_info.account.id); + } + Err(error) => { + warn!("Sent {:?} to {} failed!", msg, &player_info.account.id); + warn!("{:?}", error); + + exit = true; + } + } + } + if exit { + warn!("Long connection between \"{}\" closed.", &player_info.account.id); + break + } + } + } + + async fn background_thread(self: Arc) { + loop { + // 退出程序的监听 + if self.stop.load(SeqCst) { + tokio::time::sleep(Duration::from_secs(1)).await; + + info!("Main thread exited."); + exit(0); + } + } + } + + async fn process_debug_cli(self: Arc) { + loop { + if self.stop.load(SeqCst) { + info!("Debug console exited"); + return; + } + tokio::time::sleep(Duration::from_secs_f64(0.2)).await; + let option: Option = read_cli( + format!("SERVER {}> ", self.address.to_string()).as_str(), + "psc".to_string(), + Psc::command() + ).await; + match option { + None => {} + Some(cmd) => { + process_debug_cmd(cmd, Arc::clone(&self)); + } + } + } + } + } +} \ No newline at end of file diff --git a/core/src/pad_service/server_debug_cli.rs b/core/src/pad_service/server_debug_cli.rs new file mode 100644 index 0000000..ab9e5a5 --- /dev/null +++ b/core/src/pad_service/server_debug_cli.rs @@ -0,0 +1,204 @@ +use std::collections::HashMap; +use crate::pad_data::pad_messages::nogamepads_messages::{ControlMessage, GameMessage}; +use crate::pad_service::server::nogamepads_server::PadServer; +use clap::{Args, Parser, Subcommand}; +use std::ops::{Index}; +use std::sync::{Arc, MutexGuard, PoisonError}; +use log::{error, info}; +use crate::pad_data::pad_player_info::nogamepads_player_info::PlayerInfo; + +/// NoGamePads Server - Cli +#[derive(Parser, Debug)] +#[command(author, version, about, long_about = None)] +pub struct Psc { + #[command(subcommand)] + command: Commands, +} + +#[derive(Subcommand, Debug)] +enum Commands { + + // 清屏 + #[command(about = "Clean the screen")] + Clear, + + // 关闭服务器 + #[command(about = "Close the server")] + Stop, + + // 展示所有玩家 + #[command(about = "List all online players")] + List, + + // 展示所有封禁的玩家 + #[command(about = "List all banned players")] + Banned, + + // 检查收到的消息 + #[command(about = "Check received")] + Received(ReceivedArgs), + + // 取出一条消息 + #[command(about = "Pop a message")] + Pop(PlayerArgs), + + // 踢出玩家 + #[command(about = "Kick a player")] + Kick(PlayerArgs), + + // 封禁玩家 + #[command(about = "Ban a player")] + Ban(PlayerArgs), + + // 解封(赦免)玩家 + #[command(about = "Pardon a player")] + Pardon(PlayerArgs), + + // 激活事件触发器 + #[command(about = "Send SkinEventTrigger")] + Event(EventArgs) +} + +// 检查收到的消息 +#[derive(Args, Debug)] +struct ReceivedArgs { + + #[arg(default_value = "0")] + player: usize, + + #[arg(long)] + list: bool, +} + +/// 激活事件触发器 参数 +#[derive(Args, Debug)] +struct EventArgs { + + // 玩家序号 + #[arg(value_name = "PLAYER_INDEX")] + index: usize, + + // 事件编号 + #[arg(value_name = "CONTENT")] + message: u8, +} + +#[derive(Args, Debug)] +struct PlayerArgs { + + // 玩家序号 + #[arg(value_name = "PLAYER_INDEX")] + index: usize +} + +pub fn process_debug_cmd (cmd: Psc, server: Arc) { + match cmd.command { + + Commands::Clear => { + clearscreen::clear().expect("Failed to clear screen"); + } + + Commands::Stop => { + server.stop_server(); + } + + Commands::List => { + print_player_list(server.list_players()); + } + + Commands::Banned => { + print_player_list(server.list_players_banned()); + } + + Commands::Received(args) => { + let players = server.list_players().unwrap_or(Vec::new()); + let player = players.index(args.player.clamp(0, players.iter().count() -1)); + if args.list { + for msg in server.list_received(player) { + info!("{:?}", msg); + } + } else { + info!("Total {} messsage(s)!", server.list_received(player).iter().count()); + } + } + + Commands::Pop(args) => { + match get_player_by_index(&server, args.index) { + None => { + error!("Pup message failed : Player index \"{}\" not found!", args.index); + } + Some(player) => { + let message = server.pop_msg_or(&player, ControlMessage::Err); + info!("{:?}", message); + } + } + } + + Commands::Kick(args) => { + let player = get_player_by_index(&server, args.index); + if player.is_some() { + let player = player.unwrap(); + server.kick_player(&player); + } + } + + Commands::Ban(args) => { + let player = get_player_by_index(&server, args.index); + if player.is_some() { + let player = player.unwrap(); + server.ban_player(&player); + } + } + + Commands::Pardon(args) => { + let player = get_player_by_ban_index(&server, args.index); + if player.is_some() { + let player = player.unwrap(); + server.pardon_player(&player); + } + } + + Commands::Event(args) => { + put_to_list(server, args.index, GameMessage::SkinEventTrigger(args.message)); + } + } +} + +fn put_to_list(server: Arc, player_index: usize, message: GameMessage) { + match get_player_by_index(&server, player_index) { + None => { + error!("Put message failed : Player index \"{}\" not found!", player_index); + } + Some(player) => { + server.put_msg_to(message, &player); + } + } +} + +fn get_player_by_index(server: &Arc, index: usize) -> Option { + let list = server.list_players().unwrap_or(Vec::new()); + let max = list.iter().count(); + let index = if max > 0 { index.clamp(0, max - 1) } else { 0 }; + + let result = list.get(index).cloned(); + result +} + +fn get_player_by_ban_index(server: &Arc, index: usize) -> Option { + let list = server.list_players_banned().unwrap_or(Vec::new()); + let max = list.iter().count(); + let index = if max > 0 { index.clamp(0, max - 1) } else { 0 }; + + let result = list.get(index).cloned(); + result +} + +fn print_player_list(list: Result, PoisonError>>>) { + let list = list.unwrap_or(Vec::new()); + let mut i = 0; + for player in list { + let n = player.customize.nickname; + info!("({}){} ", i, n); + i += 1; + } +} \ No newline at end of file -- cgit